feat: Trino federation + Hadoop external tables + LLM data catalog

Trino Federation tab (3 sub-views):
- Federated: catalog landscape + a single cross-source SQL that joins
  PostgreSQL + MySQL + MongoDB (region scorecard) — the federation proof,
  computed in the background and cached (large full scans take ~2 min).
- Hadoop Lake: all federated business data materialized as external Iceberg
  tables on HDFS (iceberg.hadoop.*_ext, ~120k rows each) with live, fast
  business analytics (revenue by region/channel, top customers, HR by
  department, supply by type, telemetry averages). Includes a one-click
  "rebuild external tables" job.
- Data Dictionary: every business table + column with masked / visible PII
  badges and categories.

Backend trino_federated.py: /catalogs, /marquee(+refresh), /lake,
/materialize(+status), /dictionary. Name-based PII detection flags raw PII
in derived/lake tables as visible vs physically-masked curated layer.

LLM context: platform_context now emits a full BUSINESS DATA CATALOG section
(tables, columns, types, source row counts, federated scorecard) with exact
per-column masked/visible status, so the assistant knows the data in detail
and what is masked vs not.
This commit is contained in:
mo
2026-06-28 18:01:25 +00:00
parent 8d28695868
commit ea3e59cf9c
8 changed files with 855 additions and 5 deletions
+1 -1
View File
@@ -4,7 +4,7 @@ WORKDIR /app
RUN apt-get update && apt-get install -y --no-install-recommends curl && rm -rf /var/lib/apt/lists/*
COPY requirements.txt .
RUN pip install --no-cache-dir -r requirements.txt
COPY main.py lab_context.py agent_terminal.py workload.py node_registry.py node_ops.py topology_views.py supervisor.py approval_service.py db.py dockhand_envs.py presentation.py database_inventory.py presentation_upload.py presentation_static.py storage_s3.py elasticsearch_api.py sql_console.py hdfs_api.py ssh_terminal.py pipeline_ops.py hadoop_analytics.py agent_ops.py cdc_consumer.py movements.py dataflow.py streaming_ops.py spark_workbench.py hadoop_sql.py hdfs_kafka.py webhdfs_util.py pii_catalog.py platform_context.py hive_bench_seed.json .
COPY main.py lab_context.py agent_terminal.py workload.py node_registry.py node_ops.py topology_views.py supervisor.py approval_service.py db.py dockhand_envs.py presentation.py database_inventory.py presentation_upload.py presentation_static.py storage_s3.py elasticsearch_api.py sql_console.py hdfs_api.py ssh_terminal.py pipeline_ops.py hadoop_analytics.py agent_ops.py cdc_consumer.py movements.py dataflow.py streaming_ops.py spark_workbench.py hadoop_sql.py hdfs_kafka.py webhdfs_util.py pii_catalog.py platform_context.py trino_federated.py hive_bench_seed.json .
RUN mkdir -p /data
ENV DATABASE_URL=sqlite:////data/atc-agents.db
EXPOSE 3201
+2
View File
@@ -53,6 +53,7 @@ from dataflow import router as dataflow_router
from streaming_ops import router as streaming_router
from spark_workbench import router as spark_workbench_router
from pii_catalog import router as pii_router
from trino_federated import router as federated_router
from ssh_terminal import ssh_session
from node_registry import NODE_IDS, NODE_AGENT, NODE_REGISTRY, is_node_id
from node_ops import build_node_detail, probe_node, run_node_probe_task
@@ -757,6 +758,7 @@ app.include_router(dataflow_router)
app.include_router(streaming_router)
app.include_router(spark_workbench_router)
app.include_router(pii_router)
app.include_router(federated_router)
app.add_middleware(
CORSMiddleware,
allow_origins=["*"],
+72 -1
View File
@@ -132,13 +132,84 @@ def build_masking_section() -> str:
return "\n".join(lines)
def build_business_data_section() -> str:
"""Detailed inventory of the actual business data — tables, columns, types,
row counts and per-column masked/visible status — so the assistant knows the
data in detail and exactly what is masked vs not."""
lines: list[str] = ["=== BUSINESS DATA CATALOG (live, every detail) ==="]
# source row totals (cheap estimates)
try:
import sql_console as s
totals = {
"PostgreSQL sales_orders": s._table_row_count("postgres", "public.sales_orders"),
"MySQL employee_events": s._table_row_count("mysql", "hr.employee_events"),
"MongoDB supplychain.events": s._table_row_count("mongodb", "supplychain.events"),
}
lines.append("Source volumes (live row counts):")
for k, v in totals.items():
lines.append(f" - {k}: {v:,} rows" if isinstance(v, int) else f" - {k}: ~")
except Exception:
pass
# full column dictionary + masking per column
try:
from trino_federated import build_dictionary
d = build_dictionary()
summ = d.get("summary", {})
lines.append(
f"Catalogued tables: {summ.get('tables', 0)} — columns: {summ.get('columns', 0)}, "
f"PII columns: {summ.get('pii_columns', 0)} ({summ.get('masked_columns', 0)} masked)."
)
for t in d.get("tables", []):
lines.append(f"\n{t['engine']} · {t['fqn']}{t['desc']}")
for c in t.get("columns", []):
flag = ""
if c.get("masked"):
flag = f" [MASKED · {c.get('category')}]"
elif c.get("pii"):
flag = f" [PII visible · {c.get('category')}]"
lines.append(f" · {c['name']} ({c['type']}){flag}")
except Exception as exc:
lines.append(f"(data dictionary unavailable: {exc})")
# federated cross-source marquee (region scorecard), if computed
try:
from trino_federated import _marquee, _load_marquee
if _marquee.get("data") is None:
_load_marquee()
m = _marquee.get("data")
if m and m.get("rows"):
lines.append("\nFederated region scorecard (one Trino SQL across PostgreSQL+MySQL+MongoDB):")
for r in m["rows"][:8]:
rev = r.get("revenue")
rev_s = f"{rev:,.0f}" if isinstance(rev, (int, float)) else "?"
lines.append(
f" - {r.get('region')}: orders={r.get('orders')}, revenue={rev_s}, "
f"hr_events={r.get('hr_events')}, supply_events={r.get('supply_events')}"
)
except Exception:
pass
lines.append(
"\nData is also materialized into the Hadoop lake as external Iceberg tables "
"(iceberg.hadoop.*_ext) and exposed through one federated Trino engine "
"(catalogs: postgres_sales, mysql_hr, mongodb_supplychain, cassandra_telemetry, iceberg, kafka)."
)
return "\n".join(lines)
def build_llm_addendum() -> str:
try:
platform = build_platform_section()
except Exception as exc:
platform = f"(platform section error: {exc})"
try:
business = build_business_data_section()
except Exception as exc:
business = f"(business data section error: {exc})"
try:
masking = build_masking_section()
except Exception as exc:
masking = f"(masking section error: {exc})"
return "\n\n".join([platform, masking])
return "\n\n".join([platform, business, masking])
+400
View File
@@ -0,0 +1,400 @@
"""Trino federated business analytics + Hadoop lakehouse (external Iceberg tables).
Provides:
- /api/federated/catalogs : the federated catalog landscape (estimates, fast)
- /api/federated/marquee : a single cross-source SQL joining PG+MySQL+Mongo
(the federation proof) — computed in the background, cached
- /api/federated/lake : business analytics over the Hadoop Iceberg lake
(iceberg.hadoop external tables) — live & fast
- /api/federated/materialize: (re)build the Hadoop external tables from the sources
- /api/federated/dictionary : business data dictionary with masked / visible flags
(also fed to the LLM so it knows the data in detail)
"""
from __future__ import annotations
import json
import threading
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import Any
from fastapi import APIRouter
from fastapi.responses import JSONResponse
router = APIRouter(prefix="/api/federated", tags=["federated"])
BUSINESS_CATALOGS = ["postgres_sales", "mysql_hr", "mongodb_supplychain", "cassandra_telemetry", "iceberg"]
ALL_CATALOGS = BUSINESS_CATALOGS + ["kafka", "system"]
# Hadoop external tables (Iceberg-on-HDFS) materialized from the federated sources.
# Each: target table name in iceberg.hadoop -> source SELECT (timestamps cast to
# timestamp(6) which Iceberg requires; measures cast to double).
LAKE_SPECS: list[tuple[str, str, str]] = [
("orders_ext", "Sales orders (PostgreSQL)",
"SELECT order_id, customer_id, customer_name, customer_email, product_id, region, "
"sales_channel, order_status, currency, CAST(amount AS double) amount, "
"CAST(order_ts AS timestamp(6)) order_ts "
"FROM postgres_sales.public.sales_orders WHERE customer_name IS NOT NULL LIMIT 120000"),
("employees_ext", "Employee events (MySQL)",
"SELECT employee_id, employee_name, employee_email, employee_phone, department, role_name, "
"region, event_type, CAST(salary_change AS double) salary_change, "
"CAST(event_ts AS timestamp(6)) event_ts "
"FROM mysql_hr.hr.employee_events WHERE employee_name IS NOT NULL LIMIT 120000"),
("supply_events_ext", "Supply chain events (MongoDB)",
"SELECT event_id, type, region, source, CAST(amount AS double) amount, "
"CAST(ts AS timestamp(6)) ts FROM mongodb_supplychain.supplychain.events LIMIT 120000"),
("telemetry_ext", "Device telemetry (Cassandra)",
"SELECT device_id, CAST(metric_ts AS timestamp(6)) metric_ts, metric_type, "
"CAST(metric_value AS double) metric_value FROM cassandra_telemetry.telemetry.device_metrics LIMIT 120000"),
("customers_ext", "Distinct customers (PostgreSQL)",
"SELECT DISTINCT customer_id, customer_name, customer_email, region, sales_channel "
"FROM postgres_sales.public.sales_orders WHERE customer_name IS NOT NULL LIMIT 60000"),
("products_ext", "Distinct products (PostgreSQL)",
"SELECT DISTINCT product_id, region, sales_channel FROM postgres_sales.public.sales_orders "
"WHERE product_id IS NOT NULL LIMIT 20000"),
]
MARQUEE_SQL = (
"WITH o AS (SELECT region, count(*) orders, sum(amount) revenue "
"FROM postgres_sales.public.sales_orders GROUP BY region),\n"
" e AS (SELECT region, count(*) hr_events FROM mysql_hr.hr.employee_events GROUP BY region),\n"
" s AS (SELECT region, count(*) supply_events FROM mongodb_supplychain.supplychain.events GROUP BY region)\n"
"SELECT COALESCE(o.region, e.region, s.region) AS region,\n"
" o.orders, o.revenue, e.hr_events, s.supply_events\n"
"FROM o FULL JOIN e ON o.region = e.region\n"
" FULL JOIN s ON COALESCE(o.region, e.region) = s.region\n"
"ORDER BY o.revenue DESC NULLS LAST"
)
def _trino(sql: str, limit: int = 500) -> dict[str, Any]:
import sql_console as s
return s._run_trino(sql, limit)
def _rows_as_dicts(res: dict) -> list[dict]:
cols = res.get("columns", [])
return [{cols[i]: r[i] for i in range(min(len(cols), len(r)))} for r in res.get("rows", [])]
# ──────────────────────────────────────────────────────────────────────────────
# Marquee federated query (cross-source) — background cached
# ──────────────────────────────────────────────────────────────────────────────
MARQUEE_PATH = Path("/data/federated_marquee.json")
_MARQUEE_TTL = 1800.0
_marquee: dict[str, Any] = {"running": False, "data": None}
_marquee_lock = threading.Lock()
def _load_marquee() -> None:
try:
if MARQUEE_PATH.exists():
_marquee["data"] = json.loads(MARQUEE_PATH.read_text())
except Exception:
pass
def _marquee_worker() -> None:
try:
a = time.time()
res = _trino(MARQUEE_SQL, 50)
elapsed = int((time.time() - a) * 1000)
data = {
"ok": res.get("ok", False),
"sql": MARQUEE_SQL,
"catalogs": ["postgres_sales", "mysql_hr", "mongodb_supplychain"],
"elapsed_ms": elapsed,
"rows": _rows_as_dicts(res) if res.get("ok") else [],
"error": None if res.get("ok") else str(res.get("error", ""))[:300],
"generated_at": datetime.now(timezone.utc).isoformat(),
}
_marquee["data"] = data
try:
MARQUEE_PATH.write_text(json.dumps(data))
except Exception:
pass
finally:
_marquee["running"] = False
def _ensure_marquee(force: bool = False) -> None:
with _marquee_lock:
if _marquee["running"]:
return
data = _marquee["data"]
fresh = data and (time.time() - _ts(data.get("generated_at")) < _MARQUEE_TTL)
if fresh and not force:
return
_marquee["running"] = True
threading.Thread(target=_marquee_worker, daemon=True).start()
def _ts(iso: str | None) -> float:
if not iso:
return 0.0
try:
return datetime.fromisoformat(iso).timestamp()
except Exception:
return 0.0
@router.get("/marquee")
async def get_marquee():
if _marquee["data"] is None:
_load_marquee()
_ensure_marquee()
data = _marquee["data"]
return {
"ok": True,
"running": _marquee["running"],
"marquee": data,
"sql": MARQUEE_SQL,
}
@router.post("/marquee/refresh")
async def refresh_marquee():
_ensure_marquee(force=True)
return {"ok": True, "running": True}
# ──────────────────────────────────────────────────────────────────────────────
# Catalog landscape (fast — SHOW + planner estimates)
# ──────────────────────────────────────────────────────────────────────────────
@router.get("/catalogs")
async def get_catalogs():
import sql_console as s
cat_res = _trino("SHOW CATALOGS", 100)
catalogs = [r[0] for r in cat_res.get("rows", [])] if cat_res.get("ok") else ALL_CATALOGS
totals = {
"orders": s._table_row_count("postgres", "public.sales_orders"),
"hr_events": s._table_row_count("mysql", "hr.employee_events"),
"supply_events": s._table_row_count("mongodb", "supplychain.events"),
}
out = []
labels = {
"postgres_sales": ("PostgreSQL", "Sales / orders OLTP", "#336791"),
"mysql_hr": ("MySQL", "HR / workforce", "#00758f"),
"mongodb_supplychain": ("MongoDB", "Supply chain events", "#4db33d"),
"cassandra_telemetry": ("Cassandra", "Device telemetry", "#1287b1"),
"iceberg": ("Iceberg / Hadoop", "Lakehouse external tables", "#5b8def"),
"kafka": ("Kafka", "CDC + streaming topics", "#231f20"),
"system": ("Trino", "Engine metadata", "#dd00a1"),
}
for c in catalogs:
lbl, desc, color = labels.get(c, (c, "", "#888"))
item: dict[str, Any] = {"catalog": c, "label": lbl, "desc": desc, "color": color, "business": c in BUSINESS_CATALOGS}
if c == "postgres_sales":
item["rows"] = totals["orders"]
elif c == "mysql_hr":
item["rows"] = totals["hr_events"]
elif c == "mongodb_supplychain":
item["rows"] = totals["supply_events"]
out.append(item)
return {"ok": True, "catalogs": out, "source_totals": totals, "count": len(out)}
# ──────────────────────────────────────────────────────────────────────────────
# Hadoop lake analytics (live over the small materialized external tables)
# ──────────────────────────────────────────────────────────────────────────────
def _terms(sql: str, key_col: str, count_col: str = "c", val_col: str | None = None) -> list[dict]:
res = _trino(sql, 100)
if not res.get("ok"):
return []
cols = res.get("columns", [])
idx = {c: i for i, c in enumerate(cols)}
out = []
for r in res.get("rows", []):
item = {"key": r[idx[key_col]] if key_col in idx else None,
"count": r[idx[count_col]] if count_col in idx else 0}
if val_col and val_col in idx:
item["value"] = round(float(r[idx[val_col]] or 0), 2)
out.append(item)
return out
def _lake_table_exists(name: str) -> bool:
res = _trino(f"SELECT count(*) FROM iceberg.hadoop.{name}", 1)
return res.get("ok", False)
@router.get("/lake")
async def get_lake():
# which external tables are present
tbl_res = _trino("SHOW TABLES FROM iceberg.hadoop", 200)
tables = [r[0] for r in tbl_res.get("rows", [])] if tbl_res.get("ok") else []
table_stats = []
for t in tables:
cnt = _trino(f"SELECT count(*) FROM iceberg.hadoop.{t}", 1)
n = cnt.get("rows", [[0]])[0][0] if cnt.get("ok") else 0
table_stats.append({"table": t, "rows": n})
has_orders = "orders_ext" in tables
has_emp = "employees_ext" in tables
has_sup = "supply_events_ext" in tables
has_tel = "telemetry_ext" in tables
orders = {}
if has_orders:
orders = {
"by_region": _terms("SELECT region, count(*) c, sum(amount) rev FROM iceberg.hadoop.orders_ext GROUP BY region ORDER BY rev DESC", "region", "c", "rev"),
"by_channel": _terms("SELECT sales_channel k, count(*) c, sum(amount) rev FROM iceberg.hadoop.orders_ext GROUP BY sales_channel ORDER BY rev DESC", "k", "c", "rev"),
"by_status": _terms("SELECT order_status k, count(*) c FROM iceberg.hadoop.orders_ext GROUP BY order_status ORDER BY c DESC", "k", "c"),
"top_customers": _terms("SELECT customer_name k, count(*) c, sum(amount) rev FROM iceberg.hadoop.orders_ext GROUP BY customer_name ORDER BY rev DESC LIMIT 10", "k", "c", "rev"),
"revenue": (_trino("SELECT sum(amount), avg(amount) FROM iceberg.hadoop.orders_ext", 1).get("rows") or [[0, 0]])[0],
}
hr = {}
if has_emp:
hr = {
"by_department": _terms("SELECT department k, count(*) c FROM iceberg.hadoop.employees_ext GROUP BY department ORDER BY c DESC", "k", "c"),
"by_role": _terms("SELECT role_name k, count(*) c FROM iceberg.hadoop.employees_ext GROUP BY role_name ORDER BY c DESC LIMIT 12", "k", "c"),
}
supply = {}
if has_sup:
supply = {"by_type": _terms("SELECT type k, count(*) c FROM iceberg.hadoop.supply_events_ext GROUP BY type ORDER BY c DESC", "k", "c")}
telemetry = {}
if has_tel:
telemetry = {"by_metric": _terms("SELECT metric_type k, count(*) c, avg(metric_value) v FROM iceberg.hadoop.telemetry_ext GROUP BY metric_type ORDER BY c DESC", "k", "c", "v")}
return {
"ok": True,
"generated_at": datetime.now(timezone.utc).isoformat(),
"tables": table_stats,
"total_rows": sum(t["rows"] or 0 for t in table_stats),
"orders": orders,
"hr": hr,
"supply": supply,
"telemetry": telemetry,
}
# ──────────────────────────────────────────────────────────────────────────────
# Materialize the Hadoop external tables from the sources
# ──────────────────────────────────────────────────────────────────────────────
_mat_state: dict[str, Any] = {"running": False, "current": None, "done": [], "errors": [], "started_at": None, "finished_at": None}
_mat_lock = threading.Lock()
def _materialize_worker() -> None:
try:
for name, label, sql in LAKE_SPECS:
_mat_state["current"] = f"{name} ({label})"
try:
_trino(f"DROP TABLE IF EXISTS iceberg.hadoop.{name}", 1)
res = _trino(f"CREATE TABLE iceberg.hadoop.{name} AS {sql}", 5)
if res.get("ok"):
cnt = _trino(f"SELECT count(*) FROM iceberg.hadoop.{name}", 1)
n = cnt.get("rows", [[0]])[0][0] if cnt.get("ok") else 0
_mat_state["done"].append({"table": name, "rows": n})
else:
_mat_state["errors"].append(f"{name}: {str(res.get('error',''))[:160]}")
except Exception as exc: # noqa: BLE001
_mat_state["errors"].append(f"{name}: {str(exc)[:160]}")
finally:
_mat_state["current"] = None
_mat_state["running"] = False
_mat_state["finished_at"] = datetime.now(timezone.utc).isoformat()
@router.post("/materialize")
async def materialize():
with _mat_lock:
if _mat_state["running"]:
return {"ok": True, "already_running": True, "state": _mat_state}
_mat_state.update({"running": True, "current": "starting…", "done": [], "errors": [],
"started_at": datetime.now(timezone.utc).isoformat(), "finished_at": None})
threading.Thread(target=_materialize_worker, daemon=True).start()
return {"ok": True, "started": True, "tables": [s[0] for s in LAKE_SPECS]}
@router.get("/materialize/status")
async def materialize_status():
return {"ok": True, **_mat_state}
# ──────────────────────────────────────────────────────────────────────────────
# Business data dictionary (columns + masked/visible) — UI + LLM
# ──────────────────────────────────────────────────────────────────────────────
DICT_TABLES = [
("postgres_sales", "public", "sales_orders", "PostgreSQL", "Sales orders (OLTP, CDC source)"),
("mysql_hr", "hr", "employee_events", "MySQL", "Employee / HR events (CDC source)"),
("mongodb_supplychain", "supplychain", "events", "MongoDB", "Supply chain events (CDC source)"),
("cassandra_telemetry", "telemetry", "device_metrics", "Cassandra", "IoT device telemetry"),
("iceberg", "curated_masked", "sales_orders_masked", "Iceberg", "Curated masked layer (physically masked)"),
("iceberg", "hadoop", "orders_ext", "Hadoop", "Lake external table (orders)"),
]
_dict_cache: dict[str, Any] = {"at": 0.0, "data": None}
_DICT_TTL = 120.0
def build_dictionary() -> dict[str, Any]:
if _dict_cache["data"] and time.time() - _dict_cache["at"] < _DICT_TTL:
return _dict_cache["data"]
import sql_console as s
# masking map from the PII catalog
masked_map: dict[tuple[str, str], dict] = {}
pii_names: dict[str, str] = {} # column-name -> category (any dataset)
try:
from pii_catalog import get_pii
pii_data = get_pii()
for d in pii_data.get("datasets", []):
tname = (d.get("table") or "").split(".")[-1]
for c in d.get("pii_columns", []):
masked_map[(tname, c["name"])] = {"masked": c.get("masked"), "category": c.get("category")}
pii_names.setdefault(c["name"], c.get("category"))
except Exception:
pass
tables_out = []
for catalog, schema, table, engine, desc in DICT_TABLES:
cols_res = _trino(f'SHOW COLUMNS FROM "{catalog}"."{schema}"."{table}"', 200)
if not cols_res.get("ok"):
continue
masked_layer = schema in ("curated_masked", "curated")
cols = []
for r in cols_res.get("rows", []):
cname, ctype = r[0], r[1]
direct = masked_map.get((table, cname))
if direct:
is_pii, is_masked, cat = True, bool(direct["masked"]), direct["category"]
elif cname in pii_names:
# PII column name found in a derived/lake table — masked only if it
# is a physically-masked layer; otherwise raw PII is visible there.
is_pii, is_masked, cat = True, masked_layer, pii_names[cname]
else:
is_pii, is_masked, cat = False, False, None
cols.append({
"name": cname, "type": ctype,
"pii": is_pii, "masked": is_masked, "category": cat,
})
tables_out.append({
"catalog": catalog, "schema": schema, "table": table, "engine": engine, "desc": desc,
"fqn": f"{catalog}.{schema}.{table}",
"columns": cols,
"pii_count": sum(1 for c in cols if c["pii"]),
"masked_count": sum(1 for c in cols if c["masked"]),
})
data = {
"ok": True,
"generated_at": datetime.now(timezone.utc).isoformat(),
"tables": tables_out,
"summary": {
"tables": len(tables_out),
"columns": sum(len(t["columns"]) for t in tables_out),
"pii_columns": sum(t["pii_count"] for t in tables_out),
"masked_columns": sum(t["masked_count"] for t in tables_out),
},
}
_dict_cache.update({"at": time.time(), "data": data})
return data
@router.get("/dictionary")
async def get_dictionary():
try:
return build_dictionary()
except Exception as exc: # noqa: BLE001
return JSONResponse({"ok": False, "error": str(exc)[:300]}, status_code=502)