Skip to content

msflib.knowledge

Public modules of the msflib-knowledge package.

msflib.knowledge.actions

Database access for the canonical knowledge model.

All canonical-table access goes through these action classes; mutations use the ModelAction super() methods so ModelBase lifecycle events (knowledgeentity-create-pre-commit etc.) fire automatically.

Scoped reads take the typed :class:~msflib.knowledge.contracts.KnowledgeScope contract produced by the planner's scope rewriter — validated once at the start of the lifecycle, accessed with certainty here (no string-keyed mappings, no defensive checks). workspace_id=None renders as IS NULL (the explicit workspace-less bucket).

KnowledgeEntityAction(*, settings: SettingsBase | None = None)

Bases: ModelAction[KnowledgeEntity, KnowledgeEntityCreate, KnowledgeEntityUpdate]

resolve_by_name(session: Session, *, scope: KnowledgeScope, name: str, entity_types: list[str] | None = None, lifecycles: tuple[KnowledgeLifecycle, ...] = DEFAULT_LIFECYCLES, limit: int = 5) -> list[KnowledgeEntity]

Resolve a free-text name: exact (case-insensitive) matches first, then partial matches, deduplicated, up to limit.

find_canonical_match(session: Session, *, scope: KnowledgeScope, entity_type: str, name: str) -> KnowledgeEntity | None

Exact-match cross-fragment entity resolution within a scope.

Uses hierarchical=False: entity resolution must bind to the caller's exact scope tuple, never the hierarchical read semantics scoped_statement uses elsewhere — otherwise a conversation-private fragment could silently merge into (or match against) a workspace-general entity's evidence, exposing it beyond its conversation. See architecture doc §6.1.

Embedding-similarity resolution is a later enhancement behind the injected embedder capability.

find_duplicate_candidates(session: Session, *, scope: KnowledgeScope, entity_type: str, exclude_uids: Iterable[str] = ()) -> list[KnowledgeEntity] | None

Prefetch pool for tag-normalization duplicate-candidate detection.

Every other same-scope (hierarchical=False), same-type entity. Returns None instead of a partial list once the pool exceeds DUPLICATE_CANDIDATE_SCAN_LIMIT, so a heavily-populated type is skipped rather than compared against an arbitrary slice of itself.

promote_scope(session: Session, *, entity: KnowledgeEntity, commit: bool = False) -> KnowledgeEntity

Widen a conversation-private entity to project- (or workspace-) level visibility.

Mirrors DocumentRecordAction.promote_scope: nulls the conversation dims and stamps scope_promoted_at rather than copying the row, so re-ingestion can tell "never scoped" apart from "deliberately promoted" (see KnowledgeScopeMixin). project_id is left untouched -- promotion must never widen past the project boundary.

KnowledgeRelationAction(*, settings: SettingsBase | None = None)

Bases: ModelAction[KnowledgeRelation, KnowledgeRelationCreate, KnowledgeRelationUpdate]

archive_live_by_source(session: Session, *, scope: KnowledgeScope, source_id: str, include_subscopes: bool = False, commit: bool = False) -> list[KnowledgeRelation]

Supersede every currently-live relation asserted by source_id.

"Live" is valid_to IS NULL (the model's own assert-and-supersede contract), not a lifecycle check — lifecycle can move to archived for unrelated moderation reasons independent of temporal validity.

promote_scope(session: Session, *, relation: KnowledgeRelation, commit: bool = False) -> KnowledgeRelation

Widen a conversation-private relation to project- (or workspace-) level visibility.

See KnowledgeEntityAction.promote_scope.

KnowledgeObservationAction(*, settings: SettingsBase | None = None)

Bases: ModelAction[KnowledgeObservation, KnowledgeObservationCreate, KnowledgeObservationUpdate]

archive_live_by_source(session: Session, *, scope: KnowledgeScope, source_id: str, include_subscopes: bool = False, commit: bool = False) -> list[KnowledgeObservation]

Archive every non-archived observation asserted by source_id.

Observations have no valid_to column (evidence snapshots, not temporal facts), so lifecycle is the only liveness axis.

promote_scope(session: Session, *, observation: KnowledgeObservation, commit: bool = False) -> KnowledgeObservation

Widen a conversation-private observation to project- (or workspace-) level visibility.

See KnowledgeEntityAction.promote_scope.

normalize_tag(name: str) -> str

Lowercase form of name with everything but ASCII [a-z0-9] stripped.

Turns "ID-42" and "ID42" into the same key. A collision is only a candidate match, not proof -- see find_duplicate_candidates. Names with no ASCII letters or digits normalize to "" -- callers must never treat two empty keys as a match.

tag_prefix_key(name: str) -> str

Normalized form of name's leading tag-like run of characters.

Catches "bare tag" vs. "TAG (descriptive suffix)" -- "ID-42" and "ID-42 (Long Description)" both reduce to "id42" -- without resorting to substring containment, which would also match "ID-1" inside "ID-10". The run must contain a digit to count as tag-like, so plain names sharing a first word (e.g. "New York" / "New Jersey") don't collide. Returns "" when name has no such run; "" is never a match.

scoped_statement(model: type[ScopedModelType], *, scope: KnowledgeScope, lifecycles: tuple[KnowledgeLifecycle, ...] = (), hierarchical: bool = True, include_subscopes: bool = False) -> SelectOfScalar[ScopedModelType]

Base SELECT for a knowledge table constrained to a compiled scope.

All three knowledge models share KnowledgeScopeMixin, so the column access below is certain. workspace_id == None compiles to IS NULL and is always an equality match — tenant/workspace are not hierarchical.

project_id/conversation_id/sub_thread_id default to hierarchical (IS NULL OR = :value) per architecture doc §6.1: a row scoped to no project is workspace-general and visible to every caller; a row scoped to a project is visible only to callers whose own scope carries that same value (and a conversation nests the same way again, within a project). Pass hierarchical=False for exact-tuple matching instead — required for entity resolution (find_canonical_match), where a project- or conversation-private fragment must never merge into (or be matched against) a workspace-general (or another project's) entity. include_subscopes=True matches tenant/workspace only.

msflib.knowledge.agent

KnowledgeToolRegistry(engine: KnowledgeEngine, *, query_types: Sequence[type] = ALL_QUERY_TYPES, include_drafts: bool = False)

ToolRegistry-protocol provider of knowledge query tools.

query_tool_spec(query_type: type) -> dict[str, Any]

Build a tool spec (name, description, input schema) from a query dataclass.

tool_registry

KnowledgeToolRegistry — exposes semantic queries as agent tools.

Implements the ai_core.services.tool_registry.ToolRegistry protocol (list_tools(context=...)) structurally, without importing ai_core, so it composes through compose_tools() into build_langchain_agent().

Scope is never accepted at construction and never model-supplied: each tool reads the current invocation's scope from the live LangGraph runtime context (langgraph.runtime.get_runtime()), set per-request via agent.invoke(..., context=SomeContext(scope=...)). This is what actually lets a single registry, built once like any other ToolRegistry, serve every request with the correct caller's scope — no per-request rebuilding, no per-scope cache. If the runtime context has no scope, tools refuse to run rather than query unscoped. Tool name, description, and schema derive from each query dataclass.

KnowledgeToolRegistry(engine: KnowledgeEngine, *, query_types: Sequence[type] = ALL_QUERY_TYPES, include_drafts: bool = False)

ToolRegistry-protocol provider of knowledge query tools.

query_tool_spec(query_type: type) -> dict[str, Any]

Build a tool spec (name, description, input schema) from a query dataclass.

msflib.knowledge.compiler

EntityResolutionPass

Deduplicate entities within the fragment by (canonical type, lowercased name).

Cross-fragment resolution against already-persisted entities (exact uid/name, then embedding similarity via an injected embedder) happens at persist time in the engine — this pass only guarantees the fragment itself is internally consistent, remapping relation/observation refs onto the surviving entity.

OntologyMapPass

Rewrite entity/relation type aliases to their canonical names.

ScorePass

Assign lifecycle intent based on confidence and extraction source.

Human-provided and deterministic-parser fragments are trusted; LLM extractions below the confidence threshold land as drafts. The decision is written to each spec's typed lifecycle field for the persist step to consume; explicitly pre-set lifecycles are respected.

ValidationPass

Reject fragments with unknown types or dangling entity references.

KnowledgeCompiler(*, ontology: OntologyRegistry, passes: Sequence[CompilerPass] | None = None, draft_confidence_threshold: float = 0.7)

Runs compiler passes over a fragment before persistence.

Persistence itself (canonical rows + projections in one transaction, then eventbus publication) is owned by the engine, not the compiler.

passes

Compiler passes: validate → ontology-map → entity-resolve → score.

Each pass takes and returns a :class:KnowledgeFragment; the pipeline runs them in order. Persistence is a separate step owned by the engine (it needs a session and the projections).

ValidationPass

Reject fragments with unknown types or dangling entity references.

OntologyMapPass

Rewrite entity/relation type aliases to their canonical names.

EntityResolutionPass

Deduplicate entities within the fragment by (canonical type, lowercased name).

Cross-fragment resolution against already-persisted entities (exact uid/name, then embedding similarity via an injected embedder) happens at persist time in the engine — this pass only guarantees the fragment itself is internally consistent, remapping relation/observation refs onto the surviving entity.

ScorePass

Assign lifecycle intent based on confidence and extraction source.

Human-provided and deterministic-parser fragments are trusted; LLM extractions below the confidence threshold land as drafts. The decision is written to each spec's typed lifecycle field for the persist step to consume; explicitly pre-set lifecycles are respected.

pipeline

KnowledgeCompiler(*, ontology: OntologyRegistry, passes: Sequence[CompilerPass] | None = None, draft_confidence_threshold: float = 0.7)

Runs compiler passes over a fragment before persistence.

Persistence itself (canonical rows + projections in one transaction, then eventbus publication) is owned by the engine, not the compiler.

msflib.knowledge.config

msflib.knowledge.contracts

KnowledgeScope(tenant_id: int, workspace_id: int | None = None, project_id: str | None = None, conversation_id: str | None = None, sub_thread_id: str | None = None) dataclass

Validated scope filter values, typed to the knowledge model columns.

workspace_id=None means the explicit workspace-less bucket (rendered as IS NULL), matching the scope system's explicit-null semantics.

project_id/conversation_id/sub_thread_id are read hierarchically, not by equality: a row with project_id IS NULL is workspace-general and visible to every caller in the workspace; a row with a concrete project_id is visible only to callers whose own scope carries that same value (and, within a project, conversation_id nests the same way again). A scope with project_id=None therefore sees workspace-general knowledge only, never another project's private facts. See architecture doc §6.1.

capabilities

Injected-capability contracts.

The knowledge module has no hard dependency on msflib-ai-core; embedding and LLM extraction are supplied by the caller as callables matching these contracts (same pattern as the ai_api lazy callable factories and the ai_core IngestionPipeline's injected Embedder/Indexer).

scope

Typed compiled-scope contract.

ScopeCompiler validates a ScopeEnvelope against a profile and projects it to string-typed sink formats; per its design, domain-typed projections belong to the consuming subsystem. KnowledgeScope is that projection for this module: built exactly once by the planner's scope rewriter (after profile validation), then passed through planner → projections → actions with certain, typed attribute access — no string-keyed mappings mid-path.

KnowledgeScope(tenant_id: int, workspace_id: int | None = None, project_id: str | None = None, conversation_id: str | None = None, sub_thread_id: str | None = None) dataclass

Validated scope filter values, typed to the knowledge model columns.

workspace_id=None means the explicit workspace-less bucket (rendered as IS NULL), matching the scope system's explicit-null semantics.

project_id/conversation_id/sub_thread_id are read hierarchically, not by equality: a row with project_id IS NULL is workspace-general and visible to every caller in the workspace; a row with a concrete project_id is visible only to callers whose own scope carries that same value (and, within a project, conversation_id nests the same way again). A scope with project_id=None therefore sees workspace-general knowledge only, never another project's private facts. See architecture doc §6.1.

msflib.knowledge.deps

Vanilla KnowledgeEngine construction for host apps.

Every piece of KnowledgeEngine is optional/pluggable (see service/engine.py), so there's no single "correct" wiring, but most apps want the same default shape: canonical storage always, an AGE graph projection when running against real Postgres, and a vector-store projection via ai_core's embedding backend. get_knowledge_dependencies follows the module-deps convention used by auth/account/ai_core (see msflib.ai_core.deps.get_ai_dependencies for the closest sibling: a programmatic, non-request-scoped DependencyNamespace of plain callables, as opposed to ai_core's other factory, get_provider_dependencies, which returns FastAPI Depends-wired ones — KnowledgeEngine isn't itself request-scoped, so there's nothing to wire per-request here either).

get_knowledge_dependencies(settings: SettingsBase, *, session_factory: Callable[[], Session], vector_store_factory: Callable[[], Any] | None = None, ontology: OntologyRegistry | None = None, additional_ontology_pack_namespaces: Sequence[str] | None = None, ontology_pack_options: Mapping[str, Mapping[str, Any]] | None = None, record_dispatch_task: Any | None = None) -> DependencyNamespace

Return KnowledgeEngine building blocks other modules call programmatically.

Scopes CORE, KNOWLEDGE and AI_CORE off settings once, at namespace-construction time (mirroring get_ai_dependencies' one-time ai_settings scoping). Returns:

  • build_knowledge_engine — builds a fresh KnowledgeEngine on every call. Canonical projection always runs (works against sqlite or Postgres alike). The AGE graph projection needs a real Postgres+AGE instance, so it's only wired in when CORE.USE_SQLITE is false. The vector-store projection is gated the same way by default — its static backends (pgvector/qdrant/pinecone) all need an external service — but an explicit vector_store_factory wires it in unconditionally, since that store's contract (similarity_search/add_texts) is backend-agnostic. That's also the forward-compatible seam for routed/ failover vector stores (e.g. VectorStoreRegistryService. build_vector_store_with_failover): pass a zero-arg closure over whatever session/scope that call needs — get_vector_store's own signature, and this factory's, don't need to change either way.
  • get_knowledge_engine — build_knowledge_engine wrapped in a cache: KnowledgeEngine holds no live session (CanonicalProjection opens one per operation from session_factory), so one instance is safely a process-wide singleton, built lazily on first call. Hand this straight to ai_api's agent router (its get_knowledge_engine parameter) or to integrations.ai_api.resolve_knowledge_tool_registries.
  • ontology — the same OntologyRegistry instance wired into every engine build_knowledge_engine returns. Pass it to LLMExtractionParser/run_knowledge_extraction_worker's own ontology= param so extraction is checked against the exact schema the engine compiles and queries against, instead of rebuilding (and risking a mismatched) registry a second time.

Ontology resolution, in priority order: an explicit ontology= (a fully custom OntologyRegistry the caller built however it likes — the escape hatch for excluding msflib's built-in packs entirely) wins outright; otherwise it's built from KNOWLEDGE.ONTOLOGY_PACKS (always includes "rag_default", even if omitted from that setting — see _ensure_rag_default_pack) via default_ontology(). Pass additional_ontology_pack_namespaces=["app.knowledge.ontology.packs"] to let a host app register its own entity/relation-type packs by name; msflib's own pack namespace is always searched too, so callers never need to know or repeat it. No need to hand-roll the pack-module search or re-wire KnowledgeCompiler/projections/worker against the result separately — this and default_ontology's own additional_pack_namespaces are the intended injection seam for app-specific ontologies.

  • get_record_ingestion_service — a cached KnowledgeRecordService (msflib.knowledge.records) for queueing structured records (forms, JSON payloads, app rows). Needs record_dispatch_task: any .delay(job_id) handle, typically msflib.ingestion.build_dispatch_task(settings); honors KNOWLEDGE.RECORD_QUEUE_ENABLED.

ontology_pack_options passes through to default_ontology's own pack_options (per-pack register(registry, **options) keyword arguments -- see ontology/__init__.py), so a caller can configure a pack that takes its own keyword arguments without hand-rolling the pack resolution this function already does.

msflib.knowledge.errors

Dependency-free errors shared by the knowledge dispatch paths.

KnowledgeQueueDisabledError

Bases: QueueDisabledError

Raised when a knowledge ingestion queue is disabled by settings.

msflib.knowledge.eventbus

msflib.knowledge.integrations

Optional integrations wiring msflib.knowledge to other modules.

Each submodule pulls in one optional dependency (ai_core, documents, ai_api) and is imported directly by consumers that have that dependency installed — never through msflib.knowledge's lazy __getattr__, which must stay free of optional imports.

ai_api

Wires KnowledgeToolRegistry into ai_api's agent router.

Matches ai_api's existing get_tool_registries: Callable[[], Sequence[ToolRegistry]] factory shape exactly — no new parameter type needed on the ai_api side. KnowledgeToolRegistry itself needs no per-request scope (see its module docstring): each tool reads the caller's scope from the live LangGraph runtime context, so the registry it returns here is built once (cached by ai_api's own agent router), like any other ToolRegistry.

Takes a zero-arg get_engine getter rather than a bare KnowledgeEngine instance so the engine itself can be resolved lazily, on first tool-registry build rather than at router-mount time -- letting a caller swap the engine (e.g. a test injecting one wired with fakes) any time before that first resolution.

ai_core

Default ai_core-backed collaborators for KnowledgeEngine.

Mirrors msflib.ai_api.services.rag.resolve_default_vector_store's lazy-factory template. Knowledge owns a fixed, dedicated vector collection (unlike ai_api's configurable RETRIEVER_COLLECTION_NAME), so the collection name is hardcoded here rather than read from settings.

resolve_default_extractor(ai_settings: AICoreSettings) -> Extractor

Build the default LLM-backed Extractor for LLMExtractionParser.

Sends the prompt to the configured chat model and parses its raw content via ai_core.services.llm_json.parse_llm_json, which tolerates the syntax slips real chat models produce (curly quotes, unquoted keys, mismatched brackets) rather than a strict json.loads. No structured-output binding or schema enforcement — no such mechanism exists anywhere in ai_core's providers yet.

Static, single-config resolution — prefer resolve_registry_extractor wherever a session/scope are available (the worker path); this one remains for callers with no request/job context to route through.

resolve_registry_extractor(settings: SettingsBase, session: Session, *, scope: ScopeEnvelope, task: str = 'knowledge.extraction') -> Extractor

Registry-routed counterpart to :func:resolve_default_extractor.

Goes through ai_core's provider registry (get_ai_dependencies) instead of a single static LLM_* config, so a tenant/workspace can bind task to a specific AIProviderProfile via AICoreSettings.TASK_PROFILES and gets automatic failover across candidates (see ProviderRegistryService.build_llm_with_failover). When no profile is registered for the scope, resolution falls back to the same static config resolve_default_extractor uses, so this is a strict superset — safe to use wherever a DB session is available.

scope is taken as-is (built once by the caller, e.g. the worker's _build_scope) and passed straight through to the registry rather than decomposed and rebuilt here — scope is constructed in exactly one place per call chain so it stays easy to track/instrument end to end.

dispatch

Enqueue knowledge extraction jobs onto the shared msflib.ingestion queue.

Always defers Celery dispatch to the caller's own commit (see IngestionDispatchService.enqueue_and_dispatch's dispatch_after_commit), so a broker hiccup at enqueue time is logged and swallowed rather than raised out of an eventbus listener -- enqueue_for_record is only ever called from eventbus listeners (see integrations/documents.py).

KnowledgeQueueDisabledError

Bases: QueueDisabledError

Raised when a knowledge ingestion queue is disabled by settings.

KnowledgeIngestionDispatchService(*, dispatch: IngestionDispatchService, queue_enabled: bool = True)

enqueue_for_record(session: Session, *, document_record_id: int, target_version: int, reason: KnowledgeExtractionReason, tenant_id: int, workspace_id: int | None, commit: bool = False) -> IngestionJob

tenant_id/workspace_id must be the source DocumentRecord's own scope (see the event's producer, documents.services.pipeline), not the ambient default tenant -- a record belonging to a non-default tenant would otherwise get its extraction job (and any dead-letter entry) stamped under the wrong tenant.

documents

Wires msflib.knowledge extraction to msflib.documents' lifecycle events.

Subscribes to documents' own domain events (never the low-level drivelink hooks) so extraction only ever runs on content documents has already successfully fetched and indexed, and reason classification isn't re-derived:

  • documents.ingestion.completed (fired by the text_document ingestion stage, msflib.documents.services.pipeline) -> enqueue an IngestionJob(subject_type="knowledge_extraction") with the matching reason.
  • documents.index.deleted (fired once a DocumentRecord marked deleted commits, with a column snapshot of it) -> enqueue a delete-reason job (pure archival, no re-parsing).

reset_event_hooks(*, emitter: AppEmitter) -> None

Remove knowledge-extraction hooks and reset registration state.

register_event_hooks(*, dispatch_service: KnowledgeIngestionDispatchService, emitter: AppEmitter, force: bool = False) -> None

Register hooks that enqueue knowledge extraction jobs from documents events.

extraction_common

LLM-extraction helpers for the knowledge ingestion stages.

Both stages use default_llm_parser. Chunked extraction is applied by knowledge_extraction only; record mappers call their parser directly. Imports nothing from msflib.documents.

default_llm_parser(settings: SettingsBase, session: Session, scope: ScopeEnvelope, *, ontology: OntologyRegistry | None = None, caller: str = 'knowledge extraction') -> LLMExtractionParser

Build an LLMExtractionParser over the ai_core registry-resolved extractor.

caller only names the stage factory in the fallback warning, so the log points at whichever factory was built without ontology=.

extract_fragment_chunked(parser: LLMExtractionParser, text: str, *, provenance: Provenance, settings: SettingsBase) -> KnowledgeFragment

Parse text in context-window-sized chunks (when the chunker is available) and merge.

Whole-text extraction risks truncation on long input, and fragment_from_payload silently drops anything a truncated response doesn't return. Each chunk's fragment gets its local_ids namespaced by chunk index before merging, since every chunk is extracted independently and would otherwise reuse ids like e1 across chunks. Falls back to a single whole-text call when ai_core's chunker isn't installed.

Chunk size is derived from the extraction model's context window (AI_CORE.CHUNKING_POLICY) rather than a fixed character count, so a huge input doesn't produce a huge number of sequential LLM calls.

check_chunk_count(chunk_count: int, *, policy, source_id: str) -> None

Warn or abort before any LLM calls are made, based on chunk count.

Aborting raises ValueError so is_terminal_ingestion_error marks the job failed without retrying (unchanged input would just reproduce the same chunk count).

merge_fragments(fragments: list[KnowledgeFragment], *, provenance: Provenance, prefix: str = 'c') -> KnowledgeFragment

Concatenate independently-extracted fragments, namespacing their local ids.

Each fragment's local_ids become f"{prefix}{index}:{local_id}" and its relation/observation refs are remapped to match; uid:-prefixed external refs pass through untouched.

extractors

Pluggable, content-sniffed extractors for the knowledge_extraction stage.

A registered ContentExtractor claims raw document bytes and, when it claims one, replaces the LLM extraction pass entirely for that document. Selection is first-match-wins, evaluated before any text extraction or LLM call. See modules/knowledge/README.md ("Pluggable content extractors") for the rationale and usage.

ContentExtractor

Bases: Protocol

A specialized, content-sniffed alternative to LLM-based extraction.

name becomes Provenance.extractor for facts this extractor produces -- pick something stable and unique.

matches(*, content_bytes: bytes, mime_type: str | None, filename: str | None) -> bool

Return True iff this extractor can handle content_bytes.

Sniff the bytes themselves rather than trusting mime_type/ filename alone -- both are caller-supplied and can be wrong.

extract(content_bytes: bytes, *, document: SourceDocument, provenance: Provenance, settings: SettingsBase) -> KnowledgeFragment

Parse content_bytes into a fragment ready for KnowledgeEngine.ingest.

ContentExtractorRegistry(_entries: list[_RegisteredExtractor] = list(), _sequence: int = 0) dataclass

Ordered collection of ContentExtractors, resolved first-match-wins.

Higher priority wins; ties break by registration order (earliest first). A fresh instance (rather than always reaching for default_extractor_registry) is useful for test isolation -- make_knowledge_extraction_stage(extractor_registry=...) accepts one.

register_extractor(extractor: ContentExtractor, *, priority: int = 0, registry: ContentExtractorRegistry = default_extractor_registry) -> None

Register extractor onto registry (the process-wide default unless overridden).

models

Reason codes for documents-triggered knowledge extraction.

Extraction jobs are IngestionJob(subject_type="knowledge_extraction") rows (see integrations/pipeline.py/integrations/dispatch.py). KnowledgeExtractionReason is a knowledge-domain concept distinct from documents' own IngestionReason, which has no delete case.

pipeline

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

Wraps knowledge's business logic -- adapter-based source resolution, ContentExtractor/LLMExtractionParser/KnowledgeEngine.ingest/ archive_source -- behind the msflib.ingestion stage protocol. See modules/knowledge/README.md for the extraction pipeline and pluggable extractor design.

The LLM chunking/merging and stage error-handling helpers live in extraction_common (shared with the knowledge_record stage). A single stage does Fetch+Extract+Parse+Ingest in one call, since the resolved SourceDocument/adapter/parser state flows linearly through one job's processing and none of it is reusable by another pipeline_kind.

The stage needs a constructed KnowledgeEngine (no dependency-free default exists), so there is no top-level singleton stage instance -- the worker entrypoint must call make_knowledge_extraction_stage(engine=...) itself once it has built the engine.

make_knowledge_extraction_stage(*, engine: KnowledgeEngine, record_action: DocumentRecordAction | None = None, parser_factory: Callable[[], LLMExtractionParser] | None = None, ontology: OntologyRegistry | None = None, extractor_registry: ContentExtractorRegistry = default_extractor_registry) -> Stage

Build the knowledge_extraction stage.

ontology is only used by the default parser (ignored when parser_factory is given): pass the same OntologyRegistry the KnowledgeEngine was built with whenever it includes packs beyond default_ontology()'s rag_default -- otherwise extraction silently stays limited to that fallback vocabulary regardless of what the engine can store.

extractor_registry (see .extractors) is consulted before the LLM path; nothing registered (the default) leaves behavior unchanged. Tests should pass an isolated ContentExtractorRegistry() rather than the process-wide default_extractor_registry. See modules/knowledge/README.md for details.

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

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

msflib.knowledge.ir

Knowledge Intermediate Representation (IR).

Parsers never write to storage. Every source (LLM extraction, a structured-format parser, AAS adapter, agent memory) emits a :class:KnowledgeFragment; the compiler validates, ontology-maps, entity-resolves, and persists it.

Entity references inside a fragment use local_id (unique within the fragment). Relations and observations may reference either a local_id from the same fragment or an external persisted uid prefixed as uid:<value>.

Provenance(source_id: str, source_type: str, extracted_by: str = 'parser', extractor: str | None = None, metadata: dict[str, Any] = dict()) dataclass

Where a fragment came from and how it was produced.

ObservationSpec(entity_refs: list[str], excerpt: str | None = None, properties: dict[str, Any] = dict(), confidence: float = 1.0) dataclass

Evidence linking one or more entities to a source excerpt.

Also the reification vehicle for n-ary relations: a process step that involves several assets is one observation over all of them.

KnowledgeFragment(provenance: Provenance, entities: list[EntitySpec] = list(), relations: list[RelationSpec] = list(), observations: list[ObservationSpec] = list()) dataclass

unresolved_refs() -> set[str]

Return refs that are neither fragment local_ids nor explicit uids.

A ref is considered external (already persisted) when it matches the uid of an entity in this fragment or is prefixed with uid: followed by a non-empty value. A uid: prefix with an empty or whitespace-only suffix is malformed and remains unresolved so ValidationPass catches it instead of persistence failing later. Everything else must resolve to a fragment local_id.

msflib.knowledge.models

KnowledgeEntityCreate

Bases: SchemaBase

Internal persist DTO stamped by engine/actions with trusted scope.

These fields are not intended to be accepted directly from external callers.

KnowledgeObservation

Bases: ModelBase, KnowledgeScopeMixin, KnowledgeObservationFieldsMixin

Evidence linking entities to a source; reification vehicle for n-ary facts.

KnowledgeObservationCreate

Bases: SchemaBase

Internal persist DTO stamped by engine/actions with trusted scope.

These fields are not intended to be accepted directly from external callers.

KnowledgeRelation

Bases: ModelBase, KnowledgeScopeMixin, KnowledgeRelationFieldsMixin

Edges scope independently of their endpoints — a cross-scope edge is visible only when the caller can see the edge's own scope.

Temporal validity is assert-and-supersede: never overwrite a superseded fact, set its valid_to and insert the new assertion.

KnowledgeRelationCreate

Bases: SchemaBase

Internal persist DTO stamped by engine/actions with trusted scope.

These fields are not intended to be accepted directly from external callers.

KnowledgeRecordCreate

Bases: SchemaBase

Caller-supplied fields only; scope, submitter, version and hash go through update=.

knowledge

Canonical knowledge model — the source of truth.

Graph (Apache AGE) and vector (pgvector) representations are projections derived from these tables; agents never write to a projection directly.

Every row carries three orthogonal axes: - semantic: type + properties - visibility: scope dimensions (tenant_id, workspace_id) - lifecycle: draft / approved / archived

KnowledgeScopeMixin

Bases: SchemaBase

Visibility + lifecycle + provenance columns shared by all knowledge rows.

project_id/conversation_id/sub_thread_id are read hierarchically (IS NULL OR = :value), not by equality — see KnowledgeScope and architecture doc §6.1. scope_promoted_at mirrors msflib-documents' DocumentRecord pattern: set once a conversation-private row has been explicitly widened to workspace visibility, so a later routine re-ingest can tell "never scoped" apart from "deliberately promoted" and avoid silently re-narrowing it.

KnowledgeEntityCreate

Bases: SchemaBase

Internal persist DTO stamped by engine/actions with trusted scope.

These fields are not intended to be accepted directly from external callers.

KnowledgeRelation

Bases: ModelBase, KnowledgeScopeMixin, KnowledgeRelationFieldsMixin

Edges scope independently of their endpoints — a cross-scope edge is visible only when the caller can see the edge's own scope.

Temporal validity is assert-and-supersede: never overwrite a superseded fact, set its valid_to and insert the new assertion.

KnowledgeRelationCreate

Bases: SchemaBase

Internal persist DTO stamped by engine/actions with trusted scope.

These fields are not intended to be accepted directly from external callers.

KnowledgeObservation

Bases: ModelBase, KnowledgeScopeMixin, KnowledgeObservationFieldsMixin

Evidence linking entities to a source; reification vehicle for n-ary facts.

KnowledgeObservationCreate

Bases: SchemaBase

Internal persist DTO stamped by engine/actions with trusted scope.

These fields are not intended to be accepted directly from external callers.

records

Inline structured-record storage for KnowledgeRecordService.submit.

KnowledgeRecordCreate

Bases: SchemaBase

Caller-supplied fields only; scope, submitter, version and hash go through update=.

msflib.knowledge.ontology

OntologyRegistry()

resolve_entity_type(name: str) -> str | None

Resolve name (canonical or alias, case-insensitive alias) to canonical.

is_a(child: str, ancestor: str) -> bool

Return True if child is ancestor or transitively inherits from it.

descendants(name: str) -> set[str]

Return name plus all entity types that transitively inherit from it.

default_ontology(packs: list[str] | None = None, *, additional_pack_namespaces: Sequence[str] | None = None, pack_options: Mapping[str, Mapping[str, Any]] | None = None) -> OntologyRegistry

Build an OntologyRegistry with the requested packs registered.

Pack names map to modules exposing register(registry), searched in order under each of additional_pack_namespaces first (so a host app can shadow a built-in pack name with its own version), then always under this package's own .packs — callers never need to know or repeat msflib's internal namespace just to still get its built-in packs. This is the seam get_knowledge_dependencies's own additional_ontology_pack_namespaces threads through to here, so apps don't need to hand-roll this search themselves.

pack_options is a generic per-pack configuration seam: each pack's register(registry) is called as register(registry, **options) where options = pack_options.get(pack_name, {}). A pack that takes no extra keyword arguments (the common case, e.g. rag_default) is unaffected when no options are given for it — **{} expands to nothing. Packs that do accept keyword arguments document them on their own register().

To exclude msflib's built-ins entirely rather than merely shadow them, don't use this function: build an OntologyRegistry directly and pass it to get_knowledge_dependencies's ontology= param, which bypasses pack resolution altogether.

packs

rag_default

Default document/RAG ontology pack.

Deliberately small: generic entities and concepts extracted from documents, plus the document/chunk structure they were extracted from.

registry

Ontology registry: entity/relation type definitions with IS_A inheritance.

Separates schema (what a Pump is) from knowledge (this pump feeds that tank) and from policy (who can see it — that belongs to msflib.scope).

Single-purpose reasoning: is_a / descendants let queries like FindEntity(entity_type="Equipment") match pumps and valves. OWL-style description logic is deliberately out of scope; an RDF export would be a future projection of the canonical model.

OntologyRegistry()

resolve_entity_type(name: str) -> str | None

Resolve name (canonical or alias, case-insensitive alias) to canonical.

is_a(child: str, ancestor: str) -> bool

Return True if child is ancestor or transitively inherits from it.

descendants(name: str) -> set[str]

Return name plus all entity types that transitively inherit from it.

msflib.knowledge.parsers

build_extraction_prompt(text: str, *, ontology: OntologyRegistry, examples: Sequence[str] | None = None) -> str

Build the extraction prompt.

examples takes pre-formatted input->JSON strings, inserted verbatim before the source text — a minimal, unopinionated hook for callers who need few-shot examples on noisy source text. This is a stopgap, not the general prompt-templating story: a more holistic prompt/template system for msflib is a planned follow-up, not scoped here.

Known, not a bug, just a design limit to be aware of: downstream entity resolution (KnowledgeEntityAction.find_canonical_match) only matches by exact case-insensitive name, so near-duplicate names this prompt lets through unnormalized (e.g. "Widget Corp" vs "Widget Corp.") will still create separate canonical entities rather than merging. The instruction below reduces the common case (full name vs. abbreviation for the same entity) but not true synonyms the text never spells out consistently -- that would need fuzzy/embedding-based matching in find_canonical_match itself.

fragment_from_payload(payload: Any, *, provenance: Provenance) -> KnowledgeFragment

Assemble a fragment from an extractor payload, skipping malformed items.

llm_extraction

LLM extraction parser: text → KnowledgeFragment via an injected extractor.

The extractor is a callable (prompt: str) -> Mapping that the caller binds to their LLM with structured output (e.g. an ai_core get_llm() model wrapped with a JSON schema). This module owns prompt construction and defensive fragment assembly; it has no LLM dependency of its own.

Expected extractor payload shape::

{
  "entities": [{"local_id": "e1", "type": "Concept", "name": "...",
                "properties": {...}, "confidence": 0.9}],
  "relations": [{"source": "e1", "type": "RELATED_TO", "target": "e2",
                 "confidence": 0.8}],
  "observations": [{"entity_refs": ["e1", "e2"], "excerpt": "...",
                     "confidence": 0.9}]
}

build_extraction_prompt(text: str, *, ontology: OntologyRegistry, examples: Sequence[str] | None = None) -> str

Build the extraction prompt.

examples takes pre-formatted input->JSON strings, inserted verbatim before the source text — a minimal, unopinionated hook for callers who need few-shot examples on noisy source text. This is a stopgap, not the general prompt-templating story: a more holistic prompt/template system for msflib is a planned follow-up, not scoped here.

Known, not a bug, just a design limit to be aware of: downstream entity resolution (KnowledgeEntityAction.find_canonical_match) only matches by exact case-insensitive name, so near-duplicate names this prompt lets through unnormalized (e.g. "Widget Corp" vs "Widget Corp.") will still create separate canonical entities rather than merging. The instruction below reduces the common case (full name vs. abbreviation for the same entity) but not true synonyms the text never spells out consistently -- that would need fuzzy/embedding-based matching in find_canonical_match itself.

fragment_from_payload(payload: Any, *, provenance: Provenance) -> KnowledgeFragment

Assemble a fragment from an extractor payload, skipping malformed items.

msflib.knowledge.planner

compile_query_scope(scope: ScopeEnvelope) -> KnowledgeScope

Compile the typed scope for read plans (mandatory planner pass).

compile_write_scope(scope: ScopeEnvelope) -> KnowledgeScope

Compile the typed scope stamps for fragment persistence.

knowledge_scope_dimensions() -> set[str]

Return scope dimensions reserved for knowledge query filtering.

capabilities

Capability constants advertised by projections and required by queries.

The physical planner matches a logical plan's required capability against projection capability sets with a static preference order — no cost-based optimization until a workload proves it necessary.

logical

Logical planner: semantic query → backend-independent plan.

Rule-based and deterministic. An LLM may compose semantic queries upstream (natural language → query objects); it never participates in planning, scope injection, or backend selection.

physical

Physical planner: dispatch a logical plan to a capable projection.

Capability matching with a static preference order (registration order). Deliberately not cost-based — see the architecture doc's deferred list.

scope_rewriter

Mandatory scope pass: compile a ScopeEnvelope into a typed KnowledgeScope.

Mirrors ai_core.services.scope_metadata — one canonical projection path from scope dimensions to filter values — but the output is the typed :class:~msflib.knowledge.contracts.KnowledgeScope contract rather than a string-keyed mapping, so downstream layers (planner, actions, projections) access fields with certainty. Profile validation happens here, once, at the start of the query/write lifecycle; no re-checking mid-path.

Projections render the scope in their own dialect (SQL WHERE, Cypher property filter, pgvector metadata filter); none implements authorization itself.

knowledge_scope_dimensions() -> set[str]

Return scope dimensions reserved for knowledge query filtering.

compile_query_scope(scope: ScopeEnvelope) -> KnowledgeScope

Compile the typed scope for read plans (mandatory planner pass).

compile_write_scope(scope: ScopeEnvelope) -> KnowledgeScope

Compile the typed scope stamps for fragment persistence.

msflib.knowledge.projections

AgeGraphProjection(connection_factory: Callable[[], Any], *, graph_name: str)

AgeGraphWriter(*, graph_name: str)

graph_exists(session: Session) -> bool

Whether the AGE graph has already been created (see create_graph).

ensure_graph(session: Session) -> None

Idempotently create the AGE graph if it doesn't exist yet.

graph_summary(session: Session) -> dict[str, Any]

Diagnostics snapshot: graph existence plus vertex/edge counts.

Reuses this class's own session-init and parameter-binding fixes (_init_session, CAST(:params AS agtype) in _run_cypher) so callers — health-check/diagnostics endpoints in particular — don't need to hand-roll their own cypher() call and risk reintroducing bugs already solved here (see _run_cypher's docstring on the bind-parameter cast, and common.py's on session init).

CanonicalProjection(session_factory: SessionFactory, *, entity_action: KnowledgeEntityAction | None = None, relation_action: KnowledgeRelationAction | None = None, observation_action: KnowledgeObservationAction | None = None, ontology: OntologyRegistry | None = None)

age

AgeGraphProjection(connection_factory: Callable[[], Any], *, graph_name: str)

AgeGraphWriter(*, graph_name: str)

graph_exists(session: Session) -> bool

Whether the AGE graph has already been created (see create_graph).

ensure_graph(session: Session) -> None

Idempotently create the AGE graph if it doesn't exist yet.

graph_summary(session: Session) -> dict[str, Any]

Diagnostics snapshot: graph existence plus vertex/edge counts.

Reuses this class's own session-init and parameter-binding fixes (_init_session, CAST(:params AS agtype) in _run_cypher) so callers — health-check/diagnostics endpoints in particular — don't need to hand-roll their own cypher() call and risk reintroducing bugs already solved here (see _run_cypher's docstring on the bind-parameter cast, and common.py's on session init).

common

Shared Apache AGE helpers used by both the read projection and the writer.

AGE's cypher() SQL function requires its graph-name and query-text arguments as literal SQL text (its parser hook intercepts them at plan time), so those two pieces are string-interpolated here rather than bound as parameters. Everything else — property values, scope/lifecycle filters — is passed through the function's third agtype argument and referenced in Cypher as $name, exactly like Neo4j query parameters. Interpolated identifiers (graph name, vertex/edge labels, relation types) are validated against a strict allowlist pattern first to keep this injection-safe.

validate_identifier(value: str, *, kind: str) -> str

Guard against Cypher injection for values interpolated as literals (graph name, relation type used as an edge label) rather than bound as agtype parameters.

parse_agtype(raw: str) -> Any

Parse one agtype textual result — a JSON body with ::type annotations after every composite value's closing bracket, including nested vertices/edges inside a list (nodes(p)/relationships(p) results) — into plain Python data.

agtype_literal(value: dict[str, Any] | list[Any]) -> str

Serialize a parameter map/list for the cypher() third argument.

parse_agtype_scalar_int(raw_value: Any) -> int

Parse a scalar agtype integer result (e.g. from RETURN count(n)).

Unlike parse_agtype (composite vertex/edge JSON with ::type tags after a closing bracket), a bare scalar comes back as literal digits directly followed by ::agtype (e.g. "12::agtype") — no bracket precedes the tag, so _strip_agtype_tags won't touch it — hence the separate, simpler parse here.

matches_scope(properties: dict[str, Any], scope: KnowledgeScope, lifecycles: tuple[str, ...]) -> bool

Scope + lifecycle check applied client-side to a parsed vertex/edge's properties map.

AGE 1.5.0's Cypher grammar has no list-predicate functions (all, any, filter, list comprehensions — confirmed absent from its grammar source), so per-hop scope/lifecycle filtering over a variable-length path can't be pushed down as a WHERE clause the way the canonical projection does it. Structural matching (labels, uid, relation types, depth) still happens in Cypher; this is the client-side fallback the physical planner's docstring already anticipates for projections that can't push a predicate down.

project_id/conversation_id/sub_thread_id are checked hierarchically, not by equality, mirroring actions.scoped_statement's default (hierarchical=True): a vertex/edge with no project/conversation stamp is workspace-general and always visible; one stamped with a project or conversation is visible only to a caller scoped to that same value. See architecture doc §6.1.

relation_type_filter(relation_types: tuple[str, ...] | list[str] | str) -> str

[r:TYPE*1..3]-style label filter for a single type, or empty for any type / more than one type.

AGE 1.5.0's Cypher grammar rejects relationship-type alternation ([r:TYPE1|TYPE2]) outright -- confirmed against a live instance, ERROR: syntax error at or near "|" regardless of whether each alternative repeats the leading colon. A single type still compiles fine as a plain label filter, so that case is still pushed down; two or more types can't be expressed in the pattern at all here and must be filtered client-side by the caller instead (see matches_relation_type and its use in projection.py), the same way AGE 1.5.0's missing list-predicate functions already force scope/lifecycle/temporal filtering to happen client-side rather than in a pushed-down WHERE.

matches_relation_type(label: str, relation_types: tuple[str, ...] | list[str] | str) -> bool

Client-side counterpart to relation_type_filter for the 2+-type case it can't push down into the Cypher pattern -- see its docstring.

projection

Apache AGE projection — graph traversal in the same Postgres instance.

AGE stores the labeled property graph alongside the canonical tables. This projection is the read side (FindNeighbors / FindPath / Traverse); writes happen through :class:~msflib.knowledge.projections.age.writer.AgeGraphWriter, called directly from KnowledgeEngine.ingest so vertex/edge upserts share the canonical write's transaction (see the writer module's docstring).

Scope/lifecycle filtering over variable-length paths is applied client-side (see :func:~.common.matches_scope) rather than pushed down as a Cypher WHERE clause: AGE's grammar has no list-predicate functions (all/any/filter/list comprehensions), which is what a per-hop all(n IN nodes(p) WHERE ...) check would need — confirmed absent from AGE's grammar source and live on both 1.5.0 and 1.6.0 (the latest release as of this writing). Structural matching (labels, uid, depth) still happens in Cypher; relation-type filtering does too when there's exactly one type ([r:TYPE*1..N] compiles fine), but two or more types can't be pushed down either -- AGE also rejects relationship-type alternation ([r:TYPE1|TYPE2]) as a bare syntax error, confirmed live, so that case is filtered client-side the same way scope/lifecycle are (see :func:~.common.matches_relation_type).

AGE also has no shortestPath() (absent from its grammar source — confirmed both by inspection and live, on 1.5.0 and 1.6.0, that it's a bare syntax error). Nor is a fixed-length VLE match ([r*N..N]) a good substitute: AGE's variable-length matcher does per-path DFS enumeration with no node memoization, so even with LIMIT 1 it can blow up exponentially in branching factor before Postgres ever gets to apply the limit. FindPath instead does the search itself: a hand-rolled BFS out from source_uid, fetching one frontier level per Cypher round trip (MATCH (a)-[r]-(b) WHERE a.uid IN $frontier RETURN a.uid, r, b — AGE's grammar supports the plain IN operator, confirmed via cypher_gram.y, distinct from the unsupported list-comprehension predicates), rather than one round trip per node or one VLE call per depth. Each node is visited at most once, so this is O(explored nodes/edges) instead of exponential, and it still yields a true shortest path since BFS discovers the target at the minimum depth by construction. Should a projection ever need more than unweighted shortest-path (weighted paths, ranking, centrality), the natural next step is to swap the in-memory search here for networkx fed by the same per-level fetch — not to push more of the algorithm into Cypher.

Requires the age extension installed on the database and the age poetry extra (psycopg).

AgeGraphProjection(connection_factory: Callable[[], Any], *, graph_name: str)

writer

Apache AGE writer — keeps the graph projection in sync with canonical writes.

Called directly from KnowledgeEngine.ingest (not through the KnowledgeProjection read protocol — mirrors how canonical writes go through the action classes rather than CanonicalProjection). Upserts run through the same SQLAlchemy Session used for the canonical rows, via session.execute(text(...)), so both commit or roll back together — the "single Postgres, single transaction" principle from the architecture doc. This is why the writer takes a Session rather than the raw connection_factory the read-side projection uses (that projection runs outside any ingest transaction, from the physical planner).

Caller contract: flush the session before the first call here in a given transaction (KnowledgeEngine.ingest does this). Confirmed live against AGE 1.5.0: a cypher() MERGE+SET call issued while the ORM session still has unflushed pending state (e.g. entities resolved via find_canonical_match rather than freshly created, so nothing flushed them yet) can silently lose most of its SET-clause parameter bindings on that call. Flushing first — once, before any AGE call — avoids it reliably; per-call execution_options={"autoflush": False} on these statements does not (tried and confirmed insufficient on its own). Confirmed fixed on AGE 1.6.0 (re-ran the exact same previously-failing scenario after upgrading Postgres 14.0 → 14.23 and AGE 1.5.0 → 1.6.0 — AGE 1.6.0 rewrote MERGE property handling in cypher_merge.c) — the flush stays anyway as cheap, harmless insurance for anyone still on 1.5.0.

AgeGraphWriter(*, graph_name: str)
graph_exists(session: Session) -> bool

Whether the AGE graph has already been created (see create_graph).

ensure_graph(session: Session) -> None

Idempotently create the AGE graph if it doesn't exist yet.

graph_summary(session: Session) -> dict[str, Any]

Diagnostics snapshot: graph existence plus vertex/edge counts.

Reuses this class's own session-init and parameter-binding fixes (_init_session, CAST(:params AS agtype) in _run_cypher) so callers — health-check/diagnostics endpoints in particular — don't need to hand-roll their own cypher() call and risk reintroducing bugs already solved here (see _run_cypher's docstring on the bind-parameter cast, and common.py's on session init).

base

Projection protocol: derived representations of the canonical model.

A projection executes scoped logical plans for the capabilities it advertises. Projections never implement authorization — they receive already scope-constrained plans and either push the predicates down (SCOPE_FILTER_PUSHDOWN) or rely on executor post-filtering.

canonical

CanonicalProjection(session_factory: SessionFactory, *, entity_action: KnowledgeEntityAction | None = None, relation_action: KnowledgeRelationAction | None = None, observation_action: KnowledgeObservationAction | None = None, ontology: OntologyRegistry | None = None)

projection

Canonical (SQLModel/Postgres) projection — entity lookup, evidence, aliases.

The source of truth is also the cheapest executor for exact lookups. All database access delegates to the knowledge action classes; this projection only translates logical plans into action calls and serializes results.

CanonicalProjection(session_factory: SessionFactory, *, entity_action: KnowledgeEntityAction | None = None, relation_action: KnowledgeRelationAction | None = None, observation_action: KnowledgeObservationAction | None = None, ontology: OntologyRegistry | None = None)

vector_store

common

Shared helpers used by both the vector-store read projection and its writer.

Mirrors age/common.py's role for the AGE projection/writer pair: the entity fields carried as vector-store metadata are the same read shape used elsewhere (:data:_ENTITY_FIELDS matches age/projection.py's _ENTITY_READ_FIELDS and KnowledgeEntityRead), and the embedding text builder is the single place that decides what free text represents an entity for semantic search.

entity_embedding_text(entity: KnowledgeEntity) -> str

Free-text representation of an entity for embedding.

Entities have no dedicated description field, so this composes the typed/named identity plus any string-valued properties -- enough for semantic search without speculatively growing the canonical schema.

entity_metadata(entity: KnowledgeEntity) -> dict[str, Any]

Vector-store metadata for one entity, keyed to match the read shape.

properties is JSON-encoded rather than passed through as a nested dict: several backends this projection must stay agnostic to (Pinecone, Chroma) reject nested objects in metadata outright or silently stringify/drop them. entity_item_from_metadata reverses the encoding on read.

project_id/conversation_id/sub_thread_id are stored for the read-side client-side hierarchical filter (see vector_store/projection.py); they are not part of _ENTITY_FIELDS/KnowledgeEntityRead, so entity_item_from_metadata never surfaces them in a query result.

entity_item_from_metadata(metadata: dict[str, Any]) -> dict[str, Any]

Project stored metadata down to the same read shape other projections return (KnowledgeEntityRead-compatible field set).

projection

Vector-store projection — semantic similarity over entity/concept embeddings.

Backend-agnostic: written against the base langchain_core.vectorstores. VectorStore contract (similarity_search/add_texts, guaranteed on every backend -- pgvector, Qdrant, Pinecone, ...), not any one backend's extensions. Earlier drafts of this projection called similarity_search_by_vector/add_embeddings directly with a separately injected embedder to avoid depending on the store's own embeddings object; those methods turned out to be pgvector-specific rather than part of the universal interface, which would have broken interchangeability with other backends ai_core's factory already supports. Calling through the store's own text-based methods fixes that and, as a side benefit, keeps read and write on the exact same embedding model by construction (it's baked in once at get_vector_store(settings, collection, embeddings)), rather than risking drift against a separately injected embedder. This mirrors the precedent already in ai_core.vector_store.factory.similarity_search_scoped. No hard ai_core dependency results either way: vector_store_factory is injected, so this module never imports ai_core directly.

Scope (tenant/workspace) is pushed down as a metadata filter -- this is our own knowledge_entities collection, so we control the metadata schema and can guarantee every backend understands those two keys (see vector_store/common.py). Lifecycle and entity-type restriction are not guaranteed pushdown-able the same way (equality-only filters can't express "lifecycle in (...)" portably across backends), so those are applied client-side against an over-fetched result set -- the same compensating pattern age/projection.py uses for predicates Cypher can't push down.

project_id/conversation_id/sub_thread_id are hierarchical (IS NULL OR = :value, architecture doc §6.1), which an equality-only metadata filter can't express either, so they join lifecycle/type as a client-side check against the same over-fetched set rather than the pushed-down filter.

writer

Vector-store writer — keeps entity embeddings in sync with canonical writes.

Mirrors AgeGraphWriter's role for the graph projection: called directly from KnowledgeEngine.ingest (not through the read-side VectorStoreProjection), one upsert_entity per persisted entity.

Unlike AgeGraphWriter, this does not share the canonical SQLAlchemy Session/transaction -- the vector collection lives behind its own injected client (vector_store_factory), consistent with the "two collections, one backend" design (architecture doc §5.3) and the module's no-hard-ai_core-dependency contract. A failed vector upsert therefore does not roll back the canonical write, the same way chunk-embedding writes in ai_core's IngestionPipeline sit outside the canonical transaction.

Writes go through add_texts -- guaranteed by the base langchain_core.vectorstores.VectorStore contract on every backend, unlike add_embeddings (a pgvector-specific extension that isn't reliably present on other LangChain vector store wrappers). add_texts embeds through the store's own configured embeddings function, which keeps write and read (VectorStoreProjection, also store-driven) on the same embedding model by construction. For langchain_postgres.PGVector specifically, add_texts embeds then delegates to add_embeddings, which upserts on id conflict (ON CONFLICT DO UPDATE) -- confirmed by reading its source -- so calling this repeatedly for the same entity.uid replaces rather than duplicates. Other backends are expected to honor the same "ids is an upsert key" convention their own add_texts docs promise; this writer doesn't depend on pgvector's specific mechanism.

VectorEntityDoc(uid: str, text: str, metadata: dict[str, Any]) dataclass

Pre-computed embedding text + metadata for one entity.

Snapshotting an entity's fields into a plain dataclass -- rather than reading them off the ORM instance -- lets a caller build this before session.commit() and upsert it afterward without touching the entity again. SQLAlchemy's default expire_on_commit marks all ORM-tracked attributes stale on commit, so reading them from an after_commit hook (as KnowledgeEngine.ingest does to keep the vector store out of the canonical transaction) would otherwise trigger a refresh query per entity.

msflib.knowledge.provisioning

One-time, superuser-only provisioning for the Apache AGE + pgvector backends.

Neither CREATE EXTENSION nor the GRANTs a non-superuser app role needs on ag_catalog (and on the AGE graph's own generated schema) can be run by that role itself — Postgres requires elevated privilege for both. Nothing in the app's normal runtime path does this (AgeGraphWriter.ensure_graph is only exercised by live test fixtures), so it has to happen out-of-band, once per database, via a superuser connection.

This is deliberately not a seeder: seeders run repeatedly as the app's own runtime role and insert business data; this runs once, as a different (superuser) role, and performs DDL/GRANTs the app role must never need at request time.

provision_age_and_pgvector(superuser_dsn: str, *, app_role: str, graph_name: str, provision_pgvector: bool = True) -> None

Idempotently provision AGE (and optionally pgvector) for app_role.

superuser_dsn must authenticate as a Postgres superuser (or a role with CREATEDB/extension-owner privilege) pointed at the target database — the app's normal runtime DSN will not work here and should never need to. Safe to run repeatedly: extensions use IF NOT EXISTS, the graph is only created if absent, and GRANT is idempotent.

msflib.knowledge.queries

Semantic query objects — the public query API of the knowledge module.

Nothing outside the module sees Cypher, SQL, or a vector-store handle. Agents, routers, and services construct these objects; the planner compiles them into scoped physical plans. Each query object also doubles as an auto-generated agent tool (see agent.tool_registry), with the docstring used as the tool description.

FindEntity(uid: str | None = None, name: str | None = None, entity_type: str | None = None, limit: int = 20) dataclass

Look up knowledge entities by uid, exact name, or entity type.

Returns matching entities with their properties and provenance.

FindNeighbors(uid: str, depth: int = 1, relation_types: tuple[str, ...] = (), direction: str = 'both') dataclass

Return entities directly or transitively connected to a given entity.

Useful for questions like "what is connected to pump P-101?".

FindPath(source_uid: str, target_uid: str, max_depth: int = 6, relation_types: tuple[str, ...] = ()) dataclass

Find the shortest relationship path between two entities.

Useful for questions like "how does the feed line reach the separator?".

Traverse(start_uid: str, relation_types: tuple[str, ...], direction: str = 'out', depth: int = 3, as_of: datetime | None = None) dataclass

Walk the graph from a starting entity along given relation types.

Directional variant of FindNeighbors for process-flow questions like "what is downstream of compressor C-101?".

SearchConcepts(text: str, k: int = 8, entity_types: tuple[str, ...] = ()) dataclass

Semantic similarity search over entities and concepts.

Finds knowledge related to a natural-language description even when exact names are unknown.

FindEvidence(uid: str, limit: int = 20) dataclass

Return source evidence (observations, excerpts, provenance) for an entity.

Use this to ground or verify a claim against original sources.

ResolveAliases(name: str, entity_type: str | None = None, limit: int = 5) dataclass

Resolve a free-text name to canonical knowledge entities.

Handles synonyms, abbreviations, and tag variants (e.g. "P101" vs "Pump P-101").

QueryResult(items: list[dict] = list(), evidence: list[dict] = list(), metadata: dict = dict()) dataclass

Uniform result envelope returned by the engine for any semantic query.

msflib.knowledge.records

Structured-record ingestion: forms, JSON payloads and app rows -> knowledge.

Generic helpers (payload normalization, field selection, versions, row scope, model write hooks) live in msflib.core.records; sensitive-key handling in msflib.core.pii. Attribute access is lazy.

contracts

Contracts for structured-record ingestion: records, sources and mappers.

StructuredRecord(subject_type: str, subject_id: str, record_type: str, payload: Mapping[str, Any], scope: KnowledgeScope, version: int | None = None, account_id: int | None = None, metadata: Mapping[str, Any] = dict()) dataclass

One loaded record: filtered payload and metadata, scope from where the data lives.

RecordSource

Bases: Protocol

Loads a record from its job pointer; None means gone, so its facts are archived.

SyncableRecordSource

Bases: RecordSource, Protocol

A source whose app rows can be queued directly: snapshot(row) is what the queue sees.

RecordMapper

Bases: Protocol

Turns a StructuredRecord into a KnowledgeFragment; name is the extractor.

validate_record_subject_type(subject_type: Any) -> str

A record source's subject_type: an ingestion subject type without :.

record_source_key(subject_type: str, subject_id: str) -> str

Provenance.source_id for every fact from one record; unique per subject.

record_snapshot(subject_id: str, scope: KnowledgeScope, version: int | None, *, live: bool = True) -> IngestionSubjectSnapshot

How the job queue and the stage see a record: its id, tenant/workspace and version.

errors

Errors raised by structured-record ingestion.

RecordPayloadTooLargeError

Bases: ValueError

Raised by submit when the filtered payload exceeds RECORD_MAX_PAYLOAD_BYTES.

registry

Registries for record sources (by subject_type) and mappers (by record_type).

FunctionRecordMapper(func: Callable[..., KnowledgeFragment], *, name: str | None = None)

Adapts a plain (record, *, context) -> KnowledgeFragment callable to RecordMapper.

register_record_source(source: RecordSource, *, registry: RecordSourceRegistry | None = None, replace: bool = False) -> RecordSource

Register source on registry (the process-wide default unless overridden).

register_record_mapper(record_type: str, mapper: RecordMapper | Callable[..., KnowledgeFragment], *, registry: RecordMapperRegistry | None = None, replace: bool = False) -> RecordMapper

Register mapper for record_type; a plain function is wrapped automatically.

service

KnowledgeRecordService (inline records) and KnowledgeRecordQueue (its job queue).

KnowledgeRecordQueue(*, dispatch: IngestionDispatchService, sources: RecordSourceRegistry, enabled: bool = True, job_action: IngestionJobAction | None = None)

Bases: IngestionSubjectQueue

The knowledge_record queue; it accepts only subject types with a registered source.

KnowledgeRecordService(*, dispatch: IngestionDispatchService, sources: RecordSourceRegistry | None = None, mappers: RecordMapperRegistry | None = None, queue_enabled: bool = True, max_payload_bytes: int | None = None, job_action: IngestionJobAction | None = None, record_action: KnowledgeRecordAction | None = None)

submit(session: Session, *, record_type: str, external_id: str, payload: Any, scope: ScopeEnvelope | KnowledgeScope, fields: Sequence[str] | None = None, exclude: Sequence[str] | None = None, sensitive_fields: Sequence[str] | None = DEFAULT_SENSITIVE_FIELDS, metadata: Mapping[str, Any] | None = None, account_id: int | None = None, commit: bool = False) -> KnowledgeRecord

Store (or update) a record knowledge owns and queue it for ingestion.

scope must be server-derived, never client input. An unchanged resubmission is a no-op; a changed one bumps version.

submit_many(session: Session, *, record_type: str, items: Iterable[tuple[str, Any]], scope: ScopeEnvelope | KnowledgeScope, commit: bool = False, **options: Any) -> list[KnowledgeRecord]

submit each (external_id, payload) pair in one transaction; one job each.

retract(session: Session, *, record_type: str, external_id: str, scope: ScopeEnvelope | KnowledgeScope, purge: bool = True, commit: bool = False) -> KnowledgeRecord | None

Mark a record deleted (purge drops its payload) and queue archival of its facts.

snapshot_row(source: RecordSource | str, row: Any) -> IngestionSubjectSnapshot

row as the job queue sees it, through the registered source.

Pass it to queue.enqueue_snapshot or as sync_model_to_queue(snapshot=...).

preview(session: Session, record: StructuredRecord, *, settings: SettingsBase, parser_factory: Callable[[], LLMExtractionParser] | None = None, compiler: KnowledgeCompiler | None = None, sensitive_fields: Sequence[str] | None = DEFAULT_SENSITIVE_FIELDS) -> KnowledgeFragment

Map (and optionally compile) record without persisting anything.

sources

Built-in and app-owned record source implementations.

inline

Built-in inline source: loads records stored by KnowledgeRecordService.submit.

InlineRecordSource()

Bases: ModelRecordSource

submit already filtered the payload with the caller's policy, so it isn't re-filtered.

ensure_inline_source(registry: RecordSourceRegistry) -> None

Register InlineRecordSource unless something already claims its subject type.

model

App-owned record sources: the app's own table stays the source of truth.

ModelRecordSource (declarative, one row per record), @record_source (a loader function, with record_from_row) and BaseRecordSource (override hooks) -- from least to most custom code.

BaseRecordSource

Hook-based source: implement get_row and payload, override the rest as needed.

is_live(row: Any) -> bool

False means gone: the stage archives the record's facts.

ModelRecordSource(*, subject_type: str, model: type, record_type: str | Callable[[Any], str], payload: RowGetter | Iterable[str], fields: Sequence[str] | None = None, exclude: Sequence[str] | None = None, sensitive_fields: Sequence[str] | None = DEFAULT_SENSITIVE_FIELDS, version: str | VersionGetter | None = None, is_deleted: RowGetter | None = None, when: Callable[[Any], bool] | None = None, scope: RowScope | None = None, account: RowGetter | None = None, metadata: Callable[[Any], Mapping[str, Any]] | None = None)

Bases: BaseRecordSource

Declarative source for records stored one per row; see the README for each option.

FunctionRecordSource(subject_type: str, load: Callable[[Session, str], StructuredRecord | None], *, snapshot: Callable[[Any], IngestionSubjectSnapshot] | None = None)

A source built from a plain loader function (see record_source).

filter_record_data(payload: Any, metadata: Mapping[str, Any] | None, *, fields: Sequence[str] | None = None, exclude: Sequence[str] | None = None, sensitive_fields: Sequence[str] | None = DEFAULT_SENSITIVE_FIELDS) -> tuple[dict[str, Any], dict[str, Any]]

Normalize and filter a record's payload and metadata with the same sensitive-key policy.

record_source(subject_type: str, *, snapshot: Callable[[Any], IngestionSubjectSnapshot] | None = None, registry: RecordSourceRegistry | None = None, register: bool = True) -> Callable[[Callable[[Session, str], StructuredRecord | None]], FunctionRecordSource]

Turn a (session, subject_id) -> StructuredRecord | None loader into a source.

snapshot(row) (build it with record_snapshot) lets rows be queued directly.

record_from_row(row: Any, *, subject_type: str, record_type: str, payload: Any, subject_id: str | None = None, version: int | None = None, scope: RowScope | ScopeEnvelope | KnowledgeScope | None = None, fields: Sequence[str] | None = None, exclude: Sequence[str] | None = None, sensitive_fields: Sequence[str] | None = DEFAULT_SENSITIVE_FIELDS, account_id: int | None = None, metadata: Mapping[str, Any] | None = None) -> StructuredRecord

Build a filtered StructuredRecord from a row, for @record_source loaders.

stage

knowledge_record pipeline kind: load a structured record, map it and ingest it.

archive_record_source(engine: KnowledgeEngine, session: Session, *, source_id: str, tenant_id: int, workspace_id: int | None) -> int

Archive a record's live facts in its tenant/workspace and sub-scopes; return the count.

make_knowledge_record_stage(*, engine: KnowledgeEngine, sources: RecordSourceRegistry | None = None, mappers: RecordMapperRegistry | None = None, parser_factory: Callable[[], LLMExtractionParser] | None = None, ontology: OntologyRegistry | None = None) -> Stage

Build the knowledge_record stage; registries default to the process-wide ones.

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

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

testing

Contract check for apps implementing record sources.

check_record_source(source: SyncableRecordSource, session: Session, row: Any, *, update: Callable[[Any], None] | None = None, remove: Callable[[Any], None] | None = None, forbidden_keys: Sequence[str] = DEFAULT_SENSITIVE_FIELDS) -> StructuredRecord

Assert source round-trips row; return the loaded record.

load(snapshot(row).subject_id) must match the snapshot, and no key matching forbidden_keys may reach the payload or metadata. update must raise the version; after remove the record must load as None.

msflib.knowledge.scope_profiles

msflib.knowledge.service

ArchiveResult(relation_uids: list[str] = list(), observation_uids: list[str] = list()) dataclass

Outcome of superseding a source's previously-live facts.

IngestResult(entities: list[KnowledgeEntity] = list(), relations: list[KnowledgeRelation] = list(), observations: list[KnowledgeObservation] = list(), uid_map: dict[str, str] = dict(), resolved_existing: list[str] = list()) dataclass

Outcome of persisting one compiled fragment.

KnowledgeEngine(*, compiler: KnowledgeCompiler, projections: Sequence[KnowledgeProjection], entity_action: KnowledgeEntityAction | None = None, relation_action: KnowledgeRelationAction | None = None, observation_action: KnowledgeObservationAction | None = None, graph_writer: AgeGraphWriter | None = None, vector_writer: VectorStoreWriter | None = None)

ingest(session: Session, fragment: KnowledgeFragment, *, scope: ScopeEnvelope, commit: bool = True, supersede_source: bool = False) -> IngestResult

Compile a fragment and persist it under the caller's scope.

Entities matching an existing row (same scope, canonical type, case-insensitive name) are resolved to the existing entity instead of duplicated. Pass commit=False to participate in a larger caller transaction; ModelAction post-commit lifecycle events, the vector-store upsert, and the fragment-ingested event all fire on the eventual session.commit() rather than immediately.

Pass supersede_source=True when re-ingesting a fragment for a source that may have been extracted before (e.g. a document re-index): every currently-live relation/observation asserted by the same source_id is superseded (see archive_source) before this fragment's own facts are persisted.

archive_source(session: Session, *, source_id: str, scope: ScopeEnvelope, include_subscopes: bool = False, commit: bool = True) -> ArchiveResult

Supersede every currently-live relation/observation from source_id.

Entities are never touched: they are deduped/merged across sources (find_canonical_match), so one source's supersession must not affect another source's claim on the same entity.

If a graph_writer is configured, the AGE projection is kept in sync synchronously: archived relations already carry valid_to/ lifecycle (the same properties AgeGraphProjection already filters reads on client-side), so re-running upsert_relation on each archived relation refreshes the existing edge in place — no separate delete method needed. Vector store is untouched: it only ever holds entity embeddings, never relations, so there is nothing to desync there. include_subscopes=True archives across every project/conversation/sub-thread in the tenant/workspace.

engine

KnowledgeEngine — the service facade everything else talks to.

Reads: semantic query + ScopeEnvelope → logical plan (with mandatory scope rewriting) → physical dispatch to a projection → QueryResult with evidence.

Writes: KnowledgeFragment + ScopeEnvelope → compiler passes → persistence through the action classes (single transaction; ModelAction fires row-level lifecycle events), then graph projection upserts through the optional graph_writer (same session, same transaction), then entity-embedding upserts through the optional vector_writer (separate backend, not part of the canonical transaction — see VectorStoreWriter's module docstring) → fragment-level eventbus events.

IngestResult(entities: list[KnowledgeEntity] = list(), relations: list[KnowledgeRelation] = list(), observations: list[KnowledgeObservation] = list(), uid_map: dict[str, str] = dict(), resolved_existing: list[str] = list()) dataclass

Outcome of persisting one compiled fragment.

ArchiveResult(relation_uids: list[str] = list(), observation_uids: list[str] = list()) dataclass

Outcome of superseding a source's previously-live facts.

KnowledgeEngine(*, compiler: KnowledgeCompiler, projections: Sequence[KnowledgeProjection], entity_action: KnowledgeEntityAction | None = None, relation_action: KnowledgeRelationAction | None = None, observation_action: KnowledgeObservationAction | None = None, graph_writer: AgeGraphWriter | None = None, vector_writer: VectorStoreWriter | None = None)

ingest(session: Session, fragment: KnowledgeFragment, *, scope: ScopeEnvelope, commit: bool = True, supersede_source: bool = False) -> IngestResult

Compile a fragment and persist it under the caller's scope.

Entities matching an existing row (same scope, canonical type, case-insensitive name) are resolved to the existing entity instead of duplicated. Pass commit=False to participate in a larger caller transaction; ModelAction post-commit lifecycle events, the vector-store upsert, and the fragment-ingested event all fire on the eventual session.commit() rather than immediately.

Pass supersede_source=True when re-ingesting a fragment for a source that may have been extracted before (e.g. a document re-index): every currently-live relation/observation asserted by the same source_id is superseded (see archive_source) before this fragment's own facts are persisted.

archive_source(session: Session, *, source_id: str, scope: ScopeEnvelope, include_subscopes: bool = False, commit: bool = True) -> ArchiveResult

Supersede every currently-live relation/observation from source_id.

Entities are never touched: they are deduped/merged across sources (find_canonical_match), so one source's supersession must not affect another source's claim on the same entity.

If a graph_writer is configured, the AGE projection is kept in sync synchronously: archived relations already carry valid_to/ lifecycle (the same properties AgeGraphProjection already filters reads on client-side), so re-running upsert_relation on each archived relation refreshes the existing edge in place — no separate delete method needed. Vector store is untouched: it only ever holds entity embeddings, never relations, so there is nothing to desync there. include_subscopes=True archives across every project/conversation/sub-thread in the tenant/workspace.

hybrid_retrieval

Hybrid retrieval: vector-seed then graph-expand, composed over KnowledgeEngine.

A composition strategy over the engine's single-query facade, not a KnowledgeEngine method — this keeps the engine's surface as the stable "one semantic query in, one QueryResult out" contract, while retrieval strategies that combine multiple queries can evolve independently (e.g. a future re-ranking or multi-hop variant) without touching the engine.

hybrid_retrieve(engine: KnowledgeEngine, *, text: str, scope: ScopeEnvelope, seed_k: int = 8, expand_depth: int = 1, relation_types: tuple[str, ...] = (), direction: str = 'both', entity_types: tuple[str, ...] = (), limit: int | None = None, include_drafts: bool = False) -> QueryResult

Vector-search for seed entities, then graph-expand each seed's neighborhood.

Vector relevance is the primary order: deduplicated seeds come first, in their similarity rank, followed by all of their (deduplicated) graph neighbors — so every seed outranks every neighbor, rather than the two being pooled and re-ranked together. Empty seed results short-circuit before any graph query.