jatin gyass
initial commit
2eef9ea
Raw
History Blame Contribute Delete
2.04 kB
"""Workflow execution endpoints."""
from typing import Optional
from fastapi import APIRouter, HTTPException, BackgroundTasks
from fastapi.responses import StreamingResponse
from pydantic import BaseModel, Field
from ...agents.orchestrator import run_workflow, stream_workflow
from ...state.graph_state import create_initial_state
router = APIRouter()
class WorkflowRequest(BaseModel):
"""Workflow execution request."""
task: str = Field(..., description="Task to execute")
stream: bool = Field(False, description="Stream results in real-time")
class WorkflowResponse(BaseModel):
"""Workflow execution response."""
task_id: str
status: str
result: Optional[dict] = None
@router.post("/execute", response_model=WorkflowResponse)
async def execute_workflow(request: WorkflowRequest, background_tasks: BackgroundTasks):
"""
Execute workflow synchronously or asynchronously.
Args:
request: Workflow request payload
background_tasks: Background task runner
Returns:
WorkflowResponse: Workflow execution result
"""
try:
state = create_initial_state(request.task)
task_id = state["task_id"]
if request.stream:
# Return streaming response
return StreamingResponse(
stream_workflow(state), media_type="application/x-ndjson"
)
# Run in background
background_tasks.add_task(run_workflow, state)
return {"task_id": task_id, "status": "queued", "result": None}
except Exception as e:
raise HTTPException(status_code=400, detail=str(e)) from e
@router.get("/{task_id}/status")
async def get_workflow_status(task_id: str):
"""
Get workflow execution status.
Args:
task_id: Workflow task ID
Returns:
dict: Workflow status information
"""
# Fetch from persistent storage
return {
"task_id": task_id,
"status": "running",
"progress": 50,
"current_agent": "executor",
}