diff --git a/agents/lifecycle.py b/agents/lifecycle.py index 4c59f09..d92c142 100644 --- a/agents/lifecycle.py +++ b/agents/lifecycle.py @@ -112,6 +112,7 @@ class ExtensionService(ProjectContextMixin): evidence: dict[str, object] | None = None, source_steward_finding: StewardFinding | None = None, source_opportunity: ExplorationOpportunity | None = None, + source_roadmap_item: RoadmapItem | None = None, ) -> ExtensionCandidate: return ExtensionCandidate.objects.create( project=project, @@ -127,6 +128,7 @@ class ExtensionService(ProjectContextMixin): evidence=evidence or {}, source_steward_finding=source_steward_finding, source_opportunity=source_opportunity, + source_roadmap_item=source_roadmap_item, ) def plan_with_project_brain(self, candidate: ExtensionCandidate) -> ExtensionPlan: @@ -286,10 +288,10 @@ class EvolutionService(ProjectContextMixin): self.router = router self.bus = bus or EventBus() - def create_candidate(self, project: Project, *, target: str, objective: str, baseline_measurement: dict[str, object], desired_direction: str, target_measurement: dict[str, object] | None = None, rationale: str = "", source: str = "user", evidence: dict[str, object] | None = None, risk: str = "MEDIUM", confidence: float = 0.5, source_steward_finding: StewardFinding | None = None, source_opportunity: ExplorationOpportunity | None = None) -> EvolutionCandidate: + def create_candidate(self, project: Project, *, target: str, objective: str, baseline_measurement: dict[str, object], desired_direction: str, target_measurement: dict[str, object] | None = None, rationale: str = "", source: str = "user", evidence: dict[str, object] | None = None, risk: str = "MEDIUM", confidence: float = 0.5, source_steward_finding: StewardFinding | None = None, source_opportunity: ExplorationOpportunity | None = None, source_roadmap_item: RoadmapItem | None = None) -> EvolutionCandidate: if not baseline_measurement: raise ValidationError("EVOLVE requires a measurable baseline; route to INVESTIGATE or EXPLORE instead.") - return EvolutionCandidate.objects.create(project=project, target=target, objective=objective, baseline_measurement=baseline_measurement, desired_direction=desired_direction, target_measurement=target_measurement or {}, rationale=rationale, source=source, evidence=evidence or {}, risk=risk, confidence=confidence, source_steward_finding=source_steward_finding, source_opportunity=source_opportunity) + return EvolutionCandidate.objects.create(project=project, target=target, objective=objective, baseline_measurement=baseline_measurement, desired_direction=desired_direction, target_measurement=target_measurement or {}, rationale=rationale, source=source, evidence=evidence or {}, risk=risk, confidence=confidence, source_steward_finding=source_steward_finding, source_opportunity=source_opportunity, source_roadmap_item=source_roadmap_item) def plan_with_project_brain(self, candidate: EvolutionCandidate) -> EvolutionPlan: if not candidate.baseline_measurement: diff --git a/agents/roadmap.py b/agents/roadmap.py new file mode 100644 index 0000000..bae71c0 --- /dev/null +++ b/agents/roadmap.py @@ -0,0 +1,168 @@ +from __future__ import annotations + +import hashlib +import json +from typing import Any + +from control_plane.agents.models import ProgenySignal +from control_plane.events.bus import EventBus +from control_plane.projects.models import ( + Decision, + EvolutionCandidate, + ExtensionCandidate, + ExplorationOpportunity, + Project, + RoadmapHorizon, + RoadmapItem, + RoadmapStatus, + RoadmapTargetAction, + StewardFinding, +) +from agents.lifecycle import EvolutionService, ExtensionService, ProjectContextMixin +from agents.progeny import ProgenyService +from model_router.router import ModelCapability, ModelRequestContract, ModelRouter + + +class RoadmapService(ProjectContextMixin): + def __init__(self, router: ModelRouter | None = None, bus: EventBus | None = None) -> None: + self.router = router + self.bus = bus or EventBus() + + def upsert_item(self, project: Project, *, title: str, description: str = "", source: str = "USER", source_ref: dict[str, object] | None = None, rationale: str = "", evidence: dict[str, object] | None = None, horizon: str = RoadmapHorizon.EXPLORING, category: str = "", target_action: str = RoadmapTargetAction.NONE, scores: dict[str, object] | None = None, status: str = RoadmapStatus.PROPOSED) -> RoadmapItem: + grouping_key = self._grouping_key(project, title, category or target_action) + existing = self._find_existing(project, grouping_key, title) + values = self.score(scores or {}) + if existing: + metadata = dict(existing.metadata) + metadata["occurrences"] = int(metadata.get("occurrences", 1)) + 1 + metadata.setdefault("reinforced_by", []).append({"source": source, "source_ref": source_ref or {}}) + existing.evidence = self._merge_evidence(existing.evidence, evidence or {}) + existing.source_ref = self._merge_evidence(existing.source_ref, source_ref or {}) + existing.confidence = min(1.0, max(existing.confidence, values["confidence"]) + 0.05) + existing.composite_score = self._composite(existing) + existing.metadata = metadata + existing.save(update_fields=["evidence", "source_ref", "confidence", "composite_score", "metadata", "updated_at"]) + self.bus.publish("ROADMAP_ITEM_UPDATED", project=project, payload={"roadmap_item_id": str(existing.id), "reason": "deduplicated_reinforcement"}) + return existing + item = RoadmapItem.objects.create(project=project, title=title, description=description, source=source, source_ref=source_ref or {}, rationale=rationale, evidence=evidence or {}, horizon=horizon, category=category, status=status, target_action=target_action, grouping_key=grouping_key, value_score=values["value"], effort_score=values["effort"], risk_score=values["risk"], confidence=values["confidence"], strategic_fit=values["strategic_fit"], technical_fit=values["technical_fit"], urgency=values["urgency"]) + item.composite_score = self._composite(item) + item.priority = max(1, min(100, int(item.composite_score * 100))) + item.save(update_fields=["composite_score", "priority", "updated_at"]) + self.bus.publish("ROADMAP_ITEM_CREATED", project=project, payload={"roadmap_item_id": str(item.id), "source": source}) + return item + + def gather_candidate_items(self, project: Project) -> list[RoadmapItem]: + items: list[RoadmapItem] = [] + for opportunity in project.exploration_opportunities.exclude(status="REJECTED"): + items.append(self.upsert_item(project, title=opportunity.title, description=opportunity.description, source="EXPLORE", source_ref={"exploration_opportunity_id": str(opportunity.id)}, rationale=opportunity.rationale, evidence=opportunity.evidence, horizon=RoadmapHorizon.EXPLORING, category=opportunity.opportunity_type, target_action=opportunity.recommended_action if opportunity.recommended_action in RoadmapTargetAction.values else RoadmapTargetAction.NONE, scores={"value": opportunity.value_score, "effort": opportunity.effort_score, "risk": opportunity.risk_score, "confidence": opportunity.confidence, "strategic_fit": opportunity.strategic_fit, "technical_fit": opportunity.technical_fit})) + for finding in project.steward_findings.exclude(status__in=["RESOLVED", "DISMISSED"]): + action = finding.recommended_action if finding.recommended_action in RoadmapTargetAction.values else RoadmapTargetAction.INVESTIGATE + items.append(self.upsert_item(project, title=finding.title, description=finding.summary, source="STEWARD", source_ref={"steward_finding_id": str(finding.id)}, rationale="Steward surfaced this future project intent.", evidence=finding.evidence, horizon=RoadmapHorizon.EXPLORING, category=finding.finding_type, target_action=action, scores={"confidence": finding.confidence, "risk": 0.7 if finding.severity in ["HIGH", "CRITICAL"] else 0.4, "urgency": 0.8 if finding.severity in ["HIGH", "CRITICAL"] else 0.4})) + return items + + def review_with_project_brain(self, project: Project) -> dict[str, object]: + fallback = {"recommendations": []} + if self.router is None: + return fallback + payload = {"context": self.project_context(project), "roadmap_items": list(project.roadmap_items.values("id", "title", "description", "horizon", "status", "target_action", "value_score", "effort_score", "risk_score", "confidence", "strategic_fit", "technical_fit", "urgency", "composite_score"))} + try: + response = self.router.complete(ModelRequestContract(purpose=ModelCapability.PLANNING, model_hint="sol", project=project, prompt="Review this existing project roadmap. Return JSON with recommendations: item_id, recommendation, rationale, optional horizon/status. Do not execute work.\n" + json.dumps(payload, default=str))) + parsed = json.loads(response.content) + return parsed if isinstance(parsed, dict) else fallback + except Exception as exc: + ProgenySignal.objects.create(project=project, source="roadmap", severity="MEDIUM", failure_category="ROADMAP_PRIORITIZATION", summary="Roadmap Project Brain review failed; deterministic scoring retained.", evidence={"error": str(exc)}, grouping_key=f"roadmap:{project.id}:prioritization") + return fallback + + def apply_recommendations(self, project: Project, recommendations: dict[str, object]) -> list[RoadmapItem]: + updated: list[RoadmapItem] = [] + for rec in recommendations.get("recommendations", []): + if not isinstance(rec, dict): + continue + item_id = rec.get("item_id") + try: + item = project.roadmap_items.get(id=item_id) + except Exception: + continue + item.metadata = {**item.metadata, "project_brain_recommendations": [*item.metadata.get("project_brain_recommendations", []), rec]} + if rec.get("horizon") in RoadmapHorizon.values: + item.horizon = str(rec["horizon"]) + if rec.get("status") in RoadmapStatus.values: + item.status = str(rec["status"]) + item.save(update_fields=["horizon", "status", "metadata", "updated_at"]) + updated.append(item) + self.bus.publish("ROADMAP_ITEM_UPDATED", project=project, payload={"roadmap_item_id": str(item.id), "recommendation": rec}) + return updated + + def review_project_roadmap(self, project: Project) -> dict[str, object]: + self.gather_candidate_items(project) + recommendations = self.review_with_project_brain(project) + self.apply_recommendations(project, recommendations) + self.bus.publish("ROADMAP_REVIEW_COMPLETED", project=project, payload={"item_count": project.roadmap_items.count(), "recommendations": recommendations}) + return self.project_roadmap_view(project) + + def convert_to_extension(self, item: RoadmapItem) -> ExtensionCandidate: + candidate = ExtensionService(bus=self.bus).create_candidate(item.project, title=item.title, description=item.description, rationale=item.rationale, source="RoadmapItem", expected_value=str(item.evidence.get("expected_value", "")), affected_areas=[item.category] if item.category else [], risk=str(item.risk_score), confidence=item.confidence, evidence={"roadmap_item_id": str(item.id), **item.evidence}, source_roadmap_item=item) + item.converted_extension = candidate + item.status = RoadmapStatus.PLANNING + item.save(update_fields=["converted_extension", "status", "updated_at"]) + self.bus.publish("ROADMAP_ITEM_CONVERTED", project=item.project, payload={"roadmap_item_id": str(item.id), "extension_candidate_id": str(candidate.id)}) + return candidate + + def convert_to_evolution(self, item: RoadmapItem, *, baseline_measurement: dict[str, object], desired_direction: str = "DECREASE") -> EvolutionCandidate: + candidate = EvolutionService(bus=self.bus).create_candidate(item.project, target=item.category or item.title, objective=item.description or item.title, baseline_measurement=baseline_measurement, desired_direction=desired_direction, rationale=item.rationale, source="RoadmapItem", evidence={"roadmap_item_id": str(item.id), **item.evidence}, risk=str(item.risk_score), confidence=item.confidence, source_roadmap_item=item) + item.converted_evolution = candidate + item.status = RoadmapStatus.PLANNING + item.save(update_fields=["converted_evolution", "status", "updated_at"]) + self.bus.publish("ROADMAP_ITEM_CONVERTED", project=item.project, payload={"roadmap_item_id": str(item.id), "evolution_candidate_id": str(candidate.id)}) + return candidate + + def convert_to_investigation(self, item: RoadmapItem): + signal = ProgenySignal.objects.create(project=item.project, source="roadmap", severity="MEDIUM", failure_category=item.category or "ROADMAP_INVESTIGATION", summary=item.description or item.title, evidence={"roadmap_item_id": str(item.id), **item.evidence}, grouping_key=f"roadmap:{item.grouping_key}"[:120]) + investigation = ProgenyService(self.bus).create_smart_investigation(signal.grouping_key) + item.converted_investigation = investigation + item.status = RoadmapStatus.PLANNING + item.save(update_fields=["converted_investigation", "status", "updated_at"]) + self.bus.publish("ROADMAP_ITEM_CONVERTED", project=item.project, payload={"roadmap_item_id": str(item.id), "investigation_id": str(investigation.id)}) + return investigation + + def project_roadmap_view(self, project: Project) -> dict[str, object]: + return {horizon: [self._item(item) for item in project.roadmap_items.filter(horizon=horizon).order_by("-composite_score", "-priority", "created_at")] for horizon in RoadmapHorizon.values} + + def score(self, raw: dict[str, object]) -> dict[str, float]: + return {key: self._score_value(raw.get(key, 0.5)) for key in ["value", "effort", "risk", "confidence", "strategic_fit", "technical_fit", "urgency"]} + + def _score_value(self, value: object) -> float: + try: + score = float(value) + except (TypeError, ValueError): + return 0.5 + return max(0.0, min(1.0, score)) + + def _composite(self, item: RoadmapItem) -> float: + return (item.value_score * 0.25) + ((1 - item.effort_score) * 0.12) + ((1 - item.risk_score) * 0.12) + (item.confidence * 0.14) + (item.strategic_fit * 0.14) + (item.technical_fit * 0.11) + (item.urgency * 0.12) + + def _grouping_key(self, project: Project, title: str, category: str) -> str: + fingerprint = hashlib.sha256(f"{title.lower()}:{category.lower()}".encode("utf-8")).hexdigest()[:16] + return f"{project.id}:roadmap:{fingerprint}"[:240] + + def _find_existing(self, project: Project, grouping_key: str, title: str) -> RoadmapItem | None: + existing = project.roadmap_items.filter(grouping_key=grouping_key).first() or project.roadmap_items.filter(title__iexact=title).first() + if existing: + return existing + if project.exploration_opportunities.filter(title__iexact=title).exists() or ExtensionCandidate.objects.filter(project=project, title__iexact=title).exists() or EvolutionCandidate.objects.filter(project=project, objective__icontains=title[:80]).exists() or StewardFinding.objects.filter(project=project, title__iexact=title).exists(): + return project.roadmap_items.filter(title__iexact=title).first() + if Decision.objects.filter(project=project, decision__icontains=title[:80], decision_type__in=["REJECTED", "DEFERRED"]).exists(): + return project.roadmap_items.filter(title__iexact=title).first() + return None + + def _merge_evidence(self, current: dict[str, object], incoming: dict[str, object]) -> dict[str, object]: + merged = dict(current or {}) + for key, value in incoming.items(): + if key in merged and merged[key] != value: + merged[key] = [merged[key], value] + else: + merged[key] = value + return merged + + def _item(self, item: RoadmapItem) -> dict[str, object]: + return {"id": str(item.id), "title": item.title, "source": item.source, "rationale": item.rationale, "evidence": item.evidence, "scores": {"value": item.value_score, "effort": item.effort_score, "risk": item.risk_score, "confidence": item.confidence, "strategic_fit": item.strategic_fit, "technical_fit": item.technical_fit, "urgency": item.urgency, "composite": item.composite_score}, "status": item.status, "target_action": item.target_action, "dependencies": [str(dep.id) for dep in item.dependencies.all()], "related_items": [str(rel.id) for rel in item.related_items.all()], "conversion_lineage": {"extension_candidate_id": str(item.converted_extension_id) if item.converted_extension_id else None, "evolution_candidate_id": str(item.converted_evolution_id) if item.converted_evolution_id else None, "investigation_id": str(item.converted_investigation_id) if item.converted_investigation_id else None}, "metadata": item.metadata} diff --git a/agents/scenario_lab.py b/agents/scenario_lab.py new file mode 100644 index 0000000..8b0884a --- /dev/null +++ b/agents/scenario_lab.py @@ -0,0 +1,226 @@ +from __future__ import annotations + +import hashlib +import json +import tempfile +from collections import Counter +from pathlib import Path + +from django.utils import timezone + +from agents.lifecycle import ProjectContextMixin +from agents.progeny import ProgenyService +from agents.roadmap import RoadmapService +from control_plane.agents.models import ProgenySignal +from control_plane.events.bus import EventBus +from control_plane.projects.models import Project, Scenario, ScenarioFinding, ScenarioRun, ScenarioSuite, StewardFinding +from model_router.router import ModelCapability, ModelRequestContract, ModelRouter + + +SCENARIO_TYPES = { + "FUNCTIONAL_EDGE_CASE", + "FAILURE_INJECTION", + "DEPENDENCY_FAILURE", + "SECURITY_ADVERSARIAL", + "PERMISSION", + "CONCURRENCY", + "PERFORMANCE", + "LOAD", + "DATA_INTEGRITY", + "RECOVERY", + "USER_BEHAVIOR", + "WORKFLOW", + "AGENT_WORKFLOW", +} + + +class ScenarioValidationError(ValueError): + pass + + +class ScenarioLabService(ProjectContextMixin): + def __init__(self, router: ModelRouter | None = None, bus: EventBus | None = None) -> None: + self.router = router + self.bus = bus or EventBus() + + def create_suite(self, project: Project, *, name: str, purpose: str = "", scenarios: list[dict[str, object]] | None = None) -> ScenarioSuite: + version = (project.scenario_suites.order_by("-version").values_list("version", flat=True).first() or 0) + 1 + suite = ScenarioSuite.objects.create(project=project, name=name, version=version, purpose=purpose) + self.bus.publish("SCENARIO_SUITE_CREATED", project=project, payload={"suite_id": str(suite.id)}) + for raw in scenarios or []: + self.create_scenario(suite, raw) + return suite + + def create_scenario(self, suite: ScenarioSuite, raw: dict[str, object]) -> Scenario: + title = str(raw.get("title", raw.get("name", "Untitled scenario"))) + return Scenario.objects.create(project=suite.project, suite=suite, name=title, title=title, description=str(raw.get("description", "")), scenario_type=str(raw.get("scenario_type", "WORKFLOW")), target_component=str(raw.get("target_component", raw.get("target", "project"))), target_type=str(raw.get("target_type", "PROJECT")), target_id=str(raw.get("target_id", suite.project_id)), preconditions=list(raw.get("preconditions", [])), injected_condition=self._dict(raw.get("injected_condition", {})), expected_invariants=list(raw.get("expected_invariants", [])), success_criteria=list(raw.get("success_criteria", [])), severity=str(raw.get("severity", "MEDIUM")), source=str(raw.get("source", "USER")), definition=self._dict(raw.get("definition", {})), resource_budget=self._dict(raw.get("resource_budget", {"max_seconds": 5, "max_parallelism": 2})), metadata=self._dict(raw.get("metadata", {}))) + + def generate_scenarios(self, suite: ScenarioSuite, *, count: int = 5) -> list[Scenario]: + fallback = {"scenarios": self._fallback_scenarios(suite.project)[:count]} + payload = fallback + if self.router is not None: + try: + response = self.router.complete(ModelRequestContract(purpose=ModelCapability.PLANNING, model_hint="sol", project=suite.project, prompt="Design safe Scenario Lab candidates for this existing project. Return JSON with scenarios. Each scenario needs type, injected_condition, expected_invariants, success_criteria, and resource_budget. Do not create executable work.\n" + json.dumps(self.project_context(suite.project), default=str))) + parsed = json.loads(response.content) + if isinstance(parsed, dict) and isinstance(parsed.get("scenarios"), list): + payload = parsed + except Exception as exc: + ProgenySignal.objects.create(project=suite.project, source="scenario_lab", severity="MEDIUM", failure_category="SCENARIO_GENERATION", summary="Scenario generation failed; fallback scenarios retained.", evidence={"error": str(exc)}, grouping_key=f"scenario_lab:{suite.project_id}:generation") + return [self.create_scenario(suite, raw) for raw in payload.get("scenarios", []) if isinstance(raw, dict)] + + def validate_scenario(self, scenario: Scenario) -> bool: + reason = "" + budget = scenario.resource_budget or {} + condition = scenario.injected_condition or {} + if scenario.scenario_type not in SCENARIO_TYPES: + reason = "unsupported scenario type" + elif condition.get("destructive") is True: + reason = "destructive unsafe scenario" + elif scenario.scenario_type == "LOAD" and int(budget.get("max_parallelism", 1) or 1) > 8: + reason = "unbounded load test" + elif not condition: + reason = "missing injected condition" + elif not scenario.expected_invariants: + reason = "missing expected invariant" + elif not scenario.success_criteria: + reason = "missing observable result" + elif "max_seconds" not in budget: + reason = "missing resource budget" + duplicate = Scenario.objects.filter(project=scenario.project, scenario_type=scenario.scenario_type, title__iexact=scenario.title).exclude(id=scenario.id).first() + if duplicate: + reason = "duplicated scenario" + if reason: + scenario.status = "REJECTED" + scenario.rejection_reason = reason + scenario.save(update_fields=["status", "rejection_reason", "updated_at"]) + return False + scenario.status = "VALIDATED" + scenario.validated_at = timezone.now() + scenario.save(update_fields=["status", "validated_at", "updated_at"]) + return True + + def validate_suite(self, suite: ScenarioSuite) -> list[Scenario]: + return [scenario for scenario in suite.scenarios.all() if self.validate_scenario(scenario)] + + def freeze_suite(self, suite: ScenarioSuite) -> ScenarioSuite: + suite.status = "FROZEN" + suite.frozen_at = timezone.now() + suite.metadata = {**suite.metadata, "scenario_count": suite.scenarios.exclude(status="REJECTED").count()} + suite.save(update_fields=["status", "frozen_at", "metadata", "updated_at"]) + return suite + + def execute_suite(self, suite: ScenarioSuite, *, graph_run=None) -> list[ScenarioRun]: + runs = [] + for scenario in suite.scenarios.filter(status="VALIDATED"): + runs.append(self.execute_scenario(scenario, graph_run=graph_run)) + return runs + + def execute_scenario(self, scenario: Scenario, *, graph_run=None) -> ScenarioRun: + run = ScenarioRun.objects.create(scenario=scenario, project=scenario.project, repository_baseline=self._repository_baseline(scenario.project), graph_run=graph_run, status="RUNNING", started_at=timezone.now(), environment_metadata={"isolation": "tempdir", "canonical_repository_path": scenario.project.repository_path}) + self.bus.publish("SCENARIO_RUN_STARTED", project=scenario.project, payload={"scenario_run_id": str(run.id), "scenario_id": str(scenario.id)}) + with tempfile.TemporaryDirectory(prefix="artifex-scenario-") as tmp: + result = self._execute_mechanism(scenario, Path(tmp)) + run.status = "COMPLETE" + run.completed_at = timezone.now() + run.result = result["result"] + run.failure_evidence = result.get("failure_evidence", {}) + run.telemetry = result.get("telemetry", {}) + run.environment_metadata = {**run.environment_metadata, "workdir_removed": True} + run.save(update_fields=["status", "completed_at", "result", "failure_evidence", "telemetry", "environment_metadata", "updated_at"]) + event = "SCENARIO_FAILED" if run.result == "FAIL" else "SCENARIO_RUN_COMPLETED" + self.bus.publish(event, project=scenario.project, payload={"scenario_run_id": str(run.id), "result": run.result}) + if run.result == "FAIL": + self.create_finding(run) + return run + + def create_finding(self, run: ScenarioRun) -> ScenarioFinding: + scenario = run.scenario + category = str(scenario.injected_condition.get("failure_category", scenario.scenario_type)) + recommended_action = str(scenario.injected_condition.get("recommended_action", self._default_action(scenario))) + grouping_key = self._finding_grouping_key(scenario, category) + existing = ScenarioFinding.objects.filter(project=scenario.project, grouping_key=grouping_key, status__in=["OPEN", "ROUTED"]).first() + if existing: + existing.evidence = {**existing.evidence, "latest_run_id": str(run.id), "occurrences": int(existing.evidence.get("occurrences", 1)) + 1} + existing.save(update_fields=["evidence", "updated_at"]) + return existing + finding = ScenarioFinding.objects.create(project=scenario.project, scenario=scenario, scenario_run=run, title=f"Scenario failed: {scenario.title}", summary=str(run.failure_evidence.get("summary", scenario.description)), evidence={"scenario_run_id": str(run.id), "failure_evidence": run.failure_evidence, "occurrences": 1}, severity=scenario.severity, confidence=0.8, affected_component=scenario.target_component, failure_category=category, recommended_action=recommended_action, recommended_route=self._route_for(recommended_action), grouping_key=grouping_key, steward_policy_metadata={"monitoring_candidate": True, "scenario_type": scenario.scenario_type}) + self.bus.publish("SCENARIO_FINDING_CREATED", project=scenario.project, payload={"scenario_finding_id": str(finding.id), "recommended_action": recommended_action}) + return finding + + def route_finding(self, finding: ScenarioFinding): + action = finding.recommended_action + result = None + if action == "PROGENY": + signal = ProgenySignal.objects.create(project=finding.project, source="scenario_lab", severity=finding.severity, failure_category=finding.failure_category, summary=finding.summary, evidence={"scenario_finding_id": str(finding.id), **finding.evidence}, grouping_key=f"scenario_lab:{finding.grouping_key}"[:120]) + result = ProgenyService(self.bus).create_smart_investigation(signal.grouping_key) + finding.progeny_signal = signal + elif action in ["EXTEND", "EVOLVE", "INVESTIGATE", "NONE"]: + result = RoadmapService(bus=self.bus).upsert_item(finding.project, title=finding.title, description=finding.summary, source="SCENARIO_LAB", source_ref={"scenario_finding_id": str(finding.id)}, rationale="Scenario Lab found future project intent.", evidence=finding.evidence, horizon="NEXT", category=finding.failure_category, target_action=action if action in ["EXTEND", "EVOLVE", "INVESTIGATE"] else "NONE", scores={"confidence": finding.confidence, "risk": 0.7 if finding.severity in ["HIGH", "CRITICAL"] else 0.4, "urgency": 0.6}) + finding.roadmap_item = result + elif action == "REPAIR": + result = StewardFinding.objects.create(project=finding.project, finding_type=finding.failure_category, title=finding.title, summary=finding.summary, evidence={"scenario_finding_id": str(finding.id), **finding.evidence}, severity=finding.severity, confidence=finding.confidence, recommended_action="REPAIR", recommended_route="StewardRepair", grouping_key=f"scenario:{finding.grouping_key}"[:240]) + else: + result = RoadmapService(bus=self.bus).upsert_item(finding.project, title=finding.title, description=finding.summary, source="SCENARIO_LAB", source_ref={"scenario_finding_id": str(finding.id)}, evidence=finding.evidence, horizon="EXPLORING", category=finding.failure_category) + finding.roadmap_item = result + finding.status = "ROUTED" + finding.save(update_fields=["status", "roadmap_item", "progeny_signal", "updated_at"]) + self.bus.publish("SCENARIO_FINDING_ROUTED", project=finding.project, payload={"scenario_finding_id": str(finding.id), "recommended_action": action}) + return result + + def route_findings(self, suite: ScenarioSuite) -> list[object]: + routed = [] + for finding in ScenarioFinding.objects.filter(project=suite.project, scenario__suite=suite, status="OPEN"): + routed.append(self.route_finding(finding)) + return routed + + def coverage(self, project: Project) -> dict[str, int]: + return dict(Counter(project.scenarios.exclude(status="REJECTED").values_list("scenario_type", flat=True))) + + def summarize_suite(self, suite: ScenarioSuite) -> dict[str, object]: + runs = ScenarioRun.objects.filter(scenario__suite=suite) + return {"suite_id": str(suite.id), "status": suite.status, "coverage": self.coverage(suite.project), "results": dict(Counter(runs.values_list("result", flat=True))), "findings": list(ScenarioFinding.objects.filter(scenario__suite=suite).values("title", "recommended_action", "status", "severity", "failure_category"))} + + def _execute_mechanism(self, scenario: Scenario, workdir: Path) -> dict[str, object]: + condition = scenario.injected_condition or {} + mechanism = str(condition.get("mechanism", scenario.scenario_type)).lower() + expected = str(condition.get("expected_result", "PASS")) + if condition.get("infrastructure_failure"): + return {"result": "INFRASTRUCTURE_FAILURE", "failure_evidence": {"summary": "Scenario fixture infrastructure failed", "condition": condition}, "telemetry": {"workdir": str(workdir)}} + if mechanism not in ["test_mutation", "malformed_input", "permission_denial", "concurrency", "performance_regression", "provider_failure_replay", "workflow"]: + return {"result": "INCONCLUSIVE", "failure_evidence": {"summary": "Unsupported deterministic scenario mechanism", "mechanism": mechanism}, "telemetry": {"workdir": str(workdir)}} + if expected == "FAIL": + return {"result": "FAIL", "failure_evidence": {"summary": str(condition.get("summary", "Expected invariant failed under scenario")), "condition": condition, "invariants": scenario.expected_invariants}, "telemetry": {"mechanism": mechanism, "workdir": str(workdir)}} + if expected == "INCONCLUSIVE": + return {"result": "INCONCLUSIVE", "failure_evidence": {"summary": "Scenario did not produce observable result", "condition": condition}, "telemetry": {"mechanism": mechanism, "workdir": str(workdir)}} + return {"result": "PASS", "failure_evidence": {}, "telemetry": {"mechanism": mechanism, "workdir": str(workdir)}} + + def _fallback_scenarios(self, project: Project) -> list[dict[str, object]]: + return [ + {"title": "Malformed coder structured output", "description": "Coder returns malformed JSON and orchestration should classify rather than crash.", "scenario_type": "AGENT_WORKFLOW", "target_component": "coder", "injected_condition": {"mechanism": "malformed_input", "expected_result": "FAIL", "recommended_action": "PROGENY", "failure_category": "AGENT_WORKFLOW"}, "expected_invariants": ["Scenario Lab records finding"], "success_criteria": ["Failure is classified"], "resource_budget": {"max_seconds": 5, "max_parallelism": 1}}, + {"title": "Graph node failure recovery", "description": "Graph node reports failure evidence without crashing the lab.", "scenario_type": "RECOVERY", "target_component": "graph_runtime", "injected_condition": {"mechanism": "test_mutation", "expected_result": "PASS"}, "expected_invariants": ["Graph lineage persists"], "success_criteria": ["Run completes"], "resource_budget": {"max_seconds": 5, "max_parallelism": 1}}, + {"title": "Concurrent agent version allocation", "description": "Parallel version allocation can race.", "scenario_type": "CONCURRENCY", "target_component": "agents", "injected_condition": {"mechanism": "concurrency", "expected_result": "FAIL", "recommended_action": "EVOLVE", "failure_category": "CONCURRENCY"}, "expected_invariants": ["Uniqueness is preserved"], "success_criteria": ["Race is detected"], "resource_budget": {"max_seconds": 5, "max_parallelism": 4}}, + {"title": "Repository symlink path escape", "description": "Repository scanner must not follow symlinks outside project root.", "scenario_type": "SECURITY_ADVERSARIAL", "target_component": "repository_scanner", "injected_condition": {"mechanism": "permission_denial", "expected_result": "FAIL", "recommended_action": "REPAIR", "failure_category": "SECURITY"}, "expected_invariants": ["No path escapes root"], "success_criteria": ["Escape is blocked"], "resource_budget": {"max_seconds": 5, "max_parallelism": 1}}, + {"title": "Deterministic performance regression", "description": "Repeated project inspection exceeds threshold.", "scenario_type": "PERFORMANCE", "target_component": "project_context", "injected_condition": {"mechanism": "performance_regression", "expected_result": "FAIL", "recommended_action": "EVOLVE", "failure_category": "PERFORMANCE"}, "expected_invariants": ["Latency remains bounded"], "success_criteria": ["Regression is measured"], "resource_budget": {"max_seconds": 5, "max_parallelism": 1}}, + ] + + def _default_action(self, scenario: Scenario) -> str: + if scenario.scenario_type == "AGENT_WORKFLOW": + return "PROGENY" + if scenario.scenario_type in ["PERFORMANCE", "CONCURRENCY"]: + return "EVOLVE" + if scenario.scenario_type in ["SECURITY_ADVERSARIAL", "DATA_INTEGRITY", "RECOVERY"]: + return "REPAIR" + return "EXTEND" + + def _route_for(self, action: str) -> str: + return {"REPAIR": "StewardRepair", "EVOLVE": "RoadmapItem", "EXTEND": "RoadmapItem", "PROGENY": "ProgenySignal", "INVESTIGATE": "RoadmapItem"}.get(action, "RoadmapItem") + + def _repository_baseline(self, project: Project) -> str: + return project.repository_path or "untracked" + + def _finding_grouping_key(self, scenario: Scenario, category: str) -> str: + fingerprint = hashlib.sha256(f"{scenario.project_id}:{scenario.title.lower()}:{category}".encode("utf-8")).hexdigest()[:16] + return f"scenario:{scenario.project_id}:{fingerprint}"[:240] + + def _dict(self, value: object) -> dict[str, object]: + return value if isinstance(value, dict) else {"raw": value} diff --git a/control_plane/projects/admin.py b/control_plane/projects/admin.py index 18982c8..15e7188 100644 --- a/control_plane/projects/admin.py +++ b/control_plane/projects/admin.py @@ -14,6 +14,9 @@ from control_plane.projects.models import ( ProjectPlan, RoadmapItem, Scenario, + ScenarioFinding, + ScenarioRun, + ScenarioSuite, Task, TaskAttempt, TaskDependency, @@ -48,3 +51,6 @@ admin.site.register(Artifact) admin.site.register(RoadmapItem) admin.site.register(Finding) admin.site.register(Scenario) +admin.site.register(ScenarioSuite) +admin.site.register(ScenarioRun) +admin.site.register(ScenarioFinding) diff --git a/control_plane/projects/migrations/0006_roadmap_scenario_lab_v1.py b/control_plane/projects/migrations/0006_roadmap_scenario_lab_v1.py new file mode 100644 index 0000000..4dfbdbb --- /dev/null +++ b/control_plane/projects/migrations/0006_roadmap_scenario_lab_v1.py @@ -0,0 +1,116 @@ +import uuid + +import django.db.models.deletion +from django.db import migrations, models + + +class Migration(migrations.Migration): + dependencies = [ + ("agents", "0007_replay_arena"), + ("graph", "0004_unique_champion_graph_version"), + ("projects", "0005_extend_evolve_explore_v1"), + ] + + operations = [ + migrations.CreateModel( + name="ScenarioSuite", + fields=[ + ("id", models.UUIDField(default=uuid.uuid4, editable=False, primary_key=True, serialize=False)), + ("created_at", models.DateTimeField(auto_now_add=True)), + ("updated_at", models.DateTimeField(auto_now=True)), + ("name", models.CharField(max_length=200)), + ("version", models.PositiveIntegerField(default=1)), + ("purpose", models.TextField(blank=True)), + ("status", models.CharField(default="DRAFT", max_length=32)), + ("frozen_at", models.DateTimeField(blank=True, null=True)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("execution_graph_version", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="scenario_suites", to="graph.executiongraphversion")), + ("project", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="scenario_suites", to="projects.project")), + ], + ), + migrations.AlterField(model_name="roadmapitem", name="source", field=models.CharField(default="USER", max_length=80)), + migrations.AlterField(model_name="roadmapitem", name="status", field=models.CharField(choices=[("PROPOSED", "Proposed"), ("ACCEPTED", "Accepted"), ("PLANNING", "Planning"), ("IN_PROGRESS", "In Progress"), ("COMPLETE", "Complete"), ("DEFERRED", "Deferred"), ("REJECTED", "Rejected"), ("SUPERSEDED", "Superseded")], default="PROPOSED", max_length=32)), + migrations.AddField(model_name="roadmapitem", name="category", field=models.CharField(blank=True, max_length=80)), + migrations.AddField(model_name="roadmapitem", name="composite_score", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="confidence", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="effort_score", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="evidence", field=models.JSONField(blank=True, default=dict)), + migrations.AddField(model_name="roadmapitem", name="grouping_key", field=models.CharField(blank=True, max_length=240)), + migrations.AddField(model_name="roadmapitem", name="horizon", field=models.CharField(choices=[("NOW", "Now"), ("NEXT", "Next"), ("LATER", "Later"), ("EXPLORING", "Exploring")], default="EXPLORING", max_length=32)), + migrations.AddField(model_name="roadmapitem", name="metadata", field=models.JSONField(blank=True, default=dict)), + migrations.AddField(model_name="roadmapitem", name="rationale", field=models.TextField(blank=True)), + migrations.AddField(model_name="roadmapitem", name="risk_score", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="source_ref", field=models.JSONField(blank=True, default=dict)), + migrations.AddField(model_name="roadmapitem", name="strategic_fit", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="target_action", field=models.CharField(choices=[("EXTEND", "Extend"), ("EVOLVE", "Evolve"), ("REPAIR", "Repair"), ("INVESTIGATE", "Investigate"), ("NONE", "None")], default="NONE", max_length=32)), + migrations.AddField(model_name="roadmapitem", name="technical_fit", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="urgency", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="value_score", field=models.FloatField(default=0.5)), + migrations.AddField(model_name="roadmapitem", name="converted_extension", field=models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="converted_from_roadmap_items", to="projects.extensioncandidate")), + migrations.AddField(model_name="roadmapitem", name="converted_evolution", field=models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="converted_from_roadmap_items", to="projects.evolutioncandidate")), + migrations.AddField(model_name="roadmapitem", name="converted_investigation", field=models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="converted_from_roadmap_items", to="agents.progenyinvestigation")), + migrations.AddField(model_name="roadmapitem", name="dependencies", field=models.ManyToManyField(blank=True, related_name="dependent_roadmap_items", to="projects.roadmapitem")), + migrations.AddField(model_name="roadmapitem", name="related_items", field=models.ManyToManyField(blank=True, to="projects.roadmapitem")), + migrations.AddField(model_name="evolutioncandidate", name="source_roadmap_item", field=models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="evolution_candidates", to="projects.roadmapitem")), + migrations.AddField(model_name="scenario", name="description", field=models.TextField(blank=True)), + migrations.AddField(model_name="scenario", name="expected_invariants", field=models.JSONField(blank=True, default=list)), + migrations.AddField(model_name="scenario", name="injected_condition", field=models.JSONField(blank=True, default=dict)), + migrations.AddField(model_name="scenario", name="metadata", field=models.JSONField(blank=True, default=dict)), + migrations.AddField(model_name="scenario", name="preconditions", field=models.JSONField(blank=True, default=list)), + migrations.AddField(model_name="scenario", name="rejection_reason", field=models.TextField(blank=True)), + migrations.AddField(model_name="scenario", name="resource_budget", field=models.JSONField(blank=True, default=dict)), + migrations.AddField(model_name="scenario", name="scenario_type", field=models.CharField(default="WORKFLOW", max_length=80)), + migrations.AddField(model_name="scenario", name="severity", field=models.CharField(default="MEDIUM", max_length=32)), + migrations.AddField(model_name="scenario", name="source", field=models.CharField(default="USER", max_length=80)), + migrations.AddField(model_name="scenario", name="success_criteria", field=models.JSONField(blank=True, default=list)), + migrations.AddField(model_name="scenario", name="target_component", field=models.CharField(blank=True, max_length=240)), + migrations.AddField(model_name="scenario", name="title", field=models.CharField(blank=True, max_length=255)), + migrations.AddField(model_name="scenario", name="validated_at", field=models.DateTimeField(blank=True, null=True)), + migrations.AddField(model_name="scenario", name="suite", field=models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.CASCADE, related_name="scenarios", to="projects.scenariosuite")), + migrations.CreateModel( + name="ScenarioRun", + 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)), + ("repository_baseline", models.CharField(blank=True, max_length=255)), + ("status", models.CharField(default="PENDING", max_length=32)), + ("started_at", models.DateTimeField(blank=True, null=True)), + ("completed_at", models.DateTimeField(blank=True, null=True)), + ("environment_metadata", models.JSONField(blank=True, default=dict)), + ("result", models.CharField(blank=True, max_length=32)), + ("failure_evidence", models.JSONField(blank=True, default=dict)), + ("telemetry", models.JSONField(blank=True, default=dict)), + ("execution_graph_version", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="scenario_runs", to="graph.executiongraphversion")), + ("graph_run", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="scenario_runs", to="graph.graphrun")), + ("project", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="scenario_runs", to="projects.project")), + ("scenario", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="runs", to="projects.scenario")), + ], + ), + migrations.CreateModel( + name="ScenarioFinding", + 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)), + ("title", models.CharField(max_length=255)), + ("summary", models.TextField(blank=True)), + ("evidence", models.JSONField(blank=True, default=dict)), + ("severity", models.CharField(default="MEDIUM", max_length=32)), + ("confidence", models.FloatField(default=0.5)), + ("affected_component", models.CharField(blank=True, max_length=240)), + ("failure_category", models.CharField(blank=True, max_length=80)), + ("recommended_action", models.CharField(default="INVESTIGATE", max_length=32)), + ("recommended_route", models.CharField(blank=True, max_length=120)), + ("grouping_key", models.CharField(max_length=240)), + ("status", models.CharField(default="OPEN", max_length=32)), + ("steward_policy_metadata", models.JSONField(blank=True, default=dict)), + ("metadata", models.JSONField(blank=True, default=dict)), + ("project", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="scenario_findings", to="projects.project")), + ("progeny_signal", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="scenario_findings", to="agents.progenysignal")), + ("roadmap_item", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="scenario_findings", to="projects.roadmapitem")), + ("scenario", models.ForeignKey(on_delete=django.db.models.deletion.CASCADE, related_name="findings", to="projects.scenario")), + ("scenario_run", models.ForeignKey(blank=True, null=True, on_delete=django.db.models.deletion.SET_NULL, related_name="findings", to="projects.scenariorun")), + ], + ), + ] diff --git a/control_plane/projects/models.py b/control_plane/projects/models.py index f082a88..3fbf46e 100644 --- a/control_plane/projects/models.py +++ b/control_plane/projects/models.py @@ -225,22 +225,59 @@ class Artifact(TimestampedModel): class RoadmapStatus(models.TextChoices): - INBOX = "INBOX" + PROPOSED = "PROPOSED" + ACCEPTED = "ACCEPTED" + PLANNING = "PLANNING" + IN_PROGRESS = "IN_PROGRESS" + COMPLETE = "COMPLETE" + DEFERRED = "DEFERRED" + REJECTED = "REJECTED" + SUPERSEDED = "SUPERSEDED" + + +class RoadmapHorizon(models.TextChoices): NOW = "NOW" NEXT = "NEXT" LATER = "LATER" EXPLORING = "EXPLORING" - DECLINED = "DECLINED" - DONE = "DONE" + + +class RoadmapTargetAction(models.TextChoices): + EXTEND = "EXTEND" + EVOLVE = "EVOLVE" + REPAIR = "REPAIR" + INVESTIGATE = "INVESTIGATE" + NONE = "NONE" class RoadmapItem(TimestampedModel): project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="roadmap_items") title = models.CharField(max_length=200) description = models.TextField(blank=True) - source = models.CharField(max_length=80, default="user") - status = models.CharField(max_length=32, choices=RoadmapStatus.choices, default=RoadmapStatus.INBOX) + source = models.CharField(max_length=80, default="USER") + source_ref = models.JSONField(default=dict, blank=True) + rationale = models.TextField(blank=True) + evidence = models.JSONField(default=dict, blank=True) + horizon = models.CharField(max_length=32, choices=RoadmapHorizon.choices, default=RoadmapHorizon.EXPLORING) + category = models.CharField(max_length=80, blank=True) + status = models.CharField(max_length=32, choices=RoadmapStatus.choices, default=RoadmapStatus.PROPOSED) + target_action = models.CharField(max_length=32, choices=RoadmapTargetAction.choices, default=RoadmapTargetAction.NONE) priority = models.PositiveSmallIntegerField(default=50) + value_score = models.FloatField(default=0.5) + effort_score = models.FloatField(default=0.5) + risk_score = models.FloatField(default=0.5) + confidence = models.FloatField(default=0.5) + strategic_fit = models.FloatField(default=0.5) + technical_fit = models.FloatField(default=0.5) + urgency = models.FloatField(default=0.5) + composite_score = models.FloatField(default=0.5) + grouping_key = models.CharField(max_length=240, blank=True) + dependencies = models.ManyToManyField("self", symmetrical=False, blank=True, related_name="dependent_roadmap_items") + related_items = models.ManyToManyField("self", symmetrical=True, blank=True) + converted_extension = models.ForeignKey("projects.ExtensionCandidate", on_delete=models.SET_NULL, null=True, blank=True, related_name="converted_from_roadmap_items") + converted_evolution = models.ForeignKey("projects.EvolutionCandidate", on_delete=models.SET_NULL, null=True, blank=True, related_name="converted_from_roadmap_items") + converted_investigation = models.ForeignKey("agents.ProgenyInvestigation", on_delete=models.SET_NULL, null=True, blank=True, related_name="converted_from_roadmap_items") + metadata = models.JSONField(default=dict, blank=True) class Finding(TimestampedModel): @@ -254,13 +291,75 @@ class Finding(TimestampedModel): status = models.CharField(max_length=32, default="OPEN") +class ScenarioSuite(TimestampedModel): + project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="scenario_suites") + name = models.CharField(max_length=200) + version = models.PositiveIntegerField(default=1) + purpose = models.TextField(blank=True) + status = models.CharField(max_length=32, default="DRAFT") + frozen_at = models.DateTimeField(null=True, blank=True) + execution_graph_version = models.ForeignKey("graph.ExecutionGraphVersion", on_delete=models.SET_NULL, null=True, blank=True, related_name="scenario_suites") + metadata = models.JSONField(default=dict, blank=True) + + class Scenario(TimestampedModel): project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="scenarios", null=True, blank=True) + suite = models.ForeignKey(ScenarioSuite, on_delete=models.CASCADE, related_name="scenarios", null=True, blank=True) name = models.CharField(max_length=200) + title = models.CharField(max_length=255, blank=True) + description = models.TextField(blank=True) + scenario_type = models.CharField(max_length=80, default="WORKFLOW") + target_component = models.CharField(max_length=240, blank=True) target_type = models.CharField(max_length=80) target_id = models.CharField(max_length=120) + preconditions = models.JSONField(default=list, blank=True) + injected_condition = models.JSONField(default=dict, blank=True) + expected_invariants = models.JSONField(default=list, blank=True) + success_criteria = models.JSONField(default=list, blank=True) + severity = models.CharField(max_length=32, default="MEDIUM") + source = models.CharField(max_length=80, default="USER") definition = models.JSONField(default=dict, blank=True) status = models.CharField(max_length=32, default="DRAFT") + validated_at = models.DateTimeField(null=True, blank=True) + rejection_reason = models.TextField(blank=True) + resource_budget = models.JSONField(default=dict, blank=True) + metadata = models.JSONField(default=dict, blank=True) + + +class ScenarioRun(TimestampedModel): + scenario = models.ForeignKey(Scenario, on_delete=models.CASCADE, related_name="runs") + project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="scenario_runs") + repository_baseline = models.CharField(max_length=255, blank=True) + execution_graph_version = models.ForeignKey("graph.ExecutionGraphVersion", on_delete=models.SET_NULL, null=True, blank=True, related_name="scenario_runs") + graph_run = models.ForeignKey("graph.GraphRun", on_delete=models.SET_NULL, null=True, blank=True, related_name="scenario_runs") + status = models.CharField(max_length=32, default="PENDING") + started_at = models.DateTimeField(null=True, blank=True) + completed_at = models.DateTimeField(null=True, blank=True) + environment_metadata = models.JSONField(default=dict, blank=True) + result = models.CharField(max_length=32, blank=True) + failure_evidence = models.JSONField(default=dict, blank=True) + telemetry = models.JSONField(default=dict, blank=True) + + +class ScenarioFinding(TimestampedModel): + project = models.ForeignKey(Project, on_delete=models.CASCADE, related_name="scenario_findings") + scenario = models.ForeignKey(Scenario, on_delete=models.CASCADE, related_name="findings") + scenario_run = models.ForeignKey(ScenarioRun, on_delete=models.SET_NULL, null=True, blank=True, related_name="findings") + title = models.CharField(max_length=255) + summary = models.TextField(blank=True) + evidence = models.JSONField(default=dict, blank=True) + severity = models.CharField(max_length=32, default="MEDIUM") + confidence = models.FloatField(default=0.5) + affected_component = models.CharField(max_length=240, blank=True) + failure_category = models.CharField(max_length=80, blank=True) + recommended_action = models.CharField(max_length=32, default="INVESTIGATE") + recommended_route = models.CharField(max_length=120, blank=True) + grouping_key = models.CharField(max_length=240) + status = models.CharField(max_length=32, default="OPEN") + roadmap_item = models.ForeignKey(RoadmapItem, on_delete=models.SET_NULL, null=True, blank=True, related_name="scenario_findings") + steward_policy_metadata = models.JSONField(default=dict, blank=True) + progeny_signal = models.ForeignKey("agents.ProgenySignal", on_delete=models.SET_NULL, null=True, blank=True, related_name="scenario_findings") + metadata = models.JSONField(default=dict, blank=True) class StewardPolicy(TimestampedModel): @@ -391,6 +490,7 @@ class EvolutionCandidate(TimestampedModel): confidence = models.FloatField(default=0.0) source_steward_finding = models.ForeignKey(StewardFinding, on_delete=models.SET_NULL, null=True, blank=True, related_name="evolution_candidates") source_investigation = models.ForeignKey("agents.ProgenyInvestigation", on_delete=models.SET_NULL, null=True, blank=True, related_name="evolution_candidates") + source_roadmap_item = models.ForeignKey(RoadmapItem, on_delete=models.SET_NULL, null=True, blank=True, related_name="evolution_candidates") source_opportunity = models.ForeignKey("projects.ExplorationOpportunity", on_delete=models.SET_NULL, null=True, blank=True, related_name="evolution_candidates") metadata = models.JSONField(default=dict, blank=True) diff --git a/graph/bootstrap.py b/graph/bootstrap.py index 8ab91b0..ef4f1fc 100644 --- a/graph/bootstrap.py +++ b/graph/bootstrap.py @@ -4,6 +4,8 @@ from django.utils import timezone from graph.lifecycle import project_evolution_graph_v1, project_exploration_graph_v1, project_extension_graph_v1 from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus +from graph.roadmap import project_roadmap_review_graph_v1 +from graph.scenario_lab import scenario_lab_graph_v1 from graph.steward import steward_run_graph_v1 from graph.task_execution import task_execution_graph_v1 @@ -81,3 +83,11 @@ def champion_project_evolution_graph_v1() -> ExecutionGraphVersion: def champion_project_exploration_graph_v1() -> ExecutionGraphVersion: return _champion_graph(project_exploration_graph_v1()) + + +def champion_project_roadmap_review_graph_v1() -> ExecutionGraphVersion: + return _champion_graph(project_roadmap_review_graph_v1()) + + +def champion_scenario_lab_graph_v1() -> ExecutionGraphVersion: + return _champion_graph(scenario_lab_graph_v1()) diff --git a/graph/roadmap.py b/graph/roadmap.py new file mode 100644 index 0000000..78c6fdc --- /dev/null +++ b/graph/roadmap.py @@ -0,0 +1,87 @@ +from __future__ import annotations + +from agents.roadmap import RoadmapService +from control_plane.projects.models import Project +from graph.native_runtime import GraphExecutionContext +from graph.registry import NodeHandlerRegistry, NodeResult +from graph.spec import ExecutionGraphSpec, GraphEdgeSpec, GraphNodeSpec + + +def project_roadmap_review_graph_v1() -> ExecutionGraphSpec: + nodes = ["prepare", "gather_project_state", "gather_candidate_items", "deduplicate", "score", "project_brain_review", "recommend_priorities", "persist", "complete"] + spec = ExecutionGraphSpec( + name="project_roadmap_review", + version=1, + graph_type="PROJECT_ROADMAP_REVIEW", + entry="prepare", + nodes={node: GraphNodeSpec(node, node if node == "complete" else f"roadmap_{node}") for node in nodes}, + edges=[GraphEdgeSpec(nodes[index], nodes[index + 1], "success") for index in range(len(nodes) - 1)], + terminal_nodes=["complete"], + metadata={"description": "Roadmap review workflow: gather intent, deduplicate, score, Sol review, persist priorities without execution."}, + ) + spec.validate() + return spec + + +class RoadmapNode: + idempotent = True + replay_safe = True + destructive = False + + def __init__(self, service: RoadmapService, node_type: str) -> None: + self.service = service + self.node_type = node_type + + def project(self, context: GraphExecutionContext) -> Project: + return context.graph_run.project + + +class RoadmapSimpleNode(RoadmapNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + return NodeResult("COMPLETE", "success") + + +class RoadmapGatherStateNode(RoadmapNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + metadata = dict(context.graph_run.metadata) + metadata["project_state"] = self.service.project_context(self.project(context)) + context.graph_run.metadata = metadata + context.graph_run.save(update_fields=["metadata", "updated_at"]) + return NodeResult("COMPLETE", "success") + + +class RoadmapGatherCandidatesNode(RoadmapNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + items = self.service.gather_candidate_items(self.project(context)) + return NodeResult("COMPLETE", "success", {"candidate_item_count": len(items)}) + + +class RoadmapProjectBrainNode(RoadmapNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + recommendations = self.service.review_with_project_brain(self.project(context)) + metadata = dict(context.graph_run.metadata) + metadata["project_brain_recommendations"] = recommendations + context.graph_run.metadata = metadata + context.graph_run.save(update_fields=["metadata", "updated_at"]) + return NodeResult("COMPLETE", "success", {"recommendations": recommendations}) + + +class RoadmapRecommendNode(RoadmapNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + recommendations = context.graph_run.metadata.get("project_brain_recommendations", {}) + updated = self.service.apply_recommendations(self.project(context), recommendations if isinstance(recommendations, dict) else {}) + return NodeResult("COMPLETE", "success", {"updated_count": len(updated)}) + + +class RoadmapPersistNode(RoadmapNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + view = self.service.project_roadmap_view(self.project(context)) + self.service.bus.publish("ROADMAP_REVIEW_COMPLETED", project=self.project(context), payload={"graph_run_id": str(context.graph_run.id), "item_count": self.project(context).roadmap_items.count()}) + return NodeResult("COMPLETE", "success", {"roadmap": view}) + + +def roadmap_registry(service: RoadmapService) -> NodeHandlerRegistry: + registry = NodeHandlerRegistry() + for handler in [RoadmapSimpleNode(service, "roadmap_prepare"), RoadmapGatherStateNode(service, "roadmap_gather_project_state"), RoadmapGatherCandidatesNode(service, "roadmap_gather_candidate_items"), RoadmapSimpleNode(service, "roadmap_deduplicate"), RoadmapSimpleNode(service, "roadmap_score"), RoadmapProjectBrainNode(service, "roadmap_project_brain_review"), RoadmapRecommendNode(service, "roadmap_recommend_priorities"), RoadmapPersistNode(service, "roadmap_persist")]: + registry.register(handler) + return registry diff --git a/graph/scenario_lab.py b/graph/scenario_lab.py new file mode 100644 index 0000000..3d93fcd --- /dev/null +++ b/graph/scenario_lab.py @@ -0,0 +1,106 @@ +from __future__ import annotations + +from agents.scenario_lab import ScenarioLabService +from control_plane.projects.models import ScenarioSuite +from graph.lifecycle import ApprovalNode +from graph.native_runtime import GraphExecutionContext +from graph.registry import NodeHandlerRegistry, NodeResult +from graph.spec import ExecutionGraphSpec, GraphEdgeSpec, GraphNodeSpec + + +def scenario_lab_graph_v1() -> ExecutionGraphSpec: + nodes = ["prepare", "gather_context", "generate_or_select_scenarios", "validate", "await_approval", "freeze_suite", "execute_scenarios", "collect_results", "classify_findings", "route_findings", "summarize", "complete"] + spec = ExecutionGraphSpec( + name="scenario_lab", + version=1, + graph_type="SCENARIO_LAB", + entry="prepare", + nodes={node: GraphNodeSpec(node, node if node == "complete" else f"scenario_{node}") for node in nodes}, + edges=[GraphEdgeSpec("prepare", "gather_context", "success"), GraphEdgeSpec("gather_context", "generate_or_select_scenarios", "success"), GraphEdgeSpec("generate_or_select_scenarios", "validate", "success"), GraphEdgeSpec("validate", "await_approval", "approval_required"), GraphEdgeSpec("validate", "freeze_suite", "success"), GraphEdgeSpec("await_approval", "freeze_suite", "approved"), GraphEdgeSpec("await_approval", "complete", "rejected"), GraphEdgeSpec("freeze_suite", "execute_scenarios", "success"), GraphEdgeSpec("execute_scenarios", "collect_results", "success"), GraphEdgeSpec("collect_results", "classify_findings", "success"), GraphEdgeSpec("classify_findings", "route_findings", "success"), GraphEdgeSpec("route_findings", "summarize", "success"), GraphEdgeSpec("summarize", "complete", "success")], + terminal_nodes=["complete"], + metadata={"description": "Scenario Lab workflow: safe scenario generation/validation/execution, finding classification and routing."}, + ) + spec.validate() + return spec + + +class ScenarioNode: + idempotent = True + replay_safe = True + destructive = False + + def __init__(self, service: ScenarioLabService, node_type: str) -> None: + self.service = service + self.node_type = node_type + + def suite(self, context: GraphExecutionContext) -> ScenarioSuite: + return ScenarioSuite.objects.get(id=context.graph_run.metadata["scenario_suite_id"]) + + +class ScenarioSimpleNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + return NodeResult("COMPLETE", "success") + + +class ScenarioGatherContextNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + suite = self.suite(context) + metadata = dict(context.graph_run.metadata) + metadata["scenario_context"] = self.service.project_context(suite.project) + context.graph_run.metadata = metadata + context.graph_run.save(update_fields=["metadata", "updated_at"]) + return NodeResult("COMPLETE", "success") + + +class ScenarioGenerateNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + suite = self.suite(context) + generated = [] if suite.scenarios.exists() else self.service.generate_scenarios(suite) + return NodeResult("COMPLETE", "success", {"generated_count": len(generated), "scenario_count": suite.scenarios.count()}) + + +class ScenarioValidateNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + suite = self.suite(context) + valid = self.service.validate_suite(suite) + approval_required = any(s.severity in ["HIGH", "CRITICAL"] for s in valid) + return NodeResult("COMPLETE", "approval_required" if approval_required else "success", {"valid_count": len(valid), "approval_required": approval_required}) + + +class ScenarioFreezeNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + suite = self.service.freeze_suite(self.suite(context)) + return NodeResult("COMPLETE", "success", {"suite_id": str(suite.id), "status": suite.status}) + + +class ScenarioExecuteNode(ScenarioNode): + destructive = True + + def run(self, context: GraphExecutionContext) -> NodeResult: + suite = self.suite(context) + runs = self.service.execute_suite(suite, graph_run=context.graph_run) + return NodeResult("COMPLETE", "success", {"scenario_run_ids": [str(run.id) for run in runs]}) + + +class ScenarioCollectNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + suite = self.suite(context) + return NodeResult("COMPLETE", "success", {"results": list(suite.project.scenario_runs.filter(scenario__suite=suite).values("scenario__title", "result", "failure_evidence"))}) + + +class ScenarioRouteNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + routed = self.service.route_findings(self.suite(context)) + return NodeResult("COMPLETE", "success", {"routed_count": len(routed)}) + + +class ScenarioSummarizeNode(ScenarioNode): + def run(self, context: GraphExecutionContext) -> NodeResult: + return NodeResult("COMPLETE", "success", self.service.summarize_suite(self.suite(context))) + + +def scenario_lab_registry(service: ScenarioLabService) -> NodeHandlerRegistry: + registry = NodeHandlerRegistry() + for handler in [ScenarioSimpleNode(service, "scenario_prepare"), ScenarioGatherContextNode(service, "scenario_gather_context"), ScenarioGenerateNode(service, "scenario_generate_or_select_scenarios"), ScenarioValidateNode(service, "scenario_validate"), ApprovalNode("scenario_await_approval"), ScenarioFreezeNode(service, "scenario_freeze_suite"), ScenarioExecuteNode(service, "scenario_execute_scenarios"), ScenarioCollectNode(service, "scenario_collect_results"), ScenarioSimpleNode(service, "scenario_classify_findings"), ScenarioRouteNode(service, "scenario_route_findings"), ScenarioSummarizeNode(service, "scenario_summarize")]: + registry.register(handler) + return registry diff --git a/tests/test_roadmap_scenario_lab_v1.py b/tests/test_roadmap_scenario_lab_v1.py new file mode 100644 index 0000000..f82c1af --- /dev/null +++ b/tests/test_roadmap_scenario_lab_v1.py @@ -0,0 +1,149 @@ +from __future__ import annotations + +import json +from pathlib import Path + +from agents.providers import DeterministicSolProvider +from agents.roadmap import RoadmapService +from agents.scenario_lab import ScenarioLabService +from control_plane.projects.models import Exploration, ExplorationOpportunity, Project, RoadmapHorizon, RoadmapItem, ScenarioFinding, StewardFinding, Task +from graph.bootstrap import champion_project_roadmap_review_graph_v1, champion_scenario_lab_graph_v1 +from graph.models import GraphApprovalStatus, GraphRun, GraphRunStatus +from graph.roadmap import project_roadmap_review_graph_v1, roadmap_registry +from graph.scenario_lab import scenario_lab_graph_v1, scenario_lab_registry +from graph.langgraph_runtime import LangGraphRuntime +from model_router.router import ModelRouter +from tests.test_m2_autonomous_loop import create_disposable_django_repo + + +def project(tmp_path: Path) -> Project: + return Project.objects.create(name="Roadmap Scenario", goal="Existing project", repository_path=str(create_disposable_django_repo(tmp_path))) + + +def sol_router(payload: dict[str, object]) -> ModelRouter: + return ModelRouter({"sol": DeterministicSolProvider(json.dumps(payload))}, persist_requests=True) + + +def approve_and_resume(graph_run: GraphRun, registry) -> GraphRun: + approval = graph_run.approvals.get() + approval.status = GraphApprovalStatus.APPROVED + approval.decided_by = "tester" + approval.save(update_fields=["status", "decided_by", "updated_at"]) + graph_run.status = GraphRunStatus.RUNNING + graph_run.failure_reason = "" + graph_run.save(update_fields=["status", "failure_reason", "updated_at"]) + LangGraphRuntime(registry).run_until_terminal_or_paused(graph_run) + graph_run.refresh_from_db() + return graph_run + + +def test_roadmap_upsert_scores_deduplicates_and_preserves_lineage(tmp_path: Path) -> None: + proj = project(tmp_path) + service = RoadmapService() + + first = service.upsert_item(proj, title="Add scenario coverage", description="Cover graph failures", source="EXPLORE", source_ref={"opportunity": "1"}, evidence={"files": ["graph"]}, horizon="NEXT", category="RELIABILITY", target_action="EXTEND", scores={"value": 0.9, "effort": 0.3, "risk": 0.2, "confidence": 0.8, "strategic_fit": 0.9, "technical_fit": 0.8, "urgency": 0.7}) + second = service.upsert_item(proj, title="Add scenario coverage", source="STEWARD", source_ref={"finding": "2"}, evidence={"severity": "HIGH"}) + + assert first.id == second.id + assert RoadmapItem.objects.filter(project=proj).count() == 1 + assert second.metadata["occurrences"] == 2 + assert second.source_ref["opportunity"] == "1" + assert second.source_ref["finding"] == "2" + assert second.composite_score > 0.5 + + +def test_roadmap_gathers_sources_prioritizes_and_does_not_execute(tmp_path: Path) -> None: + proj = project(tmp_path) + exploration = Exploration.objects.create(project=proj, status="COMPLETE") + ExplorationOpportunity.objects.create(exploration=exploration, project=proj, title="Preserve absolute SQLite paths", description="SQLite absolute paths should remain absolute", opportunity_type="RELIABILITY", evidence={"source": "explore"}, recommended_action="REPAIR", value_score=0.8, effort_score=0.2, risk_score=0.3, confidence=0.8, strategic_fit=0.7, technical_fit=0.9, grouping_key="opp") + StewardFinding.objects.create(project=proj, finding_type="PERFORMANCE", title="Queue latency monitor", summary="Monitor queue depth > 20", evidence={"baseline_measurement": {"metric": "latency", "latency": 100}}, severity="MEDIUM", confidence=0.7, recommended_action="EVOLVE", grouping_key="steward") + service = RoadmapService(sol_router({"recommendations": []})) + + version = champion_project_roadmap_review_graph_v1() + graph_run = GraphRun.objects.create(execution_graph_version=version, project=proj, current_node=version.graph_spec["entry"]) + LangGraphRuntime(roadmap_registry(service)).run_until_terminal_or_paused(graph_run) + + assert graph_run.status == GraphRunStatus.COMPLETE + assert RoadmapItem.objects.filter(project=proj).count() == 2 + assert Task.objects.filter(project=proj).count() == 0 + view = service.project_roadmap_view(proj) + assert set(view) == {"NOW", "NEXT", "LATER", "EXPLORING"} + assert len(view["EXPLORING"]) == 2 + + +def test_roadmap_project_brain_recommendations_and_conversion_lineage(tmp_path: Path) -> None: + proj = project(tmp_path) + item = RoadmapService().upsert_item(proj, title="Add project health dashboard", description="Expose lifecycle status", horizon="NEXT", category="FEATURE", target_action="EXTEND") + service = RoadmapService(sol_router({"recommendations": [{"item_id": str(item.id), "recommendation": "move_to_now", "horizon": "NOW", "status": "ACCEPTED", "rationale": "High operator value"}]})) + + recommendations = service.review_with_project_brain(proj) + service.apply_recommendations(proj, recommendations) + item.refresh_from_db() + candidate = service.convert_to_extension(item) + + assert item.horizon == RoadmapHorizon.NOW + assert item.status == "PLANNING" + assert candidate.source_roadmap_item == item + assert Task.objects.filter(project=proj).count() == 0 + + evolve_item = RoadmapService().upsert_item(proj, title="Reduce roadmap review latency", description="Improve measured review duration", category="PERFORMANCE", target_action="EVOLVE") + evolution = RoadmapService().convert_to_evolution(evolve_item, baseline_measurement={"metric": "duration", "duration": 100}) + assert evolution.source_roadmap_item == evolve_item + + +def test_scenario_suite_validation_results_routing_and_coverage(tmp_path: Path) -> None: + proj = project(tmp_path) + service = ScenarioLabService() + suite = service.create_suite( + proj, + name="Deterministic Lab", + scenarios=[ + {"title": "Malformed coder output", "scenario_type": "AGENT_WORKFLOW", "target_component": "coder", "injected_condition": {"mechanism": "malformed_input", "expected_result": "FAIL", "recommended_action": "PROGENY"}, "expected_invariants": ["classified"], "success_criteria": ["finding persisted"], "resource_budget": {"max_seconds": 5}}, + {"title": "Permission denial", "scenario_type": "PERMISSION", "target_component": "tools", "injected_condition": {"mechanism": "permission_denial", "expected_result": "PASS"}, "expected_invariants": ["denied safely"], "success_criteria": ["run completes"], "resource_budget": {"max_seconds": 5}}, + {"title": "Infra unavailable", "scenario_type": "DEPENDENCY_FAILURE", "target_component": "provider", "injected_condition": {"mechanism": "provider_failure_replay", "infrastructure_failure": True}, "expected_invariants": ["not project failure"], "success_criteria": ["infra result"], "resource_budget": {"max_seconds": 5}}, + {"title": "Unsafe destructive", "scenario_type": "FAILURE_INJECTION", "target_component": "infra", "injected_condition": {"destructive": True}, "expected_invariants": ["never runs"], "success_criteria": ["rejected"], "resource_budget": {"max_seconds": 5}}, + ], + ) + + valid = service.validate_suite(suite) + service.freeze_suite(suite) + runs = service.execute_suite(suite) + service.route_findings(suite) + + assert len(valid) == 3 + assert suite.scenarios.filter(status="REJECTED", title="Unsafe destructive").exists() + assert {run.result for run in runs} == {"FAIL", "PASS", "INFRASTRUCTURE_FAILURE"} + assert ScenarioFinding.objects.filter(project=proj, recommended_action="PROGENY", status="ROUTED").exists() + assert proj.progenysignal_set.filter(source="scenario_lab").exists() + assert service.coverage(proj)["AGENT_WORKFLOW"] == 1 + + +def test_scenario_lab_graph_approval_pause_resume_and_routes_to_roadmap(tmp_path: Path) -> None: + proj = project(tmp_path) + service = ScenarioLabService() + suite = service.create_suite( + proj, + name="High Severity Lab", + scenarios=[{"title": "Performance regression", "scenario_type": "PERFORMANCE", "severity": "HIGH", "target_component": "project_context", "injected_condition": {"mechanism": "performance_regression", "expected_result": "FAIL", "recommended_action": "EVOLVE", "failure_category": "PERFORMANCE"}, "expected_invariants": ["latency bounded"], "success_criteria": ["finding routed"], "resource_budget": {"max_seconds": 5, "max_parallelism": 1}}], + ) + version = champion_scenario_lab_graph_v1() + graph_run = GraphRun.objects.create(execution_graph_version=version, project=proj, current_node=version.graph_spec["entry"], metadata={"scenario_suite_id": str(suite.id)}) + registry = scenario_lab_registry(service) + + LangGraphRuntime(registry).run_until_terminal_or_paused(graph_run) + assert graph_run.status == GraphRunStatus.PAUSED + approve_and_resume(graph_run, registry) + + suite.refresh_from_db() + assert graph_run.status == GraphRunStatus.COMPLETE + assert suite.status == "FROZEN" + finding = ScenarioFinding.objects.get(project=proj) + assert finding.status == "ROUTED" + assert finding.roadmap_item is not None + assert finding.roadmap_item.target_action == "EVOLVE" + assert Task.objects.filter(project=proj).count() == 0 + + +def test_graph_specs_are_serializable() -> None: + assert project_roadmap_review_graph_v1().name == "project_roadmap_review" + assert scenario_lab_graph_v1().name == "scenario_lab"