diff --git a/agents/judge.py b/agents/judge.py index 6180d45..b94ec26 100644 --- a/agents/judge.py +++ b/agents/judge.py @@ -10,7 +10,8 @@ class Judge: evidence: list[dict[str, object]] = [{"type": "test_status", "status": test_status}] passed = test_status == "PASS" goal = task.goal.lower() - if "health" in goal: + expects_health_endpoint = "/health" in goal or "health endpoint" in goal or "health route" in goal + if expects_health_endpoint: has_route = "/health" in diff or "path('health'" in diff or 'path("health"' in diff passed = passed and has_route and '"status": "ok"' in diff evidence.append({"type": "acceptance_check", "requirement": "health endpoint returns ok", "passed": passed}) diff --git a/agents/reviewer.py b/agents/reviewer.py index 0128a8b..d6455cb 100644 --- a/agents/reviewer.py +++ b/agents/reviewer.py @@ -24,7 +24,9 @@ class Reviewer: if not diff.strip(): status = "REJECTED" findings.append({"type": "empty_diff", "severity": "high", "message": "No implementation diff exists"}) - if "health" in task.goal.lower() and "/health" not in diff and "path('health'" not in diff and 'path("health"' not in diff: + goal = task.goal.lower() + expects_health_endpoint = "/health" in goal or "health endpoint" in goal or "health route" in goal + if expects_health_endpoint and "/health" not in diff and "path('health'" not in diff and 'path("health"' not in diff: status = "REWORK_REQUIRED" findings.append({"type": "missing_health_route", "severity": "high", "message": "Diff does not add /health"}) review = Review.objects.create( diff --git a/agents/steward.py b/agents/steward.py new file mode 100644 index 0000000..0c08a96 --- /dev/null +++ b/agents/steward.py @@ -0,0 +1,331 @@ +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] diff --git a/control_plane/projects/migrations/0004_steward_v1.py b/control_plane/projects/migrations/0004_steward_v1.py new file mode 100644 index 0000000..c35ef9f --- /dev/null +++ b/control_plane/projects/migrations/0004_steward_v1.py @@ -0,0 +1,126 @@ +import uuid + +import django.db.models.deletion +from django.db import migrations, models + + +class Migration(migrations.Migration): + dependencies = [ + ("agents", "0007_replay_arena"), + ("graph", "0004_unique_champion_graph_version"), + ("projects", "0003_commitrecord_graph_run"), + ] + + operations = [ + migrations.CreateModel( + name="StewardPolicy", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("name", models.CharField(max_length=200)), + ("enabled_checks", models.JSONField(blank=True, default=list)), + ("severity_thresholds", models.JSONField(blank=True, default=dict)), + ("auto_route_thresholds", models.JSONField(blank=True, default=dict)), + ("run_cadence", models.JSONField(blank=True, default=dict)), + ("allowed_repair_scope", models.JSONField(blank=True, default=dict)), + ("approval_requirements", models.JSONField(blank=True, default=dict)), + ("ignored_paths", models.JSONField(blank=True, default=list)), + ("budget_limits", models.JSONField(blank=True, default=dict)), + ("metadata", models.JSONField(blank=True, default=dict)), + ], + ), + migrations.CreateModel( + name="StewardEnrollment", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("status", models.CharField(default="ACTIVE", max_length=32)), + ("enrolled_at", models.DateTimeField(auto_now_add=True)), + ("last_run_at", models.DateTimeField(blank=True, null=True)), + ("next_run_at", models.DateTimeField(blank=True, null=True)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("policy", models.ForeignKey(on_delete=django.db.models.deletion.PROTECT, related_name="enrollments", to="projects.stewardpolicy")), + ("project", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="steward_enrollments", to="projects.project")), + ], + ), + migrations.CreateModel( + name="StewardRun", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("status", models.CharField(default="PENDING", max_length=32)), + ("started_at", models.DateTimeField(blank=True, null=True)), + ("completed_at", models.DateTimeField(blank=True, null=True)), + ("summary", models.TextField(blank=True)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("enrollment", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="runs", to="projects.stewardenrollment")), + ("execution_graph_version", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="steward_runs", to="graph.executiongraphversion")), + ("project", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="steward_runs", to="projects.project")), + ], + ), + migrations.CreateModel( + name="StewardCheck", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("check_type", models.CharField(max_length=80)), + ("status", models.CharField(max_length=32)), + ("evidence", models.JSONField(blank=True, default=dict)), + ("severity", models.CharField(default="INFO", max_length=32)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("steward_run", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="checks", to="projects.stewardrun")), + ], + ), + migrations.CreateModel( + name="StewardFinding", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("finding_type", models.CharField(max_length=80)), + ("title", models.CharField(max_length=255)), + ("summary", models.TextField(blank=True)), + ("evidence", models.JSONField(blank=True, default=dict)), + ("severity", models.CharField(default="INFO", max_length=32)), + ("confidence", models.FloatField(default=0.0)), + ("status", models.CharField(default="OPEN", max_length=32)), + ("recommended_action", models.CharField(default="INVESTIGATE", max_length=32)), + ("recommended_route", models.CharField(blank=True, max_length=120)), + ("grouping_key", models.CharField(max_length=240)), + ("first_seen", models.DateTimeField(blank=True, null=True)), + ("last_seen", models.DateTimeField(blank=True, null=True)), + ("occurrence_count", models.PositiveIntegerField(default=1)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("project", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="steward_findings", to="projects.project")), + ("source_check", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="findings", to="projects.stewardcheck")), + ("steward_run", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="findings", to="projects.stewardrun")), + ], + ), + migrations.AddField( + model_name="commitrecord", + name="steward_finding", + field=models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="commits", to="projects.stewardfinding"), + ), + migrations.CreateModel( + name="StewardAction", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("action_type", models.CharField(max_length=32)), + ("status", models.CharField(default="PENDING", max_length=32)), + ("requires_approval", models.BooleanField(default=False)), + ("approved_at", models.DateTimeField(blank=True, null=True)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("finding", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="actions", to="projects.stewardfinding")), + ("improvement_candidate", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="steward_actions", to="agents.improvementcandidate")), + ("investigation", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="steward_actions", to="agents.progenyinvestigation")), + ("roadmap_item", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="steward_actions", to="projects.roadmapitem")), + ("task", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="steward_actions", to="projects.task")), + ], + ), + ] diff --git a/control_plane/projects/models.py b/control_plane/projects/models.py index 3deb747..c47abd7 100644 --- a/control_plane/projects/models.py +++ b/control_plane/projects/models.py @@ -176,6 +176,9 @@ class CommitRecord(TimestampedModel): graph_run = models.ForeignKey( "graph.GraphRun", on_delete=models.SET_NULL, null=True, blank=True, related_name="commits" ) + steward_finding = models.ForeignKey( + "projects.StewardFinding", on_delete=models.SET_NULL, null=True, blank=True, related_name="commits" + ) sha = models.CharField(max_length=64) branch_name = models.CharField(max_length=255) message = models.TextField() @@ -252,3 +255,79 @@ class Scenario(TimestampedModel): target_id = models.CharField(max_length=120) definition = models.JSONField(default=dict, blank=True) status = models.CharField(max_length=32, default="DRAFT") + + +class StewardPolicy(TimestampedModel): + name = models.CharField(max_length=200) + enabled_checks = models.JSONField(default=list, blank=True) + severity_thresholds = models.JSONField(default=dict, blank=True) + auto_route_thresholds = models.JSONField(default=dict, blank=True) + run_cadence = models.JSONField(default=dict, blank=True) + allowed_repair_scope = models.JSONField(default=dict, blank=True) + approval_requirements = models.JSONField(default=dict, blank=True) + ignored_paths = models.JSONField(default=list, blank=True) + budget_limits = models.JSONField(default=dict, blank=True) + metadata = models.JSONField(default=dict, blank=True) + + +class StewardEnrollment(TimestampedModel): + project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="steward_enrollments") + status = models.CharField(max_length=32, default="ACTIVE") + policy = models.ForeignKey(StewardPolicy, on_delete=models.PROTECT, related_name="enrollments") + enrolled_at = models.DateTimeField(auto_now_add=True) + last_run_at = models.DateTimeField(null=True, blank=True) + next_run_at = models.DateTimeField(null=True, blank=True) + metadata = models.JSONField(default=dict, blank=True) + + +class StewardRun(TimestampedModel): + project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="steward_runs") + enrollment = models.ForeignKey(StewardEnrollment, on_delete=models.CASCADE, related_name="runs") + execution_graph_version = models.ForeignKey("graph.ExecutionGraphVersion", on_delete=models.SET_NULL, null=True, blank=True, related_name="steward_runs") + status = models.CharField(max_length=32, default="PENDING") + started_at = models.DateTimeField(null=True, blank=True) + completed_at = models.DateTimeField(null=True, blank=True) + summary = models.TextField(blank=True) + metadata = models.JSONField(default=dict, blank=True) + + +class StewardCheck(TimestampedModel): + steward_run = models.ForeignKey(StewardRun, on_delete=models.CASCADE, related_name="checks") + check_type = models.CharField(max_length=80) + status = models.CharField(max_length=32) + evidence = models.JSONField(default=dict, blank=True) + severity = models.CharField(max_length=32, default="INFO") + metadata = models.JSONField(default=dict, blank=True) + + +class StewardFinding(TimestampedModel): + project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="steward_findings") + steward_run = models.ForeignKey(StewardRun, on_delete=models.SET_NULL, null=True, blank=True, related_name="findings") + source_check = models.ForeignKey(StewardCheck, on_delete=models.SET_NULL, null=True, blank=True, related_name="findings") + finding_type = models.CharField(max_length=80) + title = models.CharField(max_length=255) + summary = models.TextField(blank=True) + evidence = models.JSONField(default=dict, blank=True) + severity = models.CharField(max_length=32, default="INFO") + confidence = models.FloatField(default=0.0) + status = models.CharField(max_length=32, default="OPEN") + recommended_action = models.CharField(max_length=32, default="INVESTIGATE") + recommended_route = models.CharField(max_length=120, blank=True) + grouping_key = models.CharField(max_length=240) + first_seen = models.DateTimeField(null=True, blank=True) + last_seen = models.DateTimeField(null=True, blank=True) + occurrence_count = models.PositiveIntegerField(default=1) + metadata = models.JSONField(default=dict, blank=True) + + +class StewardAction(TimestampedModel): + finding = models.ForeignKey(StewardFinding, on_delete=models.CASCADE, related_name="actions") + action_type = models.CharField(max_length=32) + status = models.CharField(max_length=32, default="PENDING") + task = models.ForeignKey(Task, on_delete=models.SET_NULL, null=True, blank=True, related_name="steward_actions") + roadmap_item = models.ForeignKey(RoadmapItem, on_delete=models.SET_NULL, null=True, blank=True, related_name="steward_actions") + investigation = models.ForeignKey("agents.ProgenyInvestigation", on_delete=models.SET_NULL, null=True, blank=True, related_name="steward_actions") + improvement_candidate = models.ForeignKey("agents.ImprovementCandidate", on_delete=models.SET_NULL, null=True, blank=True, related_name="steward_actions") + requires_approval = models.BooleanField(default=False) + approved_at = models.DateTimeField(null=True, blank=True) + metadata = models.JSONField(default=dict, blank=True) diff --git a/graph/bootstrap.py b/graph/bootstrap.py index dfa7b2e..313c21f 100644 --- a/graph/bootstrap.py +++ b/graph/bootstrap.py @@ -3,6 +3,7 @@ from __future__ import annotations from django.utils import timezone from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus +from graph.steward import steward_run_graph_v1 from graph.task_execution import task_execution_graph_v1 @@ -27,3 +28,26 @@ def champion_task_execution_graph_v1() -> ExecutionGraphVersion: version.promoted_at = timezone.now() version.save(update_fields=["status", "promoted_at"]) return version + + +def champion_steward_run_graph_v1() -> ExecutionGraphVersion: + spec = steward_run_graph_v1() + definition, _ = ExecutionGraphDefinition.objects.get_or_create( + name=spec.name, + defaults={"graph_type": spec.graph_type, "description": str(spec.metadata.get("description", ""))}, + ) + version, created = ExecutionGraphVersion.objects.get_or_create( + graph=definition, + version=spec.version, + defaults={ + "status": ExecutionGraphVersionStatus.CHAMPION, + "graph_spec": spec.to_dict(), + "metadata": {"immutable_after_use": True}, + "promoted_at": timezone.now(), + }, + ) + if not created and version.status != ExecutionGraphVersionStatus.CHAMPION: + version.status = ExecutionGraphVersionStatus.CHAMPION + version.promoted_at = timezone.now() + version.save(update_fields=["status", "promoted_at"]) + return version diff --git a/graph/native_runtime.py b/graph/native_runtime.py index 3d1773b..41938a8 100644 --- a/graph/native_runtime.py +++ b/graph/native_runtime.py @@ -106,7 +106,13 @@ class NativeGraphRuntime(GraphRuntime): metadata["interrupted_edge_result"] = result.edge_result metadata["interrupted_output"] = self._bounded(result.output_metadata or {}) graph_run.metadata = metadata - graph_run.save(update_fields=["metadata", "updated_at"]) + update_fields = ["metadata", "updated_at"] + if result.status == "PAUSED": + graph_run.status = GraphRunStatus.PAUSED + graph_run.failure_reason = result.pause_reason + update_fields.extend(["status", "failure_reason"]) + self.bus.publish("GRAPH_RUN_PAUSED", project=graph_run.project, task=graph_run.task, payload={"graph_run_id": graph_run.id, "reason": result.pause_reason}) + graph_run.save(update_fields=update_fields) graph_run.refresh_from_db() return graph_run if result.status == "PAUSED": diff --git a/graph/steward.py b/graph/steward.py new file mode 100644 index 0000000..41aeb50 --- /dev/null +++ b/graph/steward.py @@ -0,0 +1,150 @@ +from __future__ import annotations + +from django.utils import timezone + +from agents.steward import StewardService +from graph.models import GraphApproval, GraphApprovalStatus +from graph.native_runtime import GraphExecutionContext +from graph.registry import NodeHandlerRegistry, NodeResult +from graph.spec import ExecutionGraphSpec, GraphEdgeSpec, GraphNodeSpec + + +def steward_run_graph_v1() -> ExecutionGraphSpec: + spec = ExecutionGraphSpec( + name="steward_run", + version=1, + graph_type="STEWARD_RUN", + entry="prepare", + nodes={ + "prepare": GraphNodeSpec("prepare", "steward_prepare"), + "collect_signals": GraphNodeSpec("collect_signals", "steward_collect_signals"), + "run_checks": GraphNodeSpec("run_checks", "steward_run_checks"), + "normalize_findings": GraphNodeSpec("normalize_findings", "steward_normalize_findings"), + "classify": GraphNodeSpec("classify", "steward_classify"), + "route": GraphNodeSpec("route", "steward_route"), + "await_approval": GraphNodeSpec("await_approval", "steward_await_approval"), + "summarize": GraphNodeSpec("summarize", "steward_summarize"), + "complete": GraphNodeSpec("complete", "complete", {"terminal": True}), + }, + edges=[ + GraphEdgeSpec("prepare", "collect_signals", "success"), + GraphEdgeSpec("collect_signals", "run_checks", "success"), + GraphEdgeSpec("run_checks", "normalize_findings", "success"), + GraphEdgeSpec("normalize_findings", "classify", "success"), + GraphEdgeSpec("classify", "route", "success"), + GraphEdgeSpec("route", "await_approval", "approval_required"), + GraphEdgeSpec("route", "summarize", "success"), + GraphEdgeSpec("await_approval", "summarize", "approved"), + GraphEdgeSpec("await_approval", "summarize", "rejected"), + GraphEdgeSpec("summarize", "complete", "success"), + ], + terminal_nodes=["complete"], + metadata={"description": "Steward V1 governance workflow; detection, classification, routing, no code mutation."}, + ) + spec.validate() + return spec + + +class StewardNode: + idempotent = True + replay_safe = True + destructive = False + + def __init__(self, service: StewardService, node_type: str) -> None: + self.service = service + self.node_type = node_type + + def steward_run_id(self, context: GraphExecutionContext) -> str: + return str(context.graph_run.metadata["steward_run_id"]) + + +class StewardPrepareNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + return NodeResult("COMPLETE", "success", {"steward_run_id": self.steward_run_id(context)}) + + +class StewardCollectSignalsNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + return NodeResult("COMPLETE", "success") + + +class StewardRunChecksNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + from control_plane.projects.models import StewardRun + + steward_run = StewardRun.objects.get(id=self.steward_run_id(context)) + checks = self.service.run_checks(steward_run) + return NodeResult("COMPLETE", "success", {"check_count": len(checks)}) + + +class StewardNormalizeFindingsNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + from control_plane.projects.models import StewardRun + + steward_run = StewardRun.objects.get(id=self.steward_run_id(context)) + findings = self.service.normalize_findings(steward_run) + return NodeResult("COMPLETE", "success", {"finding_count": len(findings)}) + + +class StewardClassifyNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + from control_plane.projects.models import StewardRun + + steward_run = StewardRun.objects.get(id=self.steward_run_id(context)) + findings = self.service.classify_findings(steward_run) + return NodeResult("COMPLETE", "success", {"classified_count": len(findings)}) + + +class StewardRouteNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + from control_plane.projects.models import StewardRun + + steward_run = StewardRun.objects.get(id=self.steward_run_id(context)) + actions = self.service.route_findings(steward_run) + approval_required = any(action.requires_approval for action in actions) + return NodeResult("COMPLETE", "approval_required" if approval_required else "success", {"action_count": len(actions), "approval_required": approval_required}) + + +class StewardApprovalNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + node_run = context.graph_run.node_runs.filter(node_id=context.graph_run.current_node).order_by("-visit_index").first() + if GraphApproval.objects.filter(graph_run=context.graph_run, status=GraphApprovalStatus.APPROVED).exists(): + from control_plane.projects.models import StewardRun + + steward_run = StewardRun.objects.get(id=self.steward_run_id(context)) + for action in steward_run.findings.filter(actions__requires_approval=True).values_list("actions__id", flat=True).distinct(): + from control_plane.projects.models import StewardAction + + steward_action = StewardAction.objects.get(id=action) + if steward_action.task_id is None and steward_action.action_type == "REPAIR": + steward_action.task = self.service._create_repair_task(steward_action.finding) + steward_action.status = "ROUTED" + steward_action.approved_at = timezone.now() + steward_action.save(update_fields=["task", "status", "approved_at", "updated_at"]) + return NodeResult("COMPLETE", "approved") + if GraphApproval.objects.filter(graph_run=context.graph_run, status=GraphApprovalStatus.REJECTED).exists(): + return NodeResult("COMPLETE", "rejected") + GraphApproval.objects.get_or_create(graph_run=context.graph_run, node_run=node_run, reason="AWAITING_STEWARD_ROUTING_APPROVAL") + return NodeResult("PAUSED", "awaiting", pause_reason="AWAITING_STEWARD_ROUTING_APPROVAL") + + +class StewardSummarizeNode(StewardNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + from control_plane.projects.models import StewardRun + + steward_run = StewardRun.objects.get(id=self.steward_run_id(context)) + self.service.complete_run(steward_run) + return NodeResult("COMPLETE", "success", {"summary": steward_run.summary}) + + +def steward_registry(service: StewardService) -> NodeHandlerRegistry: + registry = NodeHandlerRegistry() + registry.register(StewardPrepareNode(service, "steward_prepare")) + registry.register(StewardCollectSignalsNode(service, "steward_collect_signals")) + registry.register(StewardRunChecksNode(service, "steward_run_checks")) + registry.register(StewardNormalizeFindingsNode(service, "steward_normalize_findings")) + registry.register(StewardClassifyNode(service, "steward_classify")) + registry.register(StewardRouteNode(service, "steward_route")) + registry.register(StewardApprovalNode(service, "steward_await_approval")) + registry.register(StewardSummarizeNode(service, "steward_summarize")) + return registry diff --git a/tests/test_steward_v1.py b/tests/test_steward_v1.py new file mode 100644 index 0000000..878a308 --- /dev/null +++ b/tests/test_steward_v1.py @@ -0,0 +1,153 @@ +from __future__ import annotations + +from pathlib import Path + +from agents.steward import StewardService +from control_plane.projects.models import CommitRecord, Milestone, Project, ProjectPlan, StewardAction, StewardFinding, StewardPolicy, TaskStatus +from graph.bootstrap import champion_steward_run_graph_v1 +from graph.langgraph_runtime import LangGraphRuntime +from graph.models import GraphApprovalStatus, GraphRun, GraphRunStatus +from graph.steward import steward_registry +from tests.test_m2_autonomous_loop import create_disposable_django_repo + + +def enrolled_project(tmp_path: Path, *, metadata: dict[str, object] | None = None, approval_requirements: dict[str, str] | None = None) -> tuple[Project, StewardService]: + repo = create_disposable_django_repo(tmp_path) + project = Project.objects.create(name="Stewarded", goal="Keep system healthy", repository_path=str(repo)) + policy = StewardPolicy.objects.create( + name="Steward Test Policy", + enabled_checks=["DEPENDENCY_DRIFT", "RUNTIME_CI", "SECURITY", "SECRET_EXPIRY", "PERFORMANCE"], + metadata=metadata or {}, + severity_thresholds={"performance_delta_percent": 20}, + approval_requirements=approval_requirements or {"REPAIR": "CRITICAL"}, + ) + service = StewardService() + service.enroll_project(project, policy) + return project, service + + +def run_steward_graph(project: Project, service: StewardService): + enrollment = project.steward_enrollments.get(status="ACTIVE") + version = champion_steward_run_graph_v1() + steward_run = service.start_run(enrollment, version) + graph_run = GraphRun.objects.create(execution_graph_version=version, project=project, current_node=version.graph_spec["entry"], metadata={"steward_run_id": str(steward_run.id)}) + runtime = LangGraphRuntime(steward_registry(service)) + runtime.run_until_terminal_or_paused(graph_run) + steward_run.refresh_from_db() + graph_run.refresh_from_db() + return steward_run, graph_run + + +def test_steward_enrollment_checks_and_repair_routing(tmp_path: Path) -> None: + project, service = enrolled_project( + tmp_path, + metadata={"security_findings": [{"id": "CVE-1", "severity": "MEDIUM"}], "secret_expiry_metadata": [{"name": "api_cert", "expires_at": "2026-09-01", "severity": "MEDIUM"}]}, + ) + + steward_run, graph_run = run_steward_graph(project, service) + + assert graph_run.status == GraphRunStatus.COMPLETE + assert steward_run.status == "COMPLETE" + assert steward_run.checks.count() == 5 + findings = list(steward_run.findings.order_by("finding_type")) + assert [finding.finding_type for finding in findings] == ["SECRET_EXPIRY", "SECURITY_SIGNAL"] + actions = list(StewardAction.objects.select_related("task", "finding")) + assert len(actions) == 2 + assert all(action.action_type == "REPAIR" for action in actions) + assert all(action.task and action.task.status == TaskStatus.READY for action in actions) + assert "api_cert" in actions[0].finding.evidence.get("expiring_secrets", [{}])[0].get("name", "") or "api_cert" in actions[1].finding.evidence.get("expiring_secrets", [{}])[0].get("name", "") + + +def test_steward_deduplicates_findings_across_runs(tmp_path: Path) -> None: + project, service = enrolled_project(tmp_path, metadata={"security_findings": [{"id": "CVE-1", "severity": "MEDIUM"}]}) + + first_run, _ = run_steward_graph(project, service) + second_run, _ = run_steward_graph(project, service) + + assert StewardFinding.objects.count() == 1 + finding = StewardFinding.objects.get() + assert finding.occurrence_count == 2 + assert finding.steward_run == second_run + assert first_run.findings.count() == 0 + + +def test_high_severity_repair_pauses_until_approval(tmp_path: Path) -> None: + project, service = enrolled_project(tmp_path, metadata={"security_findings": [{"id": "CVE-2", "severity": "HIGH"}]}, approval_requirements={"REPAIR": "HIGH"}) + + steward_run, graph_run = run_steward_graph(project, service) + + assert graph_run.status == GraphRunStatus.PAUSED + assert graph_run.approvals.get().status == GraphApprovalStatus.PENDING + action = StewardAction.objects.get() + assert action.status == "APPROVAL_REQUIRED" + assert action.task is None + + approval = graph_run.approvals.get() + approval.status = GraphApprovalStatus.APPROVED + approval.decided_by = "tester" + approval.save(update_fields=["status", "decided_by", "updated_at"]) + graph_run.status = GraphRunStatus.RUNNING + graph_run.save(update_fields=["status", "updated_at"]) + LangGraphRuntime(steward_registry(service)).run_until_terminal_or_paused(graph_run) + + graph_run.refresh_from_db() + steward_run.refresh_from_db() + action.refresh_from_db() + assert graph_run.status == GraphRunStatus.COMPLETE + assert steward_run.status == "COMPLETE" + assert action.status == "ROUTED" + assert action.task is not None + + +def test_steward_routes_systemic_runtime_signal_to_progeny_investigation(tmp_path: Path) -> None: + project, service = enrolled_project(tmp_path) + enrollment = project.steward_enrollments.get() + steward_run = service.start_run(enrollment) + check = service._check(steward_run, "RUNTIME_CI", "FAIL", {"events": [{"event_type": "CI_FAILED"}]}, "HIGH") + raw = service._findings_for_check(check)[0] + finding = service.upsert_finding(steward_run, check, raw) + finding.occurrence_count = 2 + finding.save(update_fields=["occurrence_count"]) + service.classify_findings(steward_run) + + action = service.route_findings(steward_run)[0] + + assert action.action_type == "INVESTIGATE" + assert action.investigation is not None + assert action.task is None + + +def test_repair_resolution_links_commit_to_steward_finding(tmp_path: Path) -> None: + project, service = enrolled_project(tmp_path, metadata={"security_findings": [{"id": "CVE-1", "severity": "MEDIUM"}]}) + run_steward_graph(project, service) + action = StewardAction.objects.select_related("task", "finding").get() + task = action.task + assert task is not None + task.status = TaskStatus.COMPLETE + task.save(update_fields=["status", "updated_at"]) + CommitRecord.objects.create(project=project, task=task, sha="abc123", branch_name="repair", message="Repair Steward finding") + + assert service.resolve_after_repair(action.finding) is True + + action.finding.refresh_from_db() + assert action.finding.status == "RESOLVED" + assert CommitRecord.objects.get(task=task).steward_finding == action.finding + + +def test_resolved_identical_findings_do_not_create_duplicate_repairs(tmp_path: Path) -> None: + project, service = enrolled_project(tmp_path, metadata={"security_findings": [{"id": "CVE-1", "severity": "MEDIUM"}]}) + run_steward_graph(project, service) + action = StewardAction.objects.select_related("task", "finding").get() + task = action.task + assert task is not None + task.status = TaskStatus.COMPLETE + task.save(update_fields=["status", "updated_at"]) + CommitRecord.objects.create(project=project, task=task, sha="abc123", branch_name="repair", message="Repair Steward finding") + service.resolve_after_repair(action.finding) + + run_steward_graph(project, service) + + assert StewardFinding.objects.count() == 1 + assert StewardAction.objects.count() == 1 + assert project.tasks.filter(task_type="repair").count() == 1 + assert StewardFinding.objects.get().status == "RESOLVED"