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
¶
msflib.documents.drivelink
¶
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.