FastAPI Integration
FlexiQRouter — pre-built APIRouter with job status, SSE progress, dead letter management.
FlexiQRouter — pre-built APIRouter with job status, SSE progress, dead letter management.
flexiq provides a pre-built APIRouter for FastAPI with endpoints for
job management, progress streaming via SSE, and dead letter queue
operations.
pip install flexiq[fastapi]This installs fastapi and pydantic as extras.
from fastapi import FastAPI
from flexiq import Queue
from flexiq.contrib.fastapi import FlexiQRouter
queue = Queue(db_path="myapp.db")
@queue.task()
def process_data(payload: dict) -> str:
return "done"
app = FastAPI()
app.include_router(FlexiQRouter(queue), prefix="/tasks")Run with:
uvicorn myapp:app --reload| Method | Path | Description |
|---|---|---|
GET | /stats | Queue statistics |
GET | /stats/queues | Per-queue statistics |
GET | /jobs/{job_id} | Job status, progress, and metadata |
GET | /jobs/{job_id}/errors | Error history for a job |
GET | /jobs/{job_id}/result | Job result (optional ?timeout=N for blocking) |
GET | /jobs/{job_id}/progress | SSE stream of progress updates |
POST | /jobs/{job_id}/cancel | Cancel a pending job |
GET | /dead-letters | List dead letter entries (paginated) |
POST | /dead-letters/{dead_id}/retry | Re-enqueue a dead letter |
GET | /health | Liveness check |
GET | /readiness | Readiness check |
GET | /resources | Worker resource status |
FlexiQRouter accepts options to control which routes are registered,
how results are serialized, and page sizes:
from fastapi import Depends, Header, HTTPException
from flexiq.contrib.fastapi import FlexiQRouter
def require_api_key(x_api_key: str = Header(...)):
if x_api_key != "secret":
raise HTTPException(status_code=403)
app.include_router(
FlexiQRouter(
queue,
include_routes={"stats", "jobs", "dead-letters", "retry-dead"},
dependencies=[Depends(require_api_key)],
sse_poll_interval=1.0,
result_timeout=5.0,
default_page_size=25,
max_page_size=200,
result_serializer=lambda v: v if isinstance(v, (str, int, float, bool, type(None))) else str(v),
),
prefix="/tasks",
)| Parameter | Type | Default | Description |
|---|---|---|---|
include_routes | set[str] | None | None | If set, only register these route names. Cannot be combined with exclude_routes. |
exclude_routes | set[str] | None | None | If set, skip these route names. Cannot be combined with include_routes. |
dependencies | Sequence[Depends] | None | None | FastAPI dependencies applied to every route (e.g. auth). |
sse_poll_interval | float | 0.5 | Seconds between SSE progress polls. |
result_timeout | float | 1.0 | Default timeout for non-blocking result fetch. |
default_page_size | int | 20 | Default page size for paginated endpoints. |
max_page_size | int | 100 | Maximum allowed page size. |
result_serializer | Callable[[Any], Any] | None | None | Custom result serializer. Receives any value, must return a JSON-serializable value. |
Valid route names: stats, jobs, job-errors, job-result,
job-progress, cancel, dead-letters, retry-dead, health,
readiness, resources, queue-stats.
The /jobs/{job_id}/result endpoint supports an optional timeout query
parameter (0–300 seconds). When timeout > 0, the request blocks until
the job completes or the timeout elapses:
# Non-blocking (default)
curl http://localhost:8000/tasks/jobs/01H5K6X.../result
# Block up to 30 seconds for the result
curl http://localhost:8000/tasks/jobs/01H5K6X.../result?timeout=30Stream real-time progress for a running job using Server-Sent Events:
import httpx
with httpx.stream("GET", "http://localhost:8000/tasks/jobs/01H5K6X.../progress") as r:
for line in r.iter_lines():
print(line)
# data: {"progress": 25, "status": "running"}
# data: {"progress": 50, "status": "running"}
# data: {"progress": 100, "status": "completed"}From the browser:
const source = new EventSource("/tasks/jobs/01H5K6X.../progress");
source.onmessage = (event) => {
const data = JSON.parse(event.data);
console.log(`Progress: ${data.progress}%`);
if (data.status === "completed" || data.status === "failed") {
source.close();
}
};The stream sends a JSON event every 0.5 seconds while the job is active, then a final event when the job reaches a terminal state.
All endpoints return validated Pydantic models with clean OpenAPI docs. You can import them for type-safe client code:
from flexiq.contrib.fastapi import (
StatsResponse,
JobResponse,
JobErrorResponse,
JobResultResponse,
CancelResponse,
DeadLetterResponse,
RetryResponse,
)from fastapi import FastAPI, Header, HTTPException, Depends
from flexiq import Queue, current_job
from flexiq.contrib.fastapi import FlexiQRouter
queue = Queue(db_path="myapp.db")
@queue.task()
def resize_image(image_url: str, sizes: list[int]) -> dict:
results = {}
for i, size in enumerate(sizes):
results[size] = do_resize(image_url, size)
current_job.update_progress(int((i + 1) / len(sizes) * 100))
return results
async def require_auth(authorization: str = Header(...)):
if not authorization.startswith("Bearer "):
raise HTTPException(401)
app = FastAPI(title="Image Service")
app.include_router(
FlexiQRouter(queue, dependencies=[Depends(require_auth)]),
prefix="/tasks",
tags=["tasks"],
)
# Start worker in a separate process:
# flexiq worker --app myapp:queue# Check job status
curl http://localhost:8000/tasks/jobs/01H5K6X... \
-H "Authorization: Bearer mytoken"
# Stream progress
curl -N http://localhost:8000/tasks/jobs/01H5K6X.../progress \
-H "Authorization: Bearer mytoken"
# Block for result (up to 60s)
curl http://localhost:8000/tasks/jobs/01H5K6X.../result?timeout=60 \
-H "Authorization: Bearer mytoken"