msflib-knowledge¶
Purpose¶
msflib-knowledge is a scoped knowledge layer for host apps. Facts (entities, relations, observations) are stored in your Postgres database with provenance and tenant and workspace scope, projected to an Apache AGE graph and a vector store, and read back through a small set of semantic query objects that double as agent tools. You feed it from documents, structured records (forms, JSON payloads, rows in your own tables) or your own code.
It does not run extraction by itself: ingestion runs as jobs on msflib-ingestion, and LLM calls go through ai-core. Reads are always scope-filtered; there is no unscoped query.
Install¶
[tool.poetry.dependencies]
msflib = { git = "https://github.com/msflib/fastapi.git", subdirectory = "core", rev = "core-v0.2.1" }
msflib-knowledge = { git = "https://github.com/msflib/fastapi.git", subdirectory = "modules/knowledge", rev = "knowledge-v0.2.1", extras = ["age", "pgvector"] }
msflib-knowledge depends on msflib-ai-core, msflib-documents, msflib-ingestion and msflib-tenancy.
| Extra | Enables |
|---|---|
age |
psycopg, needed for the Apache AGE graph projection |
pgvector |
langchain-postgres and psycopg, for the pgvector vector projection |
On SQLite (CORE.USE_SQLITE) only the canonical projection runs, so neither extra is needed for local development and tests.
The package __init__ lazily exports KnowledgeEngine, KnowledgeSettings, the three row actions, KnowledgeCompiler, KnowledgeFragment, Provenance, OntologyRegistry, default_ontology, KnowledgeToolRegistry, CanonicalProjection, LLMExtractionParser, IngestResult and provision_age_and_pgvector. The examples below import from submodules, which is equally valid.
Wiring into a host app¶
Add KnowledgeSettings to your settings class and import msflib.knowledge.models so its tables are created. Then build the engine with get_knowledge_dependencies, which wires the canonical projection always, the AGE graph projection when not on SQLite, and a vector-store projection through ai-core's embeddings backend:
from msflib.knowledge.deps import get_knowledge_dependencies
knowledge_deps = get_knowledge_dependencies(
settings,
session_factory=lambda: Session(engine),
record_dispatch_task=build_dispatch_task(settings), # only for structured records
)
knowledge = knowledge_deps.get_knowledge_engine() # built once, cached
The namespace has build_knowledge_engine (fresh engine on each call, for tests or your own lifecycle), get_knowledge_engine (the cached one), get_record_ingestion_service (a cached KnowledgeRecordService, which needs record_dispatch_task) and ontology (the registry wired into the engine).
Other parameters: vector_store_factory to supply your own vector store, ontology for a fully custom OntologyRegistry, additional_ontology_pack_namespaces for your own ontology packs, and ontology_pack_options for per-pack keyword arguments.
On Postgres the AGE projection writes cypher() statements through the same SQLAlchemy session as the canonical writes, so use a postgresql+psycopg:// URL. A different driver logs a warning. Run provision_age_and_pgvector(superuser_dsn, app_role=..., graph_name=...) once per database with a superuser connection to create the extensions, the graph and the grants; the app role cannot do this itself.
Extraction from documents¶
When msflib-documents finishes indexing a document, knowledge can extract facts from it. The hooks listen for documents.ingestion.completed, which the worker emits on its default emitter, and for documents.index.deleted, which is emitted inside web requests. Register the hooks with a dispatch service in the worker entrypoint, on an emitter the worker uses, and also in the web process for the delete event. In the worker, register the pipeline kind too.
from msflib.eventbus import AppEmitter, get_emitter
from msflib.ingestion.celery import build_dispatch_task
from msflib.ingestion.services.dispatch import IngestionDispatchService
from msflib.knowledge.integrations.dispatch import KnowledgeIngestionDispatchService
from msflib.knowledge.integrations.documents import register_event_hooks
def register_extraction_hooks(emitter):
register_event_hooks(
dispatch_service=KnowledgeIngestionDispatchService(
dispatch=IngestionDispatchService(task=build_dispatch_task(settings)),
queue_enabled=settings.scope("KNOWLEDGE").EXTRACTION_QUEUE_ENABLED,
),
emitter=emitter,
)
# worker entrypoint (there is no FastAPI app): the worker's default emitter
register_extraction_hooks(AppEmitter(get_emitter()))
# web process, after app_emitter = bind_app_emitter(app) at app creation
register_extraction_hooks(app_emitter)
# worker process
from msflib.knowledge.integrations.pipeline import (
make_knowledge_extraction_stage,
register_knowledge_extraction_pipeline,
)
register_knowledge_extraction_pipeline(stage=make_knowledge_extraction_stage(engine=knowledge))
The hooks enqueue knowledge_extraction jobs. Hooks on the web app's emitter fire for events emitted inside a request (or under with use_app_emitter(app):) but not for an emit outside a request, which is where the worker emits documents.ingestion.completed; AppEmitter(get_emitter()) registered at worker start does receive it, from any thread. The stage reads the document through its source adapter, extracts a fragment (with a registered content extractor if one claims the bytes, otherwise with the LLM), and ingests it, superseding facts from earlier versions of the same source.
Known issue (#289)
The knowledge engine registers one of its own listeners on the global emitter, which routes in an app using bind_app_emitter(app) do not reach, and events emitted in a worker or script only reach hooks registered in that process. Register the hooks in the process that emits the event, as shown above.
Known issue (#277)
knowledge_extraction is a bespoke stage, not built on ingestion_subject_stage. It has no per-subject advisory lock or fingerprint re-check, which the knowledge_record kind has, so do not assume the guarantees described for structured records apply to extraction jobs.
Structured records¶
Forms, JSON payloads and your own table rows go through the knowledge_record pipeline kind. Register it in the worker and build the service in the web process:
# worker process
from msflib.knowledge.records import make_knowledge_record_stage, register_knowledge_record_pipeline
register_knowledge_record_pipeline(
stage=make_knowledge_record_stage(engine=knowledge, ontology=knowledge_deps.ontology)
)
register_my_record_sources_and_mappers() # sources and mappers must exist in BOTH processes
# web process
records = knowledge_deps.get_record_ingestion_service()
register_my_record_sources_and_mappers()
See the structured-records example below and the long-running jobs guide for the worker and producer setup.
Agent tools¶
KnowledgeToolRegistry(engine) exposes one tool per supported query type and satisfies ai-core's ToolRegistry protocol. With ai-api, pass get_knowledge_engine=knowledge_deps.get_knowledge_engine to the router, or build the registry list yourself with resolve_knowledge_tool_registries:
from msflib.knowledge.integrations.ai_api import resolve_knowledge_tool_registries
get_tool_registries = resolve_knowledge_tool_registries(knowledge_deps.get_knowledge_engine)
Configuration¶
Settings live in the KNOWLEDGE namespace. With a nested KNOWLEDGE settings object, use KNOWLEDGE__GRAPH_BACKEND style variables. These flat aliases also work: KNOWLEDGE_GRAPH_BACKEND, KNOWLEDGE_GRAPH_NAME, KNOWLEDGE_DATABASE_URL, KNOWLEDGE_EXTRACTION_QUEUE_ENABLED, KNOWLEDGE_RECORD_QUEUE_ENABLED and KNOWLEDGE_RECORD_MAX_PAYLOAD_BYTES. When you mix KnowledgeSettings into your settings class, the field names (GRAPH_BACKEND, ...) are top-level.
Behaviour may change (#280)
The KNOWLEDGE__KEY environment variable style only works for settings composed as a field, not for subclassed hosts. This may change.
| Key | Default | Notes |
|---|---|---|
GRAPH_BACKEND |
age |
age or none. none skips the AGE projection on real Postgres. |
GRAPH_NAME |
msflib_knowledge |
AGE graph name. |
DATABASE_URL |
None |
DSN for the AGE connection. Defaults to the app database. Setting it with CORE.USE_SQLITE true raises ValueError. |
ONTOLOGY_PACKS |
["rag_default"] |
Packs to register. rag_default is always included. |
DRAFT_CONFIDENCE_THRESHOLD |
0.7 |
Declared, currently has no effect (#278). get_knowledge_dependencies does not read it (the compiler default is also tracked in #277): the engine it builds always uses the compiler default of 0.7. To change the threshold, build the engine yourself with KnowledgeCompiler(ontology=..., draft_confidence_threshold=...). |
ENTITY_RESOLUTION_SIMILARITY_THRESHOLD |
0.9 |
Declared, currently has no effect (#278). Not read by any code in the module at present. |
EXTRACTION_QUEUE_ENABLED |
True |
Whether document-completion events enqueue extraction jobs. |
RECORD_QUEUE_ENABLED |
True |
Whether structured records are queued. Independent of the previous flag. |
RECORD_MAX_PAYLOAD_BYTES |
262144 |
Maximum inline record payload after field selection. |
Retry and backoff limits for extraction and record jobs come from INGESTION.*, not from this namespace.
Key concepts¶
Canonical model. Three tables hold the facts: entities, relations and observations (KnowledgeEntity, KnowledgeRelation, KnowledgeObservation, accessed through KnowledgeEntityAction, KnowledgeRelationAction and KnowledgeObservationAction). Each row carries a type and properties, scope columns, a lifecycle (draft, approved, archived) and provenance (source_id, source_type, extracted_by, confidence). The graph and vector stores are projections of this data; you never write to them directly.
Knowledge IR and ingest. Parsers and mappers produce a KnowledgeFragment (from msflib.knowledge.ir) holding EntitySpec, RelationSpec and ObservationSpec items plus a Provenance. KnowledgeEngine.ingest(session, fragment, scope=..., commit=True, supersede_source=False) compiles it and persists it in one transaction. Entities that match an existing row (same scope, type and case-insensitive name) are resolved instead of duplicated. commit=False joins the caller's transaction. supersede_source=True first supersedes every live relation and observation from the same source_id. archive_source(...) does just that supersession; entities are never archived because other sources may rely on them.
Lifecycle. LLM fragments below the compiler's draft confidence threshold (0.7 by default) become draft. Parser and human fragments are trusted and land approved. Draft facts are hidden from queries unless include_drafts=True.
Scope. Every read goes through a scope rewriter. Scope is a ScopeEnvelope with tenant_id and workspace_id dimensions, plus optional project_id, conversation_id and sub_thread_id. The last three are read hierarchically: a row with no project_id is visible to everyone in the workspace, a row with one is visible only to callers carrying the same value. workspace_id=None means the workspace-less bucket, not a wildcard. See Scopes, tenancy and workspaces.
Queries. msflib.knowledge.queries defines FindEntity, FindNeighbors, FindPath, Traverse, SearchConcepts, FindEvidence and ResolveAliases. engine.query(query, scope=..., include_drafts=False) returns a QueryResult with items, evidence and metadata. A query is routed to a projection that supports its capability: entity lookup, evidence and alias resolution run on the canonical projection; neighbor, path and traversal queries need the AGE graph; SearchConcepts needs the vector store. engine.supports_capability(...) tells you what your wiring offers.
Ontology. An OntologyRegistry defines entity and relation types. The built-in rag_default pack provides entity types Document, Concept, Entity and Chunk, and relation types MENTIONS, PART_OF, RELATED_TO, DERIVED_FROM and POSSIBLE_DUPLICATE_OF. Add your own types with register_entity_type and register_relation_type, or ship a pack module exposing register(registry) and load it with additional_ontology_pack_namespaces.
Pluggable content extractors. register_extractor(...) in msflib.knowledge.integrations.extractors registers a ContentExtractor that sniffs raw bytes and, when it claims a document, replaces the LLM pass for it. Selection is first-match-wins by priority and the stage overrides the returned fragment's provenance, so an extractor cannot misattribute facts.
Events. knowledge-fragment-pre-ingest, knowledge-fragment-ingested, knowledge-source-archived and knowledge-record-processed (fired in the worker when a record job commits, so register its listeners in the worker process, on the emitter active there; see Extraction from documents). Row-level lifecycle events such as knowledgeentity-create-post-commit come from ModelAction. See Event bus.
Records. A record is one structured item: a payload plus its identity, scope and version. A job carries only a pointer, (subject_type, subject_id). The worker loads the record through the record source registered for that subject type, a mapper registered for its record_type turns it into a fragment, and previous facts for the record are archived before the new ones are ingested. Facts carry source_id = "<subject_type>:<subject_id>".
| Inline: knowledge stores the data | App-owned: your table is the source of truth | |
|---|---|---|
| Use when | The data has no table of its own | Your app already persists it |
| You write | A mapper and a records.submit(...) call |
A mapper, a ModelRecordSource(...) and sync_model_to_queue(...) |
| Delete | records.retract(...) |
Delete or soft-delete the row |
Inline records are stored in the knowledge_record table. create_all creates it on a fresh database; existing deployments must add it, with its two partial unique indexes on (tenant_id, workspace_id, record_type, external_id) and (tenant_id, record_type, external_id) for rows without a workspace, in their own migration before the first submit().
Examples¶
Ingest and query facts¶
Runs as written against in-memory SQLite (canonical projection only). It seeds tenant 1 explicitly: SQLite does not enforce the tenant_id foreign key, but Postgres does, and without the tenant row inserts fail there.
import msflib.knowledge.models # noqa: F401 (registers the tables)
from msflib.core.config import CoreSettings, SettingsBase
from msflib.knowledge.config import KnowledgeSettings
from msflib.knowledge.deps import get_knowledge_dependencies
from msflib.knowledge.ir import EntitySpec, KnowledgeFragment, Provenance, RelationSpec
from msflib.knowledge.queries import FindEntity
from msflib.scope import ScopeEnvelope
from msflib.tenancy.models.tenant import Tenant # noqa: F401
from msflib.tenancy.resolver import resolve_default_tenant_id
from sqlalchemy.pool import StaticPool
from sqlmodel import Session, SQLModel, create_engine
class Settings(KnowledgeSettings, CoreSettings, SettingsBase):
pass
settings = Settings() # CORE.USE_SQLITE defaults to true
engine = create_engine("sqlite://", connect_args={"check_same_thread": False}, poolclass=StaticPool)
SQLModel.metadata.create_all(engine)
with Session(engine) as session:
resolve_default_tenant_id(session) # seeds tenant 1
session.commit()
knowledge = get_knowledge_dependencies(
settings, session_factory=lambda: Session(engine)
).get_knowledge_engine()
scope = ScopeEnvelope.from_dict({"dimensions": {"tenant_id": "1", "workspace_id": "7"}})
fragment = KnowledgeFragment(
provenance=Provenance(source_id="doc-42", source_type="text", extracted_by="llm"),
entities=[
EntitySpec(local_id="e1", type="Concept", name="Pump P-101", confidence=0.9),
EntitySpec(local_id="e2", type="Concept", name="Separator V-201", confidence=0.9),
],
relations=[RelationSpec(source="e1", type="RELATED_TO", target="e2")],
)
with Session(engine) as session:
result = knowledge.ingest(session, fragment, scope=scope)
print(sorted(result.uid_map)) # ['e1', 'e2']
found = knowledge.query(FindEntity(name="Pump P-101"), scope=scope)
print([item["name"] for item in found.items]) # ['Pump P-101']
result.uid_map maps each fragment-local id to the stored uid. A caller with a different ScopeEnvelope (another tenant or workspace) gets no items from the same query.
Submit a structured record inline¶
This maps an intake form to facts, queues the job, and runs the stage in-process using the ingestion test helpers. Isolated registries keep the example from touching process-wide defaults. It runs as written against SQLite, and seeds tenant 1 for the same reason as the example above.
import msflib.knowledge.models # noqa: F401
from msflib.core.config import CoreSettings, SettingsBase
from msflib.db.sqlite import enable_savepoints
from msflib.ingestion.models import IngestionJob
from msflib.ingestion.services.dispatch import IngestionDispatchService
from msflib.ingestion.testing import FakeDispatchTask, start_stage_context
from msflib.knowledge.config import KnowledgeSettings
from msflib.knowledge.deps import get_knowledge_dependencies
from msflib.knowledge.ir import EntitySpec, KnowledgeFragment, ObservationSpec, RelationSpec
from msflib.knowledge.queries import FindEntity
from msflib.knowledge.records import (
KnowledgeRecordService,
RecordMapperRegistry,
RecordSourceRegistry,
make_knowledge_record_stage,
)
from msflib.scope import ScopeEnvelope
from msflib.tenancy.models.tenant import Tenant # noqa: F401
from msflib.tenancy.resolver import resolve_default_tenant_id
from sqlalchemy.pool import StaticPool
from sqlmodel import Session, SQLModel, create_engine
class Settings(KnowledgeSettings, CoreSettings, SettingsBase):
pass
settings = Settings()
engine = enable_savepoints(
create_engine("sqlite://", connect_args={"check_same_thread": False}, poolclass=StaticPool)
)
SQLModel.metadata.create_all(engine)
with Session(engine) as session:
resolve_default_tenant_id(session) # seeds tenant 1
session.commit()
knowledge = get_knowledge_dependencies(
settings, session_factory=lambda: Session(engine)
).get_knowledge_engine()
def intake_form(record, *, context):
submission_id = record.metadata.get("external_id", record.subject_id)
form = EntitySpec(local_id="form", type="Document", name=f"Intake form {submission_id}")
company = EntitySpec(local_id="company", type="Concept", name=record.payload["company"])
return KnowledgeFragment(
provenance=context.provenance,
entities=[form, company],
relations=[RelationSpec(source="form", type="MENTIONS", target="company")],
observations=[
ObservationSpec(entity_refs=["form", "company"], properties=dict(record.payload))
],
)
sources, mappers = RecordSourceRegistry(), RecordMapperRegistry()
mappers.register("intake_form", intake_form)
task = FakeDispatchTask()
records = KnowledgeRecordService(
dispatch=IngestionDispatchService(task=task), sources=sources, mappers=mappers
)
stage = make_knowledge_record_stage(engine=knowledge, sources=sources, mappers=mappers)
scope = ScopeEnvelope.from_dict({"dimensions": {"tenant_id": "1", "workspace_id": "7"}})
with Session(engine) as session:
records.submit(
session,
record_type="intake_form",
external_id="s-1",
payload={"company": "Acme Corp", "needs": ["pumps"], "password": "hunter2"},
fields=["company", "needs"],
scope=scope,
commit=True,
)
job = session.get(IngestionJob, task.dispatched[0])
print(job.pipeline_kind, job.reason) # knowledge_record create
stage(start_stage_context(session, job, settings)) # a worker would do this
session.commit()
found = knowledge.query(FindEntity(entity_type="Concept"), scope=scope)
print([item["name"] for item in found.items]) # ['Acme Corp']
In a real app, KnowledgeRecordService comes from get_record_ingestion_service(), the dispatch handle is build_dispatch_task(settings), and the stage runs in the Celery worker. submit only stores the filtered payload (here password is never stored, and fields= keeps only the listed keys) and queues a pointer; the job is dispatched after your transaction commits.
Notes for mappers: name each record's root entity uniquely (entities merge by scope, type and lowercase name), put answers that can change on an observation (observations from a record are replaced on reprocessing, entity properties are merged and never shrink), and use types your ontology knows. records.preview(session, record, settings=..., compiler=...) runs a mapper without persisting anything.
Ingest rows from your own table¶
Declare a source once, register it in both processes, and keep it in sync:
from functools import partial
from msflib.core.records import version_from_column
from msflib.ingestion.services.sync import sync_model_to_queue
from msflib.knowledge.records import ModelRecordSource, default_record_source_registry
form_submission_source = ModelRecordSource(
subject_type="crm.form_submission",
model=FormSubmission,
record_type="intake_form",
payload="data",
fields=["company", "needs", "notes"],
version=version_from_column("revision"),
is_deleted="deleted_at",
)
default_record_source_registry.register(form_submission_source)
sync_model_to_queue(
FormSubmission,
records.queue,
subject_type=form_submission_source.subject_type,
snapshot=partial(records.snapshot_row, form_submission_source),
emitter=get_app_emitter(app),
)
From then on, every ModelAction create, update and delete on FormSubmission queues the right job in the same transaction as your write, while the app emitter is active: in requests (including FastAPI BackgroundTasks), or in your own code wrapped in with use_app_emitter(app):. Writes from workers, scripts and threads outside that block queue nothing; call records.queue.enqueue_snapshot there. Edits queue a replace, deletes and soft deletes queue a delete, which archives the record's facts. This snippet was not run (it needs your own model); the names come from msflib.knowledge.records and msflib.ingestion.services.sync. Without a version=, an edit made while a job for the row is running can be missed.
Scope, privacy and limits for records¶
Sensitive keys (password, tokens, secrets, card numbers, ssn, iban and the like, matched case- and separator-insensitively at every depth) are stripped from payloads and metadata before storage and before the mapper sees them. File uploads in a payload are rejected; store files through msflib.documents and put the document's source_key in the payload. Inline payloads over RECORD_MAX_PAYLOAD_BYTES raise RecordPayloadTooLargeError. scope for submit() must come from your request dependencies, never from the request body.
Troubleshooting¶
| Symptom | Cause and fix |
|---|---|
ValueError: KNOWLEDGE.DATABASE_URL is set, but CORE.USE_SQLITE is True |
The AGE projection needs real Postgres. Unset the URL or set CORE.USE_SQLITE=False. |
ValueError: ... must be set when CORE.USE_SQLITE is False |
Set KNOWLEDGE.DATABASE_URL or CORE.SQLALCHEMY_DATABASE_URI. |
Warning about the psycopg2 driver |
Use a postgresql+psycopg:// URL so AGE cypher() writes work through the ORM session. |
permission denied or missing ag_catalog |
Run provision_age_and_pgvector(...) once with a superuser DSN. |
| Only a few agent tools appear | Tools are filtered by the engine's capabilities. Without AGE there are no neighbor, path or traversal tools; without a vector store there is no SearchConcepts. |
Agent tools raise "require scope on the agent's runtime context" |
The agent must run with context= carrying a scope. ai-api's agent router does this; a custom agent must too. Tools refuse to run unscoped. |
get_record_ingestion_service() needs ... record_dispatch_task |
Pass record_dispatch_task=build_dispatch_task(settings) to get_knowledge_dependencies. |
UnregisteredSubjectTypeError on enqueue |
The record source is not registered in the web process. |
| Job dead-lettered: "no record source" or "no record mapper registered" | Register sources and mappers in the worker too, then replay the dead letter. |
QueueDisabledError |
KNOWLEDGE.RECORD_QUEUE_ENABLED is false. sync_model_to_queue ignores it so saves still succeed. |
| Two different people merged into one entity | Root entity names collide. Include the record's id in the name. |
| Facts from an old version linger | Give app-owned sources a version=, or re-ingest with supersede_source=True. |
API reference¶
See the generated API reference for msflib.knowledge. modules/knowledge/README.md and modules/knowledge/STRUCTURED_RECORDS.md in the repository hold the detailed notes, including source subclassing (BaseRecordSource, @record_source), mapper and source testing (check_record_source), and the knowledge-record-processed event payload.