This repository has been archived on 2026-08-04. You can view files and clone it, but cannot push or open issues or pull requests.
orchestrator/graph.py

143 lines
4.8 KiB
Python
Raw Normal View History

from dotenv import load_dotenv
load_dotenv(".env")
from langchain_core.messages import SystemMessage, HumanMessage
from langgraph.graph import StateGraph, START, END
from langgraph.prebuilt import ToolNode
from state import OrchestratorState
from models import get_orchestrator
from agents import (
implementer_node,
reviewer_node,
researcher_node,
cross_reviewer_node,
scanner_node,
fast_coder_node,
)
from tools import get_composio_tools
ROUTER_PROMPT = """You are a task router for Sea Haven Industries. Analyze the incoming task and decide which agent should handle it.
Available agents:
- implementer: Write new code, add features, fix bugs. Use for any coding task with a clear spec.
- reviewer: Review code changes (diffs, PRs) for correctness, security, maintainability. Uses Claude.
- researcher: Look up documentation, API references, technical questions. Fast and cheap.
- cross_reviewer: Independent code review using a different AI model (GPT). Use when you want a second opinion that catches different blind spots than Claude.
- scanner: Analyze large codebases for patterns, consistency, structural issues. Uses Gemini's large context window.
- fast_coder: Quick, bounded coding for crystal-clear specs. Uses DeepSeek. Best for small, well-defined tasks.
- connector: Interact with external services (Slack, Notion, Google Drive, GitHub) — send messages, read/update pages, find files.
- done: The task is complete or doesn't need agent delegation (e.g., a simple question you can answer directly).
Respond with ONLY the agent name, nothing else. Pick the single best match."""
composio_tools = get_composio_tools()
def router_node(state: OrchestratorState) -> dict:
llm = get_orchestrator()
response = llm.invoke([
SystemMessage(content=ROUTER_PROMPT),
HumanMessage(content=state["task"]),
])
route = response.content.strip().lower()
valid = {
"implementer",
"reviewer",
"researcher",
"cross_reviewer",
"scanner",
"fast_coder",
"connector",
"done",
}
if route not in valid:
route = "researcher"
return {"route": route, "messages": [response]}
def connector_node(state: OrchestratorState) -> dict:
llm = get_orchestrator().bind_tools(composio_tools)
response = llm.invoke([
SystemMessage(
content=(
"You help interact with external services. Use the available tools to complete the task. "
"Make exactly ONE tool call, then stop. Do not chain multiple calls."
)
),
HumanMessage(content=state["task"]),
])
return {"messages": [response]}
def summarizer_node(state: OrchestratorState) -> dict:
last_msg = state["messages"][-1]
content = last_msg.content if hasattr(last_msg, "content") else str(last_msg)
if isinstance(content, list):
content = "\n".join(str(c) for c in content)
if len(content) > 2000:
content = content[:2000] + "...(truncated)"
return {"result": content}
def route_task(state: OrchestratorState) -> str:
return state["route"]
def build_graph():
graph = StateGraph(OrchestratorState)
graph.add_node("router", router_node)
graph.add_node("implementer", implementer_node)
graph.add_node("reviewer", reviewer_node)
graph.add_node("researcher", researcher_node)
graph.add_node("cross_reviewer", cross_reviewer_node)
graph.add_node("scanner", scanner_node)
graph.add_node("fast_coder", fast_coder_node)
graph.add_node("connector", connector_node)
graph.add_node("tool_executor", ToolNode(composio_tools))
graph.add_node("summarizer", summarizer_node)
graph.add_edge(START, "router")
graph.add_conditional_edges(
"router",
route_task,
{
"implementer": "implementer",
"reviewer": "reviewer",
"researcher": "researcher",
"cross_reviewer": "cross_reviewer",
"scanner": "scanner",
"fast_coder": "fast_coder",
"connector": "connector",
"done": END,
},
)
graph.add_edge("implementer", END)
graph.add_edge("reviewer", END)
graph.add_edge("researcher", END)
graph.add_edge("cross_reviewer", END)
graph.add_edge("scanner", END)
graph.add_edge("fast_coder", END)
graph.add_edge("connector", "tool_executor")
graph.add_edge("tool_executor", "summarizer")
graph.add_edge("summarizer", END)
return graph.compile()
app = build_graph()
if __name__ == "__main__":
import sys
task = " ".join(sys.argv[1:]) if len(sys.argv) > 1 else "What is the capital of France?"
result = app.invoke({"task": task, "messages": []})
print(f"\n--- Route: {result['route']} ---")
print(result.get("result", "No result"))