Skip to content

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 Quiv in 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() and shutdown() in the FastAPI lifespan so the scheduler is tied to the app process.
  • Tasks are plain functions. Define them anywhere. They only need stop_event and progress_hook in 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.