Skip to content

msflib.documents

Public modules of the msflib-documents package.

msflib.documents.actions

DocumentRecordAction

Bases: ModelAction[DocumentRecord, DocumentRecordCreate, DocumentRecordUpdate]

promote_scope(session: Session, *, record: DocumentRecord, clear_private: bool = False, commit: bool = False) -> DocumentRecord

Widen a document's scope to workspace-level by clearing its conversation_id/sub_thread_id and stamping scope_promoted_at, so we can still tell a promoted record apart from one that was never scoped and so re-ingestion won't silently narrow it back down (see DocumentRecordAction._upsert_from_source). Callers are responsible for enqueueing a metadata_change ingestion job so the vector store is re-embedded with the widened metadata.

Privacy (private_to_account_id) is orthogonal to conversation scope and is only cleared when clear_private=True is explicitly requested.

get_multi_scoped(session: Session, *, account_id: int, workspace_id: int | None = None, conversation_id: str | None = None, sub_thread_id: str | None = None, status: DocumentRecordStatus | None = None, include_deleted: bool = False, limit: int = 20, offset: int = 0) -> list[DocumentRecord]

List records visible to account_id, newest first.

workspace_id is a required_nullable scope dimension, not a regular optional filter like conversation_id/sub_thread_id: it is always applied, and None means "the caller has no workspace context" -- which must match only the catchall workspace_id IS NULL records (personal/unscoped documents), not skip the filter and leak every workspace's documents. Skipping the filter on None would let an account-only (no current workspace) request see documents scoped to workspaces it never asked for.

status (an exact match) and include_deleted are mutually exclusive filters on the same column: an explicit status wins, otherwise deleted records are excluded by default so a plain listing doesn't surface documents the account already deleted.

msflib.documents.adapters

msflib.documents.config

Optional drivelink integration for msflib-documents.

Import and call register_event_hooks from this module when both msflib-drivelink and msflib-documents are installed and you want document ingestion to be triggered automatically by drivelink node lifecycle events.

Example::

from msflib.documents.drivelink import register_event_hooks
from msflib.documents.actions import DocumentRecordAction
from msflib.documents.services import DocumentsIngestionDispatchService

register_event_hooks(
    record_action=DocumentRecordAction(),
    dispatch_service=DocumentsIngestionDispatchService(...),
    emitter=app_emitter,
    tenancy_settings=settings.scope("TENANCY"),
)

register_event_hooks(*, record_action: DocumentRecordAction, dispatch_service: DocumentsIngestionDispatchService, emitter: AppEmitter, tenancy_settings: TenancySettings, force: bool = False) -> None

Register drivelink lifecycle hooks that enqueue ingestion jobs.

When a file node is created or updated, the hook enqueues an ingestion job (commit=False, dispatch_after_commit=True under the hood -- see DocumentsIngestionDispatchService) so the Celery dispatch itself only fires once this pre-commit hook's own transaction actually commits.

tenancy_settings is the app's configured TenancySettings; pass the same instance startup seeding and request dependencies use so the hooks resolve the same default tenant row.

eventbus

reset_event_hooks(*, emitter: AppEmitter) -> None

Remove document hooks and reset registration state.

register_event_hooks(*, record_action: DocumentRecordAction, dispatch_service: DocumentsIngestionDispatchService, emitter: AppEmitter, tenancy_settings: TenancySettings, force: bool = False) -> None

Register drivelink lifecycle hooks that enqueue ingestion jobs.

When a file node is created or updated, the hook enqueues an ingestion job (commit=False, dispatch_after_commit=True under the hood -- see DocumentsIngestionDispatchService) so the Celery dispatch itself only fires once this pre-commit hook's own transaction actually commits.

tenancy_settings is the app's configured TenancySettings; pass the same instance startup seeding and request dependencies use so the hooks resolve the same default tenant row.

msflib.documents.models

DocumentIngestionCompletedEvent(document_record_id: int, target_version: int, reason: str, tenant_id: int, workspace_id: int | None) dataclass

Payload for the documents.ingestion.completed event.

Ingestion job rows themselves now live in msflib.ingestion (IngestionJob, subject_type="document"), which documents has no reason to import (it would pull the whole ingestion module into every listener). This carries just the fields external listeners actually use (e.g. msflib.knowledge.integrations.documents), keeping the event contract stable across the migration.

document

DocumentIngestionCompletedEvent(document_record_id: int, target_version: int, reason: str, tenant_id: int, workspace_id: int | None) dataclass

Payload for the documents.ingestion.completed event.

Ingestion job rows themselves now live in msflib.ingestion (IngestionJob, subject_type="document"), which documents has no reason to import (it would pull the whole ingestion module into every listener). This carries just the fields external listeners actually use (e.g. msflib.knowledge.integrations.documents), keeping the event contract stable across the migration.

msflib.documents.protocols

SourceNodeProtocol

Bases: Protocol

Structural interface for any node that can be ingested by the documents module.

msflib.drivelink.models.DrivelinkNode satisfies this protocol without any modification. Other storage backends can implement it independently.

msflib.documents.router

router(*, get_session: Callable, settings: SettingsBase, get_current_account: Callable | None = None, get_current_workspace: Callable | None = None, get_current_tenant: Callable | None = None, prefix: str = '/documents', tags: list[str] | None = None, get_policy_resolver: Callable | None = None, dispatch_task: DelayableTask | None = None) -> APIRouter

dispatch_task: a .delay(job_id)-compatible handle used to dispatch ingestion jobs to Celery. Defaults to msflib.ingestion.celery.build_dispatch_task(settings) (a producer-only Celery app -- no session_factory/task body needed, since this process only enqueues, it never executes ingestion jobs itself; see services/pipeline.py and the worker entrypoint that does execute them). Injectable for tests (pass a Mock()) or to reuse an existing Celery app instance.

get_current_tenant: resolved once at the endpoint boundary, the same way get_current_account/get_current_workspace are (see msflib.tenancy.deps.get_tenant_dependencies) -- yields a Tenant row. When omitted, endpoints that need the real tenant id fall back to resolve_default_tenant_id (today's single-tenant resolution), same as before this parameter existed.

msflib.documents.schema

msflib.documents.scope_profiles

msflib.documents.services

DocumentsQueueDisabledError

Bases: RuntimeError

Raised by enqueue_for_record when DOCUMENTS.QUEUE_ENABLED is False.

register_text_document_pipeline(registry: PipelineKindRegistry = default_pipeline_registry, *, stage: Stage | None = None, replace: bool = False) -> None

Register the text_document pipeline kind (call once per worker process).

enqueue_reindex_for_scope(session: Session, *, dispatch_service: DocumentsIngestionDispatchService, tenant_id: int, workspace_id: int | None = None, statuses: tuple[DocumentRecordStatus, ...] = (DocumentRecordStatus.indexed, DocumentRecordStatus.failed)) -> ReindexReport

Re-enqueue every in-scope, non-deleted document record for re-ingestion.

Switching a scope's vector store backend (see VectorStoreProfile/VectorStoreRegistryService) leaves documents already indexed under the old backend behind: searches against the new backend return nothing for them until they are re-ingested. This re-enqueues matching records with IngestionReason.manual_reindex; the existing ingestion worker re-runs ingest_document_content, which re-embeds and writes into whatever store the record's own scope now resolves to (see services/ingestion.py's resolve_vector_store_per_request path) -- no separate reindex pipeline is needed, just re-enqueueing plus scoped resolution.

purge_scope_from_backend(old_store: Any, *, tenant_id: str, workspace_id: int | None = None) -> bool

Delete a scope's documents from an explicitly identified old backend.

old_store must be resolved by the caller (e.g. via VectorStoreRegistryService against the previous VectorStoreProfile id) -- this never guesses "the previous config" itself, since a wrong guess would silently delete the wrong backend's data.

dispatch

Enqueue/dispatch documents' ingestion jobs onto the shared msflib.ingestion queue.

Always defers Celery dispatch to the caller's own commit (whether that's this call, via commit=True, or a later one the caller owns, e.g. a pre-commit event hook) -- see IngestionDispatchService.enqueue_and_dispatch's dispatch_after_commit.

DocumentsQueueDisabledError

Bases: RuntimeError

Raised by enqueue_for_record when DOCUMENTS.QUEUE_ENABLED is False.

loaders

register_pdf_placeholder() -> None

Optional placeholder adapter for PDF bytes.

This keeps the registry extension point explicit without forcing a PDF extraction dependency in the base module.

pipeline

text_document pipeline kind (Layer 2) for msflib.ingestion.

Wraps documents' existing, unchanged business logic -- adapter-based source resolution, ingest_document_content/delete_index_for_source_key -- behind the msflib.ingestion stage protocol.

A single stage does Fetch+Normalize+Transform+Persist+Index in one call rather than splitting into several stage functions, since each of those steps depends on state (the resolved SourceDocument, adapter) from the previous one and none of them are independently reusable by another pipeline_kind.

make_sync_document_stage(*, vector_store_factory: Callable[[], Any] | None = None, record_action: DocumentRecordAction | None = None, collection_name: str = 'documents') -> Stage

Build the text_document sync stage.

vector_store_factory/collection_name are injectable (mirroring the old router()'s get_vector_store override) so tests/alternate deployments can swap the vector store without monkeypatching internals. Everything else (max ingest size, resolve-per-request) is read live from context.settings so a worker process picks up config changes without restarting.

register_text_document_pipeline(registry: PipelineKindRegistry = default_pipeline_registry, *, stage: Stage | None = None, replace: bool = False) -> None

Register the text_document pipeline kind (call once per worker process).

reindex

enqueue_reindex_for_scope(session: Session, *, dispatch_service: DocumentsIngestionDispatchService, tenant_id: int, workspace_id: int | None = None, statuses: tuple[DocumentRecordStatus, ...] = (DocumentRecordStatus.indexed, DocumentRecordStatus.failed)) -> ReindexReport

Re-enqueue every in-scope, non-deleted document record for re-ingestion.

Switching a scope's vector store backend (see VectorStoreProfile/VectorStoreRegistryService) leaves documents already indexed under the old backend behind: searches against the new backend return nothing for them until they are re-ingested. This re-enqueues matching records with IngestionReason.manual_reindex; the existing ingestion worker re-runs ingest_document_content, which re-embeds and writes into whatever store the record's own scope now resolves to (see services/ingestion.py's resolve_vector_store_per_request path) -- no separate reindex pipeline is needed, just re-enqueueing plus scoped resolution.

purge_scope_from_backend(old_store: Any, *, tenant_id: str, workspace_id: int | None = None) -> bool

Delete a scope's documents from an explicitly identified old backend.

old_store must be resolved by the caller (e.g. via VectorStoreRegistryService against the previous VectorStoreProfile id) -- this never guesses "the previous config" itself, since a wrong guess would silently delete the wrong backend's data.