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¶
DocumentRecordis the ingestion identity of one source.source_keyis unique and typed by prefix:direct:{id}for uploads anddrivelink:{node_id}for drivelink files. Other fields:source_type,tenant_id,account_id,workspace_id, optionalconversation_id,sub_thread_idandprivate_to_account_id,current_version,latest_indexed_version,statusandlast_error.- Status:
pending(created),queued(a job is waiting),indexed(the current version is indexed),failed(retries exhausted or a terminal error, withlast_error) anddeleted. DocumentUploadNodeis 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_idis always applied as a filter andNULLmatches only records with no workspace. Conversation, thread and private fields narrow a document further, andpromotewidens it back to workspace scope. Once promoted (scope_promoted_atis set), re-ingesting the source cannot narrow the scope again. The same scope is written into the vector metadata, validated against thedocuments.ingestprofile (tenant required, workspace required but nullable). - Source adapters (
msflib.documents.adapters): aSourceAdapterresolves a record to aSourceDocumentand reads its bytes.DirectUploadAdapterand, when drivelink is installed,DrivelinkAdapterare registered inadapters.registry. Register your own withregistry.register(adapter), keyed by itssource_type. Any node type implementingSourceNodeProtocol(DrivelinkNodedoes, unchanged) can be a source. - Dispatch:
DocumentsIngestionDispatchService.enqueue_for_record(session, record=, target_version=, reason=, commit=)creates anIngestionJobfor the record.IngestionReasoniscreate,replace,manual_reindexormetadata_change. Withcommit=Falsethe Celery dispatch is deferred until the caller's own commit. - Pipeline: the
text_documentstage (sync_document_stage) resolves the record, checks the version has not been superseded, reads the bytes, extracts text withmsflib.ai_core.services.extraction.default_loader_registry, prunes old chunks for thatsource_keyand indexes the new ones. If the source is gone or the record is deleted it removes the vectors instead. Usemake_sync_document_stage(vector_store_factory=..., collection_name=...)andregister_text_document_pipeline(stage=...)to swap the vector store or collection (default collection namedocuments). - Events: after each successful stage the worker emits
documents.ingestion.completedwith aDocumentIngestionCompletedEvent(document_record_id,target_version,reason,tenant_id,workspace_id), and the drivelink hooks emitdocuments.index.deletedafter commit when a drivelink file is deleted. The knowledge module listens to these. The worker emitsdocuments.ingestion.completedon the default emitter active in the worker process, so listeners for it must be registered in the worker process onAppEmitter(get_emitter())(see Event bus and ingestion); a knowledge hook registered only in the web process never fires from a real worker.documents.index.deletedis emitted inside web requests, so that one needs the web registration. - Drivelink hooks: file nodes only. Create enqueues a
createjob. An update that changesname,path,parent_id,mime_type,versionordeleted_atenqueues a job with reasoncreate(never indexed),metadata_change(same version as the indexed one) orreplace(version advanced); a node withdeleted_atset 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 unlessforce=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
queuedand records never becomeindexed. The worker is not running, or it did not callregister_text_document_pipeline(), or it uses a differentINGESTION.BROKER_URLfrom the web process. All three are required. Check the job's status and history on the ingestion router. 503on upload, promote or delete.DOCUMENTS.QUEUE_ENABLEDis false.500 Default tenant has not been seeded. Run tenant seeding at startup, or passget_current_tenant(required if you changedDEFAULT_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_accountwas not passed, or it returnedNone.- Record is
failedwith alast_errormentioning loaders. The file type needs an extraction backend that is not installed. Install thellamaindexorcommunityextra ofmsflib-ai-core; the error names the registries that were skipped. Plain text, markdown, JSON and CSV work out of the box. - Record is
failedwith a size message. The file exceedsMAX_INGEST_BYTES. - Drivelink files never create records.
register_event_hookswas not called in the process that handles the drivelink requests, ormsflib-drivelinkis 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. WithRESOLVE_VECTOR_STORE_PER_REQUEST=trueeach document is re-embedded into the store its scope now resolves to. Passpurge_old_profile_idto delete the scope's vectors from the previous backend; purge failures are reported inerrorsand 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.