LightRT/text2sql_backend
0
1from fastapi import FastAPI , HTTPException2from src.embedding import create_embeddings3from src.graph import build_agent, AgentContext4from pydantic import BaseModel , Field5import os 6from dotenv import load_dotenv7import asyncio8from psycopg_pool import AsyncConnectionPool9from langgraph.checkpoint.postgres.aio import AsyncPostgresSaver10from contextlib import asynccontextmanager11import logging12 13load_dotenv()14 15logging.basicConfig(level=logging.INFO)16logger = logging.getLogger("text2sql")17 18DB_URI = os.getenv("DATABASE_URI")19 20@asynccontextmanager21async def lifespan(app: FastAPI):22 async with AsyncConnectionPool(23 conninfo=DB_URI,24 min_size=1,25 max_size=15,26 max_idle=300, 27 max_lifetime=1800, 28 reconnect_timeout=10,29 kwargs={"autocommit": True, "prepare_threshold": 0},30 check=AsyncConnectionPool.check_connection,31 open=False, 32 ) as pool:33 await pool.open(wait=True, timeout=15)34 checkpointer = AsyncPostgresSaver(pool)35 await checkpointer.setup()36 app.state.pool = pool37 app.state.agent = build_agent(checkpointer)38 yield39 40app = FastAPI(41 title="Text2SQL Agent API",42 description="A production-grade backend powering LangGraph agent.",43 version="1.0.0",44 lifespan=lifespan)45 46class ChatRequest(BaseModel):47 connection_url : str = Field(...)48 message: str = Field(...)49 user_id: str = Field(...)50 thread_id: str = Field(...)51 52class ChatResponse(BaseModel):53 status: str54 thread_id: str55 response: str56 57class UploadRequest(BaseModel) :58 connection_url : str = Field(...)59 user_id : str = Field(...)60 61@app.post("/upload")62async def upload_url(request : UploadRequest):63 await asyncio.to_thread(create_embeddings , request.connection_url , request.user_id)64 65 return {66 "status": "success",67 "message": "You can now chat with the agent."68 }69 70@app.post("/chat",response_model=ChatResponse)71async def chat_endpoint(request: ChatRequest):72 agent = app.state.agent 73 74 config = {"configurable": {"thread_id": request.thread_id}}75 76 try:77 result = await agent.ainvoke(78 {"messages": [{"role": "user", "content": request.message}]},79 config=config,80 context=AgentContext(user_id=request.user_id , connection_url=request.connection_url),81 )82 except Exception:83 logger.exception("Agent processing failed")84 raise HTTPException(status_code=500, detail="Agent Processing Error!")85 output_messages = result.get("messages", [])86 if not output_messages:87 raise HTTPException(status_code=500, detail="No messages returned from the agent.")88 89 return ChatResponse(90 status="success",91 thread_id=request.thread_id,92 response=output_messages[-1].content,93 )