87 lines
4.3 KiB
Python
87 lines
4.3 KiB
Python
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
|