msflib-ingestion¶
Purpose¶
msflib-ingestion is the shared background-job core for ingestion-style work: one IngestionJob table, one Celery-backed dispatch and execution path, one management API. Modules such as documents and knowledge register pipeline kinds on it, and your own app can do the same for its own data, so you do not hand-roll a polling queue per feature.
It owns the queue, retries, dead-lettering and monitoring. It does not know what a "document" or a "note" is: a job only carries a pointer, (subject_type, subject_id), and the pipeline kind you registered decides how to load and process the thing it points at.
For a step-by-step setup of a worker and a producer, see Long-running jobs and ingestion.
Install¶
[tool.poetry.dependencies]
msflib = { git = "https://github.com/msflib/fastapi.git", subdirectory = "core", rev = "core-v0.2.1" }
msflib-ingestion = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/ingestion", rev = "ingestion-v0.2.0" }
This pulls in msflib-tenancy (jobs carry a tenant_id foreign key), celery and redis. There are no extras.
msflib/ingestion/__init__.py lazily exports IngestionSettings, the three actions, IngestionDispatchService, PipelineKindRegistry, default_pipeline_registry, build_celery_app and build_dispatch_task, so from msflib.ingestion import build_dispatch_task works. Everything else is imported from its submodule, as in the examples below.
Wiring into a host app¶
Two processes share one broker, one database and one set of registered pipeline kinds:
- The FastAPI process enqueues jobs and never runs them. It builds a producer handle with
build_dispatch_task(settings)and gives it to anIngestionDispatchService(or anIngestionSubjectQueue, below). Optionally it mounts the management router. - A Celery worker process runs the jobs. Its entrypoint registers every pipeline kind, then calls
build_celery_app(settings, session_factory=...)and is started withcelery -A your_module worker.
Add IngestionSettings to your settings class, and make sure the ingestion tables are on your metadata (import msflib.ingestion.models) so your migrations or create_all create them.
from msflib.core.config import SettingsBase
from msflib.ingestion.config import IngestionSettings
class AppSettings(IngestionSettings, SettingsBase):
...
Producer side:
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
dispatch = IngestionDispatchService(task=build_dispatch_task(settings))
app.include_router(
create_ingestion_router(
get_session=get_session,
settings=settings,
mutation_role_check=require_admin, # any FastAPI dependency; see below
get_current_tenant=get_current_tenant,
get_current_workspace=get_current_workspace,
prefix="/ingestion",
)
)
Worker side:
from msflib.ingestion.celery import build_celery_app
from sqlmodel import Session
register_my_pipeline_kinds() # your code: fills the PipelineKindRegistry
app = build_celery_app(settings, session_factory=lambda: Session(engine))
mutation_role_check is required and has no permissive default. Every mutating route (trigger, retry, abort, replay, discard) depends on it. If your deployment really has no per-request auth, pass an explicit no-op such as lambda: None. get_current_tenant and get_current_workspace are optional; without a tenant dependency the router falls back to the default tenant, so list and status routes are always tenant-filtered.
Known issue (#288)
The router's default-tenant fallback looks up the tenant with slug default and ignores the host's TENANCY.DEFAULT_TENANT_SLUG. If you changed the slug, the routes return 500 "Default tenant has not been seeded". Pass get_current_tenant built from get_tenant_dependencies(session_dep=get_session, settings=settings.scope("TENANCY")).
The router validates pipeline_kind on POST /jobs/trigger against the registry in the producer process, so register your kinds there too, not only in the worker.
See Mounting routes for the general router pattern.
Configuration¶
Settings live in the INGESTION namespace. If your settings object holds IngestionSettings as a nested INGESTION field, environment variables use the double-underscore form (INGESTION__BROKER_URL). The flat names below (BROKER_URL, MAX_ATTEMPTS, ...) are accepted as aliases, and they are also the field names when you mix IngestionSettings into your settings class directly, as in the example above.
Behaviour may change (#280)
The INGESTION__KEY form only works when INGESTION is a nested field; with IngestionSettings mixed into the host class, use the flat names. This may change.
| Key | Default | Notes |
|---|---|---|
ENABLED |
True |
Declared, currently has no effect (#278). Not read by the ingestion code itself; consuming modules gate their own hooks on their own flags. The flat alias ENABLED is shared with other modules, so INGESTION__ENABLED is unreliable when several are composed as nested fields; set the field in code instead. |
BROKER_URL |
redis://localhost:6379/0 |
Must be the same in the producer and worker processes. |
RESULT_BACKEND |
None |
Optional Celery result backend. Job state lives in the database, so you normally leave this unset. |
DEFAULT_QUEUE |
ingestion |
The Celery queue jobs are sent to and the worker consumes from. |
TASK_ACKS_LATE |
True |
A job in flight when a worker dies is redelivered instead of lost. |
WORKER_PREFETCH_MULTIPLIER |
1 |
Keep at 1 with late acks. |
MAX_ATTEMPTS |
5 |
Total attempts per job, including the first. Read live on every task run. |
RETRY_BACKOFF_SECONDS |
30 |
Base delay before a retry; doubles each attempt (exponent capped at 7). |
RETRY_BACKOFF_MAX_SECONDS |
3600 |
Upper bound on the delay. |
Consuming modules add their own queue switches, for example KNOWLEDGE.EXTRACTION_QUEUE_ENABLED and KNOWLEDGE.RECORD_QUEUE_ENABLED. Retry limits are shared: they come from INGESTION.* for every pipeline kind.
Key concepts¶
Job, run and dead letter. IngestionJob has one row per unit of work, keyed by (subject_type, subject_id, pipeline_kind, target_version, reason, tenant_id, workspace_id). IngestionJobRun has one row per attempt (an audit trail with start and end times and the error text). IngestionDeadLetter holds terminally failed jobs for inspection, replay or discard.
Job status. queued, processing, done, failed, superseded. A job that failed transiently goes back to queued with next_retry_at set. failed means retries are exhausted or the error was terminal. superseded means the subject moved past the version the job was queued for, so the work was skipped. It is not an error.
Deduplication. While a job is active, its key is unique: enqueueing the same subject, version, reason and scope again returns the existing job and dispatches nothing. The key is cleared when the job reaches done, failed or superseded.
Reasons. IngestionReason (create, replace, delete, manual_reindex, metadata_change) is why the job was queued. It is stored as a string in IngestionJob.reason.
Subject, snapshot and kind. The ingestion subject is whatever subject_type and subject_id point at: a document record, a form submission, a row in your own table. A pipeline kind knows how to load it. IngestionSubjectSnapshot is what the queue and the stage need to know about a subject at one moment: subject_id, tenant_id, workspace_id, version and live (still exists). A snapshot always carries both.
Pipeline kinds and stages. A kind is an ordered chain of Stage callables in a PipelineKindRegistry. A stage is any callable (StageContext) -> StageContext with a name. The module ships default_pipeline_registry empty; consuming modules register their kinds into it, and build_celery_app(..., registry=...) and router(..., registry=...) accept a different registry.
ingestion_subject_stage. Most kinds load a subject, process it, and clean up when it is gone. msflib.ingestion.services.lifecycle.ingestion_subject_stage builds a stage that does everything around that:
- parses the job reason, and optionally restricts it with
reasons=; - runs
cleanupwhen the subject is gone, not live, or the job is adelete; adeletejob whose subject is live again in the same scope is superseded; - supersedes the job when the subject has a newer version than
target_version, or moved tenant or workspace; - takes a per-subject Postgres advisory lock (a no-op elsewhere) and reloads the subject before
processorcleanup; if its fingerprint changed meanwhile the attempt is retried; - runs each step in a savepoint and classifies a non-stage error only after the savepoint rolls back;
- fires
completed_eventonce the job commits.
You supply load, snapshot, process, and optionally cleanup, prepare (slow, read-only work that runs before the lock), fingerprint, reasons, completed_event and event_payload.
IngestionSubjectQueue. The producer side of a kind: queue.enqueue(...) takes explicit values, and queue.enqueue_snapshot(...) takes an IngestionSubjectSnapshot and works out create, replace or delete, including a first delete under the old scope when a subject moved. Both take the same per-subject lock as the stage, so an edit made while a job runs waits for it and then queues its own job instead of deduplicating into the running one. Dispatch to Celery is deferred until the caller's transaction commits, so a rollback queues nothing. Pass enabled= from the module's settings flag; a disabled queue raises QueueDisabledError. Subclass only to override check_subject.
sync_model_to_queue. msflib.ingestion.services.sync.sync_model_to_queue(Model, queue, subject_type=..., snapshot=..., emitter=...) enqueues a job on every ModelAction create, update and delete of a model, in the same transaction as the write. Writes that bypass ModelAction (raw session.add, bulk SQL) fire nothing; call queue.enqueue_snapshot yourself after those. reset_model_sync(emitter=..., pipeline_kind=..., subject_type=...) removes the hook.
Errors. Raise RetryableStageError for transient failures and TerminalStageError for bad input. Raise SupersededStageError when a stage finds its work is moot. A plain exception is classified for you by classify_stage_error. Inside ingestion_subject_stage, an error whose text is identical to the previous attempt's is also treated as terminal. QueueDisabledError and UnregisteredSubjectTypeError are raised at enqueue time.
Events. StageContext.emit_after_success(name, instance, payload) queues an event that fires after the whole job commits. Use it instead of msflib.actions.emit_after_commit inside a stage, because SQLAlchemy forbids further SQL on the session from an after_commit handler. Event listeners must be registered in the worker process, since that is where the event fires, and on the emitter active there. In the worker entrypoint, register them on the process default emitter: hooks that take an emitter= argument can be given AppEmitter(get_emitter()) (from msflib.eventbus). Listeners bound to the web app's AppEmitter are never called from a real worker; in eager mode they fire only when the task runs inside a request. See Event bus.
Examples¶
A pipeline kind for your own table¶
This registers a note_index kind, queues jobs from a snapshot, and runs the stage directly with the test helpers. It runs as written against in-memory SQLite. SQLite does not enforce the tenant_id foreign key, so the example seeds tenant 1 explicitly to behave the same on Postgres.
from msflib.core.config import SettingsBase
from msflib.db.sqlite import enable_savepoints
from msflib.ingestion.config import IngestionSettings
from msflib.ingestion.contracts.subject import IngestionSubjectSnapshot
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.ingestion.testing import FakeDispatchTask, start_stage_context
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 sqlalchemy.pool import StaticPool
from sqlmodel import Session, SQLModel, create_engine
class Note(ModelBase, table=True):
body: str = ""
version: int = 1
tenant_id: int = 1
workspace_id: int | None = None
deleted: bool = False
def load(context): # StageContext -> subject, or None when it is gone
return context.session.get(Note, int(context.job.subject_id))
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,
)
indexed = []
def process(context): # runs locked, inside a savepoint
indexed.append(context.subject.body)
def cleanup(context): # subject is gone, or the job is a delete
indexed.append(f"removed {context.job.subject_id}")
registry = PipelineKindRegistry()
registry.register_kind(
"note_index",
ingestion_subject_stage(
name="note_index", load=load, snapshot=note_snapshot, process=process, cleanup=cleanup
),
description="Index notes",
)
class Settings(IngestionSettings, SettingsBase):
pass
engine = enable_savepoints(
create_engine("sqlite://", connect_args={"check_same_thread": False}, poolclass=StaticPool)
)
SQLModel.metadata.create_all(engine)
with Session(engine) as session: # seeds tenant 1, which Note.tenant_id points at
resolve_default_tenant_id(session)
session.commit()
task = FakeDispatchTask() # records job ids instead of calling Celery
queue = IngestionSubjectQueue(
dispatch=IngestionDispatchService(task=task), pipeline_kind="note_index"
)
with Session(engine) as session:
note = Note(body="hello")
session.add(note)
session.commit()
jobs = queue.enqueue_snapshot(
session, subject_type="note", snapshot=note_snapshot(note), commit=True
)
print([(job.reason, job.status.value) for job in jobs], task.dispatched)
# [('create', 'queued')] [1]
context = start_stage_context(session, jobs[0], Settings())
registry.resolve("note_index").stages[0](context)
print(indexed) # ['hello']
note.deleted = True
session.add(note)
session.commit()
jobs = queue.enqueue_snapshot(
session, subject_type="note", snapshot=note_snapshot(note), commit=True
)
print([job.reason for job in jobs]) # ['delete']
FakeDispatchTask and start_stage_context come from msflib.ingestion.testing. start_stage_context starts an attempt the way the Celery task does and returns the StageContext a stage receives, so a stage can be unit-tested without a broker.
In the real app, register the same kind in the worker entrypoint (and in the producer process if you mount the management router), and replace FakeDispatchTask with build_dispatch_task(settings). To enqueue on every write to Note instead of calling enqueue_snapshot by hand:
from msflib.eventbus import bind_app_emitter
from msflib.ingestion.services.sync import sync_model_to_queue
app_emitter = bind_app_emitter(app) # at app creation
sync_model_to_queue(
Note, queue, subject_type="note", snapshot=note_snapshot, emitter=app_emitter
)
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): (from msflib.eventbus). A write made from a worker, a CLI script or a thread outside that block queues nothing; call queue.enqueue_snapshot there instead.
Management routes¶
With the router mounted at /ingestion:
| Method and path | What it does |
|---|---|
GET /jobs |
List jobs; filters subject_type, status, limit (1 to 100), offset. |
GET /jobs/{job_id} |
One job with its runs. |
GET /status |
Counts by status and by subject type, plus a failure rate (failed / (done + failed)). |
GET /dead-letters |
List dead letters; include_resolved also shows replayed and discarded ones. |
POST /jobs/trigger |
Enqueue a job for any subject. The payload cannot name a tenant or workspace other than the caller's (403). Unknown pipeline_kind is 422. |
POST /jobs/{job_id}/retry |
Requeue a failed job at once. Any other status is 409. |
POST /jobs/{job_id}/abort |
Cancel a queued job, or best-effort revoke a processing one. |
POST /dead-letters/{id}/replay |
Enqueue a new job from the dead letter's subject. |
POST /dead-letters/{id}/discard |
Mark the dead letter resolved without replaying. |
Everything is filtered to the caller's resolved tenant and workspace. Routes that act on one job or dead letter return 404 for rows outside it.
Troubleshooting¶
| Symptom | Cause and fix |
|---|---|
Jobs stay queued |
No worker is consuming DEFAULT_QUEUE, or the worker and producer use different BROKER_URL values. |
Worker raises KeyError: no pipeline registered for kind |
Registration is per process. Register every kind in the worker entrypoint. |
POST /jobs/trigger returns 422 Unknown pipeline_kind |
The producer process has not registered the kind. Register it there too. |
QueueDisabledError on enqueue |
The module's queue flag is off, for example KNOWLEDGE.RECORD_QUEUE_ENABLED. sync_model_to_queue ignores this error so saves still succeed. |
Job marked failed with "Celery dispatch failed" |
The broker was unreachable after the job row was committed. The job is visible and can be retried with POST /jobs/{job_id}/retry. |
UnregisteredSubjectTypeError |
The queue subclass only accepts subject types it has a loader for (the knowledge_record queue does). Register the source in the web process. |
| Edit made while a job runs is not reflected | Give the snapshot a version. Without one, the new job can deduplicate into the running one. |
| Two Celery apps in one test process corrupt each other | Build one Celery app per process. Use FakeDispatchTask for tests that only check enqueueing. |
/jobs/{job_id}/retry returns 409 |
Only failed jobs can be retried. A queued job may already have a pending retry on the broker. |
API reference¶
See the generated API reference for msflib.ingestion. modules/ingestion/README.md has the layering notes and testing instructions.
See also¶
- Long-running jobs and ingestion, a worker and producer walkthrough
- Knowledge, which registers the
knowledge_extractionandknowledge_recordkinds - Documents, which registers the
text_documentkind - Event bus