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
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 NoneThe 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:
| Scenario | Concurrent Links | Max Memory Overhead | Average Ingestion Time | UI Response Latency |
|---|---|---|---|---|
| Single link capture | 1 URL | 32 MB RSS | 1.1s | < 5 ms |
| Small batch import | 5 URLs | 48 MB RSS | 3.4s (parallel fetch) | < 8 ms |
| Bulk bookmark ingest | 25 URLs | 78 MB RSS | 14.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