| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157 |
- from __future__ import annotations
- from datetime import timedelta
- METRIC_FIELDS = (
- "job_states",
- "duration_seconds_p95",
- "failure_stage",
- "candidate_acceptance_ratio",
- "review_time_seconds_p95",
- "ontology_coverage_ratio",
- "projection_lag_seconds",
- )
- class DataResearchOperations:
- def __init__(self, repository, *, max_concurrent_jobs=10):
- if int(max_concurrent_jobs) < 1:
- raise ValueError("max_concurrent_jobs must be positive")
- self.repository = repository
- self.max_concurrent_jobs = int(max_concurrent_jobs)
- def can_start_job(self):
- return int(self.repository.running_job_count()) < self.max_concurrent_jobs
- def recover_stale_jobs(self, *, now, stale_after_seconds=900):
- before = now - timedelta(seconds=int(stale_after_seconds))
- recovered = []
- recoverable = {"extracting", "normalizing", "matching"}
- for job in self.repository.list_stale_jobs(before):
- if job.get("status") not in recoverable:
- continue
- self.repository.recover_job(
- str(job["uid"]),
- "queued",
- f"recovered stale {job['status']} job",
- )
- recovered.append(str(job["uid"]))
- return recovered
- def health(self):
- raw = dict(self.repository.aggregate_metrics() or {})
- return {
- "can_start_job": self.can_start_job(),
- "running_job_count": int(self.repository.running_job_count()),
- "max_concurrent_jobs": self.max_concurrent_jobs,
- "dead_letter_count": int(self.repository.dead_letter_count()),
- "metrics": {key: raw.get(key) for key in METRIC_FIELDS},
- }
- class DataResearchReconciler:
- def __init__(self, control_plane, graph, storage, knowledge):
- self.control_plane = control_plane
- self.graph = graph
- self.storage = storage
- self.knowledge = knowledge
- def reconcile(self, *, repair=False):
- issues = []
- active = self.graph.active_versions()
- versions = self.control_plane.published_versions()
- for version in versions:
- if active.get(version["ontology_uid"]) != version["version_uid"]:
- issues.append({"type": "projection_mismatch", **version})
- for artifact in self.control_plane.artifacts():
- if not self.storage.exists(artifact["storage_ref"]):
- issues.append({"type": "artifact_missing", **artifact})
- documents = {
- (item["object_uid"], int(item["object_version"])): item
- for item in self.control_plane.governance_documents()
- }
- for version in versions:
- key = (version["ontology_uid"], int(version.get("version", 2)))
- document = documents.get(key)
- expected = self.knowledge.expected_hash(
- version["ontology_uid"], version["version_uid"]
- )
- if document is None or document.get("content_hash") != expected:
- issues.append(
- {
- "type": "knowledge_mismatch",
- **version,
- "expected_hash": expected,
- }
- )
- repair_plan = [
- issue
- for issue in issues
- if issue["type"] in {"projection_mismatch", "knowledge_mismatch"}
- ]
- repaired = 0
- if repair:
- for issue in repair_plan:
- if issue["type"] == "projection_mismatch":
- self.graph.repair_projection(issue)
- else:
- self.knowledge.repair_document(issue)
- repaired += 1
- return {
- "repair_mode": bool(repair),
- "issues": issues,
- "repair_plan": repair_plan,
- "repaired": repaired,
- }
- class SqlAlchemyOperationsRepository:
- def __init__(self, session):
- self.session = session
- def running_job_count(self):
- from app.models.data_research import IngestionJob
- return self.session.query(IngestionJob).filter(
- IngestionJob.status.in_(("queued", "extracting", "normalizing", "matching"))
- ).count()
- def list_stale_jobs(self, before):
- from app.models.data_research import IngestionJob
- return [model.to_dict() for model in self.session.query(IngestionJob).filter(
- IngestionJob.updated_at < before,
- IngestionJob.status.in_(("extracting", "normalizing", "matching")),
- ).all()]
- def recover_job(self, uid, status, safe_error):
- from app.models.data_research import IngestionJob
- model = self.session.get(IngestionJob, str(uid))
- model.status = status
- model.last_error = safe_error
- self.session.flush()
- def aggregate_metrics(self):
- from app.models.data_research import IngestionJob
- states = {}
- for model in self.session.query(IngestionJob.status).all():
- states[model[0]] = states.get(model[0], 0) + 1
- return {
- "job_states": states,
- "duration_seconds_p95": None,
- "failure_stage": {},
- "candidate_acceptance_ratio": None,
- "review_time_seconds_p95": None,
- "ontology_coverage_ratio": None,
- "projection_lag_seconds": None,
- }
- def dead_letter_count(self):
- from sqlalchemy import text
- return self.session.execute(
- text("SELECT COUNT(*) FROM public.outbox_events WHERE status = 'dead_letter'")
- ).scalar() or 0
|