f36c8906bc
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.
315 lines
13 KiB
Python
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())
|