diff --git a/agents/progeny.py b/agents/progeny.py index e5b98c0..1e68b1f 100644 --- a/agents/progeny.py +++ b/agents/progeny.py @@ -1,11 +1,14 @@ from __future__ import annotations from dataclasses import dataclass +from datetime import datetime +from typing import Any -from control_plane.agents.models import Agent, AgentPlan, AgentVersion, BenchmarkRun, ProgenySignal, PromotionStatus +from control_plane.agents.models import Agent, AgentPlan, AgentVersion, BenchmarkRun, ImprovementCandidate, ProgenyInvestigation, ProgenySignal, PromotionStatus from control_plane.events.bus import EventBus from control_plane.events.models import EventType from control_plane.projects.models import Project, Task, Milestone +from graph.models import ExecutionGraphVersion, GraphNodeRun, GraphRun @dataclass(frozen=True) @@ -14,6 +17,20 @@ class BenchmarkDecision: metrics: dict[str, float] +@dataclass(frozen=True) +class ProgenySignalGroup: + grouping_key: str + occurrence_count: int + affected_projects: list[int] + affected_agents: list[int] + affected_graph_versions: list[int] + affected_graph_nodes: list[str] + first_seen: datetime + last_seen: datetime + severity: str + failure_category: str + + class ProgenyService: def __init__(self, bus: EventBus | None = None) -> None: self.bus = bus or EventBus() @@ -74,80 +91,329 @@ class ProgenyService: challenger.save(update_fields=["promotion_status", "updated_at"]) return BenchmarkDecision("REJECTED", challenger_metrics) - def create_reviewer_signal(self, project: Project, task: Task, milestone: Milestone, agent_version: AgentVersion, status: str, findings: list[dict[str, object]], summary: str) -> ProgenySignal: + def create_reviewer_signal( + self, + project: Project, + task: Task, + milestone: Milestone, + agent_version: AgentVersion, + status: str, + findings: list[dict[str, object]], + summary: str, + *, + graph_run: GraphRun | None = None, + graph_node_run: GraphNodeRun | None = None, + execution_graph_version: ExecutionGraphVersion | None = None, + metadata: dict[str, object] | None = None, + ) -> ProgenySignal: severity = "high" if status in ["REWORK_REQUIRED", "REJECTED"] else "info" + lineage = self._lineage(graph_run, graph_node_run, execution_graph_version) signal = ProgenySignal.objects.create( project=project, task=task, milestone=milestone, agent_version=agent_version, + graph_run=graph_run, + graph_node_run=graph_node_run, + execution_graph_version=lineage["execution_graph_version"], source="reviewer", severity=severity, failure_category=status, summary=summary, - evidence={"findings": findings}, + evidence={"findings": findings, **lineage["evidence"], **(metadata or {})}, status="OPEN", - grouping_key=f"reviewer:{task.id}:{status}", + grouping_key=self._grouping_key("reviewer", status, agent_version, lineage), model=agent_version.model, ) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "reviewer", "status": status}) return signal - def create_judge_signal(self, project: Project, task: Task, milestone: Milestone, agent_version: AgentVersion, result: str, evidence: list[dict[str, object]], summary: str) -> ProgenySignal: + def create_judge_signal( + self, + project: Project, + task: Task, + milestone: Milestone, + agent_version: AgentVersion, + result: str, + evidence: list[dict[str, object]], + summary: str, + *, + graph_run: GraphRun | None = None, + graph_node_run: GraphNodeRun | None = None, + execution_graph_version: ExecutionGraphVersion | None = None, + metadata: dict[str, object] | None = None, + ) -> ProgenySignal: severity = "high" if result == "FAIL" else "info" + lineage = self._lineage(graph_run, graph_node_run, execution_graph_version) signal = ProgenySignal.objects.create( project=project, task=task, milestone=milestone, agent_version=agent_version, + graph_run=graph_run, + graph_node_run=graph_node_run, + execution_graph_version=lineage["execution_graph_version"], source="judge", severity=severity, failure_category=result, summary=summary, - evidence={"evidence": evidence}, + evidence={"evidence": evidence, **lineage["evidence"], **(metadata or {})}, status="OPEN", - grouping_key=f"judge:{task.id}:{result}", + grouping_key=self._grouping_key("judge", result, agent_version, lineage), model=agent_version.model, ) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "judge", "result": result}) return signal - def create_model_output_signal(self, project: Project, task: Task, milestone: Milestone, agent_version: AgentVersion, error: str, raw_output: str) -> ProgenySignal: + def create_model_output_signal( + self, + project: Project, + task: Task, + milestone: Milestone, + agent_version: AgentVersion, + error: str, + raw_output: str, + *, + graph_run: GraphRun | None = None, + graph_node_run: GraphNodeRun | None = None, + execution_graph_version: ExecutionGraphVersion | None = None, + ) -> ProgenySignal: + lineage = self._lineage(graph_run, graph_node_run, execution_graph_version) signal = ProgenySignal.objects.create( project=project, task=task, milestone=milestone, agent_version=agent_version, + graph_run=graph_run, + graph_node_run=graph_node_run, + execution_graph_version=lineage["execution_graph_version"], source="model_output", severity="high", - failure_category="MALFORMED_OUTPUT", + failure_category="MODEL_OUTPUT_INVALID", summary=f"Model output malformed: {error}", - evidence={"raw_output": raw_output[:1000]}, + evidence={"raw_output": raw_output[:1000], "error": error, **lineage["evidence"]}, status="OPEN", - grouping_key=f"model_output:{task.id}", + grouping_key=self._grouping_key("model_output", "MODEL_OUTPUT_INVALID", agent_version, lineage), model=agent_version.model, ) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "model_output"}) return signal - def create_retry_exhausted_signal(self, project: Project, task: Task, milestone: Milestone, agent_version: AgentVersion, attempts: int) -> ProgenySignal: + def create_provider_signal( + self, + project: Project | None, + task: Task | None, + milestone: Milestone | None, + agent_version: AgentVersion | None, + category: str, + summary: str, + evidence: dict[str, object], + *, + severity: str = "high", + graph_run: GraphRun | None = None, + graph_node_run: GraphNodeRun | None = None, + execution_graph_version: ExecutionGraphVersion | None = None, + model: str = "", + ) -> ProgenySignal: + lineage = self._lineage(graph_run, graph_node_run, execution_graph_version) signal = ProgenySignal.objects.create( project=project, task=task, milestone=milestone, agent_version=agent_version, + graph_run=graph_run, + graph_node_run=graph_node_run, + execution_graph_version=lineage["execution_graph_version"], + source="provider", + severity=severity, + failure_category=category, + summary=summary, + evidence={**evidence, **lineage["evidence"]}, + status="OPEN", + grouping_key=self._grouping_key("provider", category, agent_version, lineage), + model=model or (agent_version.model if agent_version else ""), + ) + self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "provider", "category": category}) + return signal + + def create_retry_exhausted_signal( + self, + project: Project, + task: Task, + milestone: Milestone, + agent_version: AgentVersion, + attempts: int, + *, + graph_run: GraphRun | None = None, + graph_node_run: GraphNodeRun | None = None, + execution_graph_version: ExecutionGraphVersion | None = None, + evidence: dict[str, object] | None = None, + ) -> ProgenySignal: + lineage = self._lineage(graph_run, graph_node_run, execution_graph_version) + signal = ProgenySignal.objects.create( + project=project, + task=task, + milestone=milestone, + agent_version=agent_version, + graph_run=graph_run, + graph_node_run=graph_node_run, + execution_graph_version=lineage["execution_graph_version"], source="retry", severity="critical", - failure_category="RETRY_EXHAUSTED", + failure_category="TASK_RETRY_EXHAUSTED", summary=f"Task {task.id} exhausted {attempts} retries", - evidence={"attempts": attempts}, + evidence={"attempts": attempts, **(evidence or {}), **lineage["evidence"]}, status="OPEN", - grouping_key=f"retry:{task.id}", + grouping_key=self._grouping_key("retry", "TASK_RETRY_EXHAUSTED", agent_version, lineage), model=agent_version.model, ) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "retry"}) return signal + def create_graph_runtime_signal( + self, + graph_run: GraphRun, + category: str, + summary: str, + evidence: dict[str, object], + *, + graph_node_run: GraphNodeRun | None = None, + severity: str = "high", + ) -> ProgenySignal: + lineage = self._lineage(graph_run, graph_node_run, graph_run.execution_graph_version) + signal = ProgenySignal.objects.create( + project=graph_run.project, + task=graph_run.task, + milestone=graph_run.milestone, + graph_run=graph_run, + graph_node_run=graph_node_run, + execution_graph_version=graph_run.execution_graph_version, + source="graph_runtime", + severity=severity, + failure_category=category, + summary=summary, + evidence={**evidence, **lineage["evidence"]}, + status="OPEN", + grouping_key=self._grouping_key("graph_runtime", category, None, lineage), + ) + self.bus.publish("PROGENY_SIGNAL_CREATED", project=graph_run.project, task=graph_run.task, actor="progeny", payload={"signal_id": str(signal.id), "source": "graph_runtime", "category": category}) + return signal + + def query_inbox(self, **filters: object): + signals = ProgenySignal.objects.select_related("project", "task", "agent_version", "execution_graph_version", "graph_node_run").all() + status = filters.get("status", "OPEN") + if status: + signals = signals.filter(status=status) + if filters.get("source"): + signals = signals.filter(source=filters["source"]) + if filters.get("project"): + signals = signals.filter(project=filters["project"]) + if filters.get("agent"): + signals = signals.filter(agent_version__agent=filters["agent"]) + if filters.get("agent_version"): + signals = signals.filter(agent_version=filters["agent_version"]) + if filters.get("model"): + signals = signals.filter(model=filters["model"]) + if filters.get("execution_graph"): + signals = signals.filter(execution_graph_version__graph=filters["execution_graph"]) + if filters.get("execution_graph_version"): + signals = signals.filter(execution_graph_version=filters["execution_graph_version"]) + if filters.get("graph_node"): + signals = signals.filter(graph_node_run__node_id=filters["graph_node"]) + if filters.get("severity"): + signals = signals.filter(severity=filters["severity"]) + if filters.get("failure_category"): + signals = signals.filter(failure_category=filters["failure_category"]) + if filters.get("start"): + signals = signals.filter(created_at__gte=filters["start"]) + if filters.get("end"): + signals = signals.filter(created_at__lte=filters["end"]) + return signals.order_by("-created_at") + + def group_unresolved_signals(self, **filters: object) -> list[ProgenySignalGroup]: + signals = list(self.query_inbox(**filters)) + buckets: dict[str, list[ProgenySignal]] = {} + for signal in signals: + buckets.setdefault(signal.grouping_key or self._fallback_grouping_key(signal), []).append(signal) + groups: list[ProgenySignalGroup] = [] + severity_rank = {"info": 0, "INFO": 0, "low": 1, "medium": 2, "high": 3, "critical": 4} + for key, bucket in buckets.items(): + ordered = sorted(bucket, key=lambda signal: signal.created_at) + groups.append( + ProgenySignalGroup( + grouping_key=key, + occurrence_count=len(bucket), + affected_projects=sorted({str(signal.project_id) for signal in bucket if signal.project_id}), + affected_agents=sorted({str(signal.agent_version_id) for signal in bucket if signal.agent_version_id}), + affected_graph_versions=sorted({str(signal.execution_graph_version_id) for signal in bucket if signal.execution_graph_version_id}), + affected_graph_nodes=sorted({signal.graph_node_run.node_id for signal in bucket if signal.graph_node_run_id}), + first_seen=ordered[0].created_at, + last_seen=ordered[-1].created_at, + severity=max((signal.severity for signal in bucket), key=lambda value: severity_rank.get(value, 0)), + failure_category=ordered[-1].failure_category, + ) + ) + return sorted(groups, key=lambda group: (group.severity == "critical", group.occurrence_count, group.last_seen), reverse=True) + + def create_smart_investigation(self, grouping_key: str) -> ProgenyInvestigation: + signals = list(self.query_inbox().filter(grouping_key=grouping_key).order_by("created_at")) + if not signals: + raise ValueError(f"No open signals for grouping key {grouping_key}") + analysis = self._analyze_signals(signals) + investigation = ProgenyInvestigation.objects.create( + signal_clusters=[{"grouping_key": grouping_key, "signal_ids": [str(signal.id) for signal in signals], "occurrence_count": len(signals)}], + affected_projects=sorted({str(signal.project_id) for signal in signals if signal.project_id}), + affected_agents=sorted({str(signal.agent_version_id) for signal in signals if signal.agent_version_id}), + affected_graph_versions=sorted({str(signal.execution_graph_version_id) for signal in signals if signal.execution_graph_version_id}), + affected_nodes=sorted({signal.graph_node_run.node_id for signal in signals if signal.graph_node_run_id}), + hypotheses=analysis["hypotheses"], + recommended_target=str(analysis["target"]), + confidence=float(analysis["confidence"]), + recommended_route=str(analysis["route"]), + proposed_experiments=analysis["experiments"], + expected_impact=str(analysis["impact"]), + estimated_cost=str(analysis["cost"]), + ) + investigation.signals.set(signals) + return investigation + + def create_improvement_candidate(self, investigation: ProgenyInvestigation, hypothesis: str | None = None) -> ImprovementCandidate: + signal = investigation.signals.select_related("execution_graph_version", "agent_version").first() + target_type = investigation.recommended_target + execution_graph_version = signal.execution_graph_version if signal and target_type in {"WORKFLOW_GRAPH", "GRAPH_NODE"} else None + agent_version = signal.agent_version if signal and target_type in {"AGENT", "REVIEWER", "JUDGE"} else None + target_id = "" + target_label = target_type + if execution_graph_version is not None: + target_id = str(execution_graph_version.id) + target_label = f"{execution_graph_version.graph.name} v{execution_graph_version.version}" + elif agent_version is not None: + target_id = str(agent_version.id) + target_label = f"{agent_version.agent.name} v{agent_version.version}" + return ImprovementCandidate.objects.create( + investigation=investigation, + target_type=target_type, + target_id=target_id, + target_label=target_label, + hypothesis=hypothesis or str((investigation.hypotheses or [{}])[0].get("hypothesis", "Improve target based on Progeny investigation evidence.")), + recommended_route=investigation.recommended_route, + evidence={"investigation_id": str(investigation.id), "signal_clusters": investigation.signal_clusters}, + execution_graph_version=execution_graph_version, + agent_version=agent_version, + ) + + def list_investigations(self, status: str | None = None): + investigations = ProgenyInvestigation.objects.all() + if status: + investigations = investigations.filter(status=status) + return investigations.order_by("-created_at") + + def list_improvements(self, status: str | None = None): + candidates = ImprovementCandidate.objects.select_related("investigation", "execution_graph_version", "agent_version").all() + if status: + candidates = candidates.filter(status=status) + return candidates.order_by("-created_at") + def _score(self, version: AgentVersion, benchmark_set: list[dict[str, object]]) -> dict[str, float]: if not benchmark_set: return {"completion_rate": 0.0, "test_pass_rate": 0.0, "review_acceptance": 0.0, "tokens": 0.0, "runtime": 0.0} @@ -159,3 +425,107 @@ class ProgenyService: "tokens": float(len(version.system_contract.split())), "runtime": float(len(benchmark_set)), } + + def _lineage( + self, + graph_run: GraphRun | None, + graph_node_run: GraphNodeRun | None, + execution_graph_version: ExecutionGraphVersion | None, + ) -> dict[str, Any]: + version = execution_graph_version or (graph_run.execution_graph_version if graph_run else None) + evidence: dict[str, object] = {} + if graph_run is not None: + evidence["graph_run_id"] = graph_run.id + if version is not None: + evidence["execution_graph_version_id"] = version.id + evidence["execution_graph"] = version.graph.name + evidence["execution_graph_version"] = version.version + if graph_node_run is not None: + evidence["graph_node_run_id"] = graph_node_run.id + evidence["node_id"] = graph_node_run.node_id + evidence["node_type"] = graph_node_run.node_type + evidence["visit_index"] = graph_node_run.visit_index + return {"execution_graph_version": version, "evidence": evidence} + + def _grouping_key(self, source: str, category: str, agent_version: AgentVersion | None, lineage: dict[str, Any]) -> str: + evidence = lineage["evidence"] + graph = evidence.get("execution_graph", "no_graph") + graph_version = evidence.get("execution_graph_version", "no_version") + node_type = evidence.get("node_type", "no_node") + agent = f"agent:{agent_version.id}" if agent_version else "agent:none" + return f"{source}:{category}:{graph}:v{graph_version}:{node_type}:{agent}" + + def _fallback_grouping_key(self, signal: ProgenySignal) -> str: + node_type = signal.graph_node_run.node_type if signal.graph_node_run_id else "no_node" + graph_version = signal.execution_graph_version.version if signal.execution_graph_version_id else "no_version" + fingerprint = str(signal.evidence.get("fingerprint") or signal.evidence.get("type") or signal.summary[:80]).lower() + return f"{signal.source}:{signal.failure_category}:v{graph_version}:{node_type}:{fingerprint}" + + def _analyze_signals(self, signals: list[ProgenySignal]) -> dict[str, object]: + corpus_parts: list[str] = [] + for signal in signals: + corpus_parts.extend( + [ + signal.source, + signal.failure_category, + signal.summary, + str(signal.evidence), + signal.graph_node_run.node_type if signal.graph_node_run_id else "", + ] + ) + corpus = " ".join(corpus_parts).lower() + sources = {signal.source for signal in signals} + node_types = {signal.graph_node_run.node_type for signal in signals if signal.graph_node_run_id} + graph_versions = {str(signal.execution_graph_version_id) for signal in signals if signal.execution_graph_version_id} + target = "NO_SYSTEMIC_CHANGE" + route = "record evidence, no intervention" + confidence = 0.45 + experiments = ["Review representative signal evidence manually before changing production behavior."] + impact = "Avoid unnecessary system changes when evidence is project-specific." + cost = "low" + + if "unsupported operation" in corpus or "missing capability" in corpus: + target = "TOOL_POLICY" + route = "Progeny" + confidence = 0.82 + experiments = ["Replay affected task with candidate tool policy that grants the missing operation."] + impact = "Reduce repeated task failures caused by unavailable safe mutations." + elif "malformed json" in corpus or "model_output_invalid" in corpus: + target = "MODEL" + route = "Model Studio / Model Router" + confidence = 0.78 + experiments = ["Replay prompts with stricter response-format contract and compare valid-output rate."] + impact = "Reduce invalid model responses before they reach mutation tools." + elif "path" in corpus or "timeout" in corpus or "provider_unavailable" in corpus or "provider_timeout" in corpus or "environment" in corpus: + target = "INFRASTRUCTURE" + route = "Steward" + confidence = 0.76 + experiments = ["Replay with captured environment and provider health checks before changing agents."] + impact = "Separate environmental breakage from agent/model quality issues." + elif "graph_runtime" in sources or (len(graph_versions) == 1 and len(node_types) == 1 and len(signals) > 1): + target = "GRAPH_NODE" if node_types else "WORKFLOW_GRAPH" + route = "Progeny Graph Evolution" + confidence = 0.74 + experiments = ["Replay the cluster against a challenger graph version with adjusted node policy or transition handling."] + impact = "Reduce recurring workflow-node failures without changing task implementation agents." + cost = "medium" + elif "false" in corpus and "review" in corpus: + target = "REVIEWER" + route = "Progeny" + confidence = 0.7 + experiments = ["Replay accepted diffs against a reviewer challenger with calibrated route-detection criteria."] + impact = "Reduce false rejections while preserving quality gates." + elif signals[0].task_id and len({signal.task_id for signal in signals}) == 1 and len(signals) == 1: + target = "PROJECT_INTENT" + route = "Repair" + confidence = 0.55 + experiments = ["Inspect project-specific assertion and task acceptance criteria before system changes."] + impact = "Resolve the isolated task without overfitting global behavior." + hypothesis = { + "target": target, + "hypothesis": f"Evidence from {len(signals)} signal(s) points to {target} as the likely root-cause target.", + "supporting_evidence": [str(signal.id) for signal in signals], + "graph_nodes": sorted(node_types), + "graph_versions": sorted(graph_versions), + } + return {"target": target, "route": route, "confidence": confidence, "experiments": experiments, "impact": impact, "cost": cost, "hypotheses": [hypothesis]} diff --git a/control_plane/agents/migrations/0006_progenyinvestigation_improvementcandidate.py b/control_plane/agents/migrations/0006_progenyinvestigation_improvementcandidate.py new file mode 100644 index 0000000..6022343 --- /dev/null +++ b/control_plane/agents/migrations/0006_progenyinvestigation_improvementcandidate.py @@ -0,0 +1,53 @@ +from django.db import migrations, models +import django.db.models.deletion +import uuid + + +class Migration(migrations.Migration): + dependencies = [ + ("agents", "0005_progenysignal_graph_lineage"), + ("graph", "0004_unique_champion_graph_version"), + ] + + operations = [ + migrations.CreateModel( + name="ProgenyInvestigation", + 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="OPEN", max_length=32)), + ("signal_clusters", models.JSONField(blank=True, default=list)), + ("affected_projects", models.JSONField(blank=True, default=list)), + ("affected_agents", models.JSONField(blank=True, default=list)), + ("affected_graph_versions", models.JSONField(blank=True, default=list)), + ("affected_nodes", models.JSONField(blank=True, default=list)), + ("hypotheses", models.JSONField(blank=True, default=list)), + ("recommended_target", models.CharField(max_length=80)), + ("confidence", models.FloatField(default=0.0)), + ("recommended_route", models.CharField(max_length=120)), + ("proposed_experiments", models.JSONField(blank=True, default=list)), + ("expected_impact", models.TextField(blank=True)), + ("estimated_cost", models.CharField(blank=True, max_length=80)), + ("signals", models.ManyToManyField(blank=True, related_name="investigations", to="agents.progenysignal")), + ], + ), + migrations.CreateModel( + name="ImprovementCandidate", + 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="PROPOSED", max_length=32)), + ("target_type", models.CharField(max_length=80)), + ("target_id", models.CharField(blank=True, max_length=120)), + ("target_label", models.CharField(blank=True, max_length=240)), + ("hypothesis", models.TextField()), + ("recommended_route", models.CharField(blank=True, max_length=120)), + ("evidence", models.JSONField(blank=True, default=dict)), + ("agent_version", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="improvement_candidates", to="agents.agentversion")), + ("execution_graph_version", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="improvement_candidates", to="graph.executiongraphversion")), + ("investigation", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="improvement_candidates", to="agents.progenyinvestigation")), + ], + ), + ] diff --git a/control_plane/agents/models.py b/control_plane/agents/models.py index ee3cc39..c8c9321 100644 --- a/control_plane/agents/models.py +++ b/control_plane/agents/models.py @@ -99,3 +99,33 @@ class ProgenySignal(TimestampedModel): status = models.CharField(max_length=32, default="OPEN") grouping_key = models.CharField(max_length=120, blank=True) model = models.CharField(max_length=120, blank=True) + + +class ProgenyInvestigation(TimestampedModel): + status = models.CharField(max_length=32, default="OPEN") + signals = models.ManyToManyField(ProgenySignal, related_name="investigations", blank=True) + signal_clusters = models.JSONField(default=list, blank=True) + affected_projects = models.JSONField(default=list, blank=True) + affected_agents = models.JSONField(default=list, blank=True) + affected_graph_versions = models.JSONField(default=list, blank=True) + affected_nodes = models.JSONField(default=list, blank=True) + hypotheses = models.JSONField(default=list, blank=True) + recommended_target = models.CharField(max_length=80) + confidence = models.FloatField(default=0.0) + recommended_route = models.CharField(max_length=120) + proposed_experiments = models.JSONField(default=list, blank=True) + expected_impact = models.TextField(blank=True) + estimated_cost = models.CharField(max_length=80, blank=True) + + +class ImprovementCandidate(TimestampedModel): + status = models.CharField(max_length=32, default="PROPOSED") + target_type = models.CharField(max_length=80) + target_id = models.CharField(max_length=120, blank=True) + target_label = models.CharField(max_length=240, blank=True) + hypothesis = models.TextField() + recommended_route = models.CharField(max_length=120, blank=True) + evidence = models.JSONField(default=dict, blank=True) + investigation = models.ForeignKey(ProgenyInvestigation, on_delete=models.SET_NULL, null=True, blank=True, related_name="improvement_candidates") + execution_graph_version = models.ForeignKey("graph.ExecutionGraphVersion", on_delete=models.SET_NULL, null=True, blank=True, related_name="improvement_candidates") + agent_version = models.ForeignKey(AgentVersion, on_delete=models.SET_NULL, null=True, blank=True, related_name="improvement_candidates") diff --git a/graph/langgraph_runtime.py b/graph/langgraph_runtime.py index 0e86f15..1359c83 100644 --- a/graph/langgraph_runtime.py +++ b/graph/langgraph_runtime.py @@ -123,6 +123,13 @@ class LangGraphRuntime(GraphRuntime): graph_run.failure_reason = metadata["final_failure_reason"] graph_run.metadata = metadata graph_run.save(update_fields=["status", "failure_reason", "metadata", "updated_at"]) + try: + from agents.progeny import ProgenyService + + node_run = graph_run.node_runs.filter(node_id=node_id).order_by("-visit_index").first() + ProgenyService(self.bus).create_graph_runtime_signal(graph_run, "GRAPH_RUNTIME_ERROR", graph_run.failure_reason, {"node_id": node_id, "edge_result": edge_result}, graph_node_run=node_run) + except Exception: + pass return edge_result GraphEdgeTraversal.objects.create(graph_run=graph_run, source_node=node_id, target_node=target, condition=edge_result, result=edge_result) metadata = dict(graph_run.metadata) diff --git a/graph/native_runtime.py b/graph/native_runtime.py index 8e72930..3d1773b 100644 --- a/graph/native_runtime.py +++ b/graph/native_runtime.py @@ -123,6 +123,7 @@ class NativeGraphRuntime(GraphRuntime): graph_run.failure_reason = metadata["final_failure_reason"] graph_run.metadata = metadata graph_run.save(update_fields=["status", "metadata", "completed_at", "failure_reason", "updated_at"]) + self._create_graph_runtime_signal(graph_run, node_run, "GRAPH_NODE_FAILED", metadata["final_failure_reason"], result.failure_evidence or {}) self.bus.publish("GRAPH_RUN_FAILED", project=graph_run.project, task=graph_run.task, payload={"graph_run_id": graph_run.id, "node_id": node_id}) break next_node = self._select_next(spec, node_id, result.edge_result) @@ -134,6 +135,7 @@ class NativeGraphRuntime(GraphRuntime): graph_run.failure_reason = metadata["final_failure_reason"] graph_run.metadata = metadata graph_run.save(update_fields=["status", "metadata", "completed_at", "failure_reason", "updated_at"]) + self._create_graph_runtime_signal(graph_run, node_run, "GRAPH_RUNTIME_ERROR", graph_run.failure_reason, {"node_id": node_id, "edge_result": result.edge_result}) self.bus.publish("GRAPH_RUN_FAILED", project=graph_run.project, task=graph_run.task, payload={"graph_run_id": graph_run.id, "reason": graph_run.failure_reason}) break GraphEdgeTraversal.objects.create(graph_run=graph_run, source_node=node_id, target_node=next_node, condition=result.edge_result, result=result.edge_result) @@ -222,3 +224,11 @@ class NativeGraphRuntime(GraphRuntime): if len(text) > 20000: return {"truncated": True, "excerpt": text[:20000]} return value + + def _create_graph_runtime_signal(self, graph_run: GraphRun, node_run: GraphNodeRun, category: str, summary: str, evidence: dict[str, Any]) -> None: + try: + from agents.progeny import ProgenyService + + ProgenyService(self.bus).create_graph_runtime_signal(graph_run, category, summary, evidence, graph_node_run=node_run) + except Exception: + return None diff --git a/graph/task_nodes.py b/graph/task_nodes.py index 4769746..7772b01 100644 --- a/graph/task_nodes.py +++ b/graph/task_nodes.py @@ -4,6 +4,7 @@ from pathlib import Path from agents.coder import Coder from agents.judge import Judge +from agents.progeny import ProgenyService from agents.reviewer import Reviewer from control_plane.agents.models import AgentRole, AgentVersion from control_plane.events.bus import EventBus @@ -35,6 +36,7 @@ class TaskExecutionServices: self.coder = Coder(router) self.reviewer = Reviewer() self.judge = Judge() + self.progeny = ProgenyService(self.bus) self.tests = DeterministicTestRunner() self.test_command = test_command or ["python", "-m", "pytest"] @@ -83,6 +85,9 @@ class TaskNode: context.graph_run.metadata = metadata context.graph_run.save(update_fields=["metadata", "updated_at"]) + def current_node_run(self, context: GraphExecutionContext): + return context.graph_run.node_runs.filter(node_id=context.graph_run.current_node).order_by("-visit_index").first() + def _clear_current_failure(self, metadata: dict[str, object]) -> None: metadata.pop("current_failure_reason", None) metadata.pop("current_failure_findings", None) @@ -153,18 +158,51 @@ class CoderNode(TaskNode): metadata = self.metadata(context) attempt = TaskAttempt.objects.get(id=metadata["current_attempt_id"]) coder_version = attempt.coder - result = self.services.coder.execute(attempt.context_snapshot, self.services.tools(worktree), project=task.project, agent_version=coder_version) + try: + result = self.services.coder.execute(attempt.context_snapshot, self.services.tools(worktree), project=task.project, agent_version=coder_version) + except Exception as exc: + self.services.progeny.create_provider_signal( + task.project, + task, + task.milestone, + coder_version, + self._provider_failure_category(str(exc)), + f"Coder provider/runtime failure: {exc}", + {"error": str(exc)}, + graph_run=context.graph_run, + graph_node_run=self.current_node_run(context), + ) + raise attempt.coder_result = {"status": result.status, "summary": result.summary, "changed_files": result.changed_files, "metadata": result.metadata} attempt.save(update_fields=["coder_result", "updated_at"]) if result.status != "COMPLETE": metadata["current_failure_reason"] = "coder_failed" metadata["current_failure_findings"] = [result.summary] metadata["current_failure_node_id"] = "coder" + if "malformed json" in result.summary.lower() or "json" in result.summary.lower(): + self.services.progeny.create_model_output_signal( + task.project, + task, + task.milestone, + coder_version, + result.summary, + str(result.metadata), + graph_run=context.graph_run, + graph_node_run=self.current_node_run(context), + ) else: self._clear_current_failure(metadata) self.save_metadata(context, metadata) return NodeResult("COMPLETE", "success" if result.status == "COMPLETE" else "failure", {"coder_status": result.status, "summary": result.summary}, result.metadata.get("telemetry", {}) if isinstance(result.metadata, dict) else {}) + def _provider_failure_category(self, error: str) -> str: + lowered = error.lower() + if "timeout" in lowered or "timed out" in lowered: + return "PROVIDER_TIMEOUT" + if "unavailable" in lowered or "connection" in lowered or "configured" in lowered: + return "PROVIDER_UNAVAILABLE" + return "PROVIDER_RUNTIME_ERROR" + class RunTestsNode(TaskNode): def __init__(self, services: TaskExecutionServices) -> None: @@ -217,6 +255,18 @@ class ReviewNode(TaskNode): metadata["current_failure_reason"] = "review_failed" metadata["current_failure_findings"] = review.findings metadata["current_failure_node_id"] = "review" + self.services.progeny.create_reviewer_signal( + task.project, + task, + task.milestone, + reviewer_version, + review.status, + review.findings, + review.summary, + graph_run=context.graph_run, + graph_node_run=self.current_node_run(context), + metadata={"affected_agent_version_id": str(attempt.coder_id), "reviewer_version_id": str(reviewer_version.id), "test_status": test_run.status}, + ) else: self._clear_current_failure(metadata) self.save_metadata(context, metadata) @@ -247,6 +297,18 @@ class JudgeNode(TaskNode): metadata["current_failure_reason"] = "judge_failed" metadata["current_failure_findings"] = verification.evidence metadata["current_failure_node_id"] = "judge" + self.services.progeny.create_judge_signal( + task.project, + task, + task.milestone, + judge_version, + verification.result, + verification.evidence, + verification.summary, + graph_run=context.graph_run, + graph_node_run=self.current_node_run(context), + metadata={"affected_agent_version_id": str(attempt.coder_id), "judge_version_id": str(judge_version.id), "test_status": test_run.status}, + ) else: self._clear_current_failure(metadata) self.save_metadata(context, metadata) @@ -283,6 +345,16 @@ class RetryOrFailNode(TaskNode): return NodeResult("COMPLETE", "retry_available", {"retry_count": task.retry_count, "classification": classification}) task.status = TaskStatus.FAILED task.save(update_fields=["retry_count", "status", "updated_at"]) + self.services.progeny.create_retry_exhausted_signal( + task.project, + task, + task.milestone, + attempt.coder, + task.retry_count, + graph_run=context.graph_run, + graph_node_run=self.current_node_run(context), + evidence={"reason": reason, "classification": classification, "findings": findings}, + ) self.services.bus.publish("TASK_RETRY_EXHAUSTED", project=task.project, task=task, payload={"reason": reason, "classification": classification, "findings": findings}) self.services.bus.publish(EventType.TASK_FAILED, project=task.project, task=task, payload={"reason": reason, "will_retry": False, "classification": classification, "findings": findings}) return NodeResult("COMPLETE", "retry_exhausted", {"retry_count": task.retry_count, "classification": classification}) diff --git a/tests/test_native_graph_runtime.py b/tests/test_native_graph_runtime.py index 83b5b84..7893e1e 100644 --- a/tests/test_native_graph_runtime.py +++ b/tests/test_native_graph_runtime.py @@ -1,5 +1,6 @@ from __future__ import annotations +from control_plane.agents.models import ProgenySignal from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus, GraphEdgeTraversal, GraphRun, GraphRunStatus from graph.native_runtime import NativeGraphRuntime from graph.registry import NodeHandlerRegistry, NodeResult @@ -112,3 +113,26 @@ def test_native_graph_runtime_records_paused_runs() -> None: assert result.status == GraphRunStatus.PAUSED assert result.failure_reason == "AWAITING_APPROVAL" + + +def test_native_graph_runtime_signal_created_for_missing_edge() -> None: + spec = ExecutionGraphSpec( + name="missing_edge_graph", + version=1, + graph_type="FIXTURE", + entry="start", + nodes={"start": GraphNodeSpec("start", "start"), "done": GraphNodeSpec("done", "done")}, + edges=[GraphEdgeSpec("start", "done", "expected")], + terminal_nodes=["done"], + ) + graph_run = persist_spec(spec) + registry = NodeHandlerRegistry() + registry.register(FixedNode("start", "unexpected")) + + result = NativeGraphRuntime(registry).run_until_terminal_or_paused(graph_run) + + assert result.status == GraphRunStatus.FAILED + signal = ProgenySignal.objects.get(source="graph_runtime", failure_category="GRAPH_RUNTIME_ERROR") + assert signal.graph_run == graph_run + assert signal.graph_node_run.node_id == "start" + assert signal.evidence["edge_result"] == "unexpected" diff --git a/tests/test_progeny_investigation.py b/tests/test_progeny_investigation.py new file mode 100644 index 0000000..504ecf9 --- /dev/null +++ b/tests/test_progeny_investigation.py @@ -0,0 +1,94 @@ +from __future__ import annotations + +from agents.progeny import ProgenyService +from control_plane.agents.models import Agent, AgentVersion, ImprovementCandidate, ProgenyInvestigation, ProgenySignal, PromotionStatus +from control_plane.projects.models import Milestone, Project, ProjectPlan, Task +from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus, GraphRun +from graph.task_execution import task_execution_graph_v1 + + +def fixture_context(): + project = Project.objects.create(name="Progeny Investigation", goal="Improve Artifex") + plan = ProjectPlan.objects.create(project=project, version=1, goal="Improve Artifex") + milestone = Milestone.objects.create(project=project, plan=plan, key="P1", title="Progeny", goal="Signals") + task = Task.objects.create(project=project, milestone=milestone, task_type="implementation", goal="Investigate failure") + agent = Agent.objects.create(name="Investigation Reviewer", role="REVIEWER") + version = AgentVersion.objects.create(agent=agent, version=1, model="qwen", system_contract="Review", promotion_status=PromotionStatus.CHAMPION) + spec = task_execution_graph_v1() + definition = ExecutionGraphDefinition.objects.create(name="investigation_task_execution", graph_type=spec.graph_type) + graph_version = ExecutionGraphVersion.objects.create(graph=definition, version=1, status=ExecutionGraphVersionStatus.CHAMPION, graph_spec=spec.to_dict()) + graph_run = GraphRun.objects.create(execution_graph_version=graph_version, project=project, milestone=milestone, task=task, current_node="review") + review_node = graph_run.node_runs.create(node_id="review", node_type="review", visit_index=1) + return project, milestone, task, version, graph_version, graph_run, review_node + + +def test_smart_investigation_distinguishes_real_failure_patterns() -> None: + project, milestone, task, agent_version, graph_version, graph_run, review_node = fixture_context() + service = ProgenyService() + cases = [ + ("tool_policy", "retry", "TASK_RETRY_EXHAUSTED", "Unsupported operation: delete_file", {"classification": "missing_capability"}, "TOOL_POLICY"), + ("model_output", "model_output", "MODEL_OUTPUT_INVALID", "Provider returned malformed JSON", {"raw_output": "{invalid"}, "MODEL"), + ("infrastructure", "provider", "PROVIDER_TIMEOUT", "Spark PATH test runner environment failure", {"stderr": "python not found on PATH"}, "INFRASTRUCTURE"), + ("reviewer", "reviewer", "REJECTED", "Reviewer false rejection around health route", {"findings": [{"type": "false_rejection"}]}, "REVIEWER"), + ] + + for key, source, category, summary, evidence, expected_target in cases: + signal = ProgenySignal.objects.create( + project=project, + task=task, + milestone=milestone, + agent_version=agent_version, + graph_run=graph_run, + graph_node_run=review_node, + execution_graph_version=graph_version, + source=source, + severity="high", + failure_category=category, + summary=summary, + evidence=evidence, + grouping_key=key, + model=agent_version.model, + ) + + investigation = service.create_smart_investigation(signal.grouping_key) + + assert investigation.recommended_target == expected_target + assert investigation.signals.get() == signal + + +def test_smart_investigation_can_recommend_graph_node_and_create_graph_candidate() -> None: + project, milestone, task, agent_version, graph_version, graph_run, review_node = fixture_context() + service = ProgenyService() + for index in range(2): + ProgenySignal.objects.create( + project=project, + task=task, + milestone=milestone, + agent_version=agent_version, + graph_run=graph_run, + graph_node_run=review_node, + execution_graph_version=graph_version, + source="graph_runtime", + severity="high", + failure_category="GRAPH_RUNTIME_ERROR", + summary="Repeated failure concentrated at review node", + evidence={"node_id": "review", "index": index}, + grouping_key="graph-node-review", + model=agent_version.model, + ) + + investigation = service.create_smart_investigation("graph-node-review") + candidate = service.create_improvement_candidate( + investigation, + hypothesis="Insert static-analysis node before Reviewer to reduce REWORK_REQUIRED", + ) + + assert investigation.recommended_target == "GRAPH_NODE" + assert investigation.recommended_route == "Progeny Graph Evolution" + assert investigation.affected_graph_versions == [str(graph_version.id)] + assert investigation.affected_nodes == ["review"] + assert candidate.target_type == "GRAPH_NODE" + assert candidate.execution_graph_version == graph_version + assert candidate.hypothesis == "Insert static-analysis node before Reviewer to reduce REWORK_REQUIRED" + assert ProgenyInvestigation.objects.count() == 1 + assert ImprovementCandidate.objects.count() == 1 diff --git a/tests/test_task_execution_native_graph.py b/tests/test_task_execution_native_graph.py index 8619584..35d6492 100644 --- a/tests/test_task_execution_native_graph.py +++ b/tests/test_task_execution_native_graph.py @@ -4,6 +4,7 @@ from pathlib import Path from agents.providers import DeterministicCodingProvider from control_plane.agents.management.commands.seed_core_agents import Command as SeedAgentsCommand +from control_plane.agents.models import ProgenySignal from control_plane.events.models import Event, EventType from control_plane.projects.models import CommitRecord, Project, ProjectPlan, Milestone, Task, TaskStatus from control_plane.verification.models import Review, TestRun, Verification, VerificationResult @@ -69,6 +70,20 @@ class FirstAttemptBadProvider(DeterministicCodingProvider): return super().complete(request) +class InvalidJsonProvider(DeterministicCodingProvider): + def complete(self, request): + if "inspection phase" in request.prompt.lower(): + return super().complete(request) + return ModelResponseContract("qwen-deterministic", "{invalid", {}) + + +class TimeoutProvider(DeterministicCodingProvider): + def complete(self, request): + if "inspection phase" in request.prompt.lower(): + return super().complete(request) + raise TimeoutError("provider timed out") + + def run_graph(task: Task, *, interrupt_after: str | None = None, provider=None) -> GraphRun: SeedAgentsCommand().handle() services = TaskExecutionServices(ModelRouter({"qwen": provider or DeterministicCodingProvider()}), test_command=["python", "manage.py", "test"]) @@ -118,6 +133,10 @@ def test_native_task_execution_graph_retry_exhaustion_matches_loop_semantics(tmp assert Event.objects.filter(task=task, event_type="TASK_RETRY_EXHAUSTED").exists() assert graph_run.metadata["final_failure_reason"] == "review_failed" assert graph_run.metadata["historical_failures"][-1]["reason"] == "review_failed" + retry_signal = ProgenySignal.objects.get(task=task, failure_category="TASK_RETRY_EXHAUSTED") + assert retry_signal.graph_run == graph_run + assert retry_signal.graph_node_run.node_id == "retry_or_fail" + assert ProgenySignal.objects.filter(task=task, source="graph_runtime").count() == 0 def test_native_task_execution_graph_resume_after_coder_prevents_duplicate_commit(tmp_path: Path) -> None: @@ -142,6 +161,31 @@ def test_native_task_execution_graph_resume_after_coder_prevents_duplicate_commi assert graph_run.node_runs.filter(node_id="coder").count() == 1 +def test_model_output_invalid_creates_progeny_signal_with_graph_lineage(tmp_path: Path) -> None: + repo = create_disposable_django_repo(tmp_path) + task = create_task(repo, "Add a /health endpoint returning JSON {\"status\": \"ok\"} and add tests.", ["/health returns JSON ok", "tests pass"], max_retries=0) + + graph_run = run_graph(task, provider=InvalidJsonProvider()) + + assert graph_run.status == GraphRunStatus.FAILED + signal = ProgenySignal.objects.get(task=task, failure_category="MODEL_OUTPUT_INVALID") + assert signal.graph_run == graph_run + assert signal.graph_node_run.node_id == "coder" + assert signal.evidence["node_type"] == "coder" + + +def test_provider_timeout_creates_distinct_progeny_signal(tmp_path: Path) -> None: + repo = create_disposable_django_repo(tmp_path) + task = create_task(repo, "Add a /health endpoint returning JSON {\"status\": \"ok\"} and add tests.", ["/health returns JSON ok", "tests pass"]) + + graph_run = run_graph(task, provider=TimeoutProvider()) + + assert graph_run.status == GraphRunStatus.FAILED + signal = ProgenySignal.objects.get(task=task, source="provider", failure_category="PROVIDER_TIMEOUT") + assert signal.graph_run == graph_run + assert signal.graph_node_run.node_id == "coder" + + def test_langgraph_task_execution_graph_success_matches_native_domain_outcome(tmp_path: Path) -> None: repo = create_disposable_django_repo(tmp_path) task = create_task(repo, "Add a /health endpoint returning JSON {\"status\": \"ok\"} and add tests.", ["/health returns JSON ok", "tests pass"]) @@ -179,6 +223,10 @@ def test_langgraph_task_execution_graph_retry_exhaustion_matches_native_domain_o assert Event.objects.filter(task=task, event_type="TASK_RETRY_EXHAUSTED").exists() assert graph_run.metadata["final_failure_reason"] == "review_failed" assert graph_run.metadata["historical_failures"][-1]["reason"] == "review_failed" + reviewer_signal = ProgenySignal.objects.filter(task=task, source="reviewer", failure_category="REWORK_REQUIRED").first() + assert reviewer_signal is not None + assert reviewer_signal.graph_run == graph_run + assert reviewer_signal.graph_node_run.node_id == "review" def test_langgraph_retry_then_success_preserves_historical_failures_without_final_failure(tmp_path: Path) -> None: @@ -203,6 +251,9 @@ def test_langgraph_retry_then_success_preserves_historical_failures_without_fina "evidence": [{"type": "missing_health_route", "severity": "high", "message": "Diff does not add /health"}], } ] + reviewer_signal = ProgenySignal.objects.get(task=task, source="reviewer") + assert reviewer_signal.execution_graph_version == graph_run.execution_graph_version + assert reviewer_signal.graph_node_run.node_id == "review" def test_langgraph_resume_after_coder_tests_and_judge_does_not_duplicate_commit(tmp_path: Path) -> None: diff --git a/tests/test_v2_progeny_signals.py b/tests/test_v2_progeny_signals.py index 7044472..17d7e47 100644 --- a/tests/test_v2_progeny_signals.py +++ b/tests/test_v2_progeny_signals.py @@ -7,6 +7,8 @@ from agents.progeny import ProgenyService from control_plane.agents.models import Agent, AgentVersion, PromotionStatus from control_plane.events.models import Event from control_plane.projects.models import Project, ProjectPlan, Milestone, Task +from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus, GraphRun +from graph.task_execution import task_execution_graph_v1 class ProgenySignalTests(TestCase): @@ -38,7 +40,7 @@ class ProgenySignalTests(TestCase): self.assertEqual(signal.source, "reviewer") self.assertEqual(signal.failure_category, "REWORK_REQUIRED") self.assertEqual(signal.severity, "high") - self.assertEqual(signal.grouping_key, f"reviewer:{self.task.id}:REWORK_REQUIRED") + self.assertIn("reviewer:REWORK_REQUIRED", signal.grouping_key) event = Event.objects.get(event_type="PROGENY_SIGNAL_CREATED") self.assertEqual(event.payload["signal_id"], str(signal.id)) @@ -89,7 +91,7 @@ class ProgenySignalTests(TestCase): raw_output="{invalid json}", ) self.assertEqual(signal.source, "model_output") - self.assertEqual(signal.failure_category, "MALFORMED_OUTPUT") + self.assertEqual(signal.failure_category, "MODEL_OUTPUT_INVALID") self.assertEqual(signal.severity, "high") self.assertIn("Invalid JSON", signal.summary) @@ -105,7 +107,7 @@ class ProgenySignalTests(TestCase): attempts=3, ) self.assertEqual(signal.source, "retry") - self.assertEqual(signal.failure_category, "RETRY_EXHAUSTED") + self.assertEqual(signal.failure_category, "TASK_RETRY_EXHAUSTED") self.assertEqual(signal.severity, "critical") self.assertIn("exhausted 3 retries", signal.summary) @@ -133,6 +135,96 @@ class ProgenySignalTests(TestCase): ) self.assertEqual(signal1.grouping_key, signal2.grouping_key) + def test_signal_records_graph_lineage_when_available(self): + spec = task_execution_graph_v1() + definition = ExecutionGraphDefinition.objects.create(name="progeny_lineage", graph_type=spec.graph_type) + version = ExecutionGraphVersion.objects.create(graph=definition, version=1, status=ExecutionGraphVersionStatus.CHAMPION, graph_spec=spec.to_dict()) + graph_run = GraphRun.objects.create(execution_graph_version=version, project=self.project, milestone=self.milestone, task=self.task, current_node="review") + node_run = graph_run.node_runs.create(node_id="review", node_type="review", visit_index=2) + + signal = self.service.create_reviewer_signal( + project=self.project, + task=self.task, + milestone=self.milestone, + agent_version=self.agent_version, + status="REWORK_REQUIRED", + findings=[{"type": "missing_health_route"}], + summary="Reviewer requested rework", + graph_run=graph_run, + graph_node_run=node_run, + ) + + self.assertEqual(signal.graph_run, graph_run) + self.assertEqual(signal.execution_graph_version, version) + self.assertEqual(signal.graph_node_run, node_run) + self.assertEqual(signal.evidence["node_id"], "review") + self.assertEqual(signal.evidence["node_type"], "review") + self.assertEqual(signal.evidence["visit_index"], 2) + + def test_inbox_filters_unresolved_signals_by_graph_dimensions(self): + spec = task_execution_graph_v1() + definition = ExecutionGraphDefinition.objects.create(name="progeny_inbox", graph_type=spec.graph_type) + version = ExecutionGraphVersion.objects.create(graph=definition, version=1, status=ExecutionGraphVersionStatus.CHAMPION, graph_spec=spec.to_dict()) + graph_run = GraphRun.objects.create(execution_graph_version=version, project=self.project, milestone=self.milestone, task=self.task, current_node="judge") + node_run = graph_run.node_runs.create(node_id="judge", node_type="judge", visit_index=1) + signal = self.service.create_judge_signal( + project=self.project, + task=self.task, + milestone=self.milestone, + agent_version=self.agent_version, + result="FAIL", + evidence=[{"type": "test_status", "status": "FAIL"}], + summary="Judge failed", + graph_run=graph_run, + graph_node_run=node_run, + ) + + matches = list( + self.service.query_inbox( + source="judge", + status="OPEN", + project=self.project, + agent=self.agent, + agent_version=self.agent_version, + model="test-model", + execution_graph=definition, + execution_graph_version=version, + graph_node="judge", + severity="high", + failure_category="FAIL", + ) + ) + + self.assertEqual(matches, [signal]) + + def test_group_unresolved_signals_exposes_impact_summary(self): + spec = task_execution_graph_v1() + definition = ExecutionGraphDefinition.objects.create(name="progeny_grouping", graph_type=spec.graph_type) + version = ExecutionGraphVersion.objects.create(graph=definition, version=1, status=ExecutionGraphVersionStatus.CHAMPION, graph_spec=spec.to_dict()) + graph_run = GraphRun.objects.create(execution_graph_version=version, project=self.project, milestone=self.milestone, task=self.task, current_node="review") + node_run = graph_run.node_runs.create(node_id="review", node_type="review", visit_index=1) + for index in range(2): + self.service.create_reviewer_signal( + project=self.project, + task=self.task, + milestone=self.milestone, + agent_version=self.agent_version, + status="REWORK_REQUIRED", + findings=[{"type": "missing_health_route", "index": index}], + summary="Reviewer requested rework", + graph_run=graph_run, + graph_node_run=node_run, + ) + + groups = self.service.group_unresolved_signals(source="reviewer", failure_category="REWORK_REQUIRED") + + self.assertEqual(len(groups), 1) + self.assertEqual(groups[0].occurrence_count, 2) + self.assertEqual(groups[0].affected_projects, [str(self.project.id)]) + self.assertEqual(groups[0].affected_agents, [str(self.agent_version.id)]) + self.assertEqual(groups[0].affected_graph_versions, [str(version.id)]) + self.assertEqual(groups[0].affected_graph_nodes, ["review"]) + def test_event_lineage(self): self.service.create_reviewer_signal( project=self.project,