diff --git a/api/Dockerfile b/api/Dockerfile
index 3ac5772..89a259d 100644
--- a/api/Dockerfile
+++ b/api/Dockerfile
@@ -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
diff --git a/api/main.py b/api/main.py
index e5e7ebd..642a62a 100644
--- a/api/main.py
+++ b/api/main.py
@@ -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=["*"],
diff --git a/api/platform_context.py b/api/platform_context.py
index a32ad19..e34ae4f 100644
--- a/api/platform_context.py
+++ b/api/platform_context.py
@@ -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])
diff --git a/api/trino_federated.py b/api/trino_federated.py
new file mode 100644
index 0000000..c467be4
--- /dev/null
+++ b/api/trino_federated.py
@@ -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)
diff --git a/ui/src/App.tsx b/ui/src/App.tsx
index e105178..9f0022b 100644
--- a/ui/src/App.tsx
+++ b/ui/src/App.tsx
@@ -19,6 +19,7 @@ import { ChangesView } from './components/features/ChangesView'
import { DataFlowView } from './components/features/DataFlowView'
import { SearchView } from './components/features/SearchView'
import { DataExplorerView } from './components/features/DataExplorerView'
+import { TrinoFederationView } from './components/features/TrinoFederationView'
import { DataHubView } from './components/features/DataHubView'
import { SshTerminal } from './components/features/SshTerminal'
import { TerminalDock } from './components/features/TerminalDock'
@@ -144,6 +145,8 @@ export default function App() {
No data
+ return ( +No data
+ return ( +{label}
+{value}
+ {sub &&{sub}
} ++ One SQL engine over every source — federated business analytics, mirrored into Hadoop as external Iceberg tables +
+{marquee?.sql || m?.sql}
+ {marquee?.running && !m?.rows?.length ? (
+ | Region | +Orders | +Revenue | +HR events | +Supply events | +
|---|---|---|---|---|
| {r.region} | +{fmtNum(r.orders)} | +{fmtMoney(r.revenue)} | +{fmtNum(r.hr_events)} | +{fmtNum(r.supply_events)} | +
as of {new Date(m.generated_at).toLocaleString()} · joined live across PostgreSQL + MySQL + MongoDB
} +{m?.error || 'No result yet — click re-run.'}
+ )} ++ All federated business data materialized as external Iceberg tables on HDFS — queried live (fast). +
+ +{fmtNum(t.rows)}
+No external tables yet — click “Rebuild external tables”.
} ++ This is exactly what the assistant knows about your data — every column, its type, and whether it is masked or visible. +
+{t.engine} · {t.desc}
+