active_metadata.py 4.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154
  1. """HTTP API for active metadata discovery and field lineage."""
  2. from __future__ import annotations
  3. from flask import g, jsonify, request
  4. from app import db
  5. from app.api.meta_data import bp
  6. from app.core.meta_data.active_metadata import (
  7. ActiveMetadataConflict,
  8. ActiveMetadataError,
  9. ActiveMetadataNotFound,
  10. ActiveMetadataService,
  11. )
  12. from app.core.meta_data.active_metadata_repository import (
  13. SqlAlchemyActiveMetadataRepository,
  14. )
  15. from app.models.result import failed, success
  16. def _service():
  17. return ActiveMetadataService(SqlAlchemyActiveMetadataRepository(db.session))
  18. def _error(exc):
  19. db.session.rollback()
  20. if isinstance(exc, ActiveMetadataNotFound):
  21. return jsonify(failed(str(exc), code=404)), 404
  22. if isinstance(exc, ActiveMetadataConflict):
  23. return jsonify(failed(str(exc), code=409)), 409
  24. if isinstance(exc, ActiveMetadataError):
  25. return jsonify(failed(str(exc), code=400)), 400
  26. raise exc
  27. @bp.route("/active-metadata/plans", methods=["GET"])
  28. def list_active_metadata_plans():
  29. return jsonify(success(_service().list_plans()))
  30. @bp.route("/active-metadata/plans", methods=["POST"])
  31. def create_active_metadata_plan():
  32. try:
  33. result = _service().create_plan(
  34. request.get_json(silent=True),
  35. actor_uid=g.current_user["id"],
  36. )
  37. db.session.commit()
  38. return jsonify(success(result, "主动发现计划已创建")), 201
  39. except Exception as exc:
  40. return _error(exc)
  41. @bp.route("/active-metadata/plans/<plan_uid>/runs", methods=["POST"])
  42. def execute_active_metadata_plan(plan_uid):
  43. batch_key = str(request.headers.get("Idempotency-Key") or "").strip()
  44. if not batch_key:
  45. return jsonify(failed("Idempotency-Key is required", code=428)), 428
  46. try:
  47. result = _service().execute(
  48. plan_uid,
  49. request.get_json(silent=True),
  50. batch_key=batch_key,
  51. actor_uid=g.current_user["id"],
  52. )
  53. db.session.commit()
  54. return jsonify(success(result, "主动发现批次已完成"))
  55. except Exception as exc:
  56. return _error(exc)
  57. @bp.route("/active-metadata/plans/<plan_uid>/runs", methods=["GET"])
  58. def list_active_metadata_runs(plan_uid):
  59. try:
  60. return jsonify(success(_service().list_runs(plan_uid)))
  61. except Exception as exc:
  62. return _error(exc)
  63. @bp.route("/active-metadata/assets", methods=["GET"])
  64. def list_active_metadata_assets():
  65. try:
  66. return jsonify(
  67. success(_service().list_assets(request.args.get("source_uid")))
  68. )
  69. except Exception as exc:
  70. return _error(exc)
  71. @bp.route("/active-metadata/runs/<run_uid>/changes", methods=["GET"])
  72. def list_active_metadata_changes(run_uid):
  73. try:
  74. return jsonify(success(_service().list_changes(run_uid)))
  75. except Exception as exc:
  76. return _error(exc)
  77. @bp.route("/active-metadata/runs/<run_uid>/lineage", methods=["GET"])
  78. def list_active_metadata_lineage(run_uid):
  79. try:
  80. return jsonify(success(_service().list_lineage(run_uid)))
  81. except Exception as exc:
  82. return _error(exc)
  83. @bp.route("/active-metadata/runs/<run_uid>/health-signals", methods=["GET"])
  84. def list_active_metadata_health(run_uid):
  85. try:
  86. return jsonify(success(_service().list_health_signals(run_uid)))
  87. except Exception as exc:
  88. return _error(exc)
  89. @bp.route("/active-metadata/corrections", methods=["GET"])
  90. def list_active_metadata_corrections():
  91. try:
  92. return jsonify(
  93. success(_service().list_corrections(request.args.get("asset_uid")))
  94. )
  95. except Exception as exc:
  96. return _error(exc)
  97. @bp.route("/active-metadata/assets/<asset_uid>/corrections", methods=["POST"])
  98. def submit_active_metadata_correction(asset_uid):
  99. try:
  100. result = _service().submit_correction(
  101. asset_uid,
  102. request.get_json(silent=True),
  103. actor_uid=g.current_user["id"],
  104. )
  105. db.session.commit()
  106. return jsonify(success(result, "元数据纠错已提交")), 201
  107. except Exception as exc:
  108. return _error(exc)
  109. @bp.route(
  110. "/active-metadata/corrections/<correction_uid>/resolve",
  111. methods=["POST"],
  112. )
  113. def resolve_active_metadata_correction(correction_uid):
  114. body = request.get_json(silent=True) or {}
  115. try:
  116. result = _service().resolve_correction(
  117. correction_uid,
  118. expected_version=body.get("expected_version"),
  119. decision=body.get("decision"),
  120. resolution=body.get("resolution") or {},
  121. actor_uid=g.current_user["id"],
  122. )
  123. db.session.commit()
  124. return jsonify(success(result, "元数据纠错已处置"))
  125. except Exception as exc:
  126. return _error(exc)