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_storereturns aScopedVectorStore. It supportssimilarity_search,add_texts,add_documents,deleteandclose. The native langchain store is available asstore.raw.- Always search with
similarity_search_scopedin multi-tenant code. It builds mandatory tenant and workspace metadata filters from theScopeEnvelope. Callingstore.similarity_search(...)directly applies no filter. - Write dimension values as strings (
"1", not1), becausebuild_context_scopestringifies them. - The metadata you write must carry the dimensions the filter reads.
tenant_idandworkspace_idare 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 withoutconversation_id,sub_thread_idorprivate_to_account_idis treated as shared across the workspace, so tag anything that belongs to a single conversation, sub-thread or account. Write through theScopedVectorStore, notstore.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"], andCREATE EXTENSION vector;once on the database). Qdrant, Pinecone, Chroma, Weaviate, Milvus and FAISS are also supported (each has an extra of the same name); setVECTOR_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 coreSECRET_KEYwhen 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.mduntil 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.