from __future__ import annotations from pathlib import Path from agents.providers import DeterministicCodingProvider from control_plane.agents.management.commands.seed_core_agents import Command as SeedAgentsCommand from control_plane.events.models import Event, EventType from control_plane.projects.models import CommitRecord, Project, ProjectPlan, Milestone, Task, TaskStatus from control_plane.verification.models import Review, TestRun, Verification, VerificationResult from graph.models import ExecutionGraphDefinition, ExecutionGraphVersion, ExecutionGraphVersionStatus, GraphEdgeTraversal, GraphRun, GraphRunStatus from graph.native_runtime import NativeGraphRuntime from graph.task_execution import task_execution_graph_v1 from graph.task_nodes import TaskExecutionServices, task_execution_registry from model_router.router import ModelRouter from tests.test_m2_autonomous_loop import create_disposable_django_repo def create_task(repository_path: Path, goal: str, acceptance: list[str], *, max_retries: int = 2) -> Task: project = Project.objects.create(name=f"Graph Project {goal[:12]}", goal=goal, repository_path=str(repository_path)) plan = ProjectPlan.objects.create(project=project, version=1, goal=goal) milestone = Milestone.objects.create(project=project, plan=plan, key="G1", title="Graph", goal="Execution graph") return Task.objects.create( project=project, milestone=milestone, task_type="implementation", status=TaskStatus.RUNNING, goal=goal, acceptance_criteria=acceptance, max_retries=max_retries, ) def graph_run_for_task(task: Task) -> GraphRun: spec = task_execution_graph_v1() definition, _ = ExecutionGraphDefinition.objects.get_or_create( name=spec.name, defaults={"graph_type": spec.graph_type, "description": "Task execution graph"}, ) version, _ = ExecutionGraphVersion.objects.get_or_create( graph=definition, version=spec.version, defaults={"status": ExecutionGraphVersionStatus.CHAMPION, "graph_spec": spec.to_dict()}, ) return GraphRun.objects.create( execution_graph_version=version, project=task.project, milestone=task.milestone, feature=task.feature, task=task, current_node=spec.entry, ) def run_graph(task: Task, *, interrupt_after: str | None = None) -> GraphRun: SeedAgentsCommand().handle() services = TaskExecutionServices(ModelRouter({"qwen": DeterministicCodingProvider()}), test_command=["python", "manage.py", "test"]) runtime = NativeGraphRuntime(task_execution_registry(services)) return runtime.run_until_terminal_or_paused(graph_run_for_task(task), interrupt_after=interrupt_after) def test_native_task_execution_graph_success_matches_loop_semantics(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) task.refresh_from_db() assert graph_run.status == GraphRunStatus.COMPLETE assert task.status == TaskStatus.COMPLETE assert CommitRecord.objects.filter(task=task).count() == 1 assert TestRun.objects.get(task=task).status == "PASS" assert Review.objects.get(task=task).status == "PASS" assert Verification.objects.get(task=task).result == VerificationResult.PASS assert task.worktree.status == "CLEANED" assert GraphEdgeTraversal.objects.filter(graph_run=graph_run, source_node="judge", target_node="commit", condition="PASS").exists() def test_native_task_execution_graph_retry_exhaustion_matches_loop_semantics(tmp_path: Path) -> None: repo = create_disposable_django_repo(tmp_path) task = create_task(repo, "FORCE_BAD_IMPLEMENTATION Add a /health endpoint returning JSON {\"status\": \"ok\"} and add tests.", ["/health returns JSON ok", "tests pass"], max_retries=1) graph_run = run_graph(task) task.refresh_from_db() assert graph_run.status == GraphRunStatus.FAILED assert task.status == TaskStatus.FAILED assert task.retry_count == 2 assert CommitRecord.objects.filter(task=task).count() == 0 assert Event.objects.filter(task=task, event_type=EventType.TASK_FAILED).exists() assert Event.objects.filter(task=task, event_type="TASK_RETRY_EXHAUSTED").exists() def test_native_task_execution_graph_resume_after_coder_prevents_duplicate_commit(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"]) SeedAgentsCommand().handle() graph_run = graph_run_for_task(task) services = TaskExecutionServices(ModelRouter({"qwen": DeterministicCodingProvider()}), test_command=["python", "manage.py", "test"]) runtime = NativeGraphRuntime(task_execution_registry(services)) runtime.run_until_terminal_or_paused(graph_run, interrupt_after="coder") graph_run.refresh_from_db() assert graph_run.current_node == "coder" assert task.attempts.count() == 1 runtime.run_until_terminal_or_paused(graph_run) runtime.run_until_terminal_or_paused(graph_run) task.refresh_from_db() assert task.status == TaskStatus.COMPLETE assert CommitRecord.objects.filter(task=task).count() == 1 assert graph_run.node_runs.filter(node_id="coder").count() == 1