Wire Progeny signals to graph evidence

This commit is contained in:
Daniel Maddern 2026-08-15 17:53:52 +07:00
parent b36a1319ec
commit 2f142b59c4
10 changed files with 822 additions and 19 deletions

View file

@ -1,11 +1,14 @@
from __future__ import annotations from __future__ import annotations
from dataclasses import dataclass 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.bus import EventBus
from control_plane.events.models import EventType from control_plane.events.models import EventType
from control_plane.projects.models import Project, Task, Milestone from control_plane.projects.models import Project, Task, Milestone
from graph.models import ExecutionGraphVersion, GraphNodeRun, GraphRun
@dataclass(frozen=True) @dataclass(frozen=True)
@ -14,6 +17,20 @@ class BenchmarkDecision:
metrics: dict[str, float] 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: class ProgenyService:
def __init__(self, bus: EventBus | None = None) -> None: def __init__(self, bus: EventBus | None = None) -> None:
self.bus = bus or EventBus() self.bus = bus or EventBus()
@ -74,80 +91,329 @@ class ProgenyService:
challenger.save(update_fields=["promotion_status", "updated_at"]) challenger.save(update_fields=["promotion_status", "updated_at"])
return BenchmarkDecision("REJECTED", challenger_metrics) 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" 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( signal = ProgenySignal.objects.create(
project=project, project=project,
task=task, task=task,
milestone=milestone, milestone=milestone,
agent_version=agent_version, agent_version=agent_version,
graph_run=graph_run,
graph_node_run=graph_node_run,
execution_graph_version=lineage["execution_graph_version"],
source="reviewer", source="reviewer",
severity=severity, severity=severity,
failure_category=status, failure_category=status,
summary=summary, summary=summary,
evidence={"findings": findings}, evidence={"findings": findings, **lineage["evidence"], **(metadata or {})},
status="OPEN", status="OPEN",
grouping_key=f"reviewer:{task.id}:{status}", grouping_key=self._grouping_key("reviewer", status, agent_version, lineage),
model=agent_version.model, 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}) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "reviewer", "status": status})
return signal 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" severity = "high" if result == "FAIL" else "info"
lineage = self._lineage(graph_run, graph_node_run, execution_graph_version)
signal = ProgenySignal.objects.create( signal = ProgenySignal.objects.create(
project=project, project=project,
task=task, task=task,
milestone=milestone, milestone=milestone,
agent_version=agent_version, agent_version=agent_version,
graph_run=graph_run,
graph_node_run=graph_node_run,
execution_graph_version=lineage["execution_graph_version"],
source="judge", source="judge",
severity=severity, severity=severity,
failure_category=result, failure_category=result,
summary=summary, summary=summary,
evidence={"evidence": evidence}, evidence={"evidence": evidence, **lineage["evidence"], **(metadata or {})},
status="OPEN", status="OPEN",
grouping_key=f"judge:{task.id}:{result}", grouping_key=self._grouping_key("judge", result, agent_version, lineage),
model=agent_version.model, 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}) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "judge", "result": result})
return signal 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( signal = ProgenySignal.objects.create(
project=project, project=project,
task=task, task=task,
milestone=milestone, milestone=milestone,
agent_version=agent_version, agent_version=agent_version,
graph_run=graph_run,
graph_node_run=graph_node_run,
execution_graph_version=lineage["execution_graph_version"],
source="model_output", source="model_output",
severity="high", severity="high",
failure_category="MALFORMED_OUTPUT", failure_category="MODEL_OUTPUT_INVALID",
summary=f"Model output malformed: {error}", summary=f"Model output malformed: {error}",
evidence={"raw_output": raw_output[:1000]}, evidence={"raw_output": raw_output[:1000], "error": error, **lineage["evidence"]},
status="OPEN", 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, 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"}) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "model_output"})
return signal 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( signal = ProgenySignal.objects.create(
project=project, project=project,
task=task, task=task,
milestone=milestone, milestone=milestone,
agent_version=agent_version, 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", source="retry",
severity="critical", severity="critical",
failure_category="RETRY_EXHAUSTED", failure_category="TASK_RETRY_EXHAUSTED",
summary=f"Task {task.id} exhausted {attempts} retries", summary=f"Task {task.id} exhausted {attempts} retries",
evidence={"attempts": attempts}, evidence={"attempts": attempts, **(evidence or {}), **lineage["evidence"]},
status="OPEN", status="OPEN",
grouping_key=f"retry:{task.id}", grouping_key=self._grouping_key("retry", "TASK_RETRY_EXHAUSTED", agent_version, lineage),
model=agent_version.model, model=agent_version.model,
) )
self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "retry"}) self.bus.publish("PROGENY_SIGNAL_CREATED", project=project, task=task, actor="progeny", payload={"signal_id": str(signal.id), "source": "retry"})
return signal 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]: def _score(self, version: AgentVersion, benchmark_set: list[dict[str, object]]) -> dict[str, float]:
if not benchmark_set: if not benchmark_set:
return {"completion_rate": 0.0, "test_pass_rate": 0.0, "review_acceptance": 0.0, "tokens": 0.0, "runtime": 0.0} 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())), "tokens": float(len(version.system_contract.split())),
"runtime": float(len(benchmark_set)), "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]}

View file

@ -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")),
],
),
]

View file

@ -99,3 +99,33 @@ class ProgenySignal(TimestampedModel):
status = models.CharField(max_length=32, default="OPEN") status = models.CharField(max_length=32, default="OPEN")
grouping_key = models.CharField(max_length=120, blank=True) grouping_key = models.CharField(max_length=120, blank=True)
model = 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")

View file

@ -123,6 +123,13 @@ class LangGraphRuntime(GraphRuntime):
graph_run.failure_reason = metadata["final_failure_reason"] graph_run.failure_reason = metadata["final_failure_reason"]
graph_run.metadata = metadata graph_run.metadata = metadata
graph_run.save(update_fields=["status", "failure_reason", "metadata", "updated_at"]) 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 return edge_result
GraphEdgeTraversal.objects.create(graph_run=graph_run, source_node=node_id, target_node=target, condition=edge_result, result=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) metadata = dict(graph_run.metadata)

View file

@ -123,6 +123,7 @@ class NativeGraphRuntime(GraphRuntime):
graph_run.failure_reason = metadata["final_failure_reason"] graph_run.failure_reason = metadata["final_failure_reason"]
graph_run.metadata = metadata graph_run.metadata = metadata
graph_run.save(update_fields=["status", "metadata", "completed_at", "failure_reason", "updated_at"]) 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}) 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 break
next_node = self._select_next(spec, node_id, result.edge_result) 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.failure_reason = metadata["final_failure_reason"]
graph_run.metadata = metadata graph_run.metadata = metadata
graph_run.save(update_fields=["status", "metadata", "completed_at", "failure_reason", "updated_at"]) 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}) 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 break
GraphEdgeTraversal.objects.create(graph_run=graph_run, source_node=node_id, target_node=next_node, condition=result.edge_result, result=result.edge_result) 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: if len(text) > 20000:
return {"truncated": True, "excerpt": text[:20000]} return {"truncated": True, "excerpt": text[:20000]}
return value 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

View file

@ -4,6 +4,7 @@ from pathlib import Path
from agents.coder import Coder from agents.coder import Coder
from agents.judge import Judge from agents.judge import Judge
from agents.progeny import ProgenyService
from agents.reviewer import Reviewer from agents.reviewer import Reviewer
from control_plane.agents.models import AgentRole, AgentVersion from control_plane.agents.models import AgentRole, AgentVersion
from control_plane.events.bus import EventBus from control_plane.events.bus import EventBus
@ -35,6 +36,7 @@ class TaskExecutionServices:
self.coder = Coder(router) self.coder = Coder(router)
self.reviewer = Reviewer() self.reviewer = Reviewer()
self.judge = Judge() self.judge = Judge()
self.progeny = ProgenyService(self.bus)
self.tests = DeterministicTestRunner() self.tests = DeterministicTestRunner()
self.test_command = test_command or ["python", "-m", "pytest"] self.test_command = test_command or ["python", "-m", "pytest"]
@ -83,6 +85,9 @@ class TaskNode:
context.graph_run.metadata = metadata context.graph_run.metadata = metadata
context.graph_run.save(update_fields=["metadata", "updated_at"]) 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: def _clear_current_failure(self, metadata: dict[str, object]) -> None:
metadata.pop("current_failure_reason", None) metadata.pop("current_failure_reason", None)
metadata.pop("current_failure_findings", None) metadata.pop("current_failure_findings", None)
@ -153,18 +158,51 @@ class CoderNode(TaskNode):
metadata = self.metadata(context) metadata = self.metadata(context)
attempt = TaskAttempt.objects.get(id=metadata["current_attempt_id"]) attempt = TaskAttempt.objects.get(id=metadata["current_attempt_id"])
coder_version = attempt.coder 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.coder_result = {"status": result.status, "summary": result.summary, "changed_files": result.changed_files, "metadata": result.metadata}
attempt.save(update_fields=["coder_result", "updated_at"]) attempt.save(update_fields=["coder_result", "updated_at"])
if result.status != "COMPLETE": if result.status != "COMPLETE":
metadata["current_failure_reason"] = "coder_failed" metadata["current_failure_reason"] = "coder_failed"
metadata["current_failure_findings"] = [result.summary] metadata["current_failure_findings"] = [result.summary]
metadata["current_failure_node_id"] = "coder" 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: else:
self._clear_current_failure(metadata) self._clear_current_failure(metadata)
self.save_metadata(context, 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 {}) 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): class RunTestsNode(TaskNode):
def __init__(self, services: TaskExecutionServices) -> None: def __init__(self, services: TaskExecutionServices) -> None:
@ -217,6 +255,18 @@ class ReviewNode(TaskNode):
metadata["current_failure_reason"] = "review_failed" metadata["current_failure_reason"] = "review_failed"
metadata["current_failure_findings"] = review.findings metadata["current_failure_findings"] = review.findings
metadata["current_failure_node_id"] = "review" 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: else:
self._clear_current_failure(metadata) self._clear_current_failure(metadata)
self.save_metadata(context, metadata) self.save_metadata(context, metadata)
@ -247,6 +297,18 @@ class JudgeNode(TaskNode):
metadata["current_failure_reason"] = "judge_failed" metadata["current_failure_reason"] = "judge_failed"
metadata["current_failure_findings"] = verification.evidence metadata["current_failure_findings"] = verification.evidence
metadata["current_failure_node_id"] = "judge" 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: else:
self._clear_current_failure(metadata) self._clear_current_failure(metadata)
self.save_metadata(context, 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}) return NodeResult("COMPLETE", "retry_available", {"retry_count": task.retry_count, "classification": classification})
task.status = TaskStatus.FAILED task.status = TaskStatus.FAILED
task.save(update_fields=["retry_count", "status", "updated_at"]) 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("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}) 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}) return NodeResult("COMPLETE", "retry_exhausted", {"retry_count": task.retry_count, "classification": classification})

View file

@ -1,5 +1,6 @@
from __future__ import annotations from __future__ import annotations
from control_plane.agents.models import ProgenySignal
from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus, GraphEdgeTraversal, GraphRun, GraphRunStatus from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus, GraphEdgeTraversal, GraphRun, GraphRunStatus
from graph.native_runtime import NativeGraphRuntime from graph.native_runtime import NativeGraphRuntime
from graph.registry import NodeHandlerRegistry, NodeResult 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.status == GraphRunStatus.PAUSED
assert result.failure_reason == "AWAITING_APPROVAL" 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"

View file

@ -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

View file

@ -4,6 +4,7 @@ from pathlib import Path
from agents.providers import DeterministicCodingProvider from agents.providers import DeterministicCodingProvider
from control_plane.agents.management.commands.seed_core_agents import Command as SeedAgentsCommand 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.events.models import Event, EventType
from control_plane.projects.models import CommitRecord, Project, ProjectPlan, Milestone, Task, TaskStatus from control_plane.projects.models import CommitRecord, Project, ProjectPlan, Milestone, Task, TaskStatus
from control_plane.verification.models import Review, TestRun, Verification, VerificationResult from control_plane.verification.models import Review, TestRun, Verification, VerificationResult
@ -69,6 +70,20 @@ class FirstAttemptBadProvider(DeterministicCodingProvider):
return super().complete(request) 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: def run_graph(task: Task, *, interrupt_after: str | None = None, provider=None) -> GraphRun:
SeedAgentsCommand().handle() SeedAgentsCommand().handle()
services = TaskExecutionServices(ModelRouter({"qwen": provider or DeterministicCodingProvider()}), test_command=["python", "manage.py", "test"]) 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 Event.objects.filter(task=task, event_type="TASK_RETRY_EXHAUSTED").exists()
assert graph_run.metadata["final_failure_reason"] == "review_failed" assert graph_run.metadata["final_failure_reason"] == "review_failed"
assert graph_run.metadata["historical_failures"][-1]["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: 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 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: def test_langgraph_task_execution_graph_success_matches_native_domain_outcome(tmp_path: Path) -> None:
repo = create_disposable_django_repo(tmp_path) 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"]) 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 Event.objects.filter(task=task, event_type="TASK_RETRY_EXHAUSTED").exists()
assert graph_run.metadata["final_failure_reason"] == "review_failed" assert graph_run.metadata["final_failure_reason"] == "review_failed"
assert graph_run.metadata["historical_failures"][-1]["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: 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"}], "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: def test_langgraph_resume_after_coder_tests_and_judge_does_not_duplicate_commit(tmp_path: Path) -> None:

View file

@ -7,6 +7,8 @@ from agents.progeny import ProgenyService
from control_plane.agents.models import Agent, AgentVersion, PromotionStatus from control_plane.agents.models import Agent, AgentVersion, PromotionStatus
from control_plane.events.models import Event from control_plane.events.models import Event
from control_plane.projects.models import Project, ProjectPlan, Milestone, Task 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): class ProgenySignalTests(TestCase):
@ -38,7 +40,7 @@ class ProgenySignalTests(TestCase):
self.assertEqual(signal.source, "reviewer") self.assertEqual(signal.source, "reviewer")
self.assertEqual(signal.failure_category, "REWORK_REQUIRED") self.assertEqual(signal.failure_category, "REWORK_REQUIRED")
self.assertEqual(signal.severity, "high") 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") event = Event.objects.get(event_type="PROGENY_SIGNAL_CREATED")
self.assertEqual(event.payload["signal_id"], str(signal.id)) self.assertEqual(event.payload["signal_id"], str(signal.id))
@ -89,7 +91,7 @@ class ProgenySignalTests(TestCase):
raw_output="{invalid json}", raw_output="{invalid json}",
) )
self.assertEqual(signal.source, "model_output") 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.assertEqual(signal.severity, "high")
self.assertIn("Invalid JSON", signal.summary) self.assertIn("Invalid JSON", signal.summary)
@ -105,7 +107,7 @@ class ProgenySignalTests(TestCase):
attempts=3, attempts=3,
) )
self.assertEqual(signal.source, "retry") 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.assertEqual(signal.severity, "critical")
self.assertIn("exhausted 3 retries", signal.summary) self.assertIn("exhausted 3 retries", signal.summary)
@ -133,6 +135,96 @@ class ProgenySignalTests(TestCase):
) )
self.assertEqual(signal1.grouping_key, signal2.grouping_key) 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): def test_event_lineage(self):
self.service.create_reviewer_signal( self.service.create_reviewer_signal(
project=self.project, project=self.project,