Skip to content

AI capabilities

This guide walks through the AI building blocks in msflib-ai-core: getting a chat model and embeddings, chunking documents, storing and searching them with tenant isolation, versioning prompts, and tracking usage. Each section is short and builds on the previous one.

All imports come from subpackages. msflib.ai_core itself exports nothing.

1. Settings

AICoreSettings holds the provider, model and key fields. In a host app it is read from the AI_CORE namespace of your settings (see the environment variable rules on the module page). For scripts and tests you can construct it directly:

from msflib.ai_core.config import AICoreSettings

ai_settings = AICoreSettings(
    LLM_PROVIDER="openai",
    LLM_MODEL="gpt-4o",
    LLM_API_KEY="sk-...",
    EMBEDDING_API_KEY="sk-...",
    VECTOR_STORE_URL="postgresql+psycopg://user:password@localhost:5432/app",
)

In a host app, pass settings.scope("AI_CORE") (it works for composed and subclassed settings) or use get_ai_dependencies(settings).get_llm(); do not pass the host settings object itself, which raises AttributeError when the settings are composed. The snippets below use the standalone ai_settings.

VECTOR_STORE_URL is only needed for section 4. The default pgvector backend raises ValueError: VECTOR_STORE_URL is required for pgvector backend without it, and so do qdrant, weaviate and milvus. Pinecone needs VECTOR_STORE_API_KEY instead, and chroma and faiss need neither.

Provider names are plain strings, validated against the provider type registry. The registry mirrors the providers your installed langchain supports (for example openai, azure_openai, anthropic, groq, ollama), and each needs its package installed (see the extras table on the module page). List them with get_provider_type_registry().names() from msflib.ai_core.providers.

2. A chat model and embeddings

The next snippet needs the openai extra.

from msflib.ai_core.providers import get_embeddings, get_llm

llm = get_llm(ai_settings)
reply = llm.invoke("Summarise MSFLib in one sentence.")
print(reply.content)

embeddings = get_embeddings(ai_settings)
vector = embeddings.embed_query("hello")

get_llm and get_embeddings return the underlying langchain objects, so streaming (llm.stream(...)), batching and the rest of the langchain API work unchanged. To use another provider, change LLM_PROVIDER and LLM_MODEL, for example "anthropic" with an Anthropic model id (install the anthropic extra). Embeddings are configured separately with EMBEDDING_PROVIDER and EMBEDDING_MODEL, and Anthropic offers no embeddings.

In a host app, prefer the registry-aware variants on the dependency namespace (get_llm_from_registry(session, scope=scope)), which pick a database-backed provider profile for the caller's tenant, workspace or user before falling back to settings. See get_ai_dependencies on the module page.

3. Chunking documents

IngestionPipeline loads, cleans and chunks a document and returns the chunks with metadata.

from msflib.ai_core.services.ingestion_pipeline import IngestionPipeline

pipeline = IngestionPipeline()

result = pipeline.ingest(
    {
        "format": "markdown",
        "content": "# Refunds\n\nRefunds take 5 days.\n\n## Exceptions\n\nGift cards are final sale.",
        "source_id": "faq-1",
    }
)

for chunk in result.chunks:
    print(chunk.metadata["chunk_index"], chunk.text)

result also reports which chunking strategy ran (result.chunk_strategy) and whether it came from the built-in fallback provider (result.fallback_used; the default markdown strategy is a fallback-provider strategy, so this is True in the example above and does not signal an error).

You can pick the strategy and size per call:

long_text = "refund policy " * 400

result = pipeline.ingest(
    long_text,
    source_type="text",
    strategy_name="token_window",
    chunk_options={"chunk_size": 200, "chunk_overlap": 30},
)

Default embedder and indexer are placeholders

Without arguments, IngestionPipeline uses a deterministic stand-in embedder and an indexer that only returns ids like chunk-0. They exist for wiring and tests. For real retrieval, either pass your own embedder= and indexer= callables, or take result.chunks[i].text and store them yourself as shown in the next section. An embedder is embedder(texts: list[str]) -> list[list[float]]; an indexer is indexer(records: list[dict]) -> list[str], where each record has text, embedding and metadata.

Other constructor arguments (loaders=, cleaner=, chunking_registry=) and per-call workspace_settings= / chunking_policy= let you add loaders and tune strategy selection. Custom chunking strategies and document loaders are covered in the chunking and extraction READMEs under modules/ai_core/msflib/ai_core/services/.

4. Storing and searching with tenant isolation

This section continues from sections 1 to 3 (ai_settings, embeddings and result). It needs the pgvector extra and a Postgres database with the vector extension.

from msflib.ai_core.vector_store import get_vector_store, similarity_search_scoped
from msflib.scope import build_context_scope

store = get_vector_store(ai_settings, "faq", embeddings)

store.add_texts(
    [chunk.text for chunk in result.chunks],
    metadatas=[
        {"tenant_id": "1", "workspace_id": "7", "source_id": "faq-1"}
        for _ in result.chunks
    ],
)

scope = build_context_scope(tenant_id=1, workspace_id=7)
hits = similarity_search_scoped(store, query="How long do refunds take?", scope=scope, k=4)

Points to know:

  • get_vector_store returns a ScopedVectorStore. It supports similarity_search, add_texts, add_documents, delete and close. The native langchain store is available as store.raw.
  • Always search with similarity_search_scoped in multi-tenant code. It builds mandatory tenant and workspace metadata filters from the ScopeEnvelope. Calling store.similarity_search(...) directly applies no filter.
  • Write dimension values as strings ("1", not 1), because build_context_scope stringifies them.
  • The metadata you write must carry the dimensions the filter reads. tenant_id and workspace_id are required, exact-match string values; a scope with no workspace matches only content written with no workspace. Conversation, sub-thread and private-account dimensions are optional: content written without conversation_id, sub_thread_id or private_to_account_id is treated as shared across the workspace, so tag anything that belongs to a single conversation, sub-thread or account. Write through the ScopedVectorStore, not store.raw: on backends without native missing-key filters (Pinecone, Chroma, Milvus) it adds the __null__ markers that make untagged content visible as shared. build_retrieval_filters(scope=scope) shows exactly what a search requires.
  • Collection names are "<VECTOR_STORE_COLLECTION_PREFIX>_<name>", sanitised for the backend. The prefix organises collections; it does not isolate tenants. Isolation comes from the metadata filters.
  • The default backend is pgvector (extras = ["pgvector"], and CREATE EXTENSION vector; once on the database). Qdrant, Pinecone, Chroma, Weaviate, Milvus and FAISS are also supported (each has an extra of the same name); set VECTOR_STORE_BACKEND.

5. Versioned prompts

PromptRegistryService stores prompt versions in the database and resolves the active one by scope, in the order user, then workspace, then global.

The examples below assume session is an open SQLModel Session from your host app (for instance the one injected by your get_session dependency), and that the ai_prompt_version table exists. Import msflib.tenancy.models.tenant before SQLModel.metadata.create_all (or your migrations), because ai-core's provider and vector-store profile tables reference tenant; see the known issue on the module page.

from msflib.ai_core.services.prompt_registry import PromptRegistryService

prompts = PromptRegistryService()

prompts.register(session, "support.system", "v1", "You are a helpful assistant.", set_active=True)
prompts.register(
    session, "support.system", "v2", "You are concise.", workspace_id=5, set_active=True
)

prompts.get_active_version(session, "support.system").content
# 'You are a helpful assistant.'
prompts.get_active_version(session, "support.system", workspace_id=5).content
# 'You are concise.'

Registering the same key, version and scope twice raises PromptVersionConflictError. To roll a prompt out without a deploy, register a new version and activate it with set_active. A user-level override uses both workspace_id and account_id.

6. Usage tracking and rate limiting

from msflib.ai_core.providers import get_llm
from msflib.ai_core.tracking.usage import TokenUsageCallback

llm = get_llm(ai_settings)
usage = TokenUsageCallback(workspace_id="7")
llm.invoke("Hello", config={"callbacks": [usage]})
print(usage.prompt_tokens, usage.completion_tokens, usage.total_tokens)

The callback accumulates token counts across every call it is attached to, so create one per request. Streaming responses only report usage if the provider includes it. On every model call it also emits ai_core.token_usage on the event bus, but only if a listener is registered. Register it on the app emitter with @app_emitter.on("ai_core.token_usage") (app_emitter = bind_app_emitter(app)). It fires for model calls made inside a request (including BackgroundTasks); in threads or scripts in the web process, wrap the call in with use_app_emitter(app):, otherwise the event goes to the default emitter and the listener never runs. A Celery task has no FastAPI app: register its listener on the worker's emitter, AppEmitter(get_emitter()) (both from msflib.eventbus). The payload (workspace_id, prompt_tokens, completion_tokens, total_tokens, total_cost) carries the callback's running totals rather than per-call amounts. Use a fresh callback per request, or subtract the previous totals, before billing from it.

from msflib.ai_core.tracking.limiter import RateLimiter, RateLimitExceeded

limiter = RateLimiter(rpm=60, redis_url="redis://localhost:6379/0")

try:
    limiter.check("workspace:7")
except RateLimitExceeded:
    ...  # return HTTP 429

check raises RateLimitExceeded once a namespace makes more than rpm calls within the same clock minute (a fixed window, not a sliding one). Without a redis_url (or redis_client), counts are kept in memory and are per process, which is fine for development but not for several workers. If Redis is unreachable, check raises a Redis connection error, not RateLimitExceeded, so handle that separately. Nothing calls the limiter for you; add it where you invoke the model.

Concepts behind these pieces

  • Scopes. Provider and vector-store profiles resolve through the tiers global, tenant, workspace and user. Prompts have no tenant tier: they resolve user, then workspace, then global. Retrieval filters are built from the caller's tenant and workspace, plus the conversation, sub-thread and private-account dimensions when present. The caller's scope is a ScopeEnvelope (see Scopes, tenancy and workspaces).
  • Policies. Chunking and runtime behaviour are controlled by policy objects (ChunkingPolicySchema, MiddlewarePolicySchema) resolved from tiered settings (see Policy and Tiered configuration). A policy can select a chunking provider with a fallback chain.
  • Ingestion flow. Load, clean, chunk, embed, index. Each stage is a replaceable callable or registry entry.
  • Secrets. Provider API keys stored in profiles are encrypted with PROVIDER_SECRET_ENCRYPTION_KEY, which falls back to the core SECRET_KEY when unset. Set it explicitly in production and keep keys out of logs.

Going further

  • Agents and LangGraph: checkpointer and memory store scoping, the tools library and runtime policy are described in modules/ai_core/README.md until they are migrated into this site.
  • Long-running ingestion of uploaded files: Long-running jobs and ingestion.
  • Provider and vector-store profile management over HTTP: the routers described on the module page.