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