Skip to content

msflib.ingestion

Public modules of the msflib-ingestion package.

msflib.ingestion.actions

IngestionJobAction

Bases: ModelAction[IngestionJob, IngestionJobCreate, IngestionJobUpdate]

dedupe_key_for(*, subject_type: str, subject_id: str, pipeline_kind: str, target_version: int | None, reason: str, tenant_id: int, workspace_id: int | None) -> str staticmethod

Compute the dedupe key for a (subject, pipeline_kind, target_version, reason, scope).

Also used to restore dedupe_key when requeuing a terminally-failed job: failure_update clears it on terminal failure, so without recomputing it here a manual retry would leave the one-job-per- subject invariant unenforced until the job resolves again.

enqueue_unique(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str, reason: str, target_version: int | None = None, tenant_id: int, workspace_id: int | None = None, commit: bool = False) -> IngestionJob

Create-or-return-existing by dedupe key, race-safe under concurrent inserts.

Kept here rather than in a service: the nested-transaction + IntegrityError-fallback dance is a database write-concurrency concern, not domain policy.

enqueue_unique_created(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str, reason: str, target_version: int | None = None, tenant_id: int, workspace_id: int | None = None, commit: bool = False) -> tuple[IngestionJob, bool]

enqueue_unique that also reports whether this call created the job.

latest_for_subject(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str | None = None) -> IngestionJob | None

Most recent job for a subject, whatever its status.

count_scoped(session: Session, *, subject_type: str | None = None, status: IngestionJobStatus | None = None, tenant_id: int | None = None, workspace_id: int | None = None) -> int

Count rows matching the same filters as get_multi_scoped, for pagination totals.

status_and_subject_type_counts(session: Session, *, tenant_id: int | None = None, workspace_id: int | None = None) -> tuple[dict[IngestionJobStatus, int], dict[str, int]]

Row counts grouped by status and by subject_type, for the same tenant/workspace scope.

IngestionDeadLetterAction

Bases: ModelAction[IngestionDeadLetter, IngestionDeadLetterCreate, IngestionDeadLetterUpdate]

count_scoped(session: Session, *, subject_type: str | None = None, tenant_id: int | None = None, workspace_id: int | None = None, include_resolved: bool = False) -> int

Count rows matching the same filters as get_multi_scoped, for pagination totals.

msflib.ingestion.celery

Celery + Redis worker plumbing (Layer 0).

Real distributed workers, native retry/backoff, and a dead-letter path, replacing the DB-polling + daemon-thread mechanism documents/knowledge hand- rolled, and which other downstream apps declared as a dependency but never actually activated (celery+redis were declared, unused, deps).

task_acks_late + worker_prefetch_multiplier=1 is the Celery-native replacement for the correctness guarantee SELECT ... FOR UPDATE SKIP LOCKED + lease-TTL reclaim used to provide: an in-flight job surviving a worker crash is redelivered rather than lost.

Deployed as its own worker process (celery -A <entrypoint> worker); the consuming application supplies a module-level entrypoint that calls build_celery_app with its own settings/session factory, mirroring how documents.router()/ai_core take factories rather than assuming a single process-wide app instance.

build_celery_app(settings: SettingsBase, *, session_factory: SessionFactory, registry: PipelineKindRegistry | None = None, job_action: IngestionJobAction | None = None, run_action: IngestionJobRunAction | None = None, dead_letter_action: IngestionDeadLetterAction | None = None, app_name: str = 'msflib.ingestion') -> Celery

Build the Celery application for the ingestion worker process.

session_factory: zero-argument context-manager factory returning a Session (e.g. lambda: Session(engine)). The worker is a separate process from the FastAPI app, so no request-scoped session dependency exists here -- this mirrors documents.router()'s make_background_session override for the same reason.

build_dispatch_task(settings: SettingsBase, *, app_name: str = 'msflib.ingestion') -> _ProducerTask

Build a lightweight, .delay(job_id)-compatible dispatch handle.

For producer-only processes -- pass the result as IngestionDispatchService(task=...).

msflib.ingestion.config

msflib.ingestion.contracts

CeleryControlContract

Bases: Protocol

An object exposing .control.revoke(task_id, terminate=True).

Typically the worker-side Celery app itself, or any app sharing the same broker.

Stage

Bases: Protocol

A single Fetch/Normalize/Transform/Persist/Index/Emit step.

Implementations are plain callables (functions or callable objects) -- no base class required, only this shape.

StageContext(session: Session, job: IngestionJob, job_run: IngestionJobRun, settings: SettingsBase, state: dict[str, Any] = dict(), pending_events: list[tuple[str, Any, dict[str, Any]]] = list()) dataclass

emit_after_success(event_name: str, instance: Any = None, payload: dict[str, Any] | None = None) -> None

Queue an eventbus event to fire only once the whole job commits successfully.

Deliberately not msflib.actions.emit_after_commit (a SQLAlchemy after_commit event hook): SQLAlchemy forbids emitting further SQL on the same session from inside an after_commit handler, which breaks any listener that queries the session (reproducible, not hypothetical). The job completion handler drains this list as plain sequential code immediately after its own session.commit() returns, where session use is unrestricted.

celery

Structural contract for the optional Celery control handle passed to router().

CeleryControlContract

Bases: Protocol

An object exposing .control.revoke(task_id, terminate=True).

Typically the worker-side Celery app itself, or any app sharing the same broker.

stage

Pluggable pipeline stage protocol.

An ordered chain of stages registered per pipeline_kind (see services.registry). A stage reads/writes StageContext.state to pass data to the next stage; the session/job/job_run/settings are fixed for the whole run.

Stages raise msflib.ingestion.errors.RetryableStageError / TerminalStageError (or let the Celery task wrapper classify a plain exception via classify_stage_error) to control retry behavior.

StageContext(session: Session, job: IngestionJob, job_run: IngestionJobRun, settings: SettingsBase, state: dict[str, Any] = dict(), pending_events: list[tuple[str, Any, dict[str, Any]]] = list()) dataclass

emit_after_success(event_name: str, instance: Any = None, payload: dict[str, Any] | None = None) -> None

Queue an eventbus event to fire only once the whole job commits successfully.

Deliberately not msflib.actions.emit_after_commit (a SQLAlchemy after_commit event hook): SQLAlchemy forbids emitting further SQL on the same session from inside an after_commit handler, which breaks any listener that queries the session (reproducible, not hypothetical). The job completion handler drains this list as plain sequential code immediately after its own session.commit() returns, where session use is unrestricted.

Stage

Bases: Protocol

A single Fetch/Normalize/Transform/Persist/Index/Emit step.

Implementations are plain callables (functions or callable objects) -- no base class required, only this shape.

subject

What the lifecycle stage and IngestionSubjectQueue need to know about a job's subject.

IngestionSubjectSnapshot(subject_id: str, tenant_id: int, workspace_id: int | None = None, version: int | None = None, live: bool = True) dataclass

One ingestion subject at one moment, as the job queue sees it.

The ingestion subject is what IngestionJob.subject_type/subject_id point at (a document record, a form submission, an app row).

Holds what jobs are keyed and checked on: the id, the tenant/workspace the subject lives in, its version, and whether it still exists (live).

msflib.ingestion.errors

Typed stage error classification, driving Celery's autoretry_for directly.

Promotes msflib.core.errors.is_terminal_ingestion_error from string/type sniffing into typed exceptions that pipeline stages (Layer 2) raise on purpose, plus a classifier for stages/adapters that still just raise plain exceptions.

IngestionStageError

Bases: Exception

Base for errors raised by a pipeline stage.

RetryableStageError

Bases: IngestionStageError

Transient failure (timeout, rate limit, temporary unavailability) -- worth retrying.

TerminalStageError

Bases: IngestionStageError

Non-retryable failure (bad input, permanent 4xx, validation error) -- retrying won't help.

SupersededStageError

Bases: IngestionStageError

The job's target has been overtaken by a newer version of its subject.

Not a failure: the subject already reflects (or is queued to reflect) a later version, so this attempt's work is moot. Stages raise this explicitly (a plain/typed exception can't be distinguished from a real failure) when they detect their own supersession, mirroring the version-supersession semantics documents'/knowledge's polling queues already got right. The Celery task wrapper (celery.py) marks the job superseded rather than retrying or dead-lettering it.

QueueDisabledError

Bases: RuntimeError

Raised when enqueueing onto a pipeline kind whose queue is disabled by settings.

UnregisteredSubjectTypeError

Bases: LookupError

Raised when enqueueing a subject type no worker in this process can load.

classify_stage_error(exc: Exception) -> type[IngestionStageError]

Map an arbitrary exception raised by a stage to RetryableStageError or TerminalStageError.

Used by the Celery task wrapper (Layer 0) to normalize exceptions from stages that don't raise the typed errors directly, so autoretry_for can still key off one exception hierarchy. Delegates the actual terminal-vs-retryable call to msflib.core.errors, shared with documents/knowledge, so the rules can't drift between modules.

msflib.ingestion.models

IngestionDeadLetter

Bases: ModelBase

Archive of terminally-failed jobs, for inspection/replay/discard.

Nothing analogous exists in documents/knowledge today (exhausted jobs just sit in status=failed with no archive/replay tooling). Subject/scope fields are denormalized from the parent job at the time of dead-lettering so DLQ listing/filtering doesn't require a join, and stays inspectable even if the job row's own fields keep changing (e.g. a later manual retry) after the dead letter was written.

IngestionJob

Bases: ModelBase

One row per (subject, target_version, reason) unit of ingestion work.

subject_type/subject_id form a polymorphic key so any module can queue work here without this module importing from -- or depending on -- the modules whose subjects it queues.

IngestionJobRun

Bases: ModelBase

One row per attempt of an IngestionJob -- the audit trail.

Carries started_at/completed_at/status/error_message plus progress counters for per-attempt observability, richer than a bare attempts counter on the job itself.

job

IngestionJob

Bases: ModelBase

One row per (subject, target_version, reason) unit of ingestion work.

subject_type/subject_id form a polymorphic key so any module can queue work here without this module importing from -- or depending on -- the modules whose subjects it queues.

IngestionJobRun

Bases: ModelBase

One row per attempt of an IngestionJob -- the audit trail.

Carries started_at/completed_at/status/error_message plus progress counters for per-attempt observability, richer than a bare attempts counter on the job itself.

IngestionDeadLetter

Bases: ModelBase

Archive of terminally-failed jobs, for inspection/replay/discard.

Nothing analogous exists in documents/knowledge today (exhausted jobs just sit in status=failed with no archive/replay tooling). Subject/scope fields are denormalized from the parent job at the time of dead-lettering so DLQ listing/filtering doesn't require a join, and stays inspectable even if the job row's own fields keep changing (e.g. a later manual retry) after the dead letter was written.

msflib.ingestion.reasons

Why a job was queued; IngestionJob.reason holds one of these values.

msflib.ingestion.router

router(*, get_session: Callable, settings: SettingsBase, mutation_role_check: Callable, get_current_tenant: Callable | None = None, get_current_workspace: Callable | None = None, dispatch_task: DelayableTask | None = None, celery_control: CeleryControlContract | None = None, registry: PipelineKindRegistry | None = None, prefix: str = '/ingestion', tags: list[str] | None = None) -> APIRouter

Management API for IngestionJob: status/monitoring/audit-trail plus mutating operations.

One router covers every subject type (documents, knowledge, etc.), since they all funnel through the same IngestionJob model instead of each growing its own ops endpoints.

get_current_tenant/get_current_workspace: injectable scope resolvers, mirroring documents.router()'s get_current_account/ get_current_workspace (get_current_tenant yields a Tenant row, get_current_workspace yields a WorkspaceContract). When omitted, the tenant filter falls back to resolve_default_tenant_id (today's single-tenant resolution) rather than no filter at all -- matching every other read path in this router that resolves a tenant id.

mutation_role_check: required FastAPI dependency applied to every mutating endpoint (trigger/retry/abort/replay/discard) -- no permissive default, so a deployment with no per-request auth must pass an explicit no-op (lambda: None) rather than getting one silently.

dispatch_task: a .delay(job_id)-compatible handle used to dispatch to Celery. Defaults to celery.build_dispatch_task(settings).

celery_control: an object exposing .control.revoke(task_id, terminate=True), used by abort for cooperative cancellation of an in-flight job. Optional: without it, abort still cancels queued jobs, but a processing job is only marked for cancellation, not revoked.

registry: the PipelineKindRegistry used to validate pipeline_kind on /jobs/trigger before enqueueing. Defaults to default_pipeline_registry.

msflib.ingestion.schema

msflib.ingestion.services

IngestionDispatchService(*, task: DelayableTask, job_action: IngestionJobAction | None = None)

enqueue_and_dispatch(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str, reason: str, target_version: int | None = None, tenant_id: int, workspace_id: int | None = None, commit: bool = True, dispatch_after_commit: bool = False) -> IngestionJob

enqueue_and_dispatch_created without the created flag.

enqueue_and_dispatch_created(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str, reason: str, target_version: int | None = None, tenant_id: int, workspace_id: int | None = None, commit: bool = True, dispatch_after_commit: bool = False) -> tuple[IngestionJob, bool]

Enqueue the job row and dispatch it to Celery once it's durable.

Dispatch never fires before the row is committed -- a worker (a separate process/connection) picking up a job id its own session can't see yet, or that a rollback later undoes, is the hazard this guards against (same one msflib.actions.emit_after_commit guards for eventbus events).

  • commit=True: this call commits and dispatches immediately.
  • commit=False, dispatch_after_commit=True: the caller owns an outer transaction (e.g. a pre-commit hook); dispatch is deferred to that transaction's own session.commit(), whenever it happens.
  • commit=False, dispatch_after_commit=False (default): caller is responsible for dispatching itself after its own commit.

A dedupe hit returns the existing job without dispatching it again; the flag reports whether the job was created.

dispatch_or_mark_failed(session: Session, *, job: IngestionJob) -> None

self.task.delay(job.id), or mark the job failed and raise if that fails.

For use once a job's row is already durable and there's still a caller around to report the failure to (unlike the dispatch_after_commit path, which has none) -- shared by enqueue_and_dispatch's commit=True branch and router.retry_job, which dispatches an already-updated job directly rather than through enqueue_and_dispatch.

classify

Turn a stage failure into a typed stage error, failing fast on a repeated error.

is_repeat_of_previous_error(context: StageContext, error_text: str) -> bool

True when the previous attempt's run failed with this exact error text.

Reads IngestionJobRun history: start_attempt_update clears job.error before the stage runs.

reraise_classified(context: StageContext, exc: Exception) -> NoReturn

Re-raise exc as a typed stage error; a repeat of the previous attempt's is terminal.

dispatch

Producer-side enqueue + dispatch, decoupled from the Celery app itself.

Takes any object exposing .delay(job_id) (a real Celery Task, or a test double) rather than importing celery.build_celery_app directly -- the FastAPI process enqueueing a job doesn't need to construct the same Celery app instance the worker process runs, only something that can hand a job id to the same broker/task name.

DelayableTask

Bases: Protocol

A .delay(job_id)-compatible handle -- a real Celery Task or a test double.

DispatchFailedError

Bases: RuntimeError

Raised when handing a job to Celery fails after its row is already durable.

The job is marked failed (with dedupe_key cleared) before this is raised, so it's visible/retryable via the normal job list + retry_job instead of sitting in queued forever with no live task and nothing polling the DB to notice.

IngestionDispatchService(*, task: DelayableTask, job_action: IngestionJobAction | None = None)

enqueue_and_dispatch(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str, reason: str, target_version: int | None = None, tenant_id: int, workspace_id: int | None = None, commit: bool = True, dispatch_after_commit: bool = False) -> IngestionJob

enqueue_and_dispatch_created without the created flag.

enqueue_and_dispatch_created(session: Session, *, subject_type: str, subject_id: str, pipeline_kind: str, reason: str, target_version: int | None = None, tenant_id: int, workspace_id: int | None = None, commit: bool = True, dispatch_after_commit: bool = False) -> tuple[IngestionJob, bool]

Enqueue the job row and dispatch it to Celery once it's durable.

Dispatch never fires before the row is committed -- a worker (a separate process/connection) picking up a job id its own session can't see yet, or that a rollback later undoes, is the hazard this guards against (same one msflib.actions.emit_after_commit guards for eventbus events).

  • commit=True: this call commits and dispatches immediately.
  • commit=False, dispatch_after_commit=True: the caller owns an outer transaction (e.g. a pre-commit hook); dispatch is deferred to that transaction's own session.commit(), whenever it happens.
  • commit=False, dispatch_after_commit=False (default): caller is responsible for dispatching itself after its own commit.

A dedupe hit returns the existing job without dispatching it again; the flag reports whether the job was created.

dispatch_or_mark_failed(session: Session, *, job: IngestionJob) -> None

self.task.delay(job.id), or mark the job failed and raise if that fails.

For use once a job's row is already durable and there's still a caller around to report the failure to (unlike the dispatch_after_commit path, which has none) -- shared by enqueue_and_dispatch's commit=True branch and router.retry_job, which dispatches an already-updated job directly rather than through enqueue_and_dispatch.

execution

Job/run lifecycle domain logic for the Celery task wrapper (celery.py).

Pure functions computing what to write -- exhaustion/retry-backoff policy and job/run state transitions -- kept out of actions.py. Action classes do database reads/writes only (ModelAction.create/update already covers a plain field-set); the caller decides which fields, using these functions, then calls the action with the result.

compute_failure_outcome(*, job: IngestionJob, terminal: bool, max_attempts: int, retry_backoff_seconds: int, retry_backoff_max_seconds: int, retry_after_seconds: int | None = None) -> FailureOutcome

Decide whether a failed attempt is retryable, and if so, when.

Exponential backoff (retry_backoff_seconds * 2 ** (attempts - 1), exponent clamped to [0, 7] and the result capped at retry_backoff_max_seconds) mirrors Celery's own retry_backoff/ retry_backoff_max semantics used for autoretry_for -- see IngestionSettings. The first retry (attempts == 1) waits the base retry_backoff_seconds, not double it.

lifecycle

ingestion_subject_stage: the load / check / lock / process-or-clean-up steps of every kind.

A pipeline kind supplies how to load and snapshot its ingestion subject (what the job points at), what to do with a live one (process, optionally after an unlocked prepare) and what to do when it is gone (cleanup). The stage owns everything else:

  • the job's reason is parsed and checked against reasons;
  • a subject that is gone, not live, or targeted by a delete job is cleaned up; a delete job whose subject is live again in the same scope is superseded;
  • a newer subject version, or a subject that moved tenant/workspace, supersedes the job;
  • before any write, a per-subject Postgres advisory lock is taken and the subject reloaded; if its fingerprint changed the attempt is retried;
  • each step runs in a savepoint, and a non-stage error is classified only after the savepoint rolls back (on Postgres, classifying inside a failed transaction would itself fail);
  • completed_event fires once the job commits.

IngestionSubjectContext(stage: StageContext, reason: IngestionReason, subject: S | None = None, snapshot: IngestionSubjectSnapshot | None = None, prepared: Any = None) dataclass

Bases: Generic[S]

What prepare/process/cleanup see; subject is None in cleanup.

lock_ingestion_subject(session: Session, *, pipeline_kind: str, subject_type: str, subject_id: str) -> None

Hold a per-subject Postgres advisory lock until the transaction ends; no-op elsewhere.

run_guarded(context: StageContext, work: Callable[[], _T]) -> _T

Run work in a savepoint; classify a non-stage error after the savepoint rolls back.

ingestion_subject_stage(*, name: str, load: Callable[[StageContext], S | None], snapshot: Callable[[S], IngestionSubjectSnapshot], process: Callable[[IngestionSubjectContext[S]], Any], prepare: Callable[[IngestionSubjectContext[S]], Any] | None = None, cleanup: Callable[[IngestionSubjectContext[S]], Any] | None = None, fingerprint: Callable[[S], Any] | None = None, reasons: Collection[IngestionReason] | None = None, completed_event: str | None = None, event_payload: EventPayload = default_event_payload, lock: bool = True) -> Stage

Build a stage that runs one subject through the shared lifecycle.

load returns the subject or None when it is gone. prepare runs unlocked (put slow, read-only work there); process runs after the lock and re-check, and receives prepare's result as context.prepared. fingerprint (default: snapshot) decides whether the subject changed while the job ran.

queue

IngestionSubjectQueue: enqueue a pipeline kind's jobs; dispatch waits for the commit.

IngestionSubjectQueue(*, dispatch: IngestionDispatchService, pipeline_kind: str, enabled: bool = True, job_action: IngestionJobAction | None = None)

A pipeline kind's job queue; pass its settings flag as enabled=.

Subclass only to override check_subject (see KnowledgeRecordQueue).

latest(session: Session, *, subject_type: str, subject_id: str) -> IngestionJob | None

Most recent job of this kind for a subject, whatever its status.

enqueue(session: Session, *, subject_type: str, subject_id: str, tenant_id: int, workspace_id: int | None, target_version: int | None = None, reason: IngestionReason | str = IngestionReason.create, commit: bool = False) -> IngestionJob

Queue one job; tenant_id/workspace_id must be the subject's own scope.

Takes the subject's lock first, so a job that is running for it finishes (or stands down) before this dedupes against it.

enqueue_snapshot(session: Session, *, subject_type: str, snapshot: IngestionSubjectSnapshot, reason: IngestionReason | str | None = None, commit: bool = False) -> list[IngestionJob]

Queue the jobs that bring this kind in line with snapshot.

With no reason: a live subject is create (no earlier job) or replace; a gone one is delete, or nothing if it was never queued. A subject that moved tenant/workspace first gets a delete under the previous job's scope.

commit_and_refresh(session: Session, commit: bool, *rows: Any) -> list[Any]

Commit and refresh rows when commit is set; return them either way.

registry

Registry of pipeline stage chains, one per pipeline_kind.

Follows the register/resolve/list pattern from ai_core.services.chunking.registry.ChunkingRegistry. Empty by default -- this module owns the queue/registry mechanism only; concrete kinds (text_document, graph_extraction, structured_record) are registered by the consuming modules (documents, knowledge, others) as they migrate onto it.

sync

sync_model_to_queue: enqueue a job on every ModelAction write to a model.

reset_model_sync(*, emitter: AppEmitter, pipeline_kind: str | None = None, subject_type: str | None = None) -> None

Remove sync_model_to_queue hooks: one subject type, one kind, or all.

sync_model_to_queue(model: type | str, queue: IngestionSubjectQueue, *, subject_type: str, snapshot: Callable[[Any], IngestionSubjectSnapshot], emitter: AppEmitter, on: Iterable[ModelWriteOperation] = ('create', 'update', 'delete'), on_error: Literal['log', 'raise'] = 'log', force: bool = False) -> None

Keep queue's kind in step with model; a disabled queue is skipped silently.

See msflib.actions.on_model_write for on_error and idempotency; reset_model_sync removes the hook.

msflib.ingestion.testing

Test doubles for code that enqueues ingestion jobs.

FakeDispatchTask(dispatched: list[int] = list()) dataclass

DelayableTask double that records dispatched job ids instead of calling Celery.

start_stage_context(session: Session, job: IngestionJob, settings: SettingsBase) -> StageContext

Start an attempt the way the Celery task does and return the context a stage receives.