from __future__ import annotations import hashlib import json import subprocess from pathlib import Path from django.core.exceptions import ValidationError from django.core.serializers.json import DjangoJSONEncoder from django.utils import timezone from agents.progeny import ProgenyService from control_plane.agents.models import ImprovementCandidate, ProgenySignal from control_plane.events.bus import EventBus from control_plane.events.models import Event from control_plane.projects.models import ( CommitRecord, Milestone, Project, ProjectPlan, RoadmapItem, StewardAction, StewardCheck, StewardEnrollment, StewardFinding, StewardPolicy, StewardRun, Task, TaskStatus, ) from graph.models import ExecutionGraphVersion SEVERITY_RANK = {"INFO": 0, "LOW": 1, "MEDIUM": 2, "HIGH": 3, "CRITICAL": 4} class StewardService: def __init__(self, bus: EventBus | None = None) -> None: self.bus = bus or EventBus() def default_policy(self) -> StewardPolicy: policy, _ = StewardPolicy.objects.get_or_create( name="Steward V1 Default", defaults={ "enabled_checks": ["TEST_HEALTH", "DEPENDENCY_DRIFT", "REPOSITORY_HEALTH", "RUNTIME_CI", "SECURITY", "SECRET_EXPIRY", "PERFORMANCE"], "auto_route_thresholds": {"REPAIR": "MEDIUM"}, "approval_requirements": {"REPAIR": "HIGH", "EXTEND": "HIGH", "EVOLVE": "HIGH"}, "run_cadence": {"manual": True}, "allowed_repair_scope": {"repository": True, "infrastructure": False}, "budget_limits": {"maximum_checks": 20}, }, ) return policy def enroll_project(self, project: Project, policy: StewardPolicy | None = None) -> StewardEnrollment: if StewardEnrollment.objects.filter(project=project, status="ACTIVE").exists(): raise ValidationError("Project already has an active Steward enrollment.") enrollment = StewardEnrollment.objects.create(project=project, policy=policy or self.default_policy(), status="ACTIVE") self.bus.publish("STEWARD_PROJECT_ENROLLED", project=project, payload={"enrollment_id": str(enrollment.id)}) return enrollment def pause_project(self, project: Project) -> StewardEnrollment: enrollment = self._enrollment(project) enrollment.status = "PAUSED" enrollment.save(update_fields=["status", "updated_at"]) return enrollment def resume_project(self, project: Project) -> StewardEnrollment: enrollment = self._enrollment(project, include_paused=True) enrollment.status = "ACTIVE" enrollment.save(update_fields=["status", "updated_at"]) return enrollment def disable_project(self, project: Project) -> StewardEnrollment: enrollment = self._enrollment(project, include_paused=True) enrollment.status = "DISABLED" enrollment.save(update_fields=["status", "updated_at"]) return enrollment def start_run(self, enrollment: StewardEnrollment, execution_graph_version: ExecutionGraphVersion | None = None) -> StewardRun: run = StewardRun.objects.create(project=enrollment.project, enrollment=enrollment, execution_graph_version=execution_graph_version, status="RUNNING", started_at=timezone.now()) self.bus.publish("STEWARD_RUN_STARTED", project=enrollment.project, payload={"steward_run_id": str(run.id)}) return run def run_checks(self, steward_run: StewardRun) -> list[StewardCheck]: checks: list[StewardCheck] = [] enabled = set(steward_run.enrollment.policy.enabled_checks or []) if "TEST_HEALTH" in enabled: checks.append(self.check_test_health(steward_run)) if "DEPENDENCY_DRIFT" in enabled: checks.append(self.check_dependency_drift(steward_run)) if "REPOSITORY_HEALTH" in enabled: checks.append(self.check_repository_health(steward_run)) if "RUNTIME_CI" in enabled: checks.append(self.check_runtime_ci(steward_run)) if "SECURITY" in enabled: checks.append(self.check_security(steward_run)) if "SECRET_EXPIRY" in enabled: checks.append(self.check_secret_expiry(steward_run)) if "PERFORMANCE" in enabled: checks.append(self.check_performance(steward_run)) return checks def check_test_health(self, steward_run: StewardRun) -> StewardCheck: project = steward_run.project command = steward_run.enrollment.policy.metadata.get("test_command", ["python", "manage.py", "test"]) if not project.repository_path: return self._check(steward_run, "TEST_HEALTH", "SKIPPED", {"reason": "missing_repository_path"}, "INFO") completed = subprocess.run([str(part) for part in command], cwd=project.repository_path, capture_output=True, text=True, check=False, timeout=120) status = "PASS" if completed.returncode == 0 else "FAIL" severity = "INFO" if status == "PASS" else "HIGH" return self._check(steward_run, "TEST_HEALTH", status, {"returncode": completed.returncode, "stdout_excerpt": completed.stdout[-8000:], "stderr_excerpt": completed.stderr[-4000:]}, severity) def check_dependency_drift(self, steward_run: StewardRun) -> StewardCheck: signals = steward_run.enrollment.policy.metadata.get("dependency_signals", []) severity = max([str(item.get("severity", "INFO")) for item in signals], key=lambda value: SEVERITY_RANK.get(value, 0), default="INFO") return self._check(steward_run, "DEPENDENCY_DRIFT", "FAIL" if signals else "PASS", {"signals": signals}, severity) def check_repository_health(self, steward_run: StewardRun) -> StewardCheck: project = steward_run.project evidence: dict[str, object] = {"dirty": False, "todo_count": 0, "migration_inconsistency": False, "broken_imports": []} severity = "INFO" if project.repository_path: status = subprocess.run(["git", "status", "--short"], cwd=project.repository_path, capture_output=True, text=True, check=False) evidence["dirty"] = bool(status.stdout.strip()) ignored = set(steward_run.enrollment.policy.ignored_paths or []) todo_count = 0 for path in Path(project.repository_path).rglob("*.py"): relative = path.relative_to(project.repository_path).as_posix() if any(relative.startswith(str(prefix)) for prefix in ignored) or ".git" in path.parts: continue text = path.read_text(encoding="utf-8", errors="ignore") todo_count += text.count("TODO") + text.count("FIXME") evidence["todo_count"] = todo_count if evidence["dirty"]: severity = "MEDIUM" return self._check(steward_run, "REPOSITORY_HEALTH", "FAIL" if evidence["dirty"] else "PASS", evidence, severity) def check_runtime_ci(self, steward_run: StewardRun) -> StewardCheck: events = list(Event.objects.filter(project=steward_run.project, event_type__in=["CI_FAILED", "RUNTIME_FAILED", "TASK_FAILED"]).order_by("-created_at")[:10].values("event_type", "payload", "created_at")) events = json.loads(json.dumps(events, cls=DjangoJSONEncoder)) return self._check(steward_run, "RUNTIME_CI", "FAIL" if events else "PASS", {"events": events}, "HIGH" if events else "INFO") def check_security(self, steward_run: StewardRun) -> StewardCheck: signals = steward_run.enrollment.policy.metadata.get("security_findings", []) severity = max([str(item.get("severity", "INFO")) for item in signals], key=lambda value: SEVERITY_RANK.get(value, 0), default="INFO") return self._check(steward_run, "SECURITY", "FAIL" if signals else "PASS", {"signals": signals}, severity) def check_secret_expiry(self, steward_run: StewardRun) -> StewardCheck: expiries = steward_run.enrollment.policy.metadata.get("secret_expiry_metadata", []) redacted = [{"name": item.get("name"), "expires_at": item.get("expires_at"), "severity": item.get("severity", "MEDIUM")} for item in expiries] severity = max([str(item.get("severity", "INFO")) for item in redacted], key=lambda value: SEVERITY_RANK.get(value, 0), default="INFO") return self._check(steward_run, "SECRET_EXPIRY", "FAIL" if redacted else "PASS", {"expiring_secrets": redacted}, severity) def check_performance(self, steward_run: StewardRun) -> StewardCheck: observations = steward_run.enrollment.policy.metadata.get("performance_observations", []) regressions = [item for item in observations if float(item.get("delta_percent", 0)) > float(steward_run.enrollment.policy.severity_thresholds.get("performance_delta_percent", 20))] return self._check(steward_run, "PERFORMANCE", "FAIL" if regressions else "PASS", {"regressions": regressions}, "MEDIUM" if regressions else "INFO") def normalize_findings(self, steward_run: StewardRun) -> list[StewardFinding]: findings: list[StewardFinding] = [] for check in steward_run.checks.all(): for raw in self._findings_for_check(check): findings.append(self.upsert_finding(steward_run, check, raw)) return findings def classify_findings(self, steward_run: StewardRun) -> list[StewardFinding]: findings = list(steward_run.findings.all()) for finding in findings: action, route = self.classify(finding) finding.recommended_action = action finding.recommended_route = route finding.save(update_fields=["recommended_action", "recommended_route", "updated_at"]) return findings def route_findings(self, steward_run: StewardRun) -> list[StewardAction]: actions: list[StewardAction] = [] for finding in steward_run.findings.exclude(status__in=["ROUTED", "RESOLVED", "DISMISSED"]): if finding.recommended_action == "IGNORE": continue if StewardAction.objects.filter(finding=finding, status__in=["PENDING", "ROUTED", "APPROVAL_REQUIRED"]).exists(): continue actions.append(self.route_finding(finding, steward_run.enrollment.policy)) return actions def route_finding(self, finding: StewardFinding, policy: StewardPolicy) -> StewardAction: requires_approval = self._requires_approval(finding, policy) if finding.recommended_action == "REPAIR": task = None if requires_approval else self._create_repair_task(finding) action = StewardAction.objects.create(finding=finding, action_type="REPAIR", status="APPROVAL_REQUIRED" if requires_approval else "ROUTED", task=task, requires_approval=requires_approval) elif finding.recommended_action == "EXTEND": item = RoadmapItem.objects.create(project=finding.project, title=finding.title, description=finding.summary, source="steward") action = StewardAction.objects.create(finding=finding, action_type="EXTEND", status="APPROVAL_REQUIRED" if requires_approval else "ROUTED", roadmap_item=item, requires_approval=requires_approval) elif finding.recommended_action == "EVOLVE": candidate = ImprovementCandidate.objects.create(target_type="PROJECT", target_id=str(finding.project_id), target_label=finding.project.name, hypothesis=finding.summary, recommended_route="EvolutionCandidate", evidence={"steward_finding_id": str(finding.id), "finding_type": finding.finding_type}) action = StewardAction.objects.create(finding=finding, action_type="EVOLVE", status="APPROVAL_REQUIRED" if requires_approval else "ROUTED", improvement_candidate=candidate, requires_approval=requires_approval) elif finding.recommended_action == "INVESTIGATE": signal = ProgenySignal.objects.create(project=finding.project, source="steward", severity=finding.severity, failure_category=finding.finding_type, summary=finding.summary, evidence={"steward_finding_id": str(finding.id), **finding.evidence}, grouping_key=f"steward:{finding.grouping_key}"[:120]) investigation = ProgenyService(self.bus).create_smart_investigation(signal.grouping_key) action = StewardAction.objects.create(finding=finding, action_type="INVESTIGATE", status="APPROVAL_REQUIRED" if requires_approval else "ROUTED", investigation=investigation, requires_approval=requires_approval) else: action = StewardAction.objects.create(finding=finding, action_type="IGNORE", status="DISMISSED") finding.status = "ROUTED" if action.status != "APPROVAL_REQUIRED" else "ACKNOWLEDGED" finding.save(update_fields=["status", "updated_at"]) self.bus.publish("STEWARD_FINDING_ROUTED", project=finding.project, task=action.task, payload={"finding_id": str(finding.id), "action_id": str(action.id), "route": finding.recommended_route}) return action def complete_run(self, steward_run: StewardRun) -> StewardRun: steward_run.status = "COMPLETE" steward_run.completed_at = timezone.now() steward_run.summary = f"{steward_run.checks.count()} checks, {steward_run.findings.count()} findings" steward_run.save(update_fields=["status", "completed_at", "summary", "updated_at"]) steward_run.enrollment.last_run_at = steward_run.completed_at steward_run.enrollment.save(update_fields=["last_run_at", "updated_at"]) self.bus.publish("STEWARD_RUN_COMPLETED", project=steward_run.project, payload={"steward_run_id": str(steward_run.id), "summary": steward_run.summary}) return steward_run def resolve_after_repair(self, finding: StewardFinding) -> bool: action = finding.actions.filter(action_type="REPAIR", task__isnull=False).order_by("-created_at").first() if action is None or action.task.status != TaskStatus.COMPLETE: return False commit = CommitRecord.objects.filter(task=action.task).first() if commit is None: return False commit.steward_finding = finding commit.save(update_fields=["steward_finding", "updated_at"]) finding.status = "RESOLVED" finding.save(update_fields=["status", "updated_at"]) action.status = "RESOLVED" action.save(update_fields=["status", "updated_at"]) self.bus.publish("STEWARD_FINDING_RESOLVED", project=finding.project, task=action.task, payload={"finding_id": str(finding.id), "commit_id": str(commit.id)}) return True def dismiss_finding(self, finding: StewardFinding) -> StewardFinding: finding.status = "DISMISSED" finding.save(update_fields=["status", "updated_at"]) return finding def inspection(self, project: Project) -> dict[str, object]: enrollment = project.steward_enrollments.order_by("-created_at").first() return { "project_id": str(project.id), "status": enrollment.status if enrollment else "NOT_ENROLLED", "last_run": enrollment.last_run_at if enrollment else None, "next_run": enrollment.next_run_at if enrollment else None, "open_findings": json.loads(json.dumps(list(project.steward_findings.exclude(status__in=["RESOLVED", "DISMISSED"]).values("id", "finding_type", "severity", "recommended_action", "recommended_route", "status", "occurrence_count")), cls=DjangoJSONEncoder)), "resolved_findings": json.loads(json.dumps(list(project.steward_findings.filter(status="RESOLVED").values("id", "finding_type", "severity", "status")), cls=DjangoJSONEncoder)), "recent_checks": list(StewardCheck.objects.filter(steward_run__project=project).order_by("-created_at").values("check_type", "status", "severity")[:20]), } def _enrollment(self, project: Project, *, include_paused: bool = False) -> StewardEnrollment: statuses = ["ACTIVE", "PAUSED"] if include_paused else ["ACTIVE"] enrollment = StewardEnrollment.objects.filter(project=project, status__in=statuses).order_by("-created_at").first() if enrollment is None: raise ValidationError("Project has no Steward enrollment.") return enrollment def _check(self, steward_run: StewardRun, check_type: str, status: str, evidence: dict[str, object], severity: str) -> StewardCheck: check = StewardCheck.objects.create(steward_run=steward_run, check_type=check_type, status=status, evidence=evidence, severity=severity) self.bus.publish("STEWARD_CHECK_COMPLETED", project=steward_run.project, payload={"steward_run_id": str(steward_run.id), "check_id": str(check.id), "check_type": check_type, "status": status}) return check def _findings_for_check(self, check: StewardCheck) -> list[dict[str, object]]: if check.status in ["PASS", "SKIPPED"]: if check.check_type == "REPOSITORY_HEALTH" and int(check.evidence.get("todo_count", 0)) > 0: return [{"finding_type": "INFORMATIONAL_DRIFT", "title": "TODO/FIXME debt observed", "summary": "Repository contains TODO/FIXME markers.", "severity": "LOW", "evidence": check.evidence, "component": "repository"}] return [] if check.check_type == "TEST_HEALTH": return [{"finding_type": "TEST_REGRESSION", "title": "Deterministic tests are failing", "summary": "Project test suite failed under Steward test health check.", "severity": check.severity, "evidence": check.evidence, "component": "tests"}] if check.check_type == "DEPENDENCY_DRIFT": return [{"finding_type": "DEPENDENCY_VULNERABILITY", "title": "Dependency drift or vulnerability detected", "summary": "Dependency metadata indicates repairable drift or vulnerability.", "severity": check.severity, "evidence": check.evidence, "component": "dependencies"}] if check.check_type == "RUNTIME_CI": return [{"finding_type": "RUNTIME_CI_FAILURE", "title": "Runtime or CI failure signal observed", "summary": "Structured runtime/CI events indicate a maintenance issue.", "severity": check.severity, "evidence": check.evidence, "component": "runtime"}] if check.check_type == "SECURITY": return [{"finding_type": "SECURITY_SIGNAL", "title": "Security signal observed", "summary": "Security metadata indicates a repairable issue.", "severity": check.severity, "evidence": check.evidence, "component": "security"}] if check.check_type == "SECRET_EXPIRY": return [{"finding_type": "SECRET_EXPIRY", "title": "Secret/certificate expiry metadata observed", "summary": "A secret or certificate is approaching expiry. Secret values were not inspected.", "severity": check.severity, "evidence": check.evidence, "component": "secrets"}] if check.check_type == "PERFORMANCE": return [{"finding_type": "PERFORMANCE_REGRESSION", "title": "Performance regression signal observed", "summary": "Performance observations indicate a regression from baseline.", "severity": check.severity, "evidence": check.evidence, "component": "performance"}] if check.check_type == "REPOSITORY_HEALTH": return [{"finding_type": "REPOSITORY_DIRTY", "title": "Repository has uncommitted changes", "summary": "Repository health check found dirty state.", "severity": check.severity, "evidence": check.evidence, "component": "repository"}] return [] def upsert_finding(self, steward_run: StewardRun, check: StewardCheck, raw: dict[str, object]) -> StewardFinding: now = timezone.now() grouping_key = self._grouping_key(steward_run.project, raw) finding = StewardFinding.objects.filter(project=steward_run.project, grouping_key=grouping_key).first() if finding is None: finding = StewardFinding.objects.create(project=steward_run.project, steward_run=steward_run, source_check=check, finding_type=str(raw["finding_type"]), title=str(raw["title"]), summary=str(raw["summary"]), evidence=dict(raw.get("evidence", {})), severity=str(raw.get("severity", "INFO")), confidence=0.8, grouping_key=grouping_key, first_seen=now, last_seen=now, metadata={"component": raw.get("component", "")}) self.bus.publish("STEWARD_FINDING_CREATED", project=steward_run.project, payload={"finding_id": str(finding.id), "finding_type": finding.finding_type}) return finding finding.occurrence_count += 1 finding.last_seen = now if SEVERITY_RANK.get(str(raw.get("severity", "INFO")), 0) > SEVERITY_RANK.get(finding.severity, 0): finding.severity = str(raw["severity"]) finding.evidence = dict(raw.get("evidence", finding.evidence)) finding.steward_run = steward_run finding.source_check = check finding.save(update_fields=["occurrence_count", "last_seen", "severity", "evidence", "steward_run", "source_check", "updated_at"]) self.bus.publish("STEWARD_FINDING_UPDATED", project=steward_run.project, payload={"finding_id": str(finding.id), "occurrence_count": finding.occurrence_count}) return finding def classify(self, finding: StewardFinding) -> tuple[str, str]: if finding.finding_type in ["TEST_REGRESSION", "DEPENDENCY_VULNERABILITY", "SECURITY_SIGNAL", "SECRET_EXPIRY", "REPOSITORY_DIRTY"]: return "REPAIR", "Project DAG Repair" if finding.finding_type == "PERFORMANCE_REGRESSION": return ("REPAIR", "Project DAG Repair") if finding.severity in ["HIGH", "CRITICAL"] else ("EVOLVE", "EvolutionCandidate") if finding.finding_type == "PRODUCT_OPPORTUNITY": return "EXTEND", "RoadmapItem" if finding.finding_type == "RUNTIME_CI_FAILURE" and finding.occurrence_count > 1: return "INVESTIGATE", "Progeny Smart Investigation" if finding.finding_type == "INFORMATIONAL_DRIFT": return "IGNORE", "No execution" return "INVESTIGATE", "Progeny Smart Investigation" def _create_repair_task(self, finding: StewardFinding) -> Task: milestone = finding.project.milestones.order_by("created_at").first() if milestone is None: plan = ProjectPlan.objects.create(project=finding.project, version=finding.project.current_plan_version + 1 or 1, goal=finding.project.goal) milestone = Milestone.objects.create(project=finding.project, plan=plan, key="STEWARD", title="Steward Repairs", goal="Maintain completed system") return Task.objects.create(project=finding.project, milestone=milestone, task_type="repair", status=TaskStatus.READY, goal=f"Repair Steward finding: {finding.title}\n\nEvidence: {finding.summary}", acceptance_criteria=["Restore intended existing behavior", "Deterministic tests pass", "Add or preserve regression coverage where practical"], max_retries=2) def _requires_approval(self, finding: StewardFinding, policy: StewardPolicy) -> bool: threshold = str(policy.approval_requirements.get(finding.recommended_action, "CRITICAL")) return SEVERITY_RANK.get(finding.severity, 0) >= SEVERITY_RANK.get(threshold, 4) def _grouping_key(self, project: Project, raw: dict[str, object]) -> str: component = str(raw.get("component", "project"))[:80] evidence = str(raw.get("evidence", {}))[:1000] fingerprint = hashlib.sha256(evidence.encode("utf-8")).hexdigest()[:16] return f"{project.id}:{raw.get('finding_type')}:{component}:{fingerprint}"[:240]