operations.py 5.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157
  1. from __future__ import annotations
  2. from datetime import timedelta
  3. METRIC_FIELDS = (
  4. "job_states",
  5. "duration_seconds_p95",
  6. "failure_stage",
  7. "candidate_acceptance_ratio",
  8. "review_time_seconds_p95",
  9. "ontology_coverage_ratio",
  10. "projection_lag_seconds",
  11. )
  12. class DataResearchOperations:
  13. def __init__(self, repository, *, max_concurrent_jobs=10):
  14. if int(max_concurrent_jobs) < 1:
  15. raise ValueError("max_concurrent_jobs must be positive")
  16. self.repository = repository
  17. self.max_concurrent_jobs = int(max_concurrent_jobs)
  18. def can_start_job(self):
  19. return int(self.repository.running_job_count()) < self.max_concurrent_jobs
  20. def recover_stale_jobs(self, *, now, stale_after_seconds=900):
  21. before = now - timedelta(seconds=int(stale_after_seconds))
  22. recovered = []
  23. recoverable = {"extracting", "normalizing", "matching"}
  24. for job in self.repository.list_stale_jobs(before):
  25. if job.get("status") not in recoverable:
  26. continue
  27. self.repository.recover_job(
  28. str(job["uid"]),
  29. "queued",
  30. f"recovered stale {job['status']} job",
  31. )
  32. recovered.append(str(job["uid"]))
  33. return recovered
  34. def health(self):
  35. raw = dict(self.repository.aggregate_metrics() or {})
  36. return {
  37. "can_start_job": self.can_start_job(),
  38. "running_job_count": int(self.repository.running_job_count()),
  39. "max_concurrent_jobs": self.max_concurrent_jobs,
  40. "dead_letter_count": int(self.repository.dead_letter_count()),
  41. "metrics": {key: raw.get(key) for key in METRIC_FIELDS},
  42. }
  43. class DataResearchReconciler:
  44. def __init__(self, control_plane, graph, storage, knowledge):
  45. self.control_plane = control_plane
  46. self.graph = graph
  47. self.storage = storage
  48. self.knowledge = knowledge
  49. def reconcile(self, *, repair=False):
  50. issues = []
  51. active = self.graph.active_versions()
  52. versions = self.control_plane.published_versions()
  53. for version in versions:
  54. if active.get(version["ontology_uid"]) != version["version_uid"]:
  55. issues.append({"type": "projection_mismatch", **version})
  56. for artifact in self.control_plane.artifacts():
  57. if not self.storage.exists(artifact["storage_ref"]):
  58. issues.append({"type": "artifact_missing", **artifact})
  59. documents = {
  60. (item["object_uid"], int(item["object_version"])): item
  61. for item in self.control_plane.governance_documents()
  62. }
  63. for version in versions:
  64. key = (version["ontology_uid"], int(version.get("version", 2)))
  65. document = documents.get(key)
  66. expected = self.knowledge.expected_hash(
  67. version["ontology_uid"], version["version_uid"]
  68. )
  69. if document is None or document.get("content_hash") != expected:
  70. issues.append(
  71. {
  72. "type": "knowledge_mismatch",
  73. **version,
  74. "expected_hash": expected,
  75. }
  76. )
  77. repair_plan = [
  78. issue
  79. for issue in issues
  80. if issue["type"] in {"projection_mismatch", "knowledge_mismatch"}
  81. ]
  82. repaired = 0
  83. if repair:
  84. for issue in repair_plan:
  85. if issue["type"] == "projection_mismatch":
  86. self.graph.repair_projection(issue)
  87. else:
  88. self.knowledge.repair_document(issue)
  89. repaired += 1
  90. return {
  91. "repair_mode": bool(repair),
  92. "issues": issues,
  93. "repair_plan": repair_plan,
  94. "repaired": repaired,
  95. }
  96. class SqlAlchemyOperationsRepository:
  97. def __init__(self, session):
  98. self.session = session
  99. def running_job_count(self):
  100. from app.models.data_research import IngestionJob
  101. return self.session.query(IngestionJob).filter(
  102. IngestionJob.status.in_(("queued", "extracting", "normalizing", "matching"))
  103. ).count()
  104. def list_stale_jobs(self, before):
  105. from app.models.data_research import IngestionJob
  106. return [model.to_dict() for model in self.session.query(IngestionJob).filter(
  107. IngestionJob.updated_at < before,
  108. IngestionJob.status.in_(("extracting", "normalizing", "matching")),
  109. ).all()]
  110. def recover_job(self, uid, status, safe_error):
  111. from app.models.data_research import IngestionJob
  112. model = self.session.get(IngestionJob, str(uid))
  113. model.status = status
  114. model.last_error = safe_error
  115. self.session.flush()
  116. def aggregate_metrics(self):
  117. from app.models.data_research import IngestionJob
  118. states = {}
  119. for model in self.session.query(IngestionJob.status).all():
  120. states[model[0]] = states.get(model[0], 0) + 1
  121. return {
  122. "job_states": states,
  123. "duration_seconds_p95": None,
  124. "failure_stage": {},
  125. "candidate_acceptance_ratio": None,
  126. "review_time_seconds_p95": None,
  127. "ontology_coverage_ratio": None,
  128. "projection_lag_seconds": None,
  129. }
  130. def dead_letter_count(self):
  131. from sqlalchemy import text
  132. return self.session.execute(
  133. text("SELECT COUNT(*) FROM public.outbox_events WHERE status = 'dead_letter'")
  134. ).scalar() or 0