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 ownsession.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 ownsession.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
deletejob is cleaned up; adeletejob 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_eventfires 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.