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