Long-running jobs and ingestion¶
This guide shows how to run background work with msflib-ingestion: a Celery worker that executes jobs, a FastAPI process that enqueues them, a way to develop without a broker, how to add your own pipeline kind, and how failures, retries and dead letters behave. For the concepts and settings, see the ingestion module page.
The model in one paragraph: your FastAPI process writes a row to the IngestionJob table and, once that transaction commits, sends the job id to Celery. A worker process loads the row, looks up the pipeline kind registered for job.pipeline_kind, and runs its stages. The row, not the broker message, is the source of truth for status, so you can inspect and recover jobs from the database.
1. What you need¶
- A Redis instance (or another Celery broker) reachable from both processes, set through
INGESTION.BROKER_URL(defaultredis://localhost:6379/0). - The same database and the same ingestion tables in both processes. Import
msflib.ingestion.modelsso the tables are on your metadata; the job table has a foreign key to the tenant table, so tenancy's tables must exist too. IngestionSettingsin your settings class:
from msflib.core.config import SettingsBase
from msflib.ingestion.config import IngestionSettings
class AppSettings(IngestionSettings, SettingsBase):
...
2. The worker process¶
The worker is a separate entrypoint. It registers every pipeline kind a job might reference, then builds the Celery app. Celery finds the module-level app.
# app/ingestion_worker.py
from msflib.ingestion.celery import build_celery_app
from sqlmodel import Session
from app.core.config import settings
from app.core.db import engine
from app.pipelines import register_my_pipeline_kinds
def session_factory() -> Session:
return Session(engine)
register_my_pipeline_kinds() # fills default_pipeline_registry in this process
app = build_celery_app(settings, session_factory=session_factory)
Start it with:
celery -A app.ingestion_worker worker --loglevel=info
Notes:
session_factoryis any zero-argument callable returning a context-managerSession. A plainSession(engine)works, becauseSessionis a context manager.- Registration is per process and is not persisted. A kind registered only in the web process makes the worker raise
KeyError: no pipeline registered for kindwhen it picks up the job. - Build exactly one Celery app per process. Building several with the same task name (
msflib.ingestion.execute_job) in one process corrupts task state. - The app is configured with late acknowledgement and a prefetch multiplier of 1, so a job that was running when a worker died is redelivered.
- Modules that ship pipeline kinds export their own registration functions. The reference worker in the repository (
testsite/app/ingestion_worker.py) registerstext_documentfrommsflib.documents, andknowledge_extractionandknowledge_recordfrommsflib.knowledge. See Knowledge.
3. The FastAPI process¶
The web process only enqueues. Build one producer handle and share it:
from msflib.ingestion.celery import build_dispatch_task
from msflib.ingestion.router import router as create_ingestion_router
from msflib.ingestion.services.dispatch import IngestionDispatchService
from msflib.ingestion.services.queue import IngestionSubjectQueue
dispatch = IngestionDispatchService(task=build_dispatch_task(settings))
note_queue = IngestionSubjectQueue(
dispatch=dispatch,
pipeline_kind="note_index",
enabled=True, # pass your own settings flag here
)
build_dispatch_task needs no session factory and never imports your pipeline code. It sends by task name to DEFAULT_QUEUE.
Enqueue from request code with the queue. Dispatch is deferred until your transaction commits, so a rollback queues nothing:
@router.post("/notes")
def create_note(payload: NoteCreate, session: Session = Depends(get_session)):
note = Note(**payload.model_dump())
session.add(note)
session.flush() # assigns note.id
note_queue.enqueue_snapshot(
session, subject_type="note", snapshot=note_snapshot(note)
)
session.commit() # the job row is committed, then dispatched
return note
router here is your own APIRouter, and NoteCreate and get_session come from your host.
Or let the model do it for you: sync_model_to_queue(Note, note_queue, subject_type="note", snapshot=note_snapshot, emitter=app_emitter) queues a job on every ModelAction write to Note. Create app_emitter once when you create the app (from msflib.eventbus import bind_app_emitter, then app_emitter = bind_app_emitter(app)). The hook fires only while the app emitter is active: in requests (including FastAPI BackgroundTasks), or in your own code wrapped in with use_app_emitter(app):. Writes from workers, CLI scripts and threads outside that block queue nothing, so call note_queue.enqueue_snapshot there. Gate this registration on the same settings flag as the queue. note_snapshot is defined in section 5.
Mount the management router in the same process (section 7). It validates pipeline_kind against a registry in the web process, so call your registration function there as well.
4. Running without a broker¶
There is no built-in inline dispatcher: build_dispatch_task always sends through Celery. For development and tests you have two options that need no Redis.
Record dispatches only. msflib.ingestion.testing.FakeDispatchTask implements .delay(job_id) and records the ids. Use it to assert that your code enqueued the right jobs. start_stage_context(session, job, settings) from the same module starts an attempt and returns the StageContext, so you can call a stage function directly.
Run jobs in the same process. Build the worker app as in section 2, switch Celery to eager mode, and use the registered task as the dispatch handle. Jobs then run synchronously inside delay(), after your commit, with the full retry and dead-letter behaviour. This runs as written. It seeds tenant 1 after create_all because Note.tenant_id points at it: SQLite does not enforce the foreign key, but Postgres does.
import os
import tempfile
from msflib.core.config import SettingsBase
from msflib.ingestion.celery import TASK_NAME, build_celery_app
from msflib.ingestion.config import IngestionSettings
from msflib.ingestion.contracts.subject import IngestionSubjectSnapshot
from msflib.ingestion.errors import TerminalStageError
from msflib.ingestion.models import IngestionDeadLetter, IngestionJob
from msflib.ingestion.services.dispatch import IngestionDispatchService
from msflib.ingestion.services.lifecycle import ingestion_subject_stage
from msflib.ingestion.services.queue import IngestionSubjectQueue
from msflib.ingestion.services.registry import PipelineKindRegistry
from msflib.models import ModelBase
from msflib.tenancy.models.tenant import Tenant # noqa: F401 (registers the tenant table)
from msflib.tenancy.resolver import resolve_default_tenant_id
from sqlmodel import Session, SQLModel, create_engine, select
class Note(ModelBase, table=True):
body: str = ""
version: int = 1
tenant_id: int = 1
workspace_id: int | None = None
def load(context):
return context.session.get(Note, int(context.job.subject_id))
def snapshot(note):
return IngestionSubjectSnapshot(
subject_id=str(note.id),
tenant_id=note.tenant_id,
workspace_id=note.workspace_id,
version=note.version,
)
def process(context):
if context.subject.body == "bad":
raise TerminalStageError("cannot index this note")
print("indexed", context.subject.body)
class Settings(IngestionSettings, SettingsBase):
pass
settings = Settings()
registry = PipelineKindRegistry()
registry.register_kind(
"note_index",
ingestion_subject_stage(name="note_index", load=load, snapshot=snapshot, process=process),
)
# A file database: the eager task opens its own session while the producer's is still open.
engine = create_engine(f"sqlite:///{os.path.join(tempfile.mkdtemp(), 'app.db')}")
SQLModel.metadata.create_all(engine)
with Session(engine) as session:
resolve_default_tenant_id(session) # seeds tenant 1
session.commit()
celery_app = build_celery_app(settings, session_factory=lambda: Session(engine), registry=registry)
celery_app.conf.task_always_eager = True
celery_app.conf.task_eager_propagates = False
queue = IngestionSubjectQueue(
dispatch=IngestionDispatchService(task=celery_app.tasks[TASK_NAME]),
pipeline_kind="note_index",
)
with Session(engine) as session:
for body in ("good", "bad"):
note = Note(body=body)
session.add(note)
session.commit()
queue.enqueue_snapshot(session, subject_type="note", snapshot=snapshot(note), commit=True)
for job in session.exec(select(IngestionJob)).all():
session.refresh(job)
print(job.subject_id, job.status.value, job.attempts, job.error)
print([(d.job_id, d.reason) for d in session.exec(select(IngestionDeadLetter)).all()])
Output:
indexed good
1 done 1 None
2 failed 1 cannot index this note
[(2, 'TerminalStageError: cannot index this note')]
Eager mode in a request uses the app emitter, while a real worker uses the default emitter, so test event listeners in the process they will run in. In eager mode outside a request, only listeners on the default emitter fire. Use this only for development. Eager mode has no concurrency, no redelivery and no real backoff delays. The registry is passed explicitly here so the example does not touch default_pipeline_registry.
5. Adding a new pipeline kind¶
A kind needs a subject that can be loaded by id, a snapshot of it, and the work to do. Using the Note model from the example above:
1. Load and snapshot. load(context) returns the subject for context.job.subject_id, or None when it is gone. snapshot(subject) returns an IngestionSubjectSnapshot with the subject's id (as a string), tenant, workspace, a version (strongly recommended, see below) and live.
2. Process and clean up. process(context) handles a live subject. cleanup(context) handles a gone subject or a delete job. Both receive an IngestionSubjectContext with .session, .subject, .snapshot, .job, .reason and .settings. Do not commit inside them. They run in a savepoint, and the worker commits when the whole job succeeds, so an exception rolls back everything the step wrote.
def note_snapshot(note):
return IngestionSubjectSnapshot(
subject_id=str(note.id),
tenant_id=note.tenant_id,
workspace_id=note.workspace_id,
version=note.version,
live=not note.deleted,
)
3. Build and register the stage.
from msflib.ingestion.services.lifecycle import ingestion_subject_stage
from msflib.ingestion.services.registry import default_pipeline_registry
def register_note_pipeline():
default_pipeline_registry.register_kind(
"note_index",
ingestion_subject_stage(
name="note_index",
load=load,
snapshot=note_snapshot,
process=process,
cleanup=cleanup,
completed_event="notes.indexed", # optional; fires after the job commits
),
description="Index notes",
)
Call register_note_pipeline() in the worker entrypoint and in the web process. Event listeners for completed_event belong in the worker process, registered on the worker's emitter, AppEmitter(get_emitter()) (from msflib.eventbus); see Ingestion. Registering a kind twice raises unless you pass replace=True.
4. Slow work goes in prepare. prepare= runs before the per-subject lock is taken, so an LLM call or remote fetch does not hold it. process receives the result as context.prepared. If the subject changes while prepare runs, the attempt is retried.
5. Queue it. Create an IngestionSubjectQueue(dispatch=..., pipeline_kind="note_index") in the web process and call enqueue_snapshot or sync_model_to_queue, as in section 3.
Versions matter. A job carries the target_version it was queued for. When the stage finds a newer subject version, it supersedes the job instead of overwriting newer work with older. Without a version, an edit made while a job is running can deduplicate into the running job and be missed.
If your kind does not fit the load-and-process shape, register plain Stage callables instead: def my_stage(context: StageContext) -> StageContext, and pass them to registry.register_kind("my_kind", my_stage). Then you own the error handling. See Ingestion for the error types and emit_after_success.
6. Retries, supersession and dead letters¶
What happens to a job depends on what the stage raises.
| Outcome | Raised | Job status | Retries | Dead letter |
|---|---|---|---|---|
| Success | nothing | done |
no | no |
| Superseded | SupersededStageError |
superseded |
no | no |
| Transient failure | RetryableStageError, or a plain exception classified as retryable |
queued with next_retry_at, then failed after the last attempt |
yes, up to MAX_ATTEMPTS total |
only after the last attempt |
| Permanent failure | TerminalStageError, or a plain exception classified as terminal |
failed |
no | yes |
Each attempt adds an IngestionJobRun row with its status and error text, so the history of a flaky job is queryable.
Retry delay starts at RETRY_BACKOFF_SECONDS and doubles each attempt up to RETRY_BACKOFF_MAX_SECONDS; Celery adds jitter. Tune the limits with INGESTION.MAX_ATTEMPTS, INGESTION.RETRY_BACKOFF_SECONDS and INGESTION.RETRY_BACKOFF_MAX_SECONDS. They apply to every pipeline kind and are read when each task runs, so a settings change needs no worker restart. Inside ingestion_subject_stage, an error identical to the previous attempt's is treated as terminal rather than retried again.
When a job is terminally failed, a row is written to IngestionDeadLetter with the subject, scope, attempts and error. From there an operator can:
- replay it (
POST /ingestion/dead-letters/{id}/replay), which enqueues a new job with the original reason and target version. If an identical job is already active, the replay reuses it and does not dispatch twice. - discard it (
POST /ingestion/dead-letters/{id}/discard) to close it without action.
Alternatively POST /ingestion/jobs/{id}/retry requeues a failed job directly. The usual cause of a dead letter you can fix in code (an unregistered mapper, a missing record source) is a worker without the right registrations; fix the worker, then replay.
If the broker is down when a committed job is dispatched, the job is marked failed with a "Celery dispatch failed" error so it is visible rather than stuck in queued. Retry it once the broker is back.
7. Monitoring and operating¶
Mount the router once, in the web process:
from msflib.ingestion.celery import build_celery_app
web_celery_app = build_celery_app(settings, session_factory=session_factory) # or the worker module's `app`
app.include_router(
create_ingestion_router(
get_session=get_session,
settings=settings,
mutation_role_check=require_admin,
get_current_tenant=get_current_tenant,
get_current_workspace=get_current_workspace,
celery_control=web_celery_app, # optional, enables revoking in-flight jobs
prefix="/ingestion",
)
)
Read routes (GET /jobs, GET /jobs/{id}, GET /status, GET /dead-letters) are available to anyone your tenant and workspace dependencies let through, so add your own authentication in front of the router. The mutating routes also depend on mutation_role_check, which has no default. celery_control is any Celery app instance with .control.revoke(task_id, terminate=True). The worker module's app from section 2 is not available in the web process, so either import it or build one there with build_celery_app(settings, session_factory=...) (one per process, see section 2). Without it, abort still cancels queued jobs but only records the request for processing ones.
Watch GET /ingestion/status: by_status shows backlog (queued) and in-flight (processing) counts, and failure_rate is failed / (done + failed). A growing queued count with no processing jobs means no worker is consuming the queue. The route table is on the module page.
8. Testing your kind¶
- Call the stage directly with
start_stage_context(session, job, settings); no broker or Celery app is needed. Use an SQLite engine passed throughmsflib.db.sqlite.enable_savepointsso rollbacks inside the stage behave as on Postgres. - Assert enqueueing with
FakeDispatchTask. Itsdispatchedlist is empty until the surrounding transaction commits. - For an end-to-end test, use the eager setup from section 4. Build one Celery app for the whole test module and swap the registry or session per test, instead of rebuilding the app.
- The per-subject advisory lock is a Postgres feature and is a no-op on SQLite. Tests that cover it run against a real database.
See also¶
- Ingestion module
- Knowledge: structured-record ingestion built on this guide
- Documents
- Event bus