Skip to content

msflib-documents

Purpose

msflib-documents turns uploaded files into searchable vectors. It accepts direct uploads over HTTP, optionally follows files in drivelink, and queues each one for ingestion: bytes are fetched from storage, converted to text, chunked, embedded and written to a vector store. It tracks each source as a DocumentRecord with a status, a version and a scope (tenant, workspace, and optionally conversation, thread or private to one account).

Ingestion runs as a pipeline kind (text_document) on the shared ingestion job queue, executed by a Celery worker. The FastAPI process only enqueues jobs. Text extraction, chunking, embeddings and the vector store come from ai-core.

Install

[tool.poetry.dependencies]
msflib = { git = "https://github.com/msflib/fastapi.git", subdirectory = "core", rev = "core-v0.2.1" }
msflib-account = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/account", rev = "account-v0.2.2" }
msflib-tenancy = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/tenancy", rev = "tenancy-v0.2.0" }
msflib-ai-core = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/ai_core", rev = "ai_core-v0.2.2" }
msflib-ingestion = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/ingestion", rev = "ingestion-v0.2.0" }
msflib-documents = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/documents", rev = "documents-v0.2.1" }
# only if you want drivelink files ingested automatically:
msflib-drivelink = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/drivelink", rev = "drivelink-v0.2.0" }

msflib-documents also declares a drivelink extra (msflib-documents[drivelink]) that pulls in msflib-drivelink. Without drivelink, the drivelink adapter is simply not registered and direct uploads work normally.

Extraction is only as capable as the msflib-ai-core extras you install. The built-in fallback loader reads plain text, markdown, JSON and CSV. For PDFs and Office files add the llamaindex or community extra of msflib-ai-core (see ai-core); without one, extracting such a file fails the ingestion job. You also need an embedding provider and a vector store configured in AI_CORE (a vector store extra such as pgvector).

The ingestion queue needs a Redis broker (INGESTION.BROKER_URL, default redis://localhost:6379/0) reachable by both the web process and the worker.

Wiring into a host app

There are three parts: the router, the worker, and (optionally) the drivelink hooks.

1. Settings. Include DocumentsSettings and IngestionSettings in your settings class. The upload endpoint also needs the default tenant to exist (seed it at startup; see tenancy). Import msflib.tenancy.models.tenant, msflib.documents.models and msflib.ingestion.models before create_all, so the tables and the tenant_id foreign keys are registered.

from msflib.documents.config import DocumentsSettings
from msflib.ingestion.config import IngestionSettings

class AppSettings(DocumentsSettings, IngestionSettings, ..., CoreSettings, SettingsBase):
    pass

2. Router. router in msflib.documents.router takes keyword-only arguments.

from msflib.documents.router import router as documents_router
from msflib.documents.services.pipeline import register_text_document_pipeline

app.include_router(
    documents_router(
        get_session=get_session,
        get_current_account=get_current_account,
        get_current_workspace=get_current_workspace,   # optional
        get_current_tenant=get_current_tenant,         # optional
        settings=settings,
        prefix="/documents",
    )
)
register_text_document_pipeline()

get_current_account is optional in the signature but every route except GET /ingestion/status answers 401 without it. get_current_workspace and get_current_tenant are optional; with no tenant dependency the router falls back to the tenant with slug default and returns 500 if it has not been seeded. dispatch_task lets you inject a .delay(job_id)-compatible object instead of the default Celery producer, which is how the tests avoid needing Redis.

3. Worker. The process that runs the jobs must register the pipeline and build the Celery app:

from msflib.documents.services.pipeline import register_text_document_pipeline
from msflib.ingestion.celery import build_celery_app
from sqlmodel import Session

from your_app.db import engine   # your own engine


def session_factory() -> Session:
    return Session(engine)

register_text_document_pipeline()
app = build_celery_app(settings, session_factory=session_factory)

msflib.db.session.engine only works in hosts laid out as app.core.config, so build the session from your own engine. Run it with celery -A your_module worker. The repository's testsite/app/ingestion_worker.py is a working entry point that also registers the knowledge pipelines.

4. Drivelink hooks (optional). Call bind_app_emitter(app) when you create the app (before startup), then call this once, for example in your lifespan handler, passing get_app_emitter(app). get_app_emitter binds on first use, and binding adds middleware, which fails with RuntimeError: Cannot add middleware after an application has started inside a lifespan handler if the app was not bound earlier.

from msflib.documents.actions import DocumentRecordAction
from msflib.documents.drivelink import register_event_hooks
from msflib.documents.services import DocumentsIngestionDispatchService
from msflib.eventbus import bind_app_emitter, get_app_emitter
from msflib.ingestion.celery import build_dispatch_task
from msflib.ingestion.services.dispatch import IngestionDispatchService

app_emitter = bind_app_emitter(app)   # at app creation

record_action = DocumentRecordAction()
register_event_hooks(
    record_action=record_action,
    dispatch_service=DocumentsIngestionDispatchService(
        dispatch=IngestionDispatchService(task=build_dispatch_task(settings)),
        record_action=record_action,
    ),
    emitter=get_app_emitter(app),   # the same emitter as app_emitter
    tenancy_settings=settings.scope("TENANCY"),
)

Known issue (#288)

The router's default-tenant fallback ignores the host's TENANCY.DEFAULT_TENANT_SLUG. If you changed the slug, POST /upload returns 500 "Default tenant has not been seeded" even though your seeding used the same settings, while the drivelink hooks (given tenancy_settings) resolve your tenant. Pass get_current_tenant=get_tenant_dependencies(session_dep=get_session, settings=settings.scope("TENANCY")).get_current_tenant to the router.

Routes added by the router:

Method and path Effect
POST /upload Multipart file; optional form fields conversation_id, sub_thread_id, is_private. Stores the file, creates a record and queues a job
GET / List the caller's records in the current workspace (status, conversation_id, sub_thread_id, include_deleted, limit up to 100, offset), newest first
GET /{record_id} One record, with name, mime_type, size_bytes and a download url
GET /{record_id}/download The original bytes as an attachment
POST /{record_id}/promote Widen a conversation- or thread-scoped document to workspace scope (?clear_private=true also clears privacy) and re-index
DELETE /{record_id} Mark deleted and queue removal from the vector store
POST /reindex Re-enqueue the workspace's indexed and failed documents, for example after switching vector store backend
GET /ingestion/status Count of records by status for the caller and workspace

Access control is by owner: a record whose account_id is not the caller's gets 403, and a missing or deleted one gets 404. For job-level status (attempts, retries, history) use the ingestion module's own router with subject_type=document.

Configuration

DocumentsSettings has the DOCUMENTS namespace.

Key Default Notes
ENABLED true Declared, currently has no effect (#278). Not read by the module itself; host apps use it to decide whether to register hooks (the testsite does)
QUEUE_ENABLED true When false, anything that would enqueue a job fails. Upload, promote and delete answer 503
MAX_INGEST_BYTES 100 MiB Upload limit (HTTP 413) and the size check the worker applies before ingesting
RESOLVE_VECTOR_STORE_PER_REQUEST false When true, each document's vector store is resolved from its own tenant and workspace through the AI vector-store profiles, instead of one process-wide store

Flat aliases exist for all four keys. With subclassed settings use the flat names (QUEUE_ENABLED=false, MAX_INGEST_BYTES=...); DOCUMENTS__QUEUE_ENABLED is only read when you declare DOCUMENTS: DocumentsSettings = DocumentsSettings() as a field. Both were checked. The flat alias ENABLED is shared with the other modules that declare it, so with documents and drivelink both composed as nested fields DOCUMENTS__ENABLED=false was ignored while DRIVELINK__ENABLED=false switched off both; set the field in code instead, for example DocumentsSettings(ENABLED=False).

Behaviour may change (#280)

The NAMESPACE__KEY environment variable style only works for settings composed as a field, not for subclassed hosts. This may change.

Retries, backoff and attempt limits are not documents settings. They live in INGESTION.* (for example INGESTION.BROKER_URL, INGESTION.MAX_ATTEMPTS); see ingestion. Storage location for uploaded files uses core's CORE.STORAGE_METHOD and related keys. Embedding and vector-store selection use AI_CORE.*.

Key concepts

  • DocumentRecord is the ingestion identity of one source. source_key is unique and typed by prefix: direct:{id} for uploads and drivelink:{node_id} for drivelink files. Other fields: source_type, tenant_id, account_id, workspace_id, optional conversation_id, sub_thread_id and private_to_account_id, current_version, latest_indexed_version, status and last_error.
  • Status: pending (created), queued (a job is waiting), indexed (the current version is indexed), failed (retries exhausted or a terminal error, with last_error) and deleted.
  • DocumentUploadNode is the documents-owned source row for direct uploads: file name, storage location and owner. It is immutable; scope changes are made on the record.
  • Scope: a record is visible to its account within its workspace. workspace_id is always applied as a filter and NULL matches only records with no workspace. Conversation, thread and private fields narrow a document further, and promote widens it back to workspace scope. Once promoted (scope_promoted_at is set), re-ingesting the source cannot narrow the scope again. The same scope is written into the vector metadata, validated against the documents.ingest profile (tenant required, workspace required but nullable).
  • Source adapters (msflib.documents.adapters): a SourceAdapter resolves a record to a SourceDocument and reads its bytes. DirectUploadAdapter and, when drivelink is installed, DrivelinkAdapter are registered in adapters.registry. Register your own with registry.register(adapter), keyed by its source_type. Any node type implementing SourceNodeProtocol (DrivelinkNode does, unchanged) can be a source.
  • Dispatch: DocumentsIngestionDispatchService.enqueue_for_record(session, record=, target_version=, reason=, commit=) creates an IngestionJob for the record. IngestionReason is create, replace, manual_reindex or metadata_change. With commit=False the Celery dispatch is deferred until the caller's own commit.
  • Pipeline: the text_document stage (sync_document_stage) resolves the record, checks the version has not been superseded, reads the bytes, extracts text with msflib.ai_core.services.extraction.default_loader_registry, prunes old chunks for that source_key and indexes the new ones. If the source is gone or the record is deleted it removes the vectors instead. Use make_sync_document_stage(vector_store_factory=..., collection_name=...) and register_text_document_pipeline(stage=...) to swap the vector store or collection (default collection name documents).
  • Events: after each successful stage the worker emits documents.ingestion.completed with a DocumentIngestionCompletedEvent (document_record_id, target_version, reason, tenant_id, workspace_id), and the drivelink hooks emit documents.index.deleted after commit when a drivelink file is deleted. The knowledge module listens to these. The worker emits documents.ingestion.completed on the default emitter active in the worker process, so listeners for it must be registered in the worker process on AppEmitter(get_emitter()) (see Event bus and ingestion); a knowledge hook registered only in the web process never fires from a real worker. documents.index.deleted is emitted inside web requests, so that one needs the web registration.
  • Drivelink hooks: file nodes only. Create enqueues a create job. An update that changes name, path, parent_id, mime_type, version or deleted_at enqueues a job with reason create (never indexed), metadata_change (same version as the indexed one) or replace (version advanced); a node with deleted_at set marks the record deleted instead of enqueueing. Delete marks the record deleted. Folders are ignored. Hooks resolve the default tenant, so seed it at startup. Registering twice is a no-op unless force=True.

msflib.documents exposes DocumentRecordAction, DocumentsIngestionDispatchService, DocumentsSettings, loader_registry and register_event_hooks lazily from its package __init__.

Examples

Upload, inspect, promote and delete through the router. This ran against in-memory SQLite and local file storage, with a Mock in place of the Celery task. In a real deployment omit dispatch_task:

from unittest.mock import Mock
from fastapi import FastAPI
from fastapi.testclient import TestClient
from msflib.documents.router import router as documents_router

dispatch_task = Mock()
app = FastAPI()
app.include_router(documents_router(
    get_session=lambda: session, get_current_account=lambda: account,
    settings=settings, dispatch_task=dispatch_task,
))
c = TestClient(app)

r = c.post("/documents/upload",
           files={"file": ("notes.txt", b"hello world", "text/plain")},
           data={"conversation_id": "conv-1"})
# 200 {'record_id': 1, 'job_id': 1, 'source_key': 'direct:1', 'status': 'queued',
#      'workspace_id': None, 'conversation_id': 'conv-1', ...}
dispatch_task.delay.call_count               # 1

c.get("/documents/1").json()["status"]       # "queued"
c.get("/documents/1/download").content       # b"hello world"
c.post("/documents/1/promote").json()        # conversation_id now None, a second job queued
c.get("/documents/ingestion/status").json()  # {'module': 'documents', 'workspace_id': None, 'records': {'queued': 1}}
c.delete("/documents/1").json()["status"]    # job status, "queued"; the record is now "deleted"

The status in the delete response is the job's status. The record itself is deleted.

Ingest drivelink files automatically. With register_event_hooks registered as above, uploading a file through the drivelink router creates a document record and a job, and trashing it marks the record deleted. This ran with the same mock task:

# after drivelink upload of a.txt to a folder
[(r.source_key, r.source_type, r.status.value) for r in session.exec(select(DocumentRecord)).all()]
# [('drivelink:2', 'drivelink', 'queued')]

# after DELETE /drivelink/nodes/2 (soft delete)
# [('drivelink:2', 'deleted')]

Add a source adapter for your own storage:

from msflib.documents.adapters import registry, SourceDocument

class MyAdapter:
    source_type = "mystore"

    def resolve(self, session, record):          # return SourceDocument or None if gone
        ...
    def read_bytes(self, session, record, *, settings): ...
    def head(self, session, record, *, settings): ...

registry.register(MyAdapter())

Your own code then creates a DocumentRecord whose source_type matches the adapter (using DocumentRecordAction().create(session, data=DocumentRecordCreate(...))) and enqueues it with enqueue_for_record. The ready-made upsert helpers are upsert_from_upload_node and upsert_from_drivelink_node. This adapter snippet was not run.

Troubleshooting

  • Jobs stay queued and records never become indexed. The worker is not running, or it did not call register_text_document_pipeline(), or it uses a different INGESTION.BROKER_URL from the web process. All three are required. Check the job's status and history on the ingestion router.
  • 503 on upload, promote or delete. DOCUMENTS.QUEUE_ENABLED is false.
  • 500 Default tenant has not been seeded. Run tenant seeding at startup, or pass get_current_tenant (required if you changed DEFAULT_TENANT_SLUG, see the known issue under Wiring). The drivelink hooks raise the same error when no default tenant exists.
  • 401 Authentication required. get_current_account was not passed, or it returned None.
  • Record is failed with a last_error mentioning loaders. The file type needs an extraction backend that is not installed. Install the llamaindex or community extra of msflib-ai-core; the error names the registries that were skipped. Plain text, markdown, JSON and CSV work out of the box.
  • Record is failed with a size message. The file exceeds MAX_INGEST_BYTES.
  • Drivelink files never create records. register_event_hooks was not called in the process that handles the drivelink requests, or msflib-drivelink is not installed. Without it the drivelink adapter is not registered either.
  • Search after switching vector store backend returns nothing for old documents. Call POST /reindex. With RESOLVE_VECTOR_STORE_PER_REQUEST=true each document is re-embedded into the store its scope now resolves to. Pass purge_old_profile_id to delete the scope's vectors from the previous backend; purge failures are reported in errors and do not block the reindex.
  • Listing shows no documents although uploads succeeded. Listing filters on the caller's account and the current workspace exactly. Documents uploaded with a workspace dependency are not listed by requests made without one.

API reference

See the generated API reference for msflib.documents. The module's README.md in modules/documents has a short scope summary.

See also