Files
atc-agents/api/catalog_governance.py
mo f36c8906bc feat(governance): close the 6 data-disease gaps — DQ monitoring, ownership, access posture, lineage & observability
Adds native Command Center features (no new containers) integrated as sub-tabs
in the existing Data Explorer and Data Quality views:

- Continuous Data Quality (dq_monitor.py): live completeness/uniqueness/validity/
  freshness scorecards via Trino with rolling trends → DataQuality "Live Monitoring".
- Ownership & stewardship (catalog_governance.py): owner/steward/tier matrix,
  orphan detection, business glossary; local store best-effort synced to
  OpenMetadata (owner PATCH) → Data Explorer "Ownership".
- Access & policy posture: per-dataset compliance combining PII masking, ownership,
  live DQ and observability alerts vs data contracts → Data Explorer "Access & Policies".
- Lineage (lineage.py): staged source→CDC→Spark→S3→Iceberg→Trino→serving graph with
  live row counts and column-level PII/masking tracing → Data Explorer "Lineage".
- Observability (observability.py): volume/freshness/schema-drift monitoring with
  alerts → Data Explorer "Observability".
- Shared lake_meta.py dataset registry + bounded Trino client; fast native row-count
  and PK-indexed freshness so monitors stay cheap on 25-54M-row tables.
- LLM context (lab_context.py) enriched with DQ scores, ownership and active alerts.
2026-06-29 17:43:37 +00:00

315 lines
13 KiB
Python

"""Data governance: ownership, stewardship, glossary & posture
(Diseases #3 Ownership and #5 Governance).
Ownership/steward assignments are kept in a local store (so the feature always
works for the demo regardless of OpenMetadata ingestion state) and are
best-effort mirrored to OpenMetadata. Users/teams for the assignment dropdowns
are read from OpenMetadata when reachable. The governance *posture* view
combines, per dataset: ownership, PII/masking (pii_catalog), live data-quality
score (dq_monitor), observability alerts and a simple data-contract check.
Endpoints:
GET /api/governance/datasets -> ownership matrix (+ orphan flag)
GET /api/governance/users -> assignable users/teams
POST /api/governance/assign -> set owner/steward/team/tier (native + OM)
GET /api/governance/glossary -> business glossary terms
GET /api/governance/posture -> combined governance posture per dataset
"""
from __future__ import annotations
import json
import os
import time
from pathlib import Path
from typing import Any
import httpx
from fastapi import APIRouter
from fastapi.responses import JSONResponse
from pydantic import BaseModel
from lake_meta import DATASETS, DATASET_BY_KEY
router = APIRouter(prefix="/api/governance", tags=["governance"])
OPENMETADATA_URL = os.getenv("OPENMETADATA_URL", "").rstrip("/")
OPENMETADATA_TOKEN = os.getenv("OPENMETADATA_TOKEN", "")
OWNERS_PATH = Path(os.getenv("GOVERNANCE_OWNERS_PATH", "/data/governance_owners.json"))
CONTRACTS_PATH = Path(os.getenv("GOVERNANCE_CONTRACTS_PATH", "/data/governance_contracts.json"))
# Default data contracts (quality SLAs) per dataset — used when none stored.
DEFAULT_CONTRACTS = {
"_default": {"min_score": 80, "min_completeness": 95, "min_freshness_min": 60,
"max_critical_alerts": 0},
}
# Seed business glossary so the term list is never empty even before OM ingestion.
SEED_GLOSSARY = [
{"name": "Customer", "description": "A person or organization that places sales orders.",
"related": ["customer_id", "customer_name", "customer_email"], "domain": "Sales"},
{"name": "Order", "description": "A sales transaction with an amount, channel and status.",
"related": ["order_id", "amount", "order_status"], "domain": "Sales"},
{"name": "Revenue", "description": "Sum of order amounts over a period.",
"related": ["amount", "currency"], "domain": "Sales"},
{"name": "Employee", "description": "A member of the workforce tracked via HR events.",
"related": ["employee_id", "department", "role_name"], "domain": "People"},
{"name": "PII", "description": "Personally Identifiable Information — masked per policy.",
"related": ["customer_email", "national_id", "billing_iban"], "domain": "Governance"},
{"name": "Telemetry", "description": "Device metric readings over time.",
"related": ["device_id", "metric_type", "metric_value"], "domain": "IoT"},
]
def _headers() -> dict[str, str]:
h = {"Accept": "application/json"}
if OPENMETADATA_TOKEN:
h["Authorization"] = f"Bearer {OPENMETADATA_TOKEN}"
return h
def _load(path: Path) -> dict[str, Any]:
try:
return json.loads(path.read_text())
except Exception:
return {}
def _save(path: Path, data: dict[str, Any]) -> None:
try:
path.parent.mkdir(parents=True, exist_ok=True)
path.write_text(json.dumps(data, indent=2))
except Exception:
pass
# ── OpenMetadata best-effort ────────────────────────────────────────────────
def _om_users() -> list[dict[str, Any]]:
if not OPENMETADATA_URL:
return []
out: list[dict[str, Any]] = []
try:
with httpx.Client(timeout=8.0, verify=False) as c:
r = c.get(f"{OPENMETADATA_URL}/api/v1/users?limit=50&isBot=false", headers=_headers())
if r.status_code == 200:
for u in r.json().get("data", []):
out.append({"id": u.get("id"), "name": u.get("name"),
"display": u.get("displayName") or u.get("name"), "type": "user"})
rt = c.get(f"{OPENMETADATA_URL}/api/v1/teams?limit=50", headers=_headers())
if rt.status_code == 200:
for t in rt.json().get("data", []):
out.append({"id": t.get("id"), "name": t.get("name"),
"display": t.get("displayName") or t.get("name"), "type": "team"})
except Exception:
pass
return out
def _om_glossary() -> list[dict[str, Any]]:
if not OPENMETADATA_URL:
return []
try:
with httpx.Client(timeout=8.0, verify=False) as c:
r = c.get(f"{OPENMETADATA_URL}/api/v1/glossaryTerms?limit=100", headers=_headers())
if r.status_code == 200:
return [{"name": t.get("name"), "description": t.get("description", ""),
"domain": "OpenMetadata", "related": []}
for t in r.json().get("data", [])]
except Exception:
pass
return []
def _om_patch_owner(om_fqn: str, user: dict[str, Any]) -> bool:
"""Best-effort: set table owner in OpenMetadata via JSON-Patch."""
if not OPENMETADATA_URL or not om_fqn or not user.get("id"):
return False
try:
with httpx.Client(timeout=8.0, verify=False) as c:
g = c.get(f"{OPENMETADATA_URL}/api/v1/tables/name/{om_fqn}", headers=_headers())
if g.status_code != 200:
return False
patch = [{"op": "add", "path": "/owners/0",
"value": {"id": user["id"], "type": user.get("type", "user")}}]
h = {**_headers(), "Content-Type": "application/json-patch+json"}
p = c.patch(f"{OPENMETADATA_URL}/api/v1/tables/name/{om_fqn}",
headers=h, content=json.dumps(patch))
return p.status_code in (200, 201)
except Exception:
return False
# ── helpers pulling from sibling modules (all best-effort) ──────────────────
def _pii_summary(pii_key: str) -> dict[str, Any]:
try:
from pii_catalog import get_pii
for d in get_pii().get("datasets", []):
if d.get("key") == pii_key:
return {"pii_count": d.get("pii_count", 0),
"all_masked": d.get("all_masked", False),
"masked": sum(1 for c in d.get("pii_columns", []) if c.get("masked")),
"unmasked": sum(1 for c in d.get("pii_columns", []) if not c.get("masked"))}
except Exception:
pass
return {"pii_count": 0, "all_masked": False, "masked": 0, "unmasked": 0}
def _dq_card(key: str) -> dict[str, Any]:
try:
from dq_monitor import scorecards_view
for c in scorecards_view().get("cards", []):
if c.get("key") == key:
return {"score": c.get("score"), "issues": c.get("issues", []),
"freshness_age_min": c.get("freshness_age_min")}
except Exception:
pass
return {"score": None, "issues": []}
def _obs_alerts(key: str) -> list[dict[str, Any]]:
try:
from observability import _state, _lock
with _lock:
return [{"type": a["type"], "severity": a["severity"], "message": a["message"]}
for a in _state["alerts_active"].values() if a["dataset"] == key]
except Exception:
return []
# ── views ───────────────────────────────────────────────────────────────────
def datasets_view() -> dict[str, Any]:
owners = _load(OWNERS_PATH)
rows = []
orphans = 0
stewarded = 0
for ds in DATASETS:
rec = owners.get(ds["key"], {})
owner = rec.get("owner")
steward = rec.get("steward")
if not owner:
orphans += 1
if steward:
stewarded += 1
rows.append({
"key": ds["key"], "label": ds["label"], "engine": ds["engine"], "color": ds["color"],
"table": ds["fqtn"], "domain": rec.get("domain") or ds.get("domain"),
"owner": owner, "steward": steward, "team": rec.get("team"),
"tier": rec.get("tier"), "classification": rec.get("classification"),
"updated_at": rec.get("updated_at"),
"orphan": not owner,
"pii": _pii_summary(ds.get("pii_key", ds["key"])),
})
return {"ok": True, "datasets": rows,
"summary": {"total": len(rows), "orphans": orphans, "stewarded": stewarded,
"owned": len(rows) - orphans},
"om_connected": bool(OPENMETADATA_URL)}
class AssignRequest(BaseModel):
key: str
owner: str | None = None
steward: str | None = None
team: str | None = None
tier: str | None = None
classification: str | None = None
def posture_view() -> dict[str, Any]:
owners = _load(OWNERS_PATH)
contracts = _load(CONTRACTS_PATH)
default_c = DEFAULT_CONTRACTS["_default"]
out = []
compliant = 0
for ds in DATASETS:
key = ds["key"]
rec = owners.get(key, {})
contract = {**default_c, **(contracts.get(key, {}))}
pii = _pii_summary(ds.get("pii_key", key))
dq = _dq_card(key)
alerts = _obs_alerts(key)
crit = sum(1 for a in alerts if a["severity"] == "critical")
checks = []
score = dq.get("score")
checks.append({"name": "DQ score", "ok": score is not None and score >= contract["min_score"],
"value": score, "target": contract["min_score"]})
checks.append({"name": "Owner assigned", "ok": bool(rec.get("owner")),
"value": rec.get("owner") or "—", "target": "assigned"})
checks.append({"name": "PII masked", "ok": pii["unmasked"] == 0,
"value": f'{pii["masked"]}/{pii["pii_count"]}', "target": "all"})
checks.append({"name": "Critical alerts", "ok": crit <= contract["max_critical_alerts"],
"value": crit, "target": contract["max_critical_alerts"]})
ok = all(c["ok"] for c in checks)
if ok:
compliant += 1
out.append({
"key": key, "label": ds["label"], "engine": ds["engine"], "color": ds["color"],
"table": ds["fqtn"], "owner": rec.get("owner"), "steward": rec.get("steward"),
"tier": rec.get("tier"), "pii": pii, "dq_score": score, "issues": dq.get("issues", []),
"alerts": alerts, "contract": contract, "checks": checks, "compliant": ok,
})
return {"ok": True, "datasets": out,
"summary": {"total": len(out), "compliant": compliant,
"non_compliant": len(out) - compliant}}
def summary_for_llm() -> dict[str, Any]:
owners = _load(OWNERS_PATH)
return {
"owners": {k: {"owner": v.get("owner"), "steward": v.get("steward"), "tier": v.get("tier")}
for k, v in owners.items()},
"orphan_datasets": [d["key"] for d in DATASETS if not owners.get(d["key"], {}).get("owner")],
}
@router.get("/datasets")
async def get_datasets() -> JSONResponse:
return JSONResponse(datasets_view())
@router.get("/users")
async def get_users() -> JSONResponse:
users = _om_users()
if not users:
users = [{"id": None, "name": n, "display": n, "type": "user"}
for n in ("admin", "bart", "mo")] + \
[{"id": None, "name": "Organization", "display": "Organization", "type": "team"}]
return JSONResponse({"ok": True, "users": users, "om_connected": bool(OPENMETADATA_URL)})
@router.post("/assign")
async def assign(body: AssignRequest) -> JSONResponse:
if body.key not in DATASET_BY_KEY:
return JSONResponse({"ok": False, "error": f"unknown dataset {body.key}"}, status_code=400)
owners = _load(OWNERS_PATH)
rec = dict(owners.get(body.key, {}))
for field in ("owner", "steward", "team", "tier", "classification"):
val = getattr(body, field)
if val is not None:
rec[field] = val or None
rec["updated_at"] = time.strftime("%Y-%m-%dT%H:%M:%SZ", time.gmtime())
owners[body.key] = rec
_save(OWNERS_PATH, owners)
om_synced = False
if body.owner:
ds = DATASET_BY_KEY[body.key]
user = next((u for u in _om_users() if u.get("display") == body.owner or u.get("name") == body.owner), None)
if user:
om_synced = _om_patch_owner(ds.get("om_fqn", ""), user)
return JSONResponse({"ok": True, "key": body.key, "record": rec, "om_synced": om_synced})
@router.get("/glossary")
async def get_glossary() -> JSONResponse:
terms = _om_glossary()
source = "openmetadata"
if not terms:
terms = SEED_GLOSSARY
source = "seed"
return JSONResponse({"ok": True, "source": source, "count": len(terms), "terms": terms})
@router.get("/posture")
async def get_posture() -> JSONResponse:
return JSONResponse(posture_view())