Smart-Notes-backend / app /graph /stream_runner.py
pluto90's picture
Upload stream_runner.py
673b675 verified
Raw
History Blame Contribute Delete
2.57 kB
from app.graph.nodes.router import router_node
from app.graph.nodes.rag_agent import rag_agent_node
from app.graph.nodes.synthesizer import synthesizer_stream
from app.graph.nodes.evaluator import evaluator_node
async def stream_graph(initial_state):
"""
Streaming execution.
Emits dictionaries with keys:
event
data
"""
# ---------------- Router ----------------
yield {
"event": "status",
"data": {
"step": "router",
"message": "Routing query..."
}
}
state = router_node(initial_state)
route = state["route"]
# ---------------- Retrieval ----------------
if route == "rag":
yield {
"event": "status",
"data": {
"step": "retrieval",
"message": "Searching document..."
}
}
state = rag_agent_node(state)
elif route == "hybrid":
yield {
"event": "status",
"data": {
"step": "retrieval",
"message": "Searching document..."
}
}
else:
yield {
"event": "status",
"data": {
"step": "llm",
"message": "Using general knowledge..."
}
}
# ---------------- Generation ----------------
yield {
"event": "status",
"data": {
"step": "generation",
"message": "Generating answer..."
}
}
answer = ""
async for token in synthesizer_stream(state):
answer += token
yield {
"event": "token",
"data": token
}
state["final_answer"] = answer
# ---------------- Evaluation ----------------
yield {
"event": "status",
"data": {
"step": "evaluation",
"message": "Evaluating response..."
}
}
state = evaluator_node(state)
# ---------------- Sources ----------------
yield {
"event": "sources",
"data": state.get("sources", [])
}
# ---------------- Evaluation ----------------
yield {
"event": "evaluation",
"data": state.get("evaluation", {})
}
# ---------------- Finished ----------------
yield {
"event": "done",
"data": {
"route": state["route"],
"answer": answer
}
}