Living Document Notice
Published 2026-09-17. The evolving architecture, revisions, and connected notes for this dispatch live in the Stax Digital Garden.

Handling Staged Pipeline Ingestion in Python and FastAPI

Handling Staged Pipeline Ingestion in Python and FastAPI: High-contrast P4 paper white and amber dual-trace vector CRT macro showing five-stage ingestion pipeline with amber telemetry gates

Summary

Engineering retrospective on backend task claiming and stage isolation. Avoiding cross-job race conditions when processing concurrent link captures locally.

Under the hood, Galley runs capture, payload extraction and confidence scoring in distinct stages. When a user pastes three recipe links in quick succession, or imports a batch of bookmarks, the local backend must process these network-bound and CPU-bound tasks smoothly without blocking the local user interface.

We built Galley’s ingestion engine using Python and FastAPI, backed by an atomic SQLite task queue that prevents race conditions and manages memory cleanly.

Pipeline Architecture and State Flow

The ingestion pipeline separates work into distinct sequential stages:

[ POST /api/ingest ]
        │
        ▼ (Atomic Enqueue)
┌───────────────────────────────────────────────┐
│ SQLite Task Queue: status = 'pending'         │
└───────────────────────┬───────────────────────┘
                        │
                        ▼ (Worker Claim: status = 'processing')
┌───────────────────────────────────────────────┐
│ Ingestion Worker                              │
│   1. Fetch DOM with timeout & circuit breaker │
│   2. Extract schema JSON-LD and microdata     │
│   3. Parse ingredient AST and score fields    │
│   4. Write to staged_inbox                    │
└───────────────────────┬───────────────────────┘
                        │
                        ▼ (Atomic Complete: status = 'staged')
┌───────────────────────────────────────────────┐
│ Ready for Human Review in Review Inbox        │
└───────────────────────────────────────────────┘

Atomic Job Claiming in SQLite

To ensure multiple local worker coroutines never process the same URL simultaneously, Galley uses atomic updates with row-level state checks:

import sqlite3
from typing import Any
 
def claim_next_job(db: sqlite3.Connection, worker_id: str) -> dict[str, Any] | None:
    # Claim the oldest pending job atomically using immediate transaction
    with db:
        cursor = db.cursor()
        cursor.execute('''
            UPDATE ingestion_queue
            SET status = 'processing',
                worker_id = :worker_id,
                updated_at = CURRENT_TIMESTAMP
            WHERE id = (
                SELECT id FROM ingestion_queue
                WHERE status = 'pending'
                ORDER BY created_at ASC
                LIMIT 1
            )
            RETURNING id, url, raw_text;
        ''', {"worker_id": worker_id})
        row = cursor.fetchone()
        if row:
            return {"id": row[0], "url": row[1], "raw_text": row[2]}
        return None

The RETURNING clause introduced in modern SQLite versions allows atomic claim-and-fetch in a single round trip, eliminating the need for complex distributed locking mechanisms.

Graceful Task Cancellation and Timeouts

Web scrapers often stall when external servers throttle connections or hang during payload delivery. Galley’s worker loop wraps every network request inside an explicit timeout with circuit breaking:

import asyncio
import httpx
 
async def fetch_with_timeout(url: str, timeout_seconds: float = 8.0) -> str:
    limits = httpx.Limits(max_keepalive_connections=5, max_connections=10)
    async with httpx.AsyncClient(limits=limits, follow_redirects=True) as client:
        try:
            response = await client.get(url, timeout=timeout_seconds)
            response.raise_for_status()
            return response.text
        except (httpx.TimeoutException, httpx.HTTPError) as exc:
            raise RuntimeError(f"Ingestion fetch failed for {url}: {exc}")

If a URL fails or times out, the task state transitions to failed with the error recorded in SQLite. The background worker remains operational, immediately picking up the next pending job.

Concurrency Performance Metrics

By running extraction in an asynchronous background thread pool, the local FastAPI server maintains responsive sub-10ms UI interaction latencies even under ingestion load:

ScenarioConcurrent LinksMax Memory OverheadAverage Ingestion TimeUI Response Latency
Single link capture1 URL32 MB RSS1.1s< 5 ms
Small batch import5 URLs48 MB RSS3.4s (parallel fetch)< 8 ms
Bulk bookmark ingest25 URLs78 MB RSS14.2s (rate-limited)< 12 ms

Isolating worker tasks into discrete stages ensures that a hung network connection or malformed HTML document will never crash the core application or corrupt user data.


  • Directus Target: galley
  • Garden Source Reference: FastAPI Backend, Pipeline Workers, Concurrency Models, MOC - Culinary & Domain Workspaces, MOC - The Kitchen Chaos Factor, MOC - Bosun PKM Tools