# import asyncio # import json # import time # from fastapi import APIRouter, Form # from fastapi.responses import StreamingResponse # from app.core.mongo import conversations # from app.graph.stream_runner import stream_graph # router = APIRouter() # @router.post("/query-stream") # async def query_stream( # doc_id: str = Form(...), # question: str = Form(...) # ): # # -------------------------------------------------- # # Save user message # # -------------------------------------------------- # conversations.update_one( # {"doc_id": doc_id}, # { # "$push": { # "history": { # "role": "user", # "content": question # } # } # }, # upsert=True # ) # doc = conversations.find_one( # { # "doc_id": doc_id # } # ) # history = doc.get("history", [])[-10:] if doc else [] # history_text = "\n".join( # [ # f"{m['role'].title()}: {m['content']}" # for m in history # ] # ) # state = { # "query": question, # "doc_id": doc_id, # "history": history_text, # "route": None, # "context": None, # "sources": [], # "score": 0, # "evaluation": None, # "final_answer": None, # } # async def event_generator(): # final_answer = "" # async for event in stream_graph(state): # # ---------------------------------------- # # token # # ---------------------------------------- # if event["event"] == "token": # final_answer += event["data"] # yield ( # f"event: token\n" # f"data: {event['data']}\n\n" # ) # # ---------------------------------------- # # status # # ---------------------------------------- # elif event["event"] == "status": # yield ( # f"event: status\n" # f"data: {event['data']}\n\n" # ) # # ---------------------------------------- # # done # # ---------------------------------------- # elif event["event"] == "done": # conversations.update_one( # {"doc_id": doc_id}, # { # "$push": { # "history": { # "role": "assistant", # "content": final_answer # } # } # } # ) # payload = json.dumps(event["data"]) # yield ( # f"event: done\n" # f"data: {payload}\n\n" # ) # return StreamingResponse( # event_generator(), # media_type="text/event-stream" # ) import json from fastapi import APIRouter, Form from fastapi.responses import StreamingResponse from app.core.mongo import conversations from app.graph.stream_runner import stream_graph router = APIRouter() @router.post("/query-stream") async def query_stream( doc_id: str = Form(...), question: str = Form(...) ): conversations.update_one( {"doc_id": doc_id}, { "$push": { "history": { "role": "user", "content": question } } }, upsert=True, ) document = conversations.find_one( {"doc_id": doc_id} ) history = [] if document: history = document.get("history", [])[-10:] history_text = "\n".join( [ f"{msg['role'].title()}: {msg['content']}" for msg in history ] ) state = { "query": question, "doc_id": doc_id, "history": history_text, "route": None, "context": None, "sources": [], "score": 0, "evaluation": None, "final_answer": None, } async def event_generator(): answer = "" try: async for event in stream_graph(state): event_name = event["event"] # ---------------- Token ---------------- if event_name == "token": answer += event["data"] # ---------------- Done ---------------- elif event_name == "done": conversations.update_one( {"doc_id": doc_id}, { "$push": { "history": { "role": "assistant", "content": answer } } }, ) yield ( f"event: {event_name}\n" f"data: {json.dumps(event['data'])}\n\n" ) except Exception as e: yield ( "event: error\n" f"data: {json.dumps({'message': str(e)})}\n\n" ) return StreamingResponse( event_generator(), media_type="text/event-stream", headers={ "Cache-Control": "no-cache", "Connection": "keep-alive", "X-Accel-Buffering": "no", }, )