reconcile_governance_uids.py 2.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778
  1. from __future__ import annotations
  2. import argparse
  3. import json
  4. from typing import Any
  5. from app.core.common.identifiers import new_governance_uid
  6. GOVERNANCE_LABELS = ("BusinessDomain", "DataMeta", "DataFlow", "DataSource")
  7. def ensure_uid_constraints(session: Any) -> None:
  8. constraints = {
  9. "BusinessDomain": "business_domain_uid",
  10. "DataMeta": "data_meta_uid",
  11. "DataFlow": "data_flow_uid",
  12. "DataSource": "datasource_uid",
  13. }
  14. for label, constraint in constraints.items():
  15. session.run(
  16. f"CREATE CONSTRAINT {constraint} IF NOT EXISTS "
  17. f"FOR (n:{label}) REQUIRE n.uid IS UNIQUE"
  18. )
  19. def report_missing_uids(session: Any, limit: int = 100) -> dict[str, list[dict[str, Any]]]:
  20. report: dict[str, list[dict[str, Any]]] = {}
  21. for label in GOVERNANCE_LABELS:
  22. result = session.run(
  23. f"MATCH (n:{label}) WHERE n.uid IS NULL "
  24. "RETURN elementId(n) AS legacy_id, n.name_zh AS name_zh "
  25. "ORDER BY legacy_id LIMIT $limit",
  26. {"limit": int(limit)},
  27. )
  28. report[label] = [dict(record) for record in result]
  29. return report
  30. def backfill_missing_uids(session: Any, limit: int = 100) -> dict[str, int]:
  31. report = report_missing_uids(session, limit=limit)
  32. counts: dict[str, int] = {}
  33. for label, records in report.items():
  34. counts[label] = 0
  35. for record in records:
  36. session.run(
  37. f"MATCH (n:{label}) WHERE elementId(n) = $legacy_id AND n.uid IS NULL "
  38. "SET n.uid = $uid",
  39. {
  40. "legacy_id": str(record["legacy_id"]),
  41. "uid": new_governance_uid(),
  42. },
  43. )
  44. counts[label] += 1
  45. return counts
  46. def main() -> None:
  47. parser = argparse.ArgumentParser(description="Audit stable governance UIDs")
  48. parser.add_argument("--limit", type=int, default=100)
  49. parser.add_argument("--backfill", action="store_true")
  50. parser.add_argument("--ensure-constraints", action="store_true")
  51. args = parser.parse_args()
  52. from app.services.neo4j_driver import neo4j_driver
  53. with neo4j_driver.get_session() as session:
  54. report = report_missing_uids(session, limit=args.limit)
  55. print(json.dumps(report, ensure_ascii=False, indent=2))
  56. if args.backfill:
  57. counts = backfill_missing_uids(session, limit=args.limit)
  58. print(json.dumps({"backfilled": counts}, ensure_ascii=False))
  59. if args.ensure_constraints:
  60. ensure_uid_constraints(session)
  61. if __name__ == "__main__":
  62. main()