Bigger Applications
A larger FastAPI project spreads its code over several packages and modules. This guide shows where quiv goes in such a project. It covers one shared scheduler instance, tasks in separate files, API endpoints that control those tasks, progress updates over a WebSocket, and cancellation that lets a handler finish cleanly.
Project structure
myapp/
├── main.py # FastAPI app, lifespan, WebSocket
├── scheduler.py # Quiv instance (shared singleton)
├── tasks/
│ ├── __init__.py
│ ├── cleanup.py # DB cleanup task
│ └── report.py # Report generation task
└── routes/
├── __init__.py
└── tasks.py # Task control endpoints
1) Create the scheduler instance
Define the Quiv instance in its own module, so that every other file can import it. Do not call start() here. The FastAPI lifespan calls it.
# myapp/scheduler.py
from quiv import Quiv
scheduler = Quiv(
pool_size=4,
history_retention_seconds=7200,
timezone="America/New_York",
)
Quiv finds the asyncio event loop on first use, not at startup. This code therefore works at module level, before FastAPI or uvicorn creates a loop.
2) Create a logging context (optional)
Create a logging context that holds the trace_id, so that later log lines can carry it.
# myapp/config/logging_context.py
import asyncio
from contextvars import ContextVar, Token
from functools import wraps
import uuid
# Define a ContextVar to hold the Trace ID
trace_id_var: ContextVar[str | None] = ContextVar("trace_id", default=None)
def get_trace_id() -> str | None:
"""Returns the current Trace ID."""
return trace_id_var.get()
def get_new_trace_id() -> str:
"""Generates and returns a new Trace ID without setting it in the context."""
return str(uuid.uuid4())
def generate_trace_id(trace_id: str | None = None) -> Token:
"""Generates a new Trace ID and sets it in the context."""
if trace_id is None:
trace_id = str(uuid.uuid4())
return trace_id_var.set(trace_id)
def clear_trace_id(token: Token) -> None:
"""Clears the Trace ID from the context."""
trace_id_var.reset(token)
def with_logging_context(func):
"""Decorator to add a Trace ID to the logging context for the duration of the function call.
Args:
func: The function to decorate.
Returns:
The decorated function with Trace ID management.
"""
@wraps(func)
async def async_wrapper(*args, **kwargs):
# 1. Extract or Create
trace_id = kwargs.get("job_id") or kwargs.get("trace_id")
token = generate_trace_id(trace_id) # Falls back to uuid4() internally
try:
return await func(*args, **kwargs)
finally:
clear_trace_id(token)
@wraps(func)
def sync_wrapper(*args, **kwargs):
# 1. Extract or Create
trace_id = kwargs.get("job_id") or kwargs.get("trace_id")
token = generate_trace_id(trace_id)
try:
return func(*args, **kwargs)
finally:
clear_trace_id(token)
# Return the appropriate wrapper based on the source function
return async_wrapper if asyncio.iscoroutinefunction(func) else sync_wrapper
Make your logging adapter, or your logging handler, read the trace_id and add it to each line.
For complete code that uses a logging context, stop events, and progress hooks, read Trailarr.
3) Define tasks in separate files
Each task file imports the shared scheduler and registers its task with add_task(). A task is a plain function, either sync or async.
Cleanup task (sync, with stop event)
# myapp/tasks/cleanup.py
import logging
import threading
import time
from config.logging_context import with_logging_context
logger = logging.getLogger(__name__)
@with_logging_context
def cleanup_stale_records(
days: int,
job_id: str | None = None,
stop_event: threading.Event | None = None,
):
"""Delete records older than `days` from the database."""
# quiv injects job_id and with_logging_context decorator stores it as trace_id
# logs handler will get it using get_trace_id and adds it to all logs
# so logs from the task will be logged with that trace_id
# Attach job_id as trace context for this run
batches = 10
for batch in range(1, batches + 1):
if stop_event and stop_event.is_set():
logger.info("Cleanup cancelled at batch %d/%d", batch, batches)
return
# ... delete a batch of old records ...
time.sleep(1) # simulate work
logger.info("Cleanup finished: processed %d batches", batches)
Report task (sync, with stop event and progress hook)
# myapp/tasks/report.py
import logging
import threading
import time
from typing import Callable
from config.logging_context import with_logging_context
logger = logging.getLogger(__name__)
@with_logging_context
def generate_report(
report_type: str,
job_id: str | None = None,
stop_event: threading.Event | None = None,
progress_hook: Callable | None = None,
):
"""Generate a report with progress updates."""
steps = 5
for step in range(1, steps + 1):
if stop_event and stop_event.is_set():
logger.info("Report generation cancelled at step %d/%d", step, steps)
return
# ... do a chunk of report work ...
time.sleep(2) # simulate work
if progress_hook:
progress_hook(
step=step,
total=steps,
report_type=report_type,
)
logger.info("Report '%s' generated successfully", report_type)
4) Register tasks and wire up the lifespan
The main module registers the tasks, starts the scheduler at startup, and shuts it down at the end. Set up the progress callbacks for the WebSocket here as well.
# myapp/main.py
import asyncio
import logging
from contextlib import asynccontextmanager
from fastapi import FastAPI, WebSocket, WebSocketDisconnect
from quiv import Event
from quiv.models import Task, Job
from myapp.scheduler import scheduler
from myapp.tasks.cleanup import cleanup_stale_records
from myapp.tasks.report import generate_report
logger = logging.getLogger(__name__)
# ---------- WebSocket connection manager ----------
class ConnectionManager:
"""Track active WebSocket connections for progress broadcasts."""
def __init__(self):
self.connections: list[WebSocket] = []
async def connect(self, websocket: WebSocket):
await websocket.accept()
self.connections.append(websocket)
def disconnect(self, websocket: WebSocket):
self.connections.remove(websocket)
async def broadcast(self, message: dict):
for ws in self.connections:
try:
await ws.send_json(message)
except Exception:
pass
ws_manager = ConnectionManager()
# ---------- Progress callback ----------
async def on_report_progress(**payload):
"""Forward task progress to all connected WebSocket clients."""
logger.info("Report progress: %s", payload)
await ws_manager.broadcast({"event": "progress", "data": payload})
# ---------- Lifespan ----------
async def on_job_event(event: Event, task: Task, job: Job) -> None:
"""Forward job lifecycle events to WebSocket clients."""
payload: dict = {"event": event.value, "task": task.task_name, "job_id": job.id}
if job.duration_seconds is not None:
payload["duration_seconds"] = job.duration_seconds
if job.error_message is not None:
payload["error"] = job.error_message
await ws_manager.broadcast(payload)
@asynccontextmanager
async def lifespan(app: FastAPI):
# Register event listeners
scheduler.add_listener(Event.JOB_STARTED, on_job_event)
scheduler.add_listener(Event.JOB_COMPLETED, on_job_event)
scheduler.add_listener(Event.JOB_FAILED, on_job_event)
# Register and start tasks
scheduler.add_task(
task_name="db-cleanup",
func=cleanup_stale_records,
interval=3600,
kwargs={"days": 30},
)
scheduler.add_task(
task_name="weekly-report",
func=generate_report,
interval=604800,
delay=10,
kwargs={"report_type": "weekly-summary"},
progress_callback=on_report_progress,
)
scheduler.start()
yield
scheduler.shutdown()
app = FastAPI(lifespan=lifespan)
# ---------- WebSocket endpoint ----------
@app.websocket("/ws/progress")
async def progress_websocket(websocket: WebSocket):
await ws_manager.connect(websocket)
try:
while True:
await websocket.receive_text()
except WebSocketDisconnect:
ws_manager.disconnect(websocket)
5) API endpoints for runtime control
A separate router imports the same scheduler instance and publishes the endpoints that manage the tasks.
# myapp/routes/tasks.py
from fastapi import APIRouter, HTTPException
from myapp.scheduler import scheduler
router = APIRouter(prefix="/tasks", tags=["tasks"])
@router.post("/{task_id}/run")
def run_task_now(task_id: str):
"""Trigger a scheduled task to run immediately."""
try:
count = scheduler.run_task_immediately(task_id)
except Exception as e:
raise HTTPException(status_code=404, detail=str(e))
return {"queued": count}
@router.post("/{task_id}/pause")
def pause_task(task_id: str):
"""Pause a task by id."""
try:
scheduler.pause_task(task_id)
except Exception as e:
raise HTTPException(status_code=404, detail=str(e))
return {"status": "paused"}
@router.post("/{task_id}/resume")
def resume_task(task_id: str, delay: int = 0):
"""Resume a paused task, optionally with a delay."""
try:
scheduler.resume_task(task_id, delay=delay)
except Exception as e:
raise HTTPException(status_code=404, detail=str(e))
return {"status": "resumed"}
@router.get("/{task_id}")
def get_task(task_id: str):
"""Get a single task by id."""
try:
return scheduler.get_task(task_id)
except Exception as e:
raise HTTPException(status_code=404, detail=str(e))
@router.get("/")
def list_tasks():
"""List all scheduled tasks."""
return scheduler.get_all_tasks()
@router.get("/jobs")
def list_jobs(status: str | None = None):
"""List jobs, optionally filtered by status."""
return scheduler.get_all_jobs(status=status)
@router.post("/jobs/{job_id}/cancel")
def cancel_job(job_id: str):
"""Cancel a running job."""
cancelled = scheduler.cancel_job(job_id)
if not cancelled:
raise HTTPException(status_code=404, detail="Job not found or not running")
return {"status": "cancelled"}
An endpoint that starts work and wants to report how it went adds a one-off task and waits for its job. The endpoint holds a task_id, not a job id, so it waits with await_task:
@router.post("/refresh")
async def refresh_now():
"""Refresh the index now, and report the result."""
task_id = scheduler.add_task(task_name="refresh", func=refresh_index, run_once=True)
job = await scheduler.await_task(task_id, timeout=120)
if job.status != "completed":
raise HTTPException(status_code=502, detail=job.error_message or job.status)
return {"duration_seconds": job.duration_seconds}
The job comes back finalized, so the endpoint answers with the real outcome instead of "queued". Bound the wait with timeout. It raises the builtin TimeoutError when the deadline passes, and the job keeps running.
Task IDs in your API
add_task() returns a task_id, a UUID string. Store it in the state of your application, or return it to the client. Every later operation uses it: pause_task, resume_task, run_task_immediately, and remove_task.
Task and Job are SQLModel objects, so FastAPI serializes them as they are. You convert nothing by hand. Every datetime field is an aware UTC value: next_run_at, started_at, and ended_at. The JSON therefore ends each one with +00:00, which a browser can read and show in the timezone of the user.
Register the router in your app:
# add to myapp/main.py
from myapp.routes.tasks import router as tasks_router
app.include_router(tasks_router)
Run the full example
The examples/fastapi_app directory holds a complete version of this application that you can run. From the root of the repository:
uv run uvicorn examples.fastapi_app.main:app --reload
Then open http://127.0.0.1:8000/docs to try the API, or connect to ws://127.0.0.1:8000/ws/progress to watch the progress updates.
Key takeaways
- Single instance, shared everywhere. Create
Quivin one module and import it wherever needed. This avoids multiple schedulers and duplicate DB files. - Module-level init is safe.
Quiv()does not require a running asyncio loop at creation time. The event loop is resolved lazily when progress callbacks fire. - Lifespan owns the lifecycle. Call
start()andshutdown()in the FastAPI lifespan so the scheduler is tied to the app process. - Tasks are plain functions. Define them anywhere. They only need
stop_eventandprogress_hookin their signature if they want cancellation or progress support. See Cancellation and Progress Callbacks for detailed guides. - Progress goes through WebSocket. Async progress callbacks run on FastAPI's event loop, so they can broadcast to WebSocket clients directly. See Progress Callbacks for dispatch details.
- Event listeners for observability. Use
add_listener()to react to task and job lifecycle events. Async listeners run on the main loop, so they can broadcast to WebSocket clients alongside progress callbacks. See Event Listeners for the full event list.