Core Module¶
Package: cogniverse_core Location: libs/core/cogniverse_core/
Table of Contents¶
- Overview
- Package Structure
- Type-Safe Agent System
- AgentBase
- AgentInput / AgentOutput / AgentDeps
- A2AAgent
- Content Rails
- RLM Options
- Agent Mixins
- Registries
- Memory Management
- Query Encoding
- Event System
- Durable Execution
- Approval Interfaces
- Media Access
- Configuration
- Usage Examples
- Architecture Position
- Testing
- Cache Subsystem
- Utility Modules
- VLM Interface
- Model Loaders
- Backend Factory & Profile Validation
Overview¶
The Core package is in the Core Layer for all agent implementations in Cogniverse. It provides:
- Type-Safe Agents: Generic base classes with compile-time type checking and runtime Pydantic validation
- A2A Protocol Support: Google's Agent-to-Agent protocol for inter-agent communication
- DSPy Integration: Native DSPy module support for AI-powered agents
- Multi-Tenancy: Built-in tenant isolation for enterprise deployments
- Component Registries: Dynamic registration and discovery of agents, backends, and schemas
- Memory System: Mem0-based persistent agent memory
All concrete agent implementations (OrchestratorAgent, SearchAgent, etc.) inherit from these base classes.
Package Structure¶
cogniverse_core/conversation.py — ConversationStore: per-context conversation turns in Mem0, keyed by (tenant_id, context_id), stored verbatim and retrieved by metadata filter. get_history treats every stored row as untrusted — a row whose metadata is not a dict, or whose seq is not a number, is skipped rather than allowed to crash the read, so one malformed row can never drop a context's whole history. The agent dispatcher uses it to load/save history around each agent call so callers (the messaging gateway) need no Mem0 connection; the load is time-bounded on the reply path (CONVERSATION_LOAD_TIMEOUT_S) so a hung Mem0 degrades to no-history rather than stalling the reply, and the save runs off the reply path bounded by CONVERSATION_SAVE_TIMEOUT_S. store_turn(context_id, role, content, seq) stores the seq its writer assigned and reads order by it, so turns written by different processes read back in the writers' order, not the order the writes landed; the dispatcher takes each turn's seq from the runtime's shared conversation ledger. A turn is two appends, so the second one can fail with the first already durable: is_transient_turn_write_error(exc) types a failure the write never got a verdict for (transport, timeout, a retryable status) apart from one the backend refused, and store_missing_assistant_marker(context_id, cause, seq) writes the assistant_missing row that takes a lost reply's place, carrying the failure's type name and none of its message. get_history(context_id, max_turns) returns the newest max_turns (MAX_HISTORY_TURNS by default, every turn for None) and renders only RENDERED_TURN_ROLES (user, assistant), so a marker never reaches an agent as a reply and the user turn it belongs to reads as unanswered; get_missing_assistant_markers(context_id) returns those markers for the context, oldest first. A run cancelled before its reply stores a run_cancelled row (RUN_CANCELLED_ROLE, content RUN_CANCELLED_TEXT) in the reply's place; get_thread(context_id) returns every turn a person is shown (DISPLAYED_TURN_ROLES: user, assistant, run_cancelled), oldest first, while get_history leaves the marker out.
cogniverse_core/messaging_auth.py — InviteTokenManager (invite tokens in the system config store, canonicalized through the ConfigManager; a token is claimed for one external user and consumed with compare-and-set writes, so it is redeemed once across processes) and UserTenantMapper (Telegram-user→tenant mappings in the Mem0 system partition). Lives in core so both the messaging gateway and the runtime's registration routes share one implementation.
cogniverse_core/
├── agents/ # Agent base classes and mixins
│ ├── base.py # AgentBase[InputT, OutputT, DepsT]
│ ├── a2a_agent.py # A2AAgent with A2A protocol + DSPy
│ ├── tenant_aware_mixin.py # Multi-tenancy mixin
│ ├── a2a_mixin.py # A2AEndpointsMixin (standalone-agent agent-card endpoint)
│ ├── rails.py # Content rails (topic/safety/format) for agent I/O
│ └── rlm_options.py # RLMOptions — per-query Recursive Language Model config
├── approval/ # Human-in-the-loop approval interfaces
│ ├── interfaces.py # ApprovalStatus, ReviewItem, ReviewDecision, ApprovalBatch,
│ │ # ConfidenceExtractor, FeedbackHandler, ApprovalStorage
│ └── training_schema.py # Exact approved synthetic supervision contracts
├── registries/ # Component registries
│ ├── agent_registry.py # Agent class registration
│ ├── backend_registry.py # Backend provider registration
│ ├── schema_registry.py # Schema template registration
│ ├── schema_deploy_lease.py # Cross-process lease on the application package
│ ├── adapter_store_registry.py # AdapterStoreRegistry (entry-point auto-discovery)
│ ├── workflow_store_registry.py # WorkflowStoreRegistry (entry-point auto-discovery)
│ ├── exceptions.py # Registry exceptions
│ └── registry.py # Base registry class
├── memory/ # Memory system
│ ├── manager.py # Mem0MemoryManager (add_memory, search_memory, drop_session, lifecycle tick)
│ ├── schema.py # KnowledgeSchema, KnowledgeRegistry, Retention, Sensitivity, Pinnable, ContradictionPolicy
│ ├── provenance.py # Provenance, CitationRef, DerivationKind, ProvenanceWalker, make_provenance
│ ├── provenance_store.py # Vespa-backed provenance persistence
│ ├── contradiction.py # ContradictionDetector, ConflictSet, reconcile()
│ ├── trust.py # TrustRecord, compute_initial_trust, rank_with_trust, apply_endorsement
│ ├── federation.py # FederationService (org-trunk + tenant overlays, cross-tenant ACLs)
│ ├── pinning.py # PinService, PinQuotas, PinRecord
│ ├── lifecycle_scheduler.py # LifecycleScheduler (schema-driven periodic cleanup)
│ ├── mem0_embedder.py # DenseOnMem0Embedder (DenseOn prompt + L2-norm adapter for Mem0)
│ ├── backend_config.py # Memory backend configuration
│ ├── backend_vector_store.py # Vector store integration
│ └── _timestamps.py # Internal timestamp normalization helpers
├── common/ # Shared utilities
│ ├── cache/ # Caching subsystem (see Cache Subsystem section)
│ ├── models/ # Model loaders (see Model Loaders section)
│ ├── media/ # Media access abstraction (see Media Access section)
│ ├── utils/ # Utility functions (see Utility Modules section)
│ ├── tenant_utils.py # Re-exports foundation tenant helpers; hosts assert_tenant_exists and tenant deletion markers
│ ├── dspy_module_registry.py # Re-exports foundation DSPyModuleRegistry / DSPyOptimizerRegistry
│ ├── dynamic_dspy_mixin.py # Dynamic DSPy mixin
│ ├── health_mixin.py # Health check mixin
│ ├── agent_models.py # AgentEndpoint and shared agent models
│ └── vlm_interface.py # Vision Language Model interface
├── events/ # Real-time event notification system (see Event System section)
│ ├── types.py # Event type definitions (StatusEvent, ProgressEvent, etc.)
│ ├── queue.py # EventQueue/QueueManager protocols, TaskCancelled, per-request queue binding
│ └── backends/ # Backend implementations
│ └── memory.py # In-process EventQueue backend
├── durable/ # Durable execution for long-running workflows (see Durable Execution section)
│ ├── pipeline_checkpoint.py # PipelineCheckpoint / PipelineCheckpointStatus / PipelineCheckpointConfig
│ └── pipeline_checkpoint_storage.py # PipelineCheckpointStorage (span persist + resume lookup)
├── query/ # Query encoding utilities (see Query Encoding section)
│ └── encoders.py # QueryEncoder implementations + QueryEncoderFactory
├── factories/ # Factory classes
│ └── backend_factory.py # BackendFactory for creating backend instances
├── interfaces/ # Protocol interfaces (re-exports, no local classes)
├── validation/ # Validation utilities
│ └── profile_validator.py # ProfileValidator — backend profile config validation
├── schemas/ # Data schemas
│ └── filesystem_loader.py # FilesystemSchemaLoader — schema loading from files
└── backends/ # Backend abstractions (re-exports, no local classes)
config/ and telemetry/ backward-compatibility shims (which re-exported cogniverse_foundation.config / telemetry symbols under the cogniverse_core namespace) have been removed; import directly from cogniverse_foundation.
Type-Safe Agent System¶
The type-safe agent system uses Python generics to provide compile-time type checking and runtime validation. This is the foundation of all agents in Cogniverse.
AgentBase¶
AgentBase[InputT, OutputT, DepsT] is the abstract base class for all agents.
Type Parameters:
InputT: Agent input type (must extendAgentInput)OutputT: Agent output type (must extendAgentOutput)DepsT: Agent dependencies type (must extendAgentDeps)
Key Features:
- Generic type parameters extracted at class definition time
- Automatic Pydantic validation for inputs and outputs
- Runtime type checking
- Statistics tracking (process count, error count)
Example (illustrative pattern showing the type-safe structure - actual agents use concrete implementations):
from cogniverse_core.agents.base import AgentBase, AgentInput, AgentOutput, AgentDeps
from typing import Any, Dict, List
# Note: This example uses generic types for illustration.
# Real implementations use concrete classes from cogniverse_vespa or other backend packages.
class MySearchInput(AgentInput):
query: str
top_k: int = 10
class MySearchOutput(AgentOutput):
results: List[Dict[str, Any]]
total_count: int
class MySearchDeps(AgentDeps):
search_client: Any # In practice: VespaSearchBackend from cogniverse_vespa
embedding_model: str = "colpali"
class MySearchAgent(AgentBase[MySearchInput, MySearchOutput, MySearchDeps]):
async def _process_impl(self, input: MySearchInput) -> MySearchOutput:
# IDE autocomplete works here - input.query, input.top_k
results = await self.deps.search_client.search(
query=input.query,
limit=input.top_k
)
return MySearchOutput(results=results, total_count=len(results))
API:
| Method | Description |
|---|---|
__init__(deps: DepsT) | Initialize agent with typed dependencies |
async _process_impl(input: InputT) -> OutputT | Abstract - Implement agent logic (subclasses override) |
async process(input: InputT, stream: bool = False) | Public API - validates input and calls _process_impl |
async run(raw_input: Dict, stream: bool = False) | Run with raw dict, validates input/output |
validate_input(raw: Dict) -> InputT | Validate and convert to typed input |
validate_output(raw: Dict) -> OutputT | Validate and convert to typed output |
get_input_schema() -> Dict | Get JSON schema for input type |
get_output_schema() -> Dict | Get JSON schema for output type |
get_stats() -> Dict | Get processing statistics |
call_dspy(..., deadline=LMCallDeadline) binds the caller's deadline for the LM call it makes: the call gets the time that is left, no attempt starts once the caller stops waiting (cancelled, or the deadline passed), and past the deadline it raises LMCallDeadlineExceeded naming the bound LM's endpoint.
call_dspy(..., output_field=..., stream_view=...) emits token events with the new text in message and data={"accumulated": ..., "output_field": ...}. accumulated is always a prefix of the stripped value the module returns for that field: under a JSON adapter the listener's raw JSON is decoded (quotes and escapes resolved, nothing after the closing quote), trailing whitespace waits for the text that follows it, and stream_view narrows the text to what the agent's own post-processing of the field keeps. The streamed LM call runs on a dedicated LM stream loop, so its chunk handling never occupies the serving loop. Token streaming belongs to the agent that owns the active stream; nested agents do not emit answer tokens into their caller's stream.
A failed stream ends in one {"type": "error", "agent", "error_type", "message"} event; when the LM answered with a 4xx the event also carries status and the message names it (... (LM HTTP 413) ...).
AgentInput / AgentOutput / AgentDeps¶
These are Pydantic BaseModel subclasses that define agent interfaces:
from cogniverse_core.agents.base import AgentInput, AgentOutput, AgentDeps
class AgentInput(BaseModel):
"""Base class for all agent inputs. Undeclared fields are retained."""
model_config = ConfigDict(extra="allow")
class AgentOutput(BaseModel):
"""Base class for all agent outputs. Strict - no extra fields allowed."""
model_config = ConfigDict(extra="forbid")
class AgentDeps(BaseModel):
"""Base class for agent dependencies. Agents are tenant-agnostic at startup — tenant_id arrives per-request."""
model_config = ConfigDict(extra="allow") # Dependencies can have additional fields
Important: tenant_id is not a constructor parameter — it arrives per-request via the A2A task payload, not at agent startup. This is how multi-tenancy is enforced without coupling the agent to a specific tenant at construction time.
A2AAgent¶
A2AAgent[InputT, OutputT, DepsT] extends AgentBase with:
- A2A Protocol Support: Standard A2A spec integration via the runtime's
A2AStarletteApplication - DSPy Integration: Optional DSPy module for AI processing
- HTTP Endpoint: Agents expose a
/processendpoint (dict input, not Task objects); A2A routing handled by the runtime - Inter-Agent Communication: Agents call each other via
httpx.AsyncClientthrough the runtime'sAgentDispatcher
from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from cogniverse_core.agents.base import AgentDeps, AgentInput, AgentOutput
class MyAgent(A2AAgent[AgentInput, AgentOutput, AgentDeps]):
async def _process_impl(self, input: AgentInput) -> AgentOutput:
# Use DSPy module if available
if self.dspy_module:
result = self.dspy_module(query=input.query)
return AgentOutput(result=result.answer)
return AgentOutput(result="default")
# Create agent (deps are tenant-agnostic; tenant_id arrives per-request)
deps = AgentDeps()
config = A2AAgentConfig(
agent_name="orchestrator_agent",
agent_description="Routes queries to appropriate agents",
capabilities=["intelligent_routing", "query_analysis", "agent_orchestration"],
port=8001
)
agent = MyAgent(deps=deps, config=config)
# The A2A HTTP server is managed by the runtime's A2AStarletteApplication,
# not embedded in the agent. Agents expose a /process endpoint via the runtime.
A2A Endpoints (served by the runtime's A2AStarletteApplication, an a2a-sdk JSON-RPC 2.0 app, mounted at /a2a):
| Endpoint | Method | Description |
|---|---|---|
/a2a/ | POST | JSON-RPC 2.0 endpoint — methods message/send, message/stream (SSE), tasks/get, tasks/cancel, tasks/pushNotificationConfig/* |
/a2a/.well-known/agent-card.json | GET | Agent card per A2A spec |
/a2a/.well-known/agent.json | GET | Deprecated agent-card alias, kept for backward compatibility |
A2A requests are executed by CogniverseAgentExecutor, which dispatches through the runtime's AgentDispatcher. Separately, each agent is also reachable in-process via the agents router:
| Endpoint | Method | Description |
|---|---|---|
/agents/{agent_name}/process | POST | Process a task (dict payload, not an A2A Task object) |
/agents/{agent_name}/message | POST | Enqueue an inbound steering message (stop/constraint/interrupt) for a running agent session |
Standalone agent apps built outside the unified runtime can mix in A2AEndpointsMixin (agents/a2a_mixin.py) for a single GET /.well-known/agent-card.json endpoint plus a register_with_registry() helper that POSTs the agent's card to a central registry.
Content Rails¶
agents/rails.py provides lightweight input/output guardrails. AgentBase.process() runs the configured chains around _process_impl(), and the runtime enforces them at the gateway front door (the dispatcher's _execute_gateway_task runs input rails on the incoming query and output rails on the final response), so internal agent-to-agent calls aren't gated.
Three concrete rails (raise RailBlockedError on violation):
TopicBoundaryRail(allowed_topics, advisory=True)— flags queries that contain none of the allowed-topic keywords. Defaults to advisory (logs a warning, does not block) because coarse keyword matching would otherwise reject legitimate natural-language queries; setadvisory=Falseto hard-block (only with a deliberately broadallowed_topics).ContentSafetyRail(blocked_patterns)— blocks input/output matching any regex (e.g. prompt-injection phrases,<script>).OutputFormatRail(required_fields)— enforces required output fields + types.
Rails are configured under the rails block in config.json (enabled, input_rails, output_rails); RailsConfig.build_input_chain() / build_output_chain() compile the definitions into RailChains.
"rails": {
"enabled": true,
"input_rails": [
{"type": "topic_boundary", "params": {"advisory": true, "allowed_topics": ["video", "search"]}},
{"type": "content_safety", "params": {"blocked_patterns": ["ignore previous instructions", "<script>"]}}
],
"output_rails": []
}
RLM Options¶
agents/rlm_options.py defines RLMOptions, a per-query Pydantic config for enabling RLM (Recursive Language Model) inference — a DSPy 3.1+ paradigm that lets an LLM examine, decompose, and recursively call itself over near-infinite context instead of stuffing everything into one prompt window. It is attached to an agent's input type (e.g. SearchInput(query=..., rlm=...)) so callers can A/B test RLM against standard inference per request.
from cogniverse_core.agents.rlm_options import RLMOptions
# Explicit enable (A/B test group B)
rlm = RLMOptions(enabled=True, max_iterations=3)
# Auto-enable once context exceeds a character threshold
rlm = RLMOptions(auto_detect=True, context_threshold=50_000)
rlm.should_use_rlm(context_size=80_000) # -> True
| Field | Default | Purpose |
|---|---|---|
enabled | False | Explicitly force RLM for this query |
auto_detect | False | Enable RLM when context exceeds context_threshold |
context_threshold | 50_000 | Character threshold for auto_detect |
max_iterations | 3 (1-10) | Max REPL iterations |
max_llm_calls | 30 (1-100) | Max LLM sub-calls |
timeout_seconds | 300 (1-1800) | RLM processing timeout |
backend | "openai" | litellm backend id (openai, anthropic, litellm) |
model / api_base / api_key | None | Overrides; None falls back to the agent's model / provider defaults |
include_trajectory | False | Attach a bounded REPL trajectory to the result for auditing |
trajectory_max_entries | 32 (1-200) | Caps trajectory size when include_trajectory=True |
Agent Mixins¶
Mixins provide composable functionality that can be added to any agent:
MemoryAwareMixin¶
Adds Mem0-based persistent memory to agents:
Method Signature:
def initialize_memory(
self,
agent_name: str,
tenant_id: str,
embedder_base_url: str, # Required — DenseOn /v1/embeddings endpoint
*,
llm_model: str, # Required — no fallback default
backend_host: str = "localhost",
backend_port: int = 8080,
embedding_model: str = "lightonai/DenseOn",
llm_base_url: str = "http://localhost:11434",
llm_api_key: Optional[str] = None,
config_manager=None, # Required for schema deployment
schema_loader=None, # Required for schema templates
backend_config_port: Optional[int] = None,
auto_create_schema: bool = True,
) -> bool:
Usage Example:
from cogniverse_agents.memory_aware_mixin import MemoryAwareMixin
class MyAgent(AgentBase[...], MemoryAwareMixin):
def __init__(self, deps, config_manager, schema_loader):
super().__init__(deps)
self._config_manager = config_manager
self._schema_loader = schema_loader
# Memory is initialized per-request with tenant_id (not at construction)
async def _process_impl(self, input: MyInput) -> MyOutput:
# Initialize memory for this tenant (idempotent — only initializes once per tenant)
memory_config = self.search_config.get("memory", {})
self.initialize_memory(
agent_name="my_agent",
tenant_id=input.tenant_id, # Per-request, not from deps
embedder_base_url=memory_config.get("embedder_base_url"),
llm_model=memory_config.get("llm_model"),
backend_host=system_config.get("backend_url"),
backend_port=system_config.get("backend_port"),
embedding_model=memory_config.get("embedding_model"),
llm_base_url=memory_config.get("llm_base_url"),
config_manager=self._config_manager,
schema_loader=self._schema_loader,
backend_config_port=system_config.get("backend_config_port"),
)
# Search memories using mixin methods
context = self.get_relevant_context(query=input.query, top_k=5)
# Store new memory
self.update_memory(content=input.query, metadata={"type": "query"})
# Other available methods:
# self.is_memory_enabled() -> bool
# self.remember_success(query, result, metadata)
# self.remember_failure(query, error, metadata)
# self.clear_memory() -> bool
# self.get_memory_summary() -> Dict
TenantAwareMixin¶
Provides tenant context and isolation:
from cogniverse_core.agents.tenant_aware_mixin import TenantAwareAgentMixin
class MyAgent(AgentBase[...], TenantAwareAgentMixin):
def __init__(self, deps, ...):
super().__init__(deps)
# tenant_id is NOT set at construction — it arrives per-request
async def _process_impl(self, input: MyInput) -> MyOutput:
# tenant_id comes from the A2A task payload
tenant_id = input.tenant_id
TenantAwareAgentMixin.__init__(self, tenant_id=tenant_id)
tenant_context = self.get_tenant_context()
HealthCheckMixin¶
Adds health check capabilities:
from cogniverse_core.common.health_mixin import HealthCheckMixin
class MyAgent(AgentBase[...], HealthCheckMixin):
def get_health_status(self) -> Dict[str, Any]:
"""Override to provide custom health check logic (sync method)."""
return {
"status": "healthy",
"agent": self.__class__.__name__,
"custom_metric": self.get_custom_metric()
}
# The mixin also provides setup_health_endpoint() for FastAPI integration:
# self.setup_health_endpoint(app) # Adds GET /health endpoint
DynamicDSPyMixin¶
Provides dynamic DSPy module creation and configuration at runtime:
from cogniverse_core.common.dynamic_dspy_mixin import DynamicDSPyMixin
import dspy
class MyAgent(AgentBase[...], DynamicDSPyMixin):
def __init__(self, ...):
super().__init__(...)
# Register signatures for dynamic module creation
self.register_signature("my_task", MySignature)
def process(self, input_data):
# Get or create module based on current agent_config
module = self.get_or_create_module("my_task")
return module(query=input_data.query)
Registries¶
Registries provide dynamic component registration and discovery:
AgentRegistry¶
The AgentRegistry uses dependency injection - ConfigManager must be passed explicitly:
from cogniverse_core.registries.agent_registry import AgentRegistry
from cogniverse_foundation.config.utils import create_default_config_manager
# Create with required dependency injection
config_manager = create_default_config_manager()
registry = AgentRegistry(
tenant_id="acme",
config_manager=config_manager # Required - raises ValueError if None
)
# Register agent endpoint
from cogniverse_core.common.agent_models import AgentEndpoint
agent = AgentEndpoint(
name="search_agent",
url="http://localhost:8002",
capabilities=["video_search", "text_search"],
streams_answer_tokens=False, # True when the agent answers token by token
health_endpoint="/health",
process_endpoint="/process"
)
registry.register_agent(agent)
# Find agents by capability
search_agents = registry.find_agents_by_capability("video_search")
# Get healthy agents
healthy = registry.get_healthy_agents()
# List all registered agents
agents = registry.list_agents() # ["search_agent", ...]
# Get agent by name
agent = registry.get_agent("search_agent") # Returns AgentEndpoint or None
# Release the lazily created remote-agent connection pool at shutdown
await registry.close()
Local registry operations do not create an HTTP client. Concurrent remote health or discovery calls share one lazily constructed client, and close() detaches that client before awaiting connection-pool shutdown.
register_agent and register_agent_from_data add configured agents to this process's registry; every runtime process registers the same ones at startup. Registrations every process must serve go through a shared AgentRegistryStore passed as store= (or attached with set_store); the runtime uses RedisAgentRegistryStore:
registry = AgentRegistry(tenant_id="acme", config_manager=config_manager, store=store)
await registry.add_registration(agent) # every process
removed = await registry.remove_registration("search_agent") # False if not served
await registry.refresh() # apply the store's current registrations and removals
The served view is the configured agents overlaid with the store's registrations and removals: a registration replaces a configured agent of the same name, and removing a configured agent hides it until it is registered again. A request path calls refresh() before it reads the registry; it reads the store's version and re-reads the registrations only when that moved. Store failures raise AgentRegistryUnavailableError, and refresh() then leaves the served view as it was. Without a store a registry serves its configured agents and refuses add_registration / remove_registration.
BackendRegistry¶
from cogniverse_core.registries import BackendRegistry
# Register custom backend (implements Backend interface)
BackendRegistry.register_backend("my_backend", MyBackendClass)
# Get shared backend instance (tenant_id passed in query_dict at search time)
backend = BackendRegistry.get_search_backend(
name="my_backend",
config_manager=config_manager,
schema_loader=schema_loader
)
Resolver seam. The registry owns the lifetime of what it hands out and closes the instance on capacity eviction, on an overwriting set and on clear_instances(), so nothing may keep a handle across operations. Callers hold a zero-arg resolver instead and run each operation inside leased_backend(resolve), which resolves through the registry and holds a checkout for the block; BackendRegistry.holds(instance) reports whether the cache owns a given instance.
DSPyModuleRegistry¶
from cogniverse_core.common.dspy_module_registry import DSPyModuleRegistry
from cogniverse_foundation.config.agent_config import DSPyModuleType
import dspy
# Create module instance from a built-in module type
class QASignature(dspy.Signature):
question = dspy.InputField()
answer = dspy.OutputField()
module = DSPyModuleRegistry.create_module(
module_type=DSPyModuleType.CHAIN_OF_THOUGHT,
signature=QASignature
)
SchemaRegistry¶
Source: cogniverse_core/registries/schema_registry.py.
from cogniverse_core.registries import SchemaRegistry
# Create registry instance (requires dependencies)
registry = SchemaRegistry(
config_manager=config_manager,
backend=backend,
schema_loader=schema_loader
)
# Register deployed schema
registry.register_schema(
tenant_id="acme",
base_schema_name="video_content",
full_schema_name="video_content_acme_acme",
schema_definition=schema_definition_str,
config={"profile": "video_content"}
)
# Check if schema exists: read from its stored row on every call, so a drop
# or registration by another process is seen at once
exists = registry.schema_exists("acme", "video_content")
# Get all schemas for tenant
schemas = registry.get_tenant_schemas("acme")
tenant_deployed_schema_names(config_manager, tenant_id) is the module-level read serving-time servability uses: the base schema names the tenant has deployed, taken from the registry rows plus the pending deployment intents (a name an activation owns before its row lands). It needs no backend, and a storage read failure raises RegistryStorageError naming the tenant rather than answering with a smaller set — an outage must never read as "nothing is deployed".
DeployedSchemaNames(config_manager, refresh_after_s=DEPLOYED_SCHEMAS_REFRESH_S, max_staleness_s=DEPLOYED_SCHEMAS_MAX_STALENESS_S) answers reader(tenant_id, base_schema_name) -> bool from that read, holding each canonical tenant's deployed names in a RefreshingCache (foundation module); VespaBackend hands one to its VespaSearchBackend, which asks it before every query (see the backends module). A name in the tenant's entry answers True with no read until refresh_after_s (15 s) after the read that produced it began. Until max_staleness_s (30 s) it still answers True while one background thread (deployed-schemas-refresh) re-reads the tenant, so a search never waits on that read. Past max_staleness_s the caller reads the store itself. A name missing from the entry re-reads the store before answering False, so a deployment by any process is visible to the next call and a refusal is never served from memory. The fresh deployed names replace the entry, and a tenant with none keeps no entry. Every schema-registry row write (register_schema, unregister_schema, so deploy and delete) and every deployment-intent transition (prepare, retire, complete, recovery) calls invalidate_deployed_schema_names(tenant_id), which drops that tenant's entry from every reader in the process and detaches its read in flight. So DEPLOYED_SCHEMAS_MAX_STALENESS_S bounds only how long another process's deletion keeps answering True, and for a tenant searched continuously it is seen about DEPLOYED_SCHEMAS_REFRESH_S after it lands. Concurrent reads for one tenant share one store read; a failed read raises RegistryStorageError to each caller and caches nothing. A failed background read is logged and the entry keeps answering until max_staleness_s, after which the next call raises.
from cogniverse_core.registries.schema_registry import tenant_deployed_schema_names
deployed = tenant_deployed_schema_names(config_manager, "acme:prod")
SchemaRegistry reads persisted schemas on construction, leaves retry and backoff to its ConfigStore implementation, and wraps storage read failures once with registry context. Empty storage and store-normalized HTTP 404 results both load an empty registry.
deploy_schemas(tenant_id, base_schema_names, config=None, force=False) writes all new schema intents, deploys one application package, waits once for backend convergence, and registers every schema. It returns full names in request order; already registered schemas require no activation unless force=True. A registered schema whose stored definition differs from the definition the schema loader supplies is redeployed with the loaded one in the same package and its registry row replaced, keeping the row's config unless one is passed; a change Vespa refuses raises SchemaChangeRefusedError (a BackendDeploymentError) naming the schema and Vespa's reason, and the row keeps the definition that is live. Empty and duplicate name lists are rejected. Intents become complete only after every registration succeeds, so a partial registration leaves the whole new batch reserved for recovery. Registration runs under the backend's deployment_lease(), the lease a tenant delete reads the registry and drops the tenant's schemas under, and first re-reads the tenant's deletion marker: a tenant marked deleted after activation gets no row (TenantDeletedError), its intents stay pending and the delete drops the live schema with the tenant's others. A lease a peer holds for the whole wait raises its LeaseWaitTimeout with nothing registered. deploy_schema() delegates to this path with one name. SchemaDeploymentIntents(store) in cogniverse_core/registries/schema_deployment_intents.py reserves the full schema name under the system tenant's SCHEMA scope and schema_deployment_intents service. One key per full schema name holds the exact schema JSON, canonical tenant, base and full names, configuration, original timestamp, and pre-activation registry version. Colliding tenant names cannot reserve each other's schema. An unresolved definition stays immutable until its generation completes, including when activation was reported absent. A stale creator reuses the completed payload for the same registry version; an unresolved generation cannot be replaced merely because deletion advanced the registry version. A fresh generation requires a strictly advanced registry version and a completed prior intent; older snapshots are rejected.
SchemaDeployLease in cogniverse_core/registries/schema_deploy_lease.py (built by SchemaRegistry.deployment_lease()) serialises application-package replacement across processes and pods. The record lives under the system tenant's SCHEMA scope, schema_deploy_lease service, application key, and holds the current holder id and the hold time (DEFAULT_LEASE_SECONDS, 60 s) it was taken with; it moves only through compare_and_set_config. acquire() waits out a live holder for up to DEFAULT_WAIT_SECONDS (120 s) and then raises TimeoutError; a process takes over a holder only after watching the record's version stand still, on its own monotonic clock, for the hold time the record carries, across as many waits as that takes, and a holder treats its lease as lost once its own monotonic clock passes that hold time since its last successful claim — no timestamp is written by one node and compared on another, so clock skew cannot break mutual exclusion. renew() extends the lease and raises DeploymentLeaseLost once another holder owns it; a store failure inside renew() propagates as the store's error. release() hands the lease back and logs, rather than raises, when the store is unreachable: the package is already activated by then. A record the store did not confirm cleared is taken over at once by the process that released it, and by peers once they have watched it stand still.
A deploy holder (VespaSchemaManager.deployment_lease) renews on a background heartbeat every third of its hold for as long as it holds the lease, so a live holder is not taken over while its requests run longer than the hold; the heartbeat stops at release(). It also stops renewing once the lease has been held for MAX_TOTAL_HOLD_SECONDS (9000 s), so a holder stuck past that is taken over a hold later like a dead one. The cap covers a one-tenant, one-schema delete whose five attempts of three 310 s requests, backoff, thirteen one-page config-store visits and seven schema listings all run to their bounds (8746.25 s). Not covered: each extra visit page, extra tenant in a bulk delete or extra tombstone adds 303.75 s, and each tombstone's write adds its own time; the requests timeouts bound inactivity between bytes, not a whole call. A renewal the store refuses loses the lease, and the deploy re-checks ownership before each mutating step (session create, prepare, activate, and each backend prepare-and-activate attempt), so it stops before the next one. A step already in flight when the lease is lost cannot be recalled, because Vespa accepts no fencing token.
A record is also taken over at once when this node can prove its holder is gone. Holders are host:pidns:pid:uuid, where pidns names the PID namespace (kernel boot id and namespace inode); the probe answers "gone" only for a holder in this host's own PID namespace whose pid is no longer running, or one naming this very process that no live holder object owns any more — the state a thread that died, or a coroutine abandoned mid-deploy, leaves behind and the record itself cannot express. It never guesses about another node; a dead or partitioned holder there stops renewing, and because DEFAULT_WAIT_SECONDS (120 s) is longer than DEFAULT_LEASE_SECONDS (60 s), one wait outlasts its hold. A backend deploy (deploy_schemas), the runtime's startup metadata deploy, schema deletion and the orphan reconciler's redeploy all hold it while they enumerate the live schemas, build their package and post it, so no package is built from a snapshot another process has already moved past. Every registry read made inside the lease is strict: a failed refresh raises instead of answering from the in-memory registry, which may still hold a schema a peer deleted or lack one a peer registered since.
A package with no tenant schema carries only the metadata schemas. Before that, pyvespa added a default document type named after the application (cogniverse), which nothing registers, and every later deploy refused it as an unknown live schema. A cluster that already carries it is fixed once by dropping it as an orphan with VespaSchemaManager.delete_orphan_schemas(["cogniverse"]).
prepare(registration, grace_s=..., registry_version=0) writes the reservation conditionally; grace_s is a required keyword-only argument. reserved(live_names) (exposed as SchemaRegistry.reserved_schemas) maps each full schema name an in-flight activation owns to its exact registration: every pending intent whose schema is live, plus pending intents still inside their grace whose activation is imminent. Every process that rebuilds the application package — a schema deploy, the orphan reconciler's redeploy, tenant deletion, the runtime's startup metadata deploy — keeps these as survivors, rebuilt from the intent's definition, since the schema is live-but-unregistered for the whole convergence wait and its registration lives in another process. Vespa calls reconcile_deployment_intents(live_names) during package construction after a successful config-server enumeration. Records wait 90 seconds for normal registration. Recovery only completes schemas in that live set and supplies their definitions to the existing reconstruction path. Recovery never deploys or deletes a schema and requires no owner-death inference. Live schemas without a reconstructable definition still block deploy. Existing registered schema updates do not use name-only recovery: their live name cannot prove which definition was activated.
register_schema(..., deployment_time=..., expected_version=...) preserves the original payload and conditionally writes its canonical registry key. Identical concurrent completion is accepted; a tombstone or different newer row rejects the write. unregister_schema() persists a tombstone even for an unregistered schema, fencing pending registration writes.
complete(record) clears the active intent with revision checks and cannot replace a newer generation. An absent schema retires the intent without registering anything; the inactive record retains its definition for late activation. retire(record) applies the same rule to a deploy that failed before activation; it is conditional on the record's revision, so a record the tenant delete removed meanwhile is not written back. tenant_names(tenant_id) lists the names of a tenant's records and delete(name) removes one; the tenant delete uses them. A deploy that failed after activation (SchemaConvergenceError: the generation did not reach every service, or the new schema refused a feed, inside the budget) leaves the intent pending: the schema is live, every package built meanwhile carries it, and recovery registers it once the grace elapses. If retirement fails, the deployment error keeps its original cause and adds the retirement failure as an exception note; the durable record remains available for recovery. records() reads current journal entries; pending() lists active ones. reconcile(live_names, registered, write_registration) atomically claims at most one attempt per record per invocation. Each failed recovery raises RegistryStorageError with schema context; a fourth attempt is refused. There is no background retry loop. Journal history uses ConfigStore retention.
Schema drift migration¶
A tenant's registered definition is the schema_definition JSON in its registry row (SCHEMA scope, schema_registry service, key schema_<base>): the shipped schema file with name set to the tenant's full schema name, as it was when the tenant last deployed it. Nothing writes a tenant-specific definition; a profile's schema_config is stored with the profile and never changes the definition. A registered schema drifts when its definition, compared as parsed JSON, differs from the one the schema loader ships today with the same name; a base schema the loader does not ship has nothing to drift from. A deploy compares the stored row with the shipped definition, so whichever request first ensures a drifted schema redeploys it from the request path; the runtime's startup migration does it first.
SchemaRegistry.redeploy_drifted_schemas(should_stop=None) redeploys every drifted schema, one normal deploy per drifted tenant carrying all of that tenant's drifted schemas. Each tenant's deploy runs inside backend.deployment_lease() from the decision through the registration: which of its schemas still differ is read again from the stored rows inside the lease, so a schema another process redeployed or deleted meanwhile is not deployed, and no other deployer builds a package from the registry while the tenant's new definitions are live but not yet registered. Only the drifted schemas change: the package carries every other schema from its registry row, so other tenants' schemas and the tenant's current ones deploy as they are, a schema whose base the loader does not ship (one deployed from another schema directory) is left alone, and each re-registration keeps its row's config. A package Vespa refuses (SchemaChangeRefusedError) is deployed again one schema at a time, so the tenant's other schemas still land; a schema refused on its own is left as it is, live definition, registry row and documents alike, logged at WARNING and recorded under the system tenant (SCHEMA scope, schema_migration_refusals service, keyed by full schema name, with the SHA-256 of the definition it was refused); every later run attempts it again and records the refusal again. A run that reaches the end removes, under the deployment lease, every recorded refusal whose schema is no longer registered with a drifted definition: one that has migrated since, or was dropped. delete_tenant_refusals(config_manager, tenant_id) deletes a tenant's refusals; the tenant delete calls it once the tenant's schemas are dropped, among the rest of the tenant's rows. Neither raises: a refusal that cannot be read or deleted is logged at ERROR by tenant and schema, the run or the delete completes, and the next run removes it. A peer's deletion landing before activation drops that schema from the tenant's deploy. A tenant marked deleted (see mark_tenant_deleted), whose delete has not completed, is never redeployed: deploy_schemas refuses it with TenantDeletedError, its drifted schemas are left as they are and the other tenants still migrate. Any other error propagates, including the LeaseWaitTimeout (a TimeoutError) of a lease a peer held for the whole wait, and nothing is recorded against a tenant for it. should_stop is asked before each tenant's redeploy; once it answers True no further redeploy starts, and no refusal is removed. It returns a DriftedSchemaRedeploy: redeployed holds the full names it deployed, refused one SchemaRefusal(tenant_id, base_schema_name, schema_name, error, refused_at) per refused schema, skipped the drifted schemas a stop left, and deleted those of tenants marked deleted.
drifted_schemas(config_manager, schema_loader) lists every drifted schema as a DriftedSchema(tenant_id, base_schema_name, schema_name, refusal), ordered by tenant and schema name. refusal is the recorded refusal when it was recorded for the definition shipped now, else None: the migration has not reached the schema yet. A registry or refusal read failure raises.
AdapterStoreRegistry / WorkflowStoreRegistry¶
Thin subclasses of cogniverse_foundation.registry.EntryPointRegistry that auto-discover implementations registered via Python entry points, so cogniverse_core never imports a concrete backend package directly.
from cogniverse_core.registries import AdapterStoreRegistry, WorkflowStoreRegistry
# Implementations register under these entry-point groups in pyproject.toml:
# [project.entry-points."cogniverse.adapter.stores"]
# vespa = "cogniverse_vespa.registry.adapter_store:VespaAdapterStore"
#
# [project.entry-points."cogniverse.workflow.stores"]
# telemetry = "cogniverse_agents.workflow.telemetry_workflow_store:TelemetryWorkflowStore"
adapter_store = AdapterStoreRegistry.get(
name="vespa", config={"backend_url": "http://localhost", "backend_port": 8080}
)
workflow_store = WorkflowStoreRegistry.get(
name="telemetry", config={"telemetry_provider": telemetry_provider}
)
Query Encoding¶
query/encoders.py provides the QueryEncoder abstraction search agents use to turn a text query into the embedding shape a Vespa profile expects, plus a caching factory that picks the right encoder from configs/config.json.
from cogniverse_core.query.encoders import QueryEncoderFactory
# Cached by (model_name, inference_service, resolved service URL,
# embedding_dim) — a repeat call for the same profile reuses the
# already-loaded encoder, and two configs pointing the same service name at
# different endpoints get one encoder each.
encoder = QueryEncoderFactory.create_encoder(
profile="frame_based_colpali",
config=system_config, # SystemConfig instance, required
)
embedding = encoder.encode("find manufacturing defects")
dim = encoder.get_embedding_dim()
QueryEncoderFactory.get_supported_profiles(config=system_config)
| Class | Model family | Notes |
|---|---|---|
ColBERTQueryEncoder | ColBERT / LateOn | Per-token multi-vector; embedding_dim is required (read from schema_config.embedding_dim); supports joint query+CoT-trace encoding |
ColPaliFamilyQueryEncoder | ColPali, ColQwen, ColSmol | 320-d patch multi-vector; local or remote (inference_service_url) |
ColPaliQueryEncoder(...) / ColQwenQueryEncoder(...) | — | Thin factory functions over ColPaliFamilyQueryEncoder with model_loader="colpali"/"colqwen" |
XClipQueryEncoder | X-CLIP | Single-vector 768-d; remote via the video_embed sidecar; encodes video and text into one space |
ClapTextQueryEncoder | CLAP | Single-vector 512-d; remote via the clap_embed sidecar's /embed/text; text in the space of the stored acoustic embeddings |
QueryEncoderFactory._create_encoder_instance resolves the encoder in this order: profile_config["model_loader"] (authoritative) → model-name substring → profile-name substring, raising ValueError if none match. When a profile carries a second, acoustic embedding (audio_clap_semantic), semantic_model names its ColBERT transcript model, embedding_model names the CLAP model, and inference_services carries embedding: colbert_pylate, transcription: vllm_asr, and acoustic_embedding: clap_embed. Missing URLs for the declared services raise ValueError naming the profile and service rather than silently falling back to a local load, so a misconfigured sidecar fails loud. The profile's query encoder is the ColBERT one; the acoustic input takes QueryEncoderFactory.create_service_encoder(profile, "clap_embed", config), the cached text encoder of the service (SERVICE_TEXT_ENCODERS). A service with no text encoder, or with no configured URL, raises EncoderNotConfiguredError. Search agents (document/image) resolve their encoders through this factory, passing the merged config, so they route through the deployed sidecar exactly as the /search path does.
encoders.py also defines the two encoder fault types callers distinguish:
EncoderNotConfiguredError(aValueError) — the profile declares no usable encoder: no model name, an inference service with no configured URL, or a missingschema_config.embedding_dim.EncoderUnavailableError(aRuntimeError) — the encoder is configured but its inference service did not serve the request. Carriesprofile,service(the name the profile configures) andendpoint(the resolved sidecar URL), and chains the underlying failure.
VespaSearchBackend.search raises these when it builds or calls the encoder a profile declares, so a sidecar outage is never reported as a missing query_embeddings argument. EncoderNotConfiguredError.profile names the profile when the raiser knows it.
encoder_outage_errors()— the exception types that mean a configured encoder's service did not serve (an inference-service outage, a tripped breaker, a transport error).build_query_encoder(profile, *, config, model_name=None)— builds the profile's encoder throughQueryEncoderFactoryand raises the typed faults: a missing or incomplete setting (including a service with no URL and no in-process backend) asEncoderNotConfiguredError, an outage asEncoderUnavailableError.SharedQueryEncoder(profile, build, *, service=None)— what a text search hands the backend asquery_encoderinstead of embeddings.buildruns at most once, on the firstencode; each query text is encoded once and the embeddings, or the typed fault, are shared by every caller. An outage isEncoderUnavailableError(namingservice), a rejected query (ValueError)EncoderNotConfiguredError.
Every remote query encoder bounds its POST with model_loaders.QUERY_ENCODE_TIMEOUT_S (30s): the ColPali family through RemoteInferenceClient.query_encode_timeout_s, ColBERT on the is_query direction of /pooling, X-CLIP through embed_text, and DenseOn through RemoteOpenAIEmbedder.encode(is_query=True). The document directions keep their own budgets, DOCUMENT_ENCODE_TIMEOUT_S (120s) for a text batch and SEGMENT_EMBED_TIMEOUT_S (600s) for a video segment.
Event System¶
events/ defines the Pydantic event vocabulary, the queue protocols, and the per-request binding producers report on. The runtime's queues are Redis-backed (cogniverse_runtime.task_events, see Events Module); the in-process ones here serve a library caller outside the runtime.
from cogniverse_core.events import (
InMemoryEventQueue,
TaskCancelled,
bind_event_queue,
publish_phase,
)
queue = InMemoryEventQueue(task_id="task-123", tenant_id="acme:acme")
with bind_event_queue(queue):
await publish_phase("planning", "Creating execution plan...")
queue.cancel("operator stop")
try:
await publish_phase("execution", "Executing...")
except TaskCancelled as stopped:
print(stopped.reason) # "operator stop"
| Component | Purpose |
|---|---|
EventType, TaskState | Enums for event kind and A2A task lifecycle state |
StatusEvent, ProgressEvent, ArtifactEvent, ErrorEvent, CompleteEvent | BaseEvent subclasses emitted during processing |
EventQueue / QueueManager (Protocols) | Per-task event queue and queue-lifecycle contracts |
BaseEventQueue / BaseQueueManager | ABCs implementing the shared queue bookkeeping; a queue records the event loop it was created on (loop) so a producer on a worker thread publishes through it |
InMemoryEventQueue / InMemoryQueueManager (events/backends/memory.py) | In-process implementation; get_queue_manager() / reset_queue_manager() manage the process-wide singleton |
CancellationToken (events/queue.py) | Cooperative cancellation signal threaded through a running task |
bind_event_queue(queue) / current_event_queue() | Bind the queue a request reports to (a ContextVar, inherited by the tasks and threads the request starts), never held on a shared agent |
publish_phase(phase, message, check_cancelled=True) | Publish a working StatusEvent on the bound queue (no-op without one), then raise TaskCancelled when the task was cancelled |
raise_if_cancelled(queue=None) / TaskCancelled | Stop a producer at a boundary once its task was cancelled; carries task_id and reason |
AgentBase.report_phase(phase, message) is emit_progress (the streaming caller's progress dict) plus publish_phase.
Durable Execution¶
durable/ provides checkpoint + resume for long-running workflows — the DSPy optimization / auto-eval jobs (optimization_cli modes, job_executor) that run as single-container Argo pods and today lose all progress when the pod is killed. A checkpoint captures the stage pipeline's progress so a restarted run continues from the last completed stage instead of re-running expensive compile()s.
This is distinct from WorkflowStore (the offline workflow-intelligence learning corpus).
from cogniverse_core.durable import (
PipelineCheckpoint,
PipelineCheckpointStatus,
PipelineCheckpointStorage,
)
storage = PipelineCheckpointStorage(
grpc_endpoint="localhost:4317",
http_endpoint="http://localhost:6006",
tenant_id="acme:acme",
)
await storage.save_checkpoint(checkpoint) # persisted as a telemetry span
latest = await storage.get_latest_checkpoint(workflow_id) # None if unknown; raises on outage
for phase in latest.pending_phases(): # stages still to run
...
await storage.mark_status(checkpoint_id, PipelineCheckpointStatus.COMPLETED)
| Component | Purpose |
|---|---|
PipelineCheckpoint | Resumable stage-pipeline state: phases, phase_index, completed_units, cursor, metadata, resume_count. pending_phases() / completed_unit_keys() drive resume (a failed unit re-runs; a completed one is skipped). |
PipelineCheckpointStatus | Run lifecycle: ACTIVE, COMPLETED, FAILED. |
PipelineCheckpointConfig | Telemetry project name + retention windows. |
PipelineCheckpointStorage | Persists checkpoints as pipeline_checkpoint telemetry spans; get_latest_checkpoint finds the most recent by workflow_id (status transitions recorded as span annotations). Reads raise on backend outage rather than reading as "no checkpoint". |
Approval Interfaces¶
approval/interfaces.py defines domain-agnostic human-in-the-loop review contracts. They live in cogniverse_core (rather than in one implementation package) specifically so both cogniverse_agents and cogniverse_synthetic can depend on the same interfaces without depending on each other.
approval/training_schema.py is the shared validator for signed synthetic examples consumed by optimization and finetuning. It accepts only the canonical agent types, rejects surrounding whitespace rather than normalizing labels, and checks semantic supervision: entity relationships reference exact entity text, profile selections belong to the offered profile set with known modality and complexity labels, and query enhancements contain the complete observed output.
from cogniverse_core.approval import (
ApprovalBatch,
ApprovalStatus,
ApprovalStorage,
ConfidenceExtractor,
FeedbackHandler,
ReviewDecision,
ReviewItem,
)
item = ReviewItem(item_id="q1", data={"query": "..."}, confidence=0.42)
batch = ApprovalBatch(batch_id="b1", items=[item])
batch.pending_review # pending and regenerated items awaiting human review
batch.approval_rate # (auto_approved + approved) / total
| Type | Role |
|---|---|
ApprovalStatus | auto_approved / pending_review / approved / rejected / regenerated; a regenerated replacement remains in ApprovalBatch.pending_review until a reviewer resolves it |
ReviewItem / ReviewDecision / ApprovalBatch | Dataclasses carrying domain-specific data + review outcome |
ConfidenceExtractor (ABC) | extract(data) -> float — domain-specific confidence scoring |
FeedbackHandler (ABC) | process_rejection(item, decision) -> Optional[ReviewItem] — regenerate on rejection |
ApprovalStorage (ABC) | save_batch / get_batch / update_item / get_pending_batches — persistence backend |
Media Access¶
common/media/ is the URI → local-Path resolver used by both ingestion (writing source_url) and evaluation (reading frames for the visual judge). It is tenant-scoped and fsspec-backed for network schemes.
from cogniverse_core.common.media import MediaConfig, MediaLocator
config = MediaConfig.from_dict({
"default_uri_scheme": "file",
"cache": {"max_bytes_gb": 50, "ttl_days": 7},
"backends": {"s3": {"endpoint_url": "http://minio:9000"}},
})
locator = MediaLocator(tenant_id="acme", config=config)
path = locator.localize("s3://bucket/videos/clip.mp4") # downloads + caches, returns local Path
locator.exists("pvc://videos/clip.mp4")
stat = locator.stat("https://example.com/clip.mp4") # MediaStat(size, etag, ...)
| Scheme | Behavior |
|---|---|
file:// / bare path | Short-circuits to the local path, no copy |
pvc://<volume>/<rest> | Resolves to <config.pvc_mount_root>/<volume>/<rest> |
s3://, http(s)://, gs://, az:// | Fetched via fsspec and cached locally through MediaCache |
MediaCache (common/media/cache.py) is a content-addressed, tenant-scoped local cache with atomic staged writes (os.replace) and LRU eviction by atime once max_bytes is exceeded. A running byte total keeps under-budget puts walk-free; the tree is walked only on the first put, when over budget, or when a TTL sweep is due (at most once per TTL period, so an expired entry lingers at most one extra period).
Memory Management¶
The memory system uses Mem0 for persistent, tenant-isolated agent memory:
from cogniverse_core.memory.manager import Mem0MemoryManager
# Get memory manager (singleton per tenant via __new__)
memory = Mem0MemoryManager(tenant_id="acme")
# The manager keeps no backend handle: `_resolve_backend` re-resolves through
# BackendRegistry per operation, leased for the operation's duration. The same
# resolver goes into Mem0's vector_store config as `backend_resolver`, so
# BackendVectorStore resolves and leases per insert/search/get/update/list.
# Initialize with required parameters
from cogniverse_core.memory.schema import build_default_registry
memory.initialize(
backend_host="localhost",
backend_port=8080,
llm_model="openai/google/gemma-4-e4b-it",
embedding_model="lightonai/DenseOn",
llm_base_url="http://localhost:11434",
embedder_base_url="http://localhost:8000", # Required: DenseOn /v1 endpoint
config_manager=config_manager,
schema_loader=schema_loader,
backend_config_port=19071, # Optional, defaults to 19071
base_schema_name="agent_memories", # Optional
auto_create_schema=True, # Optional
# Optional but load-bearing: when set, add_memory enforces schema
# provenance + auto-attaches initial trust; get_relevant_context
# applies trust ranking and per-schema contradiction reconciliation.
knowledge_registry=build_default_registry(),
)
# Add memories
memory.add_memory(
content="RAG is Retrieval-Augmented Generation...",
tenant_id="acme",
agent_name="search_agent",
metadata={"topic": "ml_concepts"}
)
# Search memories
results = memory.search_memory(
query="retrieval augmented generation",
tenant_id="acme",
agent_name="search_agent",
top_k=5
)
# Get all memories for agent
all_memories = memory.get_all_memories(
tenant_id="acme",
agent_name="search_agent"
)
# Delete memory
memory.delete_memory(
memory_id="mem_xyz",
tenant_id="acme",
agent_name="search_agent"
)
# Count live and archived memories over the whole partition; a backend
# failure raises rather than counting zero
memory.get_memory_stats(tenant_id="acme", agent_name="search_agent")
# {"total_memories": 12, "archived_memories": 2, "enabled": True,
# "tenant_id": "acme", "agent_name": "search_agent"}
Mem0MemoryManager.tenant_partition_schema_exists(tenant_id) reports whether the tenant partition schema is deployed. get_all_memories() uses it to return an empty result for schema-less tenants before any store read runs.
DenseOn embedder adapter¶
The manager configures Mem0's embedder with the cogniverse_denseon provider (memory/mem0_embedder.py), not Mem0's stock openai provider. DenseOnMem0Embedder wraps RemoteOpenAIEmbedder so the DenseOn sentence-transformers prompt prefix (query: / document:) and normalize_embeddings=True that the model ships with are restored client-side over /v1/embeddings. Mem0's memory_action selects the prompt: search embeds a query, add/update (and any unspecified action) embed a document. Without this adapter, stored memory vectors would silently change value and drift Mem0's closeness ranking against existing agent_memories rows.
Federation: org trunk + tenant overlays¶
FederationService lets multiple tenants under the same org share a trunk of knowledge while overlaying tenant-specific facts.
from cogniverse_core.memory.federation import FederationService
from cogniverse_core.memory.schema import Pinnable, build_default_registry
svc = FederationService(
memory_manager_factory=Mem0MemoryManager, # per-tenant singleton
registry=build_default_registry(),
)
# Read: tenant rows + org-trunk rows, deduped by metadata.subject_key
# (tenant overlay wins on collision); each row tagged with
# _federation_origin = "tenant" | "org_trunk".
rows = svc.federated_get_all(tenant_id="acme:production", agent_name="search_agent")
# Promote: copy a tenant memory into the org trunk so siblings see it.
result = svc.promote_to_org_trunk(
source_tenant_id="acme:production",
source_memory=memory,
actor_role=Pinnable.TENANT_ADMIN,
actor_id="admin_alpha",
)
Storage: the org trunk lives under a dedicated tenant_id {org}:_org_trunk (Mem0+Vespa already isolate per-tenant_id, so no new backend wiring required). Promotion stamps promoted_from_tenant, promoted_by, and promoted_by_role into the new record's metadata for history.
| Schema sensitivity | Promotion outcome |
|---|---|
tenant_private | Refused (FederationDeniedError) — even org admins cannot escalate. |
org_shared | Promotable when actor's role is at-or-above pinnable_by floor. |
global_shared | Reserved for a future cross-org channel. |
ACLs are enforced at query time: federated_get_all only ever reads from the caller's tenant + that tenant's org trunk, so cross-org leakage is structurally prevented.
Contradiction detection + reconciliation¶
Two memories about the same subject can disagree. The ContradictionDetector groups candidate memories by metadata.subject_key (set by the writing agent) and emits a ConflictSet per subject_key that has more than one distinct content signature.
from cogniverse_core.memory.contradiction import (
ContradictionDetector,
reconcile,
)
from cogniverse_core.memory.schema import ContradictionPolicy
detector = ContradictionDetector()
conflicts = detector.detect(memories)
# Each ConflictSet carries subject_key + conflicting_memory_ids.
# Reconcile at retrieval time per the schema's contradiction_policy:
visible = reconcile(memories, ContradictionPolicy.TRUST_RANKED)
| Policy | Effect on conflicting members |
|---|---|
latest_wins | Keep the highest created_at member only |
trust_ranked | Keep the highest trust_score × confidence member |
preserve_both | Keep all members and tag each with metadata['disputed'] = True |
Memories without a subject_key pass through unchanged — the detector has no way to know what they are claims about. Conflict sets persist under sentinel agent_name="_conflict_store" (matching the pinning service's sentinel-agent pattern) so they don't pollute normal-agent search results; a future ContradictionReconciliationAgent will consume them.
Trust / source ranking¶
Each memory carries a TrustRecord derived at write time from the schema's default_trust and the provenance's derivation_kind. Trust ages slowly (≈0.5 pt/day above the initial baseline), is bumped by explicit user/admin endorsements, and is composed with relevance and confidence at retrieval time.
from cogniverse_core.memory.trust import (
TrustRecord,
apply_endorsement,
attach_trust_to_metadata,
compute_initial_trust,
rank_with_trust,
)
trust = compute_initial_trust(schema, provenance=prov)
metadata = attach_trust_to_metadata(metadata, trust)
# After a tenant admin endorses the memory:
new_trust = apply_endorsement(trust, "tenant_admin") # +0.10
# Persist by re-writing the memory with attach_trust_to_metadata(meta, new_trust).
# Retrieval ranking: relevance × trust × confidence
ranked = rank_with_trust(search_results)
| Derivation kind | Trust weight |
|---|---|
direct_ingest | × 1.20 |
user_assert | × 1.10 |
extraction | × 1.00 |
summarization | × 0.90 |
synthesis | × 0.85 |
agent_inference | × 0.70 |
| Endorser | Δ trust |
|---|---|
user | +0.05 |
tenant_admin | +0.10 |
org_admin | +0.20 |
Decay floor is initial_score — a memory never loses more trust than it had originally gained above the schema baseline. Endorsements raise the effective ceiling without changing the floor; federation and contradiction reconciliation both plug into this scoring loop.
Provenance + citation graph¶
Every memory write carries a Provenance record describing where the content came from. The schema's provenance_required flag gates writes that omit it. ProvenanceStore(backend_resolver=..., tenant_id=...) takes the resolver, not a backend instance, and leases per read/write.
from cogniverse_core.memory.provenance import (
CitationRef,
DerivationKind,
ProvenanceWalker,
attach_to_metadata,
make_provenance,
)
prov = make_provenance(
written_by="agent:search_agent",
derivation_kind=DerivationKind.SYNTHESIS,
confidence=0.85,
derived_from=[
CitationRef.external("https://wiki/source-a"),
CitationRef.memory("m_prior_answer"),
],
trace_id=current_trace_id, # optional Phoenix trace id
)
memory_manager.add_memory(
content="Synthesised claim citing two sources",
tenant_id="acme",
agent_name="search_agent",
metadata=attach_to_metadata({"kind": "synthesis_fact"}, prov),
infer=False,
)
# Read path: walk the citation chain back to primary sources.
graph = ProvenanceWalker(memory_manager).walk(memory_id, tenant_id="acme")
for node in graph.nodes:
print(node.depth, node.memory_id, node.content_excerpt)
for src in graph.primary_sources:
print("source:", src.ref_kind, src.ref_id)
# Explicit repair reads the primary's canonical metadata and upserts the
# stable prov-<tenant>-<memory> row. It raises on a concurrent primary change.
row_id = memory_manager.repair_provenance(memory_id)
Storage: provenance is attached in-band in metadata["provenance"] on the memory record and indexed as one row per memory in the per-tenant provenance Vespa schema (ProvenanceStore). Walker traversal batches each BFS level into one indexed query; cycle and depth limits (max_depth, max_nodes) protect against runaway chains.
The primary and indexed records form one verified persistence contract. ProvenanceStore.attach() requires an exact one-document feed result and raises ProvenanceWriteError for rejected or unresolved writes. A newly created Mem0 ADD is removed when its indexed write fails, together with any indexed row that write left behind. An existing UPDATE is retained and the error identifies its memory id for explicit repair; it is never deleted as compensation.
ProvenanceWalker compares the in-band declaration with the indexed row before returning a graph. Missing, malformed, mismatched, or orphaned indexed state raises ProvenanceConsistencyError; a storage outage raises rather than becoming absence. A memory reference is a normal terminal leaf only when both the primary and indexed row are absent, or when a present primary legitimately declares no provenance. Citation reads never repair state as a side effect.
Mem0MemoryManager.repair_provenance() is the explicit idempotent repair path. Supported manager writes and repair take a per-canonical-tenant lease backed by the configuration store's conditional writes, so separate processes and manager instances share one ownership boundary. Repair reattaches the primary's canonical provenance under the stable row id and verifies the primary and index before releasing that lease. Its success linearizes at the verified primary read while lease ownership excludes supported writers.
The lease covers only operations that have a primary and an indexed row to keep consistent. add_memory resolves the requested provenance from the caller's metadata before taking anything, and skips the store lease when there is none — so a conversation turn and an agent remember, which carry no provenance, never take a cluster-wide per-tenant mutex and never hold one across mem0's extraction pass. They skip the in-process lock in front of the lease too, so they run concurrently with each other. The lease a thread holds is recorded per thread, so only that thread's writes are fenced on it. A leased write ensures the tenant's memory and provenance schemas can be fed before it takes the lock and lease: that first feed in a process can deploy a missing or drifted schema, and inside the lease the deploy would outlast a hold sized for a memory write. update_memory always takes the lease: whether the stored primary declares provenance can change between any unleased read of it and the write. An update that leaves the primary without provenance also deletes its indexed row. Deletion, repair, archiving and admin restore always take it; restore reads the primary inside the lease, so it cannot revert a leased update or archive with a stale copy, and a retention archive re-reads the metadata inside the lease and writes only the metadata. The search path's last_accessed bump runs outside the lease from the search's own snapshot, so it skips provenance-bearing hits and writes only the metadata of the rest, never their text; a metadata change landing between the search and the bump (an archive) can still be overwritten. The sweeps (clear_agent_memory, cleanup_with_schema, drop_session) acquire per row: one lease for a whole retention pass excluded every other writer for the tenant until the last of a 205-row clear was gone.
Its sizing is a memory write's, not a package activation's: PROVENANCE_LEASE_SECONDS (60 s) covers a primary write, its read-back and the indexed feed with margin for an extraction pass, and PROVENANCE_WAIT_SECONDS (75 s) is deliberately longer than the hold. That inequality is the recoverability contract: a holder on another node leaves no liveness proof this node can read, so waiting is the only way to outlast a record it stopped renewing, and a wait shorter than the hold makes a leaked record permanently unrecoverable. A holder that died on this host, or one naming this very process with no live holder object behind it, is taken over at once by the lease's own liveness probe without waiting at all.
A lease has a finite hold time, so holding one across a whole operation is not exclusivity by itself: a boundary call that stalls past expiry lets a peer take over while the original holder is still inside. Ownership is therefore revalidated before every primary/index mutation and again before a successful return, and renewed once the lease has aged past half its hold time so a long operation stays owned. A holder whose own monotonic clock has passed the hold time is fenced with DeploymentLeaseLost: it neither writes nor acknowledges success after a peer becomes entitled to the lease.
Supported hard deletion enters the same ownership scope, so a delete cannot land inside repair's verification window and leave it reporting success for a primary that is already gone. delete_memory, namespace clearing, schema retention cleanup and session drop all remove the primary and its stable indexed row together; a row left behind would be an orphan that every later citation read rejects. A failure at either boundary raises instead of reporting the memory fully deleted.
ProvenanceStore.delete does not deploy a schema in order to delete from it. When the tenant's provenance schema has never been deployed, there can be no indexed row for it, so delete returns idempotently. The delete itself goes through VespaBackend.delete_live_document, a Document v1 delete that never builds an ingestion client — a client's cache miss redeploys a missing, tombstoned or drifted schema, which from the request path is a full application-package redeploy, and a namespace clear that fans out one delete per memory turns that into a redeploy race between tenants. Vespa answers a delete from a document type it does not have (never deployed, or removed by a peer) with success, so the row's absence stays idempotent; while a removal is still reaching the content nodes it answers "Unknown document type" for that type, which reads and deletes also take as absence. A schema-registry lookup failure still raises rather than being read as "no schema": only a clean, successful "not deployed" answer short-circuits the delete. The known cost is that a tenant whose provenance schema exists but is briefly unreadable at the moment of the check also skips its indexed row, leaving an orphan for that memory rather than deleting it.
update_memory returns False for an update that did not happen. Once the primary has been rewritten it can no longer say that truthfully, so any failure after that point propagates instead of collapsing into False — a DeploymentLeaseLost as itself, anything else as a ProvenanceWriteError with the memory id and the failure as its cause: the content is changed and the index disagrees with it, which the caller has to see to retry or repair.
Each indexed row also stores a SHA-256 digest of the primary content and canonical provenance. A raw writer outside the manager/lease contract may mutate after repair's linearization point; the next citation read detects the digest mismatch and raises rather than serving the row as consistent.
Upgrading from before the digest field. Rows indexed before this field existed read back with an empty primary_digest, which can never equal a SHA-256. Those rows are treated as legacy: the digest comparison is skipped and every other consistency check still applies, so a pre-upgrade memory stays readable instead of reporting torn provenance for data that is intact. No migration or bulk sweep is required — the next attach or an explicit repair_provenance() writes a real digest under the row's stable id, after which the row is checked like any other. primary_provenance_digest always returns 64 hex characters, so an empty digest is unambiguously a legacy row and never a real mismatch. The provenance schema gained the primary_digest field in the same change; a tenant's provenance schema registered before it is redeployed by the schema drift migration. Concurrent external changes inside repair are retried up to the requested bound and then raise ProvenanceRepairConflictError.
Contradiction detection and trust ranking both read this provenance graph to score conflicting claims.
Pinning service¶
PinService lets users, tenant admins, and org admins pin memories so they survive lifecycle cleanup, trust decay, and any future curator pass.
from cogniverse_core.memory.pinning import PinQuotas, PinService
from cogniverse_core.memory.schema import Pinnable, build_default_registry
registry = build_default_registry()
quotas = PinQuotas.from_tenant_config(tenant_config) # honours admin overrides
service = PinService(memory_manager, registry, quotas=quotas)
# Pin a tenant_instruction memory as tenant admin.
service.pin(
target_memory_id=memory_id,
target_kind="tenant_instruction",
pinned_by=Pinnable.TENANT_ADMIN,
actor_id="admin_alpha",
tenant_id="acme",
)
# Inspect.
service.is_pinned(memory_id, "acme") # bool
service.list_pins("acme") # all PinRecord across roles
service.quota_used(Pinnable.USER, "acme") # int
# Unpin (authority rules: user can only unpin their own; tenant admin can
# unpin user/tenant pins; org admin can unpin anything).
service.unpin(
target_memory_id=memory_id,
requester=Pinnable.TENANT_ADMIN,
actor_id="admin_alpha",
tenant_id="acme",
)
Quota defaults — user 50, tenant_admin 500, org_admin unlimited. Override per-tenant via TenantConfig.metadata["pin_quota"] = {"user": N, ...}. An org admin can also override an existing pin from a lower role (the previous pin record is dropped before the new one persists).
Pin records live under a sentinel agent_name="_pinning" so they never pollute normal-agent search results. The schema registry's validate_pin_authority(role) enforces the per-kind floor before any write hits Vespa.
Note on the
actor_idkwarg. The public API takesactor_id, but the persisted metadata key ispin_actor_id. Mem0 treatsactor_idas a promoted-payload key (it's lifted out ofmetadataand into the top-level payload on insert), which would erase the pin's persisted history on round-trip. Other writers storing per-actor identifiers in memory metadata should avoid the bare keyactor_idfor the same reason — pick a namespaced key (pin_actor_id,endorsement_actor_id, etc.).
Knowledge schema registry¶
Each kind of memory carries a KnowledgeSchema describing its retention, sensitivity, pin authority, provenance requirement, contradiction policy, and default trust. The registry is the single source of truth the pinning service, the lifecycle scheduler, and provenance handling all read from.
from cogniverse_core.memory.schema import (
KnowledgeRegistry,
KnowledgeSchema,
Pinnable,
Retention,
SchemaViolationError,
build_default_registry,
)
registry = build_default_registry() # seeded with conversation_turn,
# learned_strategy, tenant_instruction,
# external_doc, entity_fact, kg_node, kg_edge
schema = registry.get("entity_fact")
# Defaults are conservative when the kind is unregistered:
# permanent + tenant_private + provenance_required.
# Validation gates the write before any Vespa I/O:
schema.validate_write(provenance=my_provenance, pinned_by=Pinnable.USER)
# Raises SchemaViolationError when:
# * provenance is required but missing or has empty derived_from
# * pin requester's role is below the schema's pinnable_by floor
Register custom kinds at boot — replace=True is required to overwrite an existing kind so accidental redefinition fails loudly.
Strategy decay¶
StrategyLearner tracks confirmation_count + last_confirmed_at on each strategy. When the dedup search finds a near-duplicate (Jaccard overlap above 0.9):
- The existing record is deleted and a fresh copy is added with
confirmation_countbumped by one andcreated_at/last_confirmed_atreset to now. - The bumped record carries an accumulated
trace_countso retrieval reflects the total weight of evidence behind the strategy.
At retrieval, StrategyLearner.rank_strategies_with_decay downweights strategies with confirmation_count < 3 AND age > 14 days by a factor of 0.5 — they sink in the result list instead of competing equally with high-confirmation strategies.
The schema cleanup hook (_retire_unconfirmed_strategy, registered for learned_strategy in build_default_registry) deletes records that stay under 3 confirmations for more than 30 days. Pinned strategies are filtered out by LifecycleScheduler.pin_lookup before the hook runs, so admin-promoted strategies are immune to retirement.
Schema-driven lifecycle¶
The runtime ticks the lifecycle in schema-driven mode only — every memory's metadata["kind"] is looked up in the KnowledgeRegistry and that schema's retention policy decides whether to delete it. There is no bulk-age fallback; a memory without a registered kind falls back to the registry's safe default (permanent) and is never auto-deleted. Numeric created_at values accept signed epoch seconds, milliseconds, microseconds, or nanoseconds; normalization preserves dates before 1970.
Retention | Behaviour |
|---|---|
PERMANENT | Never auto-deleted. |
EPHEMERAL_SESSION | Event-driven: cleared by Mem0MemoryManager.drop_session(session_id, registry). Two HTTP endpoints reach it: DELETE /admin/tenants/{tenant_id}/sessions/{session_id} (single tenant) and POST /admin/sessions/{session_id}/close (fan-out across every warm tenant of every runtime worker process — the gateway's logout / disconnect hook). Writes MUST carry metadata["session_id"] (promoted to a fast-search Vespa field at insert time, so drop_session filters the store server-side instead of scanning the tenant's memories) or the schema rejects them, AND the kind's pinnable_by must be Pinnable.NOBODY (the schema constructor refuses any other value, since pinning a session memory and then losing it on session close would be a foot-gun). The default registry ships kind="session_scratch" for this. |
EPHEMERAL_DAYS(N) | Soft-deleted (metadata.archived=true) when created_at is older than N days; hard-deleted at 2N days. Restoreable via the admin restore endpoint inside the soft-delete window. |
SCHEMA_DRIVEN | Defers to the schema's cleanup_hook callable. |
Pinned memory ids (from PinService.list_pins) are filtered out before any deletion attempt by the periodic scheduler — pins always win there. drop_session is the user's explicit "end my session" signal and is not pin-gated; the EPHEMERAL_SESSION + Pinnable.NOBODY schema invariant makes the question moot — there can never be a pinned session memory.
The dispatcher auto-stamps metadata["session_id"] on every memory-aware agent's writes for the duration of one request via AgentDispatcher._scoped_session(agent, session_id), which calls MemoryAwareMixin.set_session_id(...) on entry and clears it on exit. Agent code never has to thread session_id through manually; a caller-supplied session_id in metadata always wins over the dispatcher stamp.
The tick summary returned by LifecycleScheduler.tick_once() reports {"tenants": {tenant_id: {kind: deleted_count} | "schema absent" | "error: <ExceptionName>"}, "total_deleted": int}. Soft-delete events appear under the {kind}:archived key; hard-deletes appear under {kind} directly. Schema-less tenants are skipped after a schema-exists probe and logged at INFO.
Scheduled lifecycle cleanup¶
Memories accumulate; the runtime ships a LifecycleScheduler that runs schema-driven cleanup on each warm tenant on a periodic tick. It is started in the FastAPI lifespan and stopped on shutdown.
| Env var | Default | Effect |
|---|---|---|
COGNIVERSE_MEMORY_LIFECYCLE_DISABLED | unset | When 1/true/yes, the scheduler does not start. |
COGNIVERSE_MEMORY_LIFECYCLE_INTERVAL | 3600 | Seconds between ticks. |
Per-tenant errors during a tick are recorded in the run summary but do not abort the run; offending tenants are visible via the scheduler's last_run_summary for operator inspection.
Configuration¶
Configuration management for agents and system settings:
from cogniverse_foundation.config.utils import create_default_config_manager
from cogniverse_core.common.tenant_utils import parse_tenant_id, get_tenant_storage_path
# Create config manager (reads from configs/config.json and environment)
config_manager = create_default_config_manager()
# Get global system configuration (no tenant_id argument)
system_config = config_manager.get_system_config()
# Tenant utility functions
org_id, tenant_name = parse_tenant_id("acme:production") # Returns (org_id, tenant_name)
storage_path = get_tenant_storage_path(base_dir="data", tenant_id="acme")
# Apply overrides
tenant_config.update({
"max_concurrent_requests": 100,
"embedding_model": "colpali-v2"
})
Usage Examples¶
Creating a Complete Agent¶
from cogniverse_core.agents.base import AgentBase, AgentInput, AgentOutput, AgentDeps
from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from typing import List, Optional
# 1. Define types
class SummarizerInput(AgentInput):
text: str
max_length: int = 100
style: str = "concise"
class SummarizerOutput(AgentOutput):
summary: str
word_count: int
key_points: List[str]
class SummarizerDeps(AgentDeps):
model_name: str = "gpt-4"
temperature: float = 0.7
# 2. Implement agent
class SummarizerAgent(A2AAgent[SummarizerInput, SummarizerOutput, SummarizerDeps]):
async def _process_impl(self, input: SummarizerInput) -> SummarizerOutput:
# Use DSPy module if available
if self.dspy_module:
result = self.dspy_module(
text=input.text,
max_length=input.max_length,
style=input.style
)
return SummarizerOutput(
summary=result.summary,
word_count=len(result.summary.split()),
key_points=result.key_points
)
# Fallback logic
summary = input.text[:input.max_length] + "..."
return SummarizerOutput(
summary=summary,
word_count=len(summary.split()),
key_points=[]
)
# 3. Run agent
if __name__ == "__main__":
deps = SummarizerDeps() # No tenant_id — it's per-request
config = A2AAgentConfig(
agent_name="summarizer_agent",
agent_description="Summarizes text content",
capabilities=["text_summarization"],
port=8002
)
agent = SummarizerAgent(deps=deps, config=config)
# The A2A HTTP server is managed by the runtime's A2AStarletteApplication.
# Register agents with the runtime and use the /a2a endpoint.
Calling Between Agents¶
Inter-agent calls are made via httpx.AsyncClient (not via self.call_agent()). The orchestrator dispatches to agents through AgentDispatcher and CogniverseAgentExecutor wired in the runtime:
import httpx
class OrchestratorAgent(A2AAgent[OrchestratorInput, OrchestratorOutput, OrchestratorDeps]):
async def _process_impl(self, input: OrchestratorInput) -> OrchestratorOutput:
async with httpx.AsyncClient() as client:
# Call search agent /process endpoint
search_resp = await client.post(
"http://localhost:8001/process",
json={"query": input.query, "top_k": 10}
)
search_result = search_resp.json()
# Call summarizer agent /process endpoint
summary_resp = await client.post(
"http://localhost:8002/process",
json={"text": str(search_result["results"]), "max_length": 200}
)
summary_result = summary_resp.json()
return OrchestratorOutput(
search_results=search_result["results"],
summary=summary_result["summary"]
)
Architecture Position¶
flowchart TB
subgraph AppLayer["<span style='color:#000'>Application Layer</span>"]
Runtime["<span style='color:#000'>cogniverse-runtime (FastAPI)</span>"]
end
subgraph ImplLayer["<span style='color:#000'>Implementation Layer</span>"]
Agents["<span style='color:#000'>cogniverse-agents<br/>(OrchestratorAgent, SearchAgent)</span>"]
Vespa["<span style='color:#000'>cogniverse-vespa<br/>(Vespa backend)</span>"]
Synthetic["<span style='color:#000'>cogniverse-synthetic<br/>(data gen)</span>"]
end
subgraph CoreLayer["<span style='color:#000'>Core Layer</span>"]
Core["<span style='color:#000'>cogniverse-core ◄─ YOU ARE HERE<br/>AgentBase, A2AAgent, Registries, Memory, Config</span>"]
Evaluation["<span style='color:#000'>cogniverse-evaluation</span>"]
Telemetry["<span style='color:#000'>cogniverse-telemetry-phoenix</span>"]
end
subgraph FoundationLayer["<span style='color:#000'>Foundation Layer</span>"]
SDK["<span style='color:#000'>cogniverse-sdk (interfaces)</span>"]
Foundation["<span style='color:#000'>cogniverse-foundation<br/>(config, telemetry base)</span>"]
end
AppLayer --> ImplLayer
ImplLayer --> CoreLayer
CoreLayer --> FoundationLayer
style AppLayer fill:#90caf9,stroke:#1565c0,color:#000
style Runtime fill:#90caf9,stroke:#1565c0,color:#000
style ImplLayer fill:#ffcc80,stroke:#ef6c00,color:#000
style Agents fill:#ffcc80,stroke:#ef6c00,color:#000
style Vespa fill:#ffcc80,stroke:#ef6c00,color:#000
style Synthetic fill:#ffcc80,stroke:#ef6c00,color:#000
style CoreLayer fill:#ce93d8,stroke:#7b1fa2,color:#000
style Core fill:#ce93d8,stroke:#7b1fa2,color:#000
style Evaluation fill:#ce93d8,stroke:#7b1fa2,color:#000
style Telemetry fill:#ce93d8,stroke:#7b1fa2,color:#000
style FoundationLayer fill:#a5d6a7,stroke:#388e3c,color:#000
style SDK fill:#a5d6a7,stroke:#388e3c,color:#000
style Foundation fill:#a5d6a7,stroke:#388e3c,color:#000 Testing¶
# Run all core tests (agents and common utilities)
JAX_PLATFORM_NAME=cpu uv run pytest tests/agents/ tests/common/ -v
# Run specific test categories
uv run pytest tests/agents/unit/ -v
uv run pytest tests/common/unit/ -v
uv run pytest tests/agents/integration/ -v
# Run with coverage
uv run pytest tests/agents/ tests/common/ --cov=cogniverse_core --cov-report=html
Test Categories:
tests/agents/unit/- Unit tests for base classes, mixins, registriestests/agents/integration/- Integration tests with multiple componentstests/common/- Tests for shared utilities
Cache Subsystem¶
Location: common/cache/
The cache subsystem provides tiered caching for embeddings and pipeline artifacts.
CacheBackend (base.py)¶
Abstract base class for cache backends:
class CacheBackend(ABC):
async def get(self, key: str) -> Optional[Any]: ...
async def set(self, key: str, value: Any, ttl: Optional[int] = None) -> bool: ...
async def delete(self, key: str) -> bool: ...
async def exists(self, key: str) -> bool: ...
async def clear(self, pattern: Optional[str] = None) -> int: ...
async def get_stats(self) -> Dict[str, Any]: ...
CacheManager (base.py)¶
Manages multiple cache backends with tiered caching:
from cogniverse_core.common.cache import CacheManager, CacheConfig
config = CacheConfig(
backends=[
{"backend_type": "structured_filesystem", "priority": 0},
],
default_ttl=3600,
enable_compression=True,
serialization_format="pickle" # or "json", "msgpack"
)
manager = CacheManager(config)
await manager.set("key", value, ttl=3600)
result = await manager.get("key")
PipelineArtifactCache (pipeline_cache.py)¶
Caches video processing pipeline artifacts:
from cogniverse_core.common.cache import PipelineArtifactCache
cache = PipelineArtifactCache(
cache_manager,
ttl=604800, # 7 days
profile="video_colpali_mv_frame"
)
# Cache keyframes
await cache.set_keyframes(
video_path="video.mp4",
keyframes_metadata={"keyframes": [...]},
keyframe_images={"0": image_array},
strategy="similarity",
threshold=0.999
)
# Retrieve with optional image loading
metadata = await cache.get_keyframes(
video_path="video.mp4",
strategy="similarity",
load_images=True
)
Artifact families: keyframes (get/set_keyframes), audio transcripts (get/set_transcript), frame descriptions (get/set_descriptions), per-segment frames (get/set_segment_frames), and single-vector segmentation results (get/set_segmentation — boundary math + transcript alignment keyed by strategy params and a transcript fingerprint; the single-vector segmentation strategy hands the pipeline cache to SingleVectorVideoProcessor so repeated segmentation of the same video/params serves from cache).
CacheBackendRegistry (registry.py)¶
Plugin registry for cache backends:
from cogniverse_core.common.cache import CacheBackendRegistry
# Register custom backend
CacheBackendRegistry.register("redis", RedisCacheBackend)
# Create from config
backend = CacheBackendRegistry.create({"backend_type": "structured_filesystem", ...})
# List registered backends
backends = CacheBackendRegistry.list_backends() # ["structured_filesystem", ...]
Built-in CacheBackend implementations (common/cache/backends/):
| Backend | Registered name | Storage |
|---|---|---|
StructuredFilesystemBackend | structured_filesystem | Local filesystem, keyed directory layout |
S3CacheBackend | s3 | S3-compatible object storage (works with MinIO) |
Utility Modules¶
Location: common/utils/
retry.py - Retry with Exponential Backoff¶
from cogniverse_core.common.utils.retry import retry_with_backoff, RetryConfig
config = RetryConfig(
max_attempts=3,
initial_delay=1.0,
max_delay=60.0,
exponential_base=2.0,
jitter=True, # Prevents thundering herd
exceptions=(ConnectionError, TimeoutError)
)
@retry_with_backoff(config=config)
def fetch_data():
return requests.get(url)
# With callbacks
@retry_with_backoff(
on_retry=lambda e, attempt: logger.warning(f"Retry {attempt}: {e}"),
on_failure=lambda e: logger.error(f"Failed: {e}")
)
def process_item(item):
return api.process(item)
async_polling.py - Semantic Wait Functions¶
from cogniverse_core.common.utils.async_polling import wait_for_retry_backoff
# Wait with backoff for retries
wait_for_retry_backoff(
attempt=2,
base_delay=1.0,
max_delay=60.0,
exponential=True
)
Other Utilities¶
| Module | Purpose |
|---|---|
output_manager.py | Manage output directories and artifacts |
async_bridge.py | Run a coroutine to completion from sync code |
async_polling.py | Async polling with backoff for polled waits |
VLM Interface¶
Location: common/vlm_interface.py
Vision Language Model interface using DSPy for visual content analysis.
from cogniverse_core.common.vlm_interface import VLMInterface
vlm = VLMInterface(
config_manager=config_manager,
tenant_id="acme"
)
# Visual analysis
result = await vlm.analyze_visual_content(
image_paths=["frame1.jpg", "frame2.jpg"],
query="Find manufacturing defects"
)
# Returns: {descriptions, themes, key_objects, insights, relevance_score}
DSPy Signatures:
| Signature | Purpose |
|---|---|
VisualAnalysisSignature | Basic analysis (descriptions, themes, objects, insights, relevance_score) |
Model Loaders¶
Location: common/models/
ModelLoader (model_loaders.py)¶
Abstract base class for model loaders:
from cogniverse_core.common.models import ModelLoader
class CustomLoader(ModelLoader):
def load_model(self) -> Tuple[Any, Any]:
# Load and return (model, processor)
pass
# Auto device detection
device = loader.get_device() # "cuda", "mps", or "cpu"
dtype = loader.get_dtype() # bfloat16 for CUDA, float32 otherwise
ModelLoaderFactory (model_loaders.py)¶
Factory for creating model loaders based on the model_loader config key:
from cogniverse_core.common.models import ModelLoaderFactory
# Config must contain "model_loader" key — raises ValueError if missing.
# A "remote_inference_url" selects the remote loader; ColQwen3/Tomoro
# (model_type qwen3_vl) is remote-only and must always carry one.
loader = ModelLoaderFactory.create_loader(
model_name="TomoroAI/tomoro-colqwen3-embed-4b",
config={
"model_loader": "colpali",
"embedding_type": "multi_vector",
"remote_inference_url": "http://localhost:8000",
},
logger=logger,
)
model, processor = loader.load_model()
# Capability probe: True when the model has no in-process loader (must be
# served via the vLLM sidecar). Gate on this instead of catching the
# in-process load's RuntimeError.
from cogniverse_core.common.models import is_remote_only_model
is_remote_only_model("TomoroAI/tomoro-colqwen3-embed-4b") # True
Loader Registry:
model_loader key | Local Loader | Remote Loader |
|---|---|---|
colpali | ColPaliModelLoader | RemoteColPaliLoader |
colqwen | ColQwenModelLoader | RemoteColPaliLoader |
xclip | — (remote only) | RemoteXClipLoader |
colbert | ColBERTModelLoader | RemoteColBERTLoader |
whisper | — | RemoteWhisperLoader |
Remote loaders are selected when remote_inference_url is set in config. RemoteWhisperLoader's wrapper transcribes through the chunked client below and returns text, language, duration and segments. RemoteXClipLoader forwards its exact configured model name on every video-segment request, so the remote service can reject requests for any other checkpoint instead of silently selecting a default. Remote loader cache keys include a one-way fingerprint of the exact credential snapshot used to construct the client. A custom endpoint key or rotated Modal environment key therefore cannot reuse a client authenticated with another key.
ColQwen3/Tomoro is remote-only. TomoroAI/tomoro-colqwen3-embed-4b (architecture qwen3_vl) has no in-process loader: the pinned transformers (4.56.2, capped by pylate) lacks qwen3_vl support and colpali_engine mis-maps it to idefics3. Constructing the local ColPaliModelLoader/ColQwenModelLoader (or the corresponding query encoder) for such a model without remote_inference_url raises a clear RuntimeError directing the operator to serve via vLLM and set inference_service_url (profile inference_services.embedding).
Chunked Whisper transcription (whisper_transcription.py)¶
transcribe_in_chunks(samples, transcribe_chunk, *, language, source, logger) sends 16 kHz mono PCM16 audio to an OpenAI-compatible Whisper endpoint chunk by chunk and merges the answers into full_text, language, duration and segments (times from the start of the audio). Every remote Whisper client goes through it: AudioProcessor, AudioAnalysisAgent and RemoteWhisperLoader. transcribe_chunk(chunk, language, timestamps, temperature) sends one request with response_format(timestamps) and sampling_fields(temperature) (temperature, and seed = SAMPLING_SEED).
split_for_whisper(samples)cuts where vLLM's Whisper server cuts a long file: audio of at most 30 s is one chunk; longer audio is cut every 30 s at the start of the quietest 0.1 s window in the chunk's last second.- Each chunk is asked for
verbose_json(timings). vLLM builds that transcript only from text between two adjacent timestamp tokens: a decode with no such pair comes back empty, text after the last pair is dropped, and a decode ending on a pair before the end of the chunk leaves the rest undecoded. A timed answer whose segments run to the chunk's end (reaches_chunk_end: the last end within one 0.02 s timestamp step of it, since chunk lengths are no multiple of the step) lost nothing and is kept as it is. Otherwise the chunk is also asked forjson, which keeps the whole decode, so a transcription costs one or two ASR requests per 30 s chunk plus one per answer asked again. align_text(text, timed, duration, *, no_space=False)times thejsonwords with theverbose_jsonsegments. Words are compared case- and punctuation-blind (characters forjaandzh). A json word matching a timed word takes its segment, and a timed word the json answer lacks stays in its segment. Where the answers word the same stretch differently, the rendering with more words is kept (the json one on a tie), so a json decode that stops early or collapses cannot replace timed text. A json word the timed answer lacks gets a segment spanning the untimed gap it falls in (before the first segment, between two that do not meet, after the last), or else joins the segment before it. With no timed segments the text is one segment spanning the chunk. No segment runs past the chunk's duration, so none reaches into the next chunk.- An answer is unusable when it is a repetition loop (
compression_ratio(text)aboveGARBLED_COMPRESSION_RATIO, 2.4) or empty for a chunk whose loudest 25 ms frame reachesSILENCE_FLOOR_DBFS(-60); a timed answer is also unusable without segments. An unusable answer is asked again at the next ofFALLBACK_TEMPERATURES(0.0, 0.2, 0.4, 0.6, 0.8, 1.0;TRANSCRIBE_ATTEMPTSis 6). When neither ajsonanswer nor a timed answer running to the end is usable, the chunk raisesGarbledTranscriptError(source,chunk_index,start_s,end_s,compression_ratios,attempts) if any looped, elseEmptyTranscriptError(source,chunk_index,start_s,end_s,loudest_frame_dbfs,attempts). When noverbose_jsonanswer is usable, the text spans the chunk. A silent chunk may come back empty. A request that fails is not repeated. - With no language named, the first chunk's timed answer names it and every later request is sent in it, as the server does for a whole file. Chunk texts join with a space, or with nothing for
jaandzh. lenient_chunk_answer(body, chunk)(the processor's and the loader's parser) and the agent's strict parser clamp a segment time past the chunk's duration to it withclamp_to_duration, which logs the original value at DEBUG unless it is past by less than half a timestamp step (vLLM's0.02 * nrounds 29.4 to 29.400000000000002): Whisper times text into the padding after short audio (the live server gave 29.98 s on an 18.77 s chunk).- Limits: the loop check measures a whole answer, so a short repetition inside an otherwise ordinary answer is kept (the live server's "DR. DR. DR. DR. SOUDOS, ..." sat in a json answer whose ratio was 1.76). Text that both answers miss over the same stretch, or that a timed answer running to the chunk's end skipped, is not detected.
decode_audio(path)decodes any container's first audio stream to 16 kHz mono PCM16 (pyav);pcm16_wav_samplesandwav_bytesconvert to and from WAV.
ColBERTModelLoader (model_loaders.py)¶
Loads ColBERT late-interaction models via PyLate for document and audio semantic embeddings:
from cogniverse_core.common.models import ColBERTModelLoader
loader = ColBERTModelLoader(
model_name="lightonai/LateOn",
config={"model_loader": "colbert"},
logger=logger,
)
model, _ = loader.load_model()
# model is a pylate.models.ColBERT instance
# Returns 128-dim per-token multi-vector embeddings
RemoteColBERTLoader (model_loaders.py)¶
Serves the ColBERT contract from the PyLate service's /pooling endpoint (cogniverse_cli/modal_inference/servers/pylate.py) instead of an in-process pylate model. load_model() returns a ColBERTRemoteWrapper whose .encode(texts, is_query=...) mirrors pylate.models.ColBERT.encode(): the wrapper sends raw text plus is_query, and the service applies PyLate's own query expansion over masked padding positions and document punctuation skiplist, so the per-token matrices come back unchanged.
from cogniverse_core.common.models.model_loaders import RemoteColBERTLoader
loader = RemoteColBERTLoader(
model_name="lightonai/LateOn",
config={"remote_inference_url": "http://localhost:8000"},
logger=logger,
)
model, _ = loader.load_model()
doc_tokens = model.encode(["Vespa stores token embeddings."], is_query=False)[0]
RemoteInferenceClient (model_loaders.py)¶
Client for remote model inference providers:
from cogniverse_core.common.models.model_loaders import RemoteInferenceClient
client = RemoteInferenceClient(
endpoint_url="http://localhost:8080",
api_key="..."
)
# Process images with retry logic
result = client.process_images(
images=["image1.jpg", pil_image]
)
Supported providers:
- Infinity (ColPali and similar models)
- Modal (custom deployed models)
- Custom REST APIs
RemoteGlinerClient (model_loaders.py)¶
RemoteGlinerClient preserves the GLiNER predict_entities(text, labels, threshold) contract over /predict_entities. Modal endpoints use COGNIVERSE_INFERENCE_API_KEY; non-Modal endpoints use api_key. Cache keys include a one-way fingerprint of the credential snapshot used to construct the client.
RemoteClapClient (model_loaders.py)¶
RemoteClapClient calls the clap_embed sidecar's /embed/text and /embed/audio and returns the 512-d CLAP vector, raising with the URL and cause on a transport or HTTP failure and rejecting any response whose vector is not exactly 512 floats. The runtime audio embedding generator delegates its remote CLAP calls to it.
SemanticEmbedder (semantic_embedder.py)¶
Text embedder used by memory/dedup code paths that need a plain sentence embedding (not the multi-vector ColBERT/ColPali contract). It calls a remote OpenAI-compatible /v1/embeddings endpoint; no model is loaded in-process. With no endpoint it raises SemanticEmbedderNotConfiguredError naming the settings.
from cogniverse_core.common.models.semantic_embedder import get_semantic_embedder
# Resolution order: explicit remote_url/model_name arg -> the entrypoint's
# default (COGNIVERSE_SEMANTIC_EMBED_URL, else INFERENCE_SERVICE_URLS["denseon"];
# COGNIVERSE_SEMANTIC_EMBED_MODEL, else DenseOn) -> SemanticEmbedderNotConfiguredError.
embedder = get_semantic_embedder()
vectors = embedder.encode(["find manufacturing defects"], is_query=True) # (1, D)
# Authenticated remote endpoint; the credential is validated, forwarded on
# every request, and represented only by a digest in the shared cache key.
modal_embedder = get_semantic_embedder(
remote_url="https://denseon.modal.run",
)
| Class | Backend |
|---|---|
RemoteOpenAIEmbedder | Authenticated HTTP client for an OpenAI-compatible /v1/embeddings server; Modal URLs use COGNIVERSE_INFERENCE_API_KEY, non-Modal URLs may accept an exact bearer Authorization mapping, and DenseOn's query:/document: prompt prefixes plus L2 normalization are applied client-side |
Remote instances are cached module-level by (remote_url, model, credential fingerprint). Concurrent agents therefore share only clients for the same endpoint, model, and authentication context; reset_semantic_embedder_cache() clears them for tests. Connection and timeout errors preserve their original requests exception type while adding the exact model and /v1/embeddings endpoint; bearer values are never included in that context.
Backend Factory & Profile Validation¶
BackendFactory (factories/backend_factory.py)¶
Centralizes backend construction so the Backend ↔ SchemaRegistry circular dependency is resolved in one place instead of scattered ad hoc wiring:
from cogniverse_core.factories.backend_factory import BackendFactory
backend = BackendFactory.create_backend_with_dependencies(
backend_class=MyBackendClass,
backend_config=backend_config,
config_manager=config_manager,
schema_loader=schema_loader,
backend_init_config={"profile": "video_colpali_mv_frame"},
)
Sequence: construct the backend (without a registry) -> build or reuse a SchemaRegistry -> inject it into the backend -> call backend.initialize() -> inject the registry into the backend's internal schema manager if present.
ProfileValidator (validation/profile_validator.py)¶
Validates a BackendProfileConfig before it is created or updated — schema template existence, importable strategy classes, and consistent profile settings — so a broken profile fails at admin-API time, not at first search.
from cogniverse_core.validation.profile_validator import ProfileValidator
validator = ProfileValidator(config_manager, schema_templates_dir=Path("configs/schemas"))
errors = validator.validate_profile(profile, tenant_id="acme", is_update=False)
VALID_EMBEDDING_TYPES = ["multi_vector", "single_vector"] are the accepted embedding types. Profile types are the types of the shipped profiles in configs/config.json backend.profiles. model_loader must be one of model_loaders.EMBEDDING_MODEL_LOADERS (colbert, colpali, colqwen, xclip), and is required for every type whose shipped profiles all name one, so a profile ingestion cannot embed with is refused at creation. process_type, when set, must be one of unified_config.PROCESS_TYPES (direct_video, frame_based, video_chunks), and extra_config keys may not name a profile field.
FilesystemSchemaLoader (schemas/filesystem_loader.py)¶
Source: cogniverse_core/schemas/filesystem_loader.py.
Loads Vespa schema template files (.sd content + JSON metadata) from disk for SchemaRegistry and BackendFactory to deploy; implements the SchemaLoader interface from cogniverse_sdk. Its source is ("filesystem", <resolved directory>), so two loaders over one directory are interchangeable to the backend registry's cache-hit check.
Related Documentation¶
- Agents Module - Concrete agent implementations (OrchestratorAgent, SearchAgent, etc.)
- Multi-Agent Interactions - A2A protocol flows
- SDK Architecture - Package structure
- Creating Agents Tutorial - Step-by-step agent creation
Summary: The Core module provides the type-safe foundation for all agents in Cogniverse. AgentBase[InputT, OutputT, DepsT] ensures compile-time type checking and runtime validation, while A2AAgent adds A2A protocol support and DSPy integration (HTTP server endpoints are managed by the runtime's A2AStarletteApplication). Mixins provide composable functionality for memory, multi-tenancy, and health checks.