Skip to content

Agents Module

Package: cogniverse_agents (Implementation Layer) Location: libs/agents/cogniverse_agents/


Table of Contents

  1. Module Overview
  2. Package Structure
  3. Core Agents
  4. SearchAgent
  5. GatewayAgent
  6. OrchestratorAgent (A2A Entry Point)
  7. ProfileSelectionAgent
  8. EntityExtractionAgent
  9. SearchAgent (Ensemble Mode)
  10. DetailedReportAgent
  11. DocumentAgent
  12. ImageSearchAgent
  13. SummarizerAgent
  14. AudioAnalysisAgent
  15. TextAnalysisAgent
  16. SearchAgent (Refactored)
  17. QueryEnhancementAgent
  18. CodingAgent
  19. DeepResearchAgent
  20. Knowledge Agents
  21. AuditExplanationAgent
  22. KnowledgeSummarizationAgent
  23. TemporalReasoningAgent
  24. FederatedQueryAgent
  25. CrossTenantComparisonAgent
  26. KnowledgeGraphTraversalAgent
  27. MultiDocumentSynthesisAgent
  28. ContradictionReconciliationAgent
  29. CitationTracingAgent
  30. Agent Architecture
  31. Multi-Tenant Integration
  32. Usage Examples
  33. Streaming API
  34. RLM Inference (Recursive Language Models)
  35. Testing
  36. Audit Checklist
  37. Real-Time Event Notifications
  38. Approval Workflow System
  39. Inference System

Module Overview

The Agents package (cogniverse-agents) provides concrete agent implementations for the Cogniverse multi-agent AI platform. The architecture supports any agent type - content understanding agents ship by default, but web browsing, code analysis, and domain-specific agents can be integrated via the same AgentBase/A2AAgent base classes. All agents are tenant-aware and integrate with the core SDK packages.

Key Agents

Search and Orchestration Agents:

  1. GatewayAgent - Query entry point: GLiNER-based triage classifying queries as simple or complex (<100ms, no LLM)
  2. OrchestratorAgent - Autonomous A2A orchestrator: DSPy planning, parallel execution, cross-modal fusion
  3. SearchAgent - Multi-modal video search (ColPali, X-CLIP)
  4. ProfileSelectionAgent - LLM-based intelligent backend profile selection and ensemble composition
  5. EntityExtractionAgent - Named entity extraction with DSPy primary path; GLiNER + SpaCy fallback (PERSON, ORGANIZATION, CONCEPT, PLACE, EVENT, TECHNOLOGY; verbatim query spans)
  6. SearchAgent - Enhanced with ensemble mode and RRF fusion for multi-profile queries
  7. DetailedReportAgent - Comprehensive report generation with VLM visual analysis
  8. DocumentAgent - Dual-strategy document search (visual ColPali + text semantic)
  9. ImageSearchAgent - Image similarity search using ColPali embeddings
  10. SummarizerAgent - Intelligent summarization with thinking phase
  11. AudioAnalysisAgent - Audio search with Whisper transcription
  12. TextAnalysisAgent - Runtime-configurable text analysis with DSPy
  13. SearchAgent (Refactored) - Simplified video search with unified service
  14. QueryEnhancementAgent - Enhances queries with synonyms, context, and related terms to improve search recall
  15. CodingAgent - Iterative code generation with semantic code search (LateOn-Code-edge) and sandboxed execution
  16. DeepResearchAgent - Multi-step research: decomposes queries, dispatches parallel searches, synthesizes a cited report

Knowledge Management Agents:

  1. MultiDocumentSynthesisAgent - Synthesize and reconcile claims across multiple source documents with citations
  2. KnowledgeGraphTraversalAgent - Walk the knowledge graph to find entity chains and relationship sub-graphs
  3. CrossTenantComparisonAgent - Compare knowledge views across multiple tenants (org-admin scoped)
  4. ContradictionReconciliationAgent - Surface and resolve ConflictSet entries in the contradiction store
  5. CitationTracingAgent - Walk provenance chains back to primary sources for a given answer
  6. TemporalReasoningAgent - Answer questions about how knowledge changed over time
  7. FederatedQueryAgent - Retrieve knowledge from both tenant overlay and org trunk in one call
  8. KnowledgeSummarizationAgent - Summarize a knowledge slice (by subject / kind / time window) with citations
  9. AuditExplanationAgent - Explain why an answer was produced by tracing provenance, trust, and contradictions

Design Principles

  • Tenant-Agnostic at Startup: Agents boot without tenant_id — it arrives per-request in A2A task payload
  • Memory-Enabled: Integration with Mem0 via MemoryAwareMixin from cogniverse_agents
  • Base Class Inheritance: Extend A2AAgent[InputT, OutputT, DepsT] from cogniverse_core with type-safe generics
  • DSPy 3.0 Integration: A2A protocol + DSPy modules for optimization
  • Streaming Support: OpenAI-style stream=True parameter for progressive results
  • Production-Ready: Health checks, graceful degradation, telemetry

Extensibility

The agent architecture is not limited to content understanding. The AgentBase and A2AAgent base classes support any agent type:

  • Web Browsing Agents: Research, scraping, monitoring
  • Code Agents: Analysis, generation, refactoring
  • Data Agents: Database queries, API integrations
  • Communication Agents: Email, Slack, notifications
  • Domain-Specific Agents: Legal, medical, financial analysis

To add a custom agent, implement A2AAgent[InputT, OutputT, DepsT] and register with the AgentRegistry. All agents automatically gain tenant isolation, memory, telemetry, and DSPy optimization capabilities.

Package Dependencies

# Agents package depends on:
from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from cogniverse_core.agents.base import AgentBase, AgentDeps, AgentInput, AgentOutput
from cogniverse_core.agents.tenant_aware_mixin import TenantAwareAgentMixin
from cogniverse_agents.memory_aware_mixin import MemoryAwareMixin
from cogniverse_core.common.health_mixin import HealthCheckMixin
from cogniverse_foundation.telemetry.manager import TelemetryManager
from cogniverse_foundation.config.unified_config import SystemConfig

Package Structure

graph TD
    Root["<span style='color:#000'><b>cogniverse_agents/</b></span>"]

    Root --> Init["<span style='color:#000'>__init__.py</span>"]
    Root --> GatewayAgent["<span style='color:#000'><b>gateway_agent.py</b><br/>A2A Entry Point</span>"]
    Root --> OrchestratorAgent["<span style='color:#000'><b>orchestrator_agent.py</b><br/>A2A Orchestrator</span>"]
    Root --> AudioAgent["<span style='color:#000'>audio_analysis_agent.py</span>"]
    Root --> DocAgent["<span style='color:#000'>document_agent.py</span>"]
    Root --> ImageAgent["<span style='color:#000'>image_search_agent.py</span>"]
    Root --> SummarizerAgent["<span style='color:#000'>summarizer_agent.py</span>"]
    Root --> More["<span style='color:#000'>... (knowledge agents + other agent files)</span>"]

    Root --> RoutingDir["<span style='color:#000'><b>routing/</b><br/>13 files</span>"]
    RoutingDir --> RoutingInit["<span style='color:#000'>__init__.py</span>"]
    RoutingDir --> RelExtract["<span style='color:#000'>relationship_extraction_tools.py</span>"]
    RoutingDir --> DspyRouter["<span style='color:#000'>dspy_relationship_router.py</span>"]
    RoutingDir --> DspySigs["<span style='color:#000'>dspy_routing_signatures.py</span>"]
    RoutingDir --> MoreRouting["<span style='color:#000'>... (utility files)</span>"]

    Root --> SearchDir["<span style='color:#000'><b>search/</b><br/>9 files</span>"]
    SearchDir --> SearchInit["<span style='color:#000'>__init__.py</span>"]
    SearchDir --> MMRerank["<span style='color:#000'>multi_modal_reranker.py</span>"]
    SearchDir --> HybridRerank["<span style='color:#000'>hybrid_reranker.py</span>"]
    SearchDir --> LearnedRerank["<span style='color:#000'>learned_reranker.py</span>"]
    SearchDir --> RerankersDir["<span style='color:#000'>rerankers/</span>"]

    Root --> OrchDir["<span style='color:#000'><b>orchestrator/</b><br/>2 files</span>"]
    OrchDir --> OrchInit["<span style='color:#000'>__init__.py</span>"]
    OrchDir --> SuffContext["<span style='color:#000'>sufficient_context_signature.py</span>"]

    Root --> OptDir["<span style='color:#000'><b>optimizer/</b><br/>5 files</span>"]
    OptDir --> OptInit["<span style='color:#000'>__init__.py</span>"]
    OptDir --> ArtifactMgr["<span style='color:#000'>artifact_manager.py</span>"]
    OptDir --> DspyAgentOpt["<span style='color:#000'>dspy_agent_optimizer.py</span>"]
    OptDir --> SigVariants["<span style='color:#000'>signature_variants.py</span>"]
    OptDir --> StrategyLearner["<span style='color:#000'>strategy_learner.py</span>"]

    Root --> InferenceDir["<span style='color:#000'><b>inference/</b><br/>6 files</span>"]
    InferenceDir --> RlmInf["<span style='color:#000'>rlm_inference.py</span>"]
    InferenceDir --> InstrumentedRlm["<span style='color:#000'>instrumented_rlm.py</span>"]

    Root --> ApprovalDir["<span style='color:#000'><b>approval/</b><br/>4 files</span>"]
    ApprovalDir --> ApprovalStorage["<span style='color:#000'>approval_storage.py</span>"]
    ApprovalDir --> HumanApproval["<span style='color:#000'>human_approval_agent.py</span>"]
    ApprovalDir --> ApprovalOrch["<span style='color:#000'>orchestrator.py</span>"]

    Root --> GraphDir["<span style='color:#000'><b>graph/</b><br/>12 files</span>"]
    GraphDir --> GraphMgr["<span style='color:#000'>graph_manager.py</span>"]
    GraphDir --> ClaimExtract["<span style='color:#000'>claim_extractor.py</span>"]

    Root --> WikiDir["<span style='color:#000'><b>wiki/</b><br/>3 files</span>"]
    WikiDir --> WikiMgr["<span style='color:#000'>wiki_manager.py</span>"]

    Root --> MixinsDir["<span style='color:#000'><b>mixins/</b><br/>2 files</span>"]
    MixinsDir --> RlmMixin["<span style='color:#000'>rlm_aware_mixin.py</span>"]

    Root --> ToolsDir["<span style='color:#000'><b>tools/</b></span>"]
    Root --> WorkflowDir["<span style='color:#000'><b>workflow/</b></span>"]

    style Root fill:#ce93d8,stroke:#7b1fa2,color:#000
    style GatewayAgent fill:#ffcc80,stroke:#ef6c00,color:#000
    style OrchestratorAgent fill:#ffcc80,stroke:#ef6c00,color:#000
    style RoutingDir fill:#81d4fa,stroke:#0288d1,color:#000
    style SearchDir fill:#81d4fa,stroke:#0288d1,color:#000
    style OrchDir fill:#81d4fa,stroke:#0288d1,color:#000
    style OptDir fill:#81d4fa,stroke:#0288d1,color:#000
    style InferenceDir fill:#81d4fa,stroke:#0288d1,color:#000
    style ApprovalDir fill:#81d4fa,stroke:#0288d1,color:#000
    style GraphDir fill:#81d4fa,stroke:#0288d1,color:#000
    style WikiDir fill:#81d4fa,stroke:#0288d1,color:#000
    style MixinsDir fill:#81d4fa,stroke:#0288d1,color:#000
    style ToolsDir fill:#81d4fa,stroke:#0288d1,color:#000
    style WorkflowDir fill:#81d4fa,stroke:#0288d1,color:#000

The wiki implementation is the libs/agents/cogniverse_agents/wiki/ subpackage; wiki_manager.py owns page persistence and lint reporting.

Total Files: 96 Python files (35 at top level + 61 in subdirectories)

Key Agent Files (all at top level):

  • gateway_agent.py: A2A entry point with GLiNER classification
  • orchestrator_agent.py: A2A orchestrator with DSPy planning

Core Agents

1. SearchAgent

Purpose: Text-to-video search with ColPali and X-CLIP embeddings Constructor: SearchAgent(deps: SearchAgentDeps, schema_loader=None, config_manager=None, port: int = 8002)

Backend resolution: the agent holds no backend. _get_backend() is the resolver seam — it resolves through BackendRegistry on every call — and _search_backend(query_dict) runs one search with that instance leased for the call. AudioAnalysisAgent and SearchService use the same two methods.

Recorded searches: every text search the agent runs records a search_service.search span (search_span, the span SearchService.search records) under the request's tenant, with the user's query as query, the rewrite it searched as enhanced_query, and the profile, strategy (default when none was requested), top_k and result rows. A single-profile or relationship-aware search records one span, an ensemble one per profile, and a multi-query fusion one span holding the fused results. The evaluation views and the optimization framework read these spans; a failure to record the result rows is logged and never fails the search.

Profile resolution: profiles, models and encoder services come from the config of SearchAgentDeps.tenant_id — the tenant's own profiles merged over the system's — and from the system tenant when it is unset. The dispatcher caches one agent per tenant and profile.

Query encoding: constructing the agent builds no encoder. A text search (plain, relationship-aware, query variants, ensemble) hands the backend a SharedQueryEncoder instead of embeddings; the backend encodes through it only when the resolved ranking strategy needs embeddings, so bm25_only never builds or calls the encoder. query_encoder builds the active profile's encoder once, on first use, under a lock; media queries (an image or video file) go through it via content_processor. Encoder failures surface as EncoderNotConfiguredError and EncoderUnavailableError, which /agents/{name}/process answers as typed 500 and 503 bodies.

Multi-Modal Support

flowchart LR
    Input["<span style='color:#000'>Input</span>"] --> TextQuery["<span style='color:#000'>Text Query</span>"]
    Input --> VideoFile["<span style='color:#000'>Video File</span>"]
    Input --> ImageFile["<span style='color:#000'>Image File</span>"]

    TextQuery --> ColPaliEncode["<span style='color:#000'>ColPali Text Encoder</span>"]
    TextQuery --> X-CLIPEncode["<span style='color:#000'>X-CLIP Text Encoder</span>"]

    VideoFile --> FrameExtract["<span style='color:#000'>Extract Frames<br/>1 FPS</span>"]
    FrameExtract --> VideoEncode["<span style='color:#000'>Encode Frames<br/>ColPali/X-CLIP</span>"]

    ImageFile --> ImageEncode["<span style='color:#000'>Encode Image<br/>ColPali/X-CLIP</span>"]

    ColPaliEncode --> VespaSearch["<span style='color:#000'>Vespa Search<br/>Schema: video_frames_{tenant_id}</span>"]
    X-CLIPEncode --> VespaSearch
    VideoEncode --> VespaSearch
    ImageEncode --> VespaSearch

    VespaSearch --> Rerank["<span style='color:#000'>Rerank Results<br/>Relationship boost</span>"]
    Rerank --> Results["<span style='color:#000'>Ranked Results</span>"]

    style Input fill:#90caf9,stroke:#1565c0,color:#000
    style TextQuery fill:#90caf9,stroke:#1565c0,color:#000
    style VideoFile fill:#90caf9,stroke:#1565c0,color:#000
    style ImageFile fill:#90caf9,stroke:#1565c0,color:#000
    style ColPaliEncode fill:#81d4fa,stroke:#0288d1,color:#000
    style X-CLIPEncode fill:#81d4fa,stroke:#0288d1,color:#000
    style FrameExtract fill:#ffcc80,stroke:#ef6c00,color:#000
    style VideoEncode fill:#81d4fa,stroke:#0288d1,color:#000
    style ImageEncode fill:#81d4fa,stroke:#0288d1,color:#000
    style VespaSearch fill:#90caf9,stroke:#1565c0,color:#000
    style Rerank fill:#ffcc80,stroke:#ef6c00,color:#000
    style Results fill:#a5d6a7,stroke:#388e3c,color:#000

Class Definition

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_core.agents.a2a_agent import A2AAgent
from cogniverse_agents.memory_aware_mixin import MemoryAwareMixin
from cogniverse_agents.mixins.rlm_aware_mixin import RLMAwareMixin

class SearchAgent(
    RLMAwareMixin,
    MemoryAwareMixin,
    A2AAgent[SearchInput, SearchOutput, SearchAgentDeps],
):
    """
    Type-safe generic multi-modal search agent with full A2A protocol support.
    """

    def __init__(
        self,
        deps: SearchAgentDeps,
        schema_loader=None,   # REQUIRED
        config_manager=None,
        port: int = 8002,
    ):
        """
        Initialize generic search agent with typed dependencies.

        Args:
            deps: SearchAgentDeps with backend configuration
            schema_loader: SchemaLoader instance (REQUIRED)
            config_manager: ConfigManager instance (optional, creates default if None)
            port: A2A server port

        Raises:
            ValueError: If schema_loader is None
        """
        ...

Key Methods

search_by_text(query, *, tenant_id, modality="video", top_k=10, **kwargs) -> List[Dict]

Text search. This method is synchronous (not async). tenant_id is required per-request.

def search_by_text(
    self,
    query: str,
    *,
    tenant_id: str,             # Required per-request
    modality: str = "video",
    top_k: int = 10,
    **kwargs,
) -> List[Dict[str, Any]]:
    """
    Search content using text query.

    Args:
        query: Text search query
        tenant_id: Tenant identifier (required)
        modality: Content modality to search (video/image/text/audio/document)
        top_k: Number of results to return
        **kwargs: Additional search parameters (ranking, etc.)

    Returns:
        List of search results
    """
    ...

Usage Example

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from cogniverse_foundation.config.utils import create_default_config_manager
from pathlib import Path

config_manager = create_default_config_manager()
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))
deps = SearchAgentDeps()

# Create agent — schema_loader is required
agent = SearchAgent(deps=deps, schema_loader=schema_loader, config_manager=config_manager)

# Search (synchronous) — tenant_id required per-request
results = agent.search_by_text(
    query="cooking tutorial",
    tenant_id="acme",
    modality="video",
    top_k=10,
)

Note: SearchAgent provides text-to-video search only. For video-to-video similarity search, use SearchAgent which supports this via _search_by_video():

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from cogniverse_foundation.config.utils import create_default_config_manager
from pathlib import Path

# SearchAgent supports video-to-video search
# Note: schema_loader is REQUIRED, config_manager is optional (will create default if None)
schema_loader = FilesystemSchemaLoader(Path("schemas"))
config_manager = create_default_config_manager()
deps = SearchAgentDeps()  # No tenant_id at construction
search_agent = SearchAgent(deps=deps, schema_loader=schema_loader, config_manager=config_manager)

# Video-to-video similarity search (internal method)
# Note: video_bytes must be defined (e.g., from reading a file)
with open("query_video.mp4", "rb") as f:
    video_bytes = f.read()

results = search_agent._search_by_video(
    video_data=video_bytes,
    filename="query_video.mp4",
    tenant_id="acme",
    modality="video",
    top_k=10,
)

Multi-Tenant Search Flow

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_foundation.config.utils import create_default_config_manager

config_manager = create_default_config_manager()
deps = SearchAgentDeps()

# ONE agent serves ALL tenants — tenant_id is per-request
agent = SearchAgent(deps=deps, config_manager=config_manager, schema_loader=schema_loader)

# Tenant A: acme
results_acme = agent.search_by_text("cooking videos", tenant_id="acme")
# Searches schema: video_colpali_smol500_mv_frame_acme
# Only acme's videos returned

# Tenant B: startup (same agent instance)
results_startup = agent.search_by_text("cooking videos", tenant_id="startup")
# Searches schema: video_colpali_smol500_mv_frame_startup
# Only startup's videos returned

# Physical isolation via Vespa schema naming — no cross-tenant data access possible

Gateway: GatewayAgent

Location: libs/agents/cogniverse_agents/gateway_agent.py Purpose: Query entry point — classifies queries as simple or complex using GLiNER entity detection, then routes without an LLM call Base: A2AAgent[GatewayInput, GatewayOutput, GatewayDeps] Port: 8014 (standalone A2A server default; cogniverse_runtime runs GatewayAgent in-process on the runtime's own port, 8000, rather than as a separate service) Target latency: <100ms (no LLM, model inference only)

What It Does

GatewayAgent is the first agent to handle every incoming query. It uses GLiNER zero-shot NER to detect the content modality and generation type, then decides whether the query is simple (direct to a specialist agent) or complex (forward to OrchestratorAgent for multi-step planning). A deterministic MODALITY_KEYWORDS fallback (literal word matching) is merged into _classify_modality() so queries GLiNER misses or where the remote service times out can still route correctly via keyword detection alone. The confidence assigned to a keyword-only match (KEYWORD_MODALITY_CONFIDENCE = 0.5) is above fast_path_confidence_threshold (0.4), so obvious single-modal keyword queries stay on the simple fast path. The runtime response also copies the applied fast_path_confidence_threshold and gliner_threshold into gateway.

Input / Output

class GatewayInput(AgentInput):
    query: str          # User query to classify and route
    tenant_id: Optional[str]  # Per-request tenant identifier

class GatewayOutput(AgentOutput):
    query: str          # Original query (passed through)
    complexity: Literal["simple", "complex"]
    modality: Literal["video", "text", "audio", "image", "document", "both"]
    generation_type: Literal["raw_results", "summary", "detailed_report"]
    routed_to: str      # Target agent name, or "orchestrator_agent"
    confidence: float   # min(modality_confidence, generation_confidence)
    fast_path_confidence_threshold: float  # fast-path threshold snapshot
    gliner_threshold: float  # GLiNER threshold snapshot
    reasoning: str      # Human-readable explanation of the routing decision
    entity_extraction_failed: bool  # True on a GLiNER outage (keyword-only
                                    # classification); distinguishes a sidecar
                                    # outage from genuine low confidence

6 Complexity Signals

_is_complex() returns True (routes to orchestrator) when any of these hold:

# Signal Example
1 Low modality confidence confidence < fast_path_confidence_threshold (default 0.4) — neither GLiNER nor the deterministic MODALITY_KEYWORDS fallback resolved a modality
2 Multiple modalities detected Query spans both video and document content (modality == "both")
3 Generation type is detailed_report Always needs search → analyze → write pipeline
4 Analysis/synthesis verbs present "analyze", "compare", "evaluate", "synthesize", etc.
5 Multi-step markers present "then", "after that", "first...next", "followed by", etc.
6 Compound query (multiple clauses) 3+ commas or 2+ " and " separators in query

Routing Logic

Simple queries are resolved via SIMPLE_ROUTE_MAP — a lookup on (modality, generation_type):

SIMPLE_ROUTE_MAP = {
    ("video", "raw_results"):     "search_agent",
    ("document", "raw_results"):  "document_agent",
    ("audio", "raw_results"):     "audio_analysis_agent",
    ("image", "raw_results"):     "image_search_agent",
    ("video", "summary"):         "summarizer_agent",
    ("video", "detailed_report"): "detailed_report_agent",
    # ... all (modality, generation_type) combinations
}

Complex queries always route to "orchestrator_agent" regardless of modality.

GLiNER Labels

GatewayAgent uses 7 focused labels (experimentally tuned — 7 labels yield average top score 0.56 vs 0.41 with 21 labels):

  • Modality labels: video_content, text_information, audio_content, image_content, document_content
  • Generation labels: summary_request, detailed_report_request

Telemetry

Emits a cogniverse.gateway span for every classified query through the canonical span contract (record_span_io): input.value is the query, operation is gateway, and output.value is a JSON object:

complexity       — "simple" | "complex"
modality         — detected modality
generation_type  — detected generation type
routed_to        — target agent name
confidence       — float 0.0–1.0

It also emits a cogniverse.routing span whose output.value carries the routing decision (chosen_agent, recommended_agent, confidence, reasoning, complexity, modality, generation_type), read by RoutingEvaluator.

Configuration

class GatewayDeps(AgentDeps):
    gliner_model_name: str = "urchade/gliner_large-v2.1"
    gliner_threshold: float = 0.3           # entity detection confidence floor
    gliner_inference_url: Optional[str] = None  # gliner sidecar URL; None = in-process
    fast_path_confidence_threshold: float = 0.4  # min confidence for simple routing

deps = GatewayDeps(fast_path_confidence_threshold=0.4)
gateway = GatewayAgent(deps=deps, port=8014)

Usage Example

from cogniverse_agents.gateway_agent import GatewayAgent, GatewayDeps, GatewayInput

deps = GatewayDeps()
gateway = GatewayAgent(deps=deps)

result = await gateway._process_impl(
    GatewayInput(query="Find videos about transformer architecture", tenant_id="acme")
)
# result.complexity == "simple"
# result.modality == "video"
# result.routed_to == "search_agent"

result = await gateway._process_impl(
    GatewayInput(query="Compare and analyze transformer vs LSTM architectures across papers and videos")
)
# result.complexity == "complex"
# result.routed_to == "orchestrator_agent"
# result.reasoning == "Orchestrator needed: multiple modalities detected; analysis keywords: analyze, compare"

3. OrchestratorAgent (A2A Entry Point)

Location: libs/agents/cogniverse_agents/orchestrator_agent.py Purpose: Central orchestration entry point — plans and executes multi-agent pipelines via A2A Base: MemoryAwareMixin, A2AAgent[OrchestratorInput, OrchestratorOutput, OrchestratorDeps] Port: 8013

Architecture

flowchart TB
    Client["<span style='color:#000'>Client</span>"] -->|HTTP POST /tasks/send| Orchestrator["<span style='color:#000'>OrchestratorAgent<br/>(DSPy Planner)</span>"]

    Orchestrator -->|A2A| QE["<span style='color:#000'>QueryEnhancementAgent</span>"]
    Orchestrator -->|A2A| EE["<span style='color:#000'>EntityExtractionAgent</span>"]
    Orchestrator -->|A2A| PS["<span style='color:#000'>ProfileSelectionAgent</span>"]
    Orchestrator -->|A2A| SA["<span style='color:#000'>SearchAgent</span>"]
    Orchestrator -->|A2A| SU["<span style='color:#000'>SummarizerAgent</span>"]

    style Client fill:#90caf9,stroke:#1565c0,color:#000
    style Orchestrator fill:#ce93d8,stroke:#7b1fa2,color:#000
    style QE fill:#ffcc80,stroke:#ef6c00,color:#000
    style EE fill:#ffcc80,stroke:#ef6c00,color:#000
    style PS fill:#ffcc80,stroke:#ef6c00,color:#000
    style SA fill:#a5d6a7,stroke:#388e3c,color:#000
    style SU fill:#ffcc80,stroke:#ef6c00,color:#000

Implementation

The orchestrator uses DSPy for planning and AgentRegistry for discovery:

from cogniverse_agents.orchestrator_agent import (
    OrchestratorAgent,
    OrchestratorDeps,
    OrchestratorInput,
    OrchestratorOutput,
    close_orchestrator_http_client,
)
from cogniverse_core.registries.agent_registry import AgentRegistry

# Construction — tenant-agnostic, no env vars
registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
deps = OrchestratorDeps()
orchestrator = OrchestratorAgent(
    deps=deps, registry=registry, config_manager=config_manager, port=8013
)

# Processing — tenant_id and session_id arrive per-request
result = await orchestrator._process_impl(
    OrchestratorInput(
        query="Show me machine learning videos",
        tenant_id="acme_corp",
        session_id="sess-uuid",
    )
)

# Result contains: plan_steps, parallel_groups, agent_results, final_output, execution_summary

# Release the current event loop's shared HTTP connection pool at shutdown.
await close_orchestrator_http_client()

Key Features

  • DSPy Planning: OrchestrationModule uses dspy.ChainOfThought to plan agent sequences
  • Parallel Execution: Steps can run in parallel groups (e.g., entity extraction + query enhancement)
  • Agent Discovery: Uses AgentRegistry.find_agents_by_capability() for dynamic agent lookup
  • Agent Calls: Executes agents through one shared httpx.AsyncClient per event loop; close_orchestrator_http_client() removes and closes the current loop's client and raises RuntimeError if shutdown fails
  • Graceful Degradation: Captures agent failures without stopping the pipeline

4. ProfileSelectionAgent

Location: libs/agents/cogniverse_agents/profile_selection_agent.py Purpose: LLM-based intelligent backend profile selection and ensemble composition Base Classes: A2AAgent[ProfileSelectionInput, ProfileSelectionOutput, ProfileSelectionDeps] Port: 8011 (standalone A2A server default; runs in-process on port 8000 in cogniverse_runtime)

Overview

ProfileSelectionAgent uses small language models (SmolLM 3B or Qwen 2.5 3B) via DSPy to intelligently select which backend search profiles to use for a given query. It can recommend single profile searches or ensemble searches with multiple profiles.

Key Capabilities:

  • Analyze query complexity (entities, relationships, keywords)
  • Match query characteristics to profile strengths
  • Recommend single profile or ensemble mode
  • Provide reasoning and confidence scores

Architecture

flowchart TB
    Query["<span style='color:#000'>User Query</span>"] --> ProfileAgent["<span style='color:#000'>ProfileSelectionAgent</span>"]

    ProfileAgent --> Features["<span style='color:#000'>Extract Query Features</span>"]
    Features --> Entities["<span style='color:#000'>Entity Count</span>"]
    Features --> Relationships["<span style='color:#000'>Relationship Count</span>"]
    Features --> Keywords["<span style='color:#000'>Visual/Temporal Keywords</span>"]
    Features --> Length["<span style='color:#000'>Query Length</span>"]

    Entities --> LLM["<span style='color:#000'>SmolLM 3B via DSPy</span>"]
    Relationships --> LLM
    Keywords --> LLM
    Length --> LLM

    LLM --> Decision{"<span style='color:#000'>LLM Reasoning</span>"}
    Decision --> Profiles["<span style='color:#000'>Selected Profiles</span>"]
    Decision --> Confidence["<span style='color:#000'>Confidence Score</span>"]
    Decision --> UseEnsemble["<span style='color:#000'>Use Ensemble Flag</span>"]

    Profiles --> Response["<span style='color:#000'>ProfileSelectionOutput</span>"]
    Confidence --> Response
    UseEnsemble --> Response

    Response --> Orchestrator["<span style='color:#000'>OrchestratorAgent</span>"]

    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style ProfileAgent fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Features fill:#ffcc80,stroke:#ef6c00,color:#000
    style Entities fill:#ffcc80,stroke:#ef6c00,color:#000
    style Relationships fill:#ffcc80,stroke:#ef6c00,color:#000
    style Keywords fill:#ffcc80,stroke:#ef6c00,color:#000
    style Length fill:#ffcc80,stroke:#ef6c00,color:#000
    style LLM fill:#81d4fa,stroke:#0288d1,color:#000
    style Decision fill:#ffcc80,stroke:#ef6c00,color:#000
    style Profiles fill:#a5d6a7,stroke:#388e3c,color:#000
    style Confidence fill:#a5d6a7,stroke:#388e3c,color:#000
    style UseEnsemble fill:#a5d6a7,stroke:#388e3c,color:#000
    style Response fill:#a5d6a7,stroke:#388e3c,color:#000
    style Orchestrator fill:#ce93d8,stroke:#7b1fa2,color:#000

DSPy Signature

# libs/agents/cogniverse_agents/profile_selection_agent.py

import dspy

class ProfileSelectionSignature(dspy.Signature):
    """Select optimal backend profile based on query analysis"""

    # Inputs
    query: str = dspy.InputField(
        desc="User query to analyze"
    )
    available_profiles: str = dspy.InputField(
        desc="Comma-separated list of available profiles"
    )

    # Outputs
    selected_profile: str = dspy.OutputField(
        desc="Best matching profile name"
    )
    confidence: str = dspy.OutputField(
        desc="Confidence score 0.0-1.0"
    )
    reasoning: str = dspy.OutputField(
        desc="Explanation for profile selection"
    )
    query_intent: ProfileQueryIntent = dspy.OutputField(
        desc="Detected intent: multi_modal_search, video_search, image_search, text_search, audio_search, document_search, relationship_aware_search, ensemble_search, code_search, wiki_search"
    )
    modality: str = dspy.OutputField(
        desc="Target modality: audio, code, document, image, text, video, wiki"
    )
    complexity: Literal["simple", "medium", "complex"] = dspy.OutputField(
        desc="Query complexity: simple, medium, complex"
    )

Class Definition

from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from cogniverse_core.agents.base import AgentDeps, AgentInput, AgentOutput
from cogniverse_agents.memory_aware_mixin import MemoryAwareMixin
from pydantic import BaseModel, Field
from typing import List
import dspy

class ProfileCandidate(BaseModel):
    """Candidate profile with score"""
    profile_name: str
    score: float = Field(ge=0.0, le=1.0, description="Confidence score")
    reasoning: str = Field(description="Why this profile was selected")

class ProfileSelectionInput(AgentInput):
    """Input for profile selection."""
    query: str = Field(..., description="Query to analyze")
    available_profiles: Optional[List[str]] = Field(None, description="Available profiles to choose from")

class ProfileSelectionOutput(AgentOutput):
    """Output from profile selection."""
    query: str = Field(..., description="Original query")
    selected_profile: str = Field(..., description="Selected profile")
    confidence: float = Field(0.0, ge=0.0, le=1.0, description="Confidence score")
    reasoning: str = Field("", description="Selection reasoning")
    query_intent: str = Field(
        "",
        description="Detected query intent: multi_modal_search, video_search, image_search, text_search, audio_search, document_search, relationship_aware_search, ensemble_search, code_search, wiki_search",
    )
    modality: str = Field("video", description="Target modality")
    complexity: str = Field("simple", description="Query complexity")
    alternatives: List[ProfileCandidate] = Field(default_factory=list, description="Alternative profiles")

class ProfileSelectionDeps(AgentDeps):
    """Dependencies for profile selection agent (tenant-agnostic at startup)."""
    available_profiles: List[str] = Field(
        default_factory=lambda: [
            "video_colpali_smol500_mv_frame",
            "video_colqwen_omni_mv_chunk_30s",
            "video_xclip_sv_chunk_6s",
            "video_xclip_sv_chunk_6s",
        ],
        description="Default available profiles (must match config.json backend.profiles)",
    )

class ProfileSelectionAgent(
    MemoryAwareMixin,
    A2AAgent[ProfileSelectionInput, ProfileSelectionOutput, ProfileSelectionDeps],
):
    """
    Intelligent profile selection agent using DSPy and small LLMs.

    Uses SmolLM 3B or Qwen 2.5 3B to analyze queries and recommend
    optimal backend profiles for search.
    """

    def __init__(self, deps: ProfileSelectionDeps, port: int = 8011):
        """
        Initialize profile selection agent (tenant-agnostic).

        Args:
            deps: ProfileSelectionDeps with available_profiles config
            port: A2A HTTP server port
        """
        # Create DSPy module
        selection_module = ProfileSelectionModule()

        config = A2AAgentConfig(
            agent_name="profile_selection_agent",
            agent_description="Type-safe profile selection with LLM-based reasoning",
            capabilities=[
                "profile_selection",
                "query_analysis",
                "modality_detection",
                "intent_classification",
                "profile_ranking",
            ],
            port=port,
            version="1.0.0",
        )
        super().__init__(deps=deps, config=config, dspy_module=selection_module)

    async def _process_impl(self, input: ProfileSelectionInput) -> ProfileSelectionOutput:
        """Type-safe profile selection."""
        query = input.query
        profiles = input.available_profiles or self.deps.available_profiles
        profiles_str = ", ".join(profiles) if isinstance(profiles, list) else profiles

        result = self.dspy_module.forward(query=query, available_profiles=profiles_str)
        try:
            confidence = float(result.confidence)
        except (ValueError, AttributeError):
            confidence = 0.5

        return ProfileSelectionOutput(
            query=query,
            selected_profile=result.selected_profile,
            confidence=confidence,
            reasoning=result.reasoning,
            query_intent=result.query_intent,
            modality=result.modality,
            complexity=result.complexity,
            alternatives=[],
        )

Key Methods

_process_impl(input: ProfileSelectionInput) -> ProfileSelectionOutput

Main processing method (required by AgentBase).

async def _process_impl(
    self,
    input: ProfileSelectionInput
) -> ProfileSelectionOutput:
    """
    Process profile selection request with typed input/output.

    Args:
        input: Typed input with query and optional available_profiles

    Returns:
        ProfileSelectionOutput with selected profile and reasoning
    """
    query = input.query
    profiles = input.available_profiles or self.deps.available_profiles

    # Convert profiles list to comma-separated string for DSPy
    profiles_str = ", ".join(profiles) if isinstance(profiles, list) else profiles

    # Select profile using DSPy LLM reasoning
    result = self.dspy_module.forward(query=query, available_profiles=profiles_str)

    # Parse and return typed output
    return ProfileSelectionOutput(
        query=query,
        selected_profile=result.selected_profile,
        confidence=confidence,
        reasoning=result.reasoning,
        query_intent=result.query_intent,
        modality=result.modality,
        complexity=result.complexity,
        alternatives=[],
    )

Configuration

When the runtime handles profile_selection_agent, it derives the candidate set with tenant_usable_profile_names(ConfigManager, tenant_id). A profile is usable only when both halves hold: its embedding service resolves to a URL in SystemConfig.inference_service_urls, and this tenant's schema for it is deployed (read from the schema registry). Neither half alone can serve a search. When nothing is usable the raised error names each held-back profile under missing inference services= or undeployed schemas=. tenant_profile_servability(ConfigManager, tenant_id) returns every configured profile as a ProfileServability(name, profile, state) row in selection order, and servable_tenant_profiles is its servable subset — the set GET /search/profiles advertises. A schema-registry outage raises rather than reporting the tenant's schemas as undeployed. ProfileSelectionDeps.available_profiles is only a standalone fallback for local construction or tests without tenant state.

# Standalone fallback; the runtime injects tenant-usable profiles instead.
deps = ProfileSelectionDeps()

Decision Criteria

Profile Output Fields:

  • selected_profile — single best-matching profile name
  • query_intent — intent classification: text_search, video_search, image_search, etc.
  • modality — target modality: video, image, text, audio
  • complexity — literal query complexity: simple, medium, complex
  • alternatives — top 3 alternative profiles with scores and reasoning

API Usage

A2A Task Message:

{
  "id": "task_001",
  "messages": [
    {
      "role": "user",
      "parts": [
        {
          "type": "text",
          "text": "show me robots playing soccer in tournaments"
        },
        {
          "type": "data",
          "data": {
            "available_profiles": ["colpali", "xclip", "qwen"]
          }
        }
      ]
    }
  ]
}

Response:

{
  "id": "task_001",
  "messages": [
    {
      "role": "assistant",
      "parts": [
        {
          "type": "data",
          "data": {
            "status": "success",
            "agent": "profile_selection_agent",
            "query": "show me robots playing soccer in tournaments",
            "selected_profile": "video_colpali_smol500_mv_frame",
            "confidence": 0.85,
            "reasoning": "Query involves visual content (robots, soccer) requiring visual embedding model. ColPali excels at frame-based visual understanding.",
            "query_intent": "video_search",
            "modality": "video",
            "complexity": "medium",
            "alternatives": [
              {"profile_name": "video_colqwen_omni_mv_chunk_30s", "score": 0.7, "reasoning": "Alternative for video modality"}
            ]
          }
        }
      ]
    }
  ]
}

5. EntityExtractionAgent

Location: libs/agents/cogniverse_agents/entity_extraction_agent.py Purpose: Fast entity and relationship extraction for query enhancement Base Classes: A2AAgent[EntityExtractionInput, EntityExtractionOutput, EntityExtractionDeps] Port: 8010 (standalone A2A server default; runs in-process on port 8000 in cogniverse_runtime)

Overview

EntityExtractionAgent extracts named and unnamed entities and relationships from user queries with DSPy as the primary path. Its signature's entities output is a list[EntityMention]: each item is exactly {text, type}, text a verbatim span of the query and type the EntityType literal — PERSON, ORGANIZATION, CONCEPT, PLACE, EVENT, or TECHNOLOGY (ENTITY_TYPES is the same set). An entity span is its head noun with the modifiers that precede it — adjectives, compound-noun modifiers and numbers, without a leading article — and ends at that head noun: a participial phrase, relative clause or prepositional phrase after it belongs to no entity, and a noun inside one is extracted as its own entity. Role nouns like man, woman, people, and biker map to PERSON, physical things map to CONCEPT, settings map to PLACE, camera and screen map to TECHNOLOGY, and crash maps to EVENT. Named universities, companies, institutions, teams, and agencies map to ORGANIZATION even when the name is also a place name or something happens at them; cities, venues, and a campus or building named as the setting map to PLACE. Teaching sessions map to EVENT and informational resources to CONCEPT as complete noun phrases, typed by their head noun (lecture notes are a CONCEPT resource), retaining descriptive adjectives and compound-noun modifiers while excluding leading articles and everything after the head noun; an unmodified session or resource noun is emitted bare. Named programming languages, libraries, frameworks, and tools map to TECHNOLOGY by their bare name, extracted separately from any session or resource phrase that mentions them; activity words such as programming or training never join the span and never make a subject of study an EVENT. Fields of study and topics map to CONCEPT. A query with no session or resource noun yields no EVENT or resource entity. The prompt scans left to right and emits entities once in source order, without prioritizing proper names or grouping by type. The reasoning field forms complete spans before classification. Signature examples cover span boundaries, session-first, resource-first, subject-only, mixed session-and-subject, bare-resource, and organization-versus-city queries, retaining modifiers and keeping a repeated software subject at its first occurrence. The demo-less prompt must fit the optimizer teacher's input allowance (its context window minus the completion it reserves) for every committed ground-truth query; the teacher raises on a signature that overflows instead of truncating it. When the LM call fails, it falls back to GLiNER NER + SpaCy dependency analysis. If the fallback is unavailable or also fails, the request raises.

EntityExtractionModule binds StructuredJSONAdapter (cogniverse_foundation.dspy) in its own forward, so every caller — the served agent and the optimizer's bootstrap teacher — sends the same prompt and the same response_format: a json_schema derived from the signature's output fields — reasoning a string, entities an array of EntityMention objects whose text (string) and type (enum of the six types) are both required — every object additionalProperties: false, strict: true. The serving engine's guided decoding therefore cannot answer with an object that omits entities, an item without both keys, an extra key, or a type outside the vocabulary; an engine that ignored the schema raises AdapterParseError instead of yielding a silently wrong extraction.

The optimizer persists the compiled module under ("model", "entity_extraction"); its demos carry entities as lists of {text, type} objects. The agent loads an artifact only when its saved instructions and field prefixes/descriptions match the live signature (signature_contract_mismatch); any other artifact is refused with artifact_load_status = "signature_mismatch" and the agent serves the base module until the optimizer compiles one against the live signature.

artifact_load_status is one of no_telemetry, no_artifact, signature_mismatch, loaded, store_unavailable, error. store_unavailable means the artefact store could not answer, so what the tenant has promoted is unknown — distinct from no_artifact, which means the store answered and holds nothing. The dispatcher's per-request overlay (resolve_artefact_for_request) carries the same field: on a store outage it returns {"served_from": "default", "prompts": None, "version": None, "variant_id": ..., "artifact_load_status": "store_unavailable", "variant_lookup_status": ..., "error": ...} so a request served on defaults during an outage is distinguishable from one served on defaults because nothing was promoted. The signature-variant lookup reads the admin config store, which is wired independently, and reports itself in variant_lookup_status: when only that read fails the canary/variant decision still runs on the selection this replica last cached, and the status says the selection is unconfirmed.

Because any DSPy failure falls through to the GLiNER path, the cogniverse.entity_extraction span names which failure it was. entity_extraction.fallback_reason is schema_refused when an AdapterParseError appears anywhere in the cause chain, request_rejected:<status> when the chain carries a RoutedLMCallFailed other than UpstreamUnavailable (the routed LM's own status) or an exception with a 4xx status_code (the engine refused the request — a body it would not accept, a credential, a quota), grounding_failed when the LM answered and no entity it returned is a span of the raw query, and lm_unavailable otherwise — including an UpstreamUnavailable, a timeout or a refused connection, whose synthetic 408 / 500 is never read as a refusal. The exception is in entity_extraction.fallback_error. A span that served the DSPy answer carries neither attribute.

The DSPy path is prompted with the memory-augmented query and grounded against the raw one, so a mention contributed by tenant instructions or remembered context is not a span of the query the caller sent. Each such mention is dropped on its own and the rest of the answer still serves: entity_extraction.grounding_dropped_count holds how many were dropped and entity_extraction.grounding_dropped the text:type of each. Both are absent when nothing was dropped. grounding_failed is the state only when nothing survived.

extractor_unavailable is the same vocabulary's value for GLiNER itself not answering -- an unprovisioned sidecar, a missing dependency, an unreachable inference service. GLiNERRelationshipExtractor.extract_entities raises GLiNEREntityExtractionUnavailableError carrying the model and the inference URL, and the two routing consumers that keep serving on their own fallback carry that value rather than an empty result: ComposableQueryAnalysisModule sets fallback_reason / fallback_model / fallback_inference_url on the prediction it returns, and RelationshipExtractorTool. extract_comprehensive_relationships the same three keys on its result dict. query_structure reports the spaCy parse and is never a failure marker.

Key Capabilities:

  • Primary path: DSPy ChainOfThought entity extraction
  • Fallback path: GLiNER entity extraction + SpaCy relationship extraction
  • Entity type classification using PERSON, ORGANIZATION, CONCEPT, PLACE, EVENT, and TECHNOLOGY
  • Relationship extraction between entities (subject-relation-object triples)
  • Per-entity GLiNER confidence on the fallback path (path_used: "fast") only; DSPy-path entities (path_used: "dspy") have no confidence key
  • Confidence per relationship
  • Dominant entity type detection
  • Typed output enforced by the engine's schema; span validity via entity_is_valid_for_query(text, entity_type, query)
  • Telemetry span emission (cogniverse.entity_extraction)

Architecture

flowchart LR
    Query["<span style='color:#000'>User Query</span>"] --> EntityAgent["<span style='color:#000'>EntityExtractionAgent</span>"]

    EntityAgent --> DSPy["<span style='color:#000'>DSPy ChainOfThought<br/>(primary path)</span>"]
    EntityAgent --> FastPath["<span style='color:#000'>GLiNER + SpaCy<br/>(fallback path)</span>"]

    FastPath --> Entities["<span style='color:#000'>List[Entity]</span>"]
    FastPath --> Rels["<span style='color:#000'>List[Relationship]</span>"]
    DSPy --> ValidateEntities["<span style='color:#000'>_validated_entities()</span>"]
    ValidateEntities --> Entities

    Entities --> Output["<span style='color:#000'>EntityExtractionOutput</span>"]
    Rels --> Output

    Output --> Orchestrator["<span style='color:#000'>OrchestratorAgent</span>"]

    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style EntityAgent fill:#ce93d8,stroke:#7b1fa2,color:#000
    style FastPath fill:#a5d6a7,stroke:#388e3c,color:#000
    style DSPy fill:#81d4fa,stroke:#0288d1,color:#000
    style ValidateEntities fill:#ffcc80,stroke:#ef6c00,color:#000
    style Entities fill:#a5d6a7,stroke:#388e3c,color:#000
    style Rels fill:#a5d6a7,stroke:#388e3c,color:#000
    style Output fill:#a5d6a7,stroke:#388e3c,color:#000
    style Orchestrator fill:#ce93d8,stroke:#7b1fa2,color:#000

Class Definition

import json
import dspy
from pydantic import BaseModel, Field
from typing import Any, Dict, List, Optional

from cogniverse_agents.memory_aware_mixin import MemoryAwareMixin
from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from cogniverse_core.agents.base import AgentDeps, AgentInput, AgentOutput

class Entity(BaseModel):
    """Extracted entity with type and metadata."""
    text: str = Field(description="Entity text as a verbatim span of the query")
    type: str = Field(
        description="Entity type: PERSON, ORGANIZATION, CONCEPT, PLACE, EVENT, or TECHNOLOGY"
    )
    confidence: Optional[float] = Field(
        default=None,
        exclude_if=lambda value: value is None,
        description=(
            "GLiNER score 0-1, set on the fast path only. The DSPy path's "
            "schema carries no score, so its entities serialize without this key"
        ),
    )
    context: str = Field(default="", description="Surrounding context")

EntityType = Literal["CONCEPT", "EVENT", "ORGANIZATION", "PERSON", "PLACE", "TECHNOLOGY"]
ENTITY_TYPES = frozenset(get_args(EntityType))

class EntityMention(BaseModel):
    """One entity as the extraction signature's output schema carries it."""
    model_config = ConfigDict(extra="forbid")
    text: str = Field(description="Verbatim span of the query")
    type: EntityType

class EntityExtractionSignature(dspy.Signature):
    query: str = dspy.InputField(desc="User query to analyze")
    entities: list[EntityMention] = dspy.OutputField(
        desc="Entities in order of first appearance, each a verbatim query span with its type"
    )

class Relationship(BaseModel):
    """Extracted relationship between entities."""
    subject: str = Field(description="Source entity")
    relation: str = Field(description="Relationship type")
    object: str = Field(description="Target entity")
    confidence: float = Field(default=0.5, description="Confidence 0-1")

class EntityExtractionInput(AgentInput):
    """Type-safe input for entity extraction."""
    query: str = Field(..., description="Query to extract entities from")
    tenant_id: Optional[str] = Field(None, description="Tenant identifier")

class EntityExtractionOutput(AgentOutput):
    """Type-safe output from entity extraction."""
    query: str = Field(..., description="Original query")
    entities: List[Entity] = Field(default_factory=list, description="Extracted entities")
    relationships: List[Relationship] = Field(default_factory=list, description="Extracted relationships")
    entity_count: int = Field(0, description="Number of entities found")
    has_entities: bool = Field(False, description="Whether entities were found")
    dominant_types: List[str] = Field(default_factory=list, description="Most common entity types")
    path_used: str = Field("dspy", description="Extraction path: dspy or fast")

class EntityExtractionDeps(AgentDeps):
    """Dependencies for entity extraction agent (tenant-agnostic at startup)."""
    gliner_model_name: Optional[str] = Field(
        None,
        description=(
            "GLiNER model identifier for the fallback path. None resolves to "
            "DEFAULT_GLINER_MODEL in GLiNERRelationshipExtractor."
        ),
    )
    gliner_inference_url: Optional[str] = Field(
        None,
        description=(
            "Optional remote GLiNER service URL (GLiNER inference service). "
            "When set, the fallback path posts to this endpoint instead of "
            "loading gliner in-process — required on slim runtime images."
        ),
    )

class EntityExtractionAgent(
    MemoryAwareMixin,
    A2AAgent[EntityExtractionInput, EntityExtractionOutput, EntityExtractionDeps],
):
    def __init__(self, deps: EntityExtractionDeps, port: int = 8010):
        extraction_module = EntityExtractionModule()
        config = A2AAgentConfig(
            agent_name="entity_extraction_agent",
            agent_description="Type-safe entity extraction from user queries",
            capabilities=[
                "entity_extraction",
                "named_entity_recognition",
                "entity_classification",
                "query_understanding",
            ],
            port=port,
            version="1.0.0",
        )
        super().__init__(deps=deps, config=config, dspy_module=extraction_module)

        # GLiNER + SpaCy for the fallback path.
        self._gliner_extractor = None
        self._spacy_analyzer = None
        self._initialize_extractors()

Key Methods

_process_impl(input: EntityExtractionInput) -> EntityExtractionOutput

Main processing method (required by AgentBase). Uses DSPy primary routing:

  1. DSPy primary path: _extract_dspy_path(prompt_query) runs first and validates the typed mentions the adapter parsed.
  2. GLiNER + SpaCy fallback: _extract_fast_path(query) runs only when the LM call raises. It extracts entities with GLiNER and relationships with SpaCy.

If the fallback is unavailable or also raises, the request fails with an error that names both failures. The path_used output is "dspy" for the primary path and "fast" for the fallback path.

After extraction, emits a cogniverse.entity_extraction telemetry span through the canonical span contract (record_span_io): input.value is the query, output.value is a JSON object {entities, relationships, entity_count, relationship_count, path_used}, and operation is entity_extraction. The entity-extraction optimizer reads the (query -> entities) training pair back via read_span_io.

_validated_entities(mentions: List[EntityMention], query: str) -> List[Entity]

Turns the schema-typed mentions into served entities. The schema already fixes each mention's keys and type; this checks what it cannot:

  • text (stripped) must be non-empty and appear in the query, case-insensitively (entity_is_valid_for_query); the served text is the query's own characters for that span
  • a repeated (text.casefold(), type) pair keeps its first mention
  • entities are ordered by where their span starts in the query, an enclosing span before a shorter one starting at the same place
  • dropped mentions are logged with the offending text and type
  • confidence is unset, so the served entity has no confidence key: the DSPy path produces no score

The agent never fabricates a replacement span for invalid output.

Configuration

EntityExtractionDeps has two optional fields that control the GLiNER fallback path:

  • gliner_model_name — GLiNER model identifier (None resolves to the default in GLiNERRelationshipExtractor).
  • gliner_inference_url — URL of the remote GLiNER inference service. When set, the fallback path POSTs to it instead of loading GLiNER in-process — required on slim runtime images that omit the heavy torch stack. The canonical server is cogniverse_cli.modal_inference.servers.gliner.

The DSPy primary path is scoped per-call via dspy.context(lm=...) using create_dspy_lm() from the centralized llm_config.

from cogniverse_foundation.config.llm_factory import create_dspy_lm

# Remote GLiNER inference service (production):
deps = EntityExtractionDeps(gliner_inference_url="http://cogniverse-gliner:8080")
agent = EntityExtractionAgent(deps=deps, port=8010)

# In-process GLiNER (development/testing):
deps = EntityExtractionDeps()
agent = EntityExtractionAgent(deps=deps, port=8010)

# DSPy LM is scoped per-call, not configured globally:
lm = create_dspy_lm(llm_endpoint_config)
with dspy.context(lm=lm):
    result = agent.extraction_module(query="machine learning videos")

API Usage

A2A Task Message:

{
  "id": "task_002",
  "messages": [
    {
      "role": "user",
      "parts": [
        {
          "type": "text",
          "text": "Python programming with TensorFlow for deep learning"
        }
      ]
    }
  ]
}

Response:

{
  "id": "task_002",
  "messages": [
    {
      "role": "assistant",
      "parts": [
        {
          "type": "data",
          "data": {
            "status": "success",
            "agent": "entity_extraction_agent",
            "query": "Python programming with TensorFlow for deep learning",
            "entities": [
              {"text": "Python", "type": "TECHNOLOGY", "context": "Python programming with TensorFlow f"},
              {"text": "TensorFlow", "type": "TECHNOLOGY", "context": "Python programming with TensorFlow for deep learning"},
              {"text": "deep learning", "type": "CONCEPT", "context": "ogramming with TensorFlow for deep learning"}
            ],
            "relationships": [],
            "entity_count": 3,
            "has_entities": true,
            "dominant_types": ["TECHNOLOGY", "CONCEPT"],
            "path_used": "dspy"
          }
        }
      ]
    }
  ]
}

Only the GLiNER fallback scores entities: with "path_used": "fast" every entity also carries its GLiNER confidence (0-1).


OrchestratorAgent (Multi-Agent Workflow Detail)

Location: libs/agents/cogniverse_agents/orchestrator_agent.py Purpose: Multi-agent workflow coordination with planning and action phases Base Classes: MemoryAwareMixin, A2AAgent[OrchestratorInput, OrchestratorOutput, OrchestratorDeps] Port: 8013

Overview

OrchestratorAgent coordinates complex multi-agent workflows by dividing execution into two phases: 1. Planning Phase: Parallel execution of ProfileSelectionAgent and EntityExtractionAgent 2. Action Phase: Sequential or parallel execution of SearchAgent based on planning results

Key Capabilities:

  • Two-phase workflow orchestration (planning → action)
  • Parallel agent execution in planning phase
  • Agent discovery via AgentRegistry
  • Result aggregation and metadata tracking
  • Multi-tenant support

Architecture

sequenceDiagram
    participant User
    participant Orch as OrchestratorAgent
    participant Registry as AgentRegistry
    participant Prof as ProfileSelectionAgent
    participant Entity as EntityExtractionAgent
    participant Search as SearchAgent

    User->>Orch: POST /tasks/send<br/>{query: "robots playing soccer"}

    Note over Orch: PLANNING PHASE (Parallel)

    par Profile Selection
        Orch->>Registry: GET /agents/by-capability/profile_selection
        Registry-->>Orch: [ProfileSelectionAgent @ :8011]
        Orch->>Prof: POST /tasks/send
        Prof-->>Orch: {selected_profile, query_intent, modality, confidence}
    and Entity Extraction
        Orch->>Registry: GET /agents/by-capability/entity_extraction
        Registry-->>Orch: [EntityExtractionAgent @ :8010]
        Orch->>Entity: POST /tasks/send
        Entity-->>Orch: {entities, entity_count, has_entities, dominant_types}
    end

    Note over Orch: Planning Complete (~150-200ms)
    Note over Orch: ACTION PHASE

    Orch->>Registry: GET /agents/by-capability/search
    Registry-->>Orch: [SearchAgent @ :8002]

    Orch->>Search: POST /tasks/send<br/>{query, profile, modality}
    Search-->>Orch: {results, metadata}

    Orch->>Orch: Aggregate results + metadata
    Orch-->>User: {results, planning_time, search_time}

Class Definition

from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from cogniverse_core.agents.base import AgentDeps, AgentInput, AgentOutput
from pydantic import Field
from typing import Any, Dict, Optional
import asyncio

class OrchestratorInput(AgentInput):
    """Type-safe input for orchestration."""
    query: str = Field(..., description="Query to orchestrate")
    tenant_id: str = Field(..., description="Tenant identifier (per-request, required)")
    session_id: Optional[str] = Field(default=None, description="Session identifier (per-request)")

class OrchestratorOutput(AgentOutput):
    """Type-safe output from orchestration."""
    query: str = Field(..., description="Original query")
    workflow_id: str = Field("", description="Unique workflow identifier")
    plan_steps: List[Dict[str, Any]] = Field(default_factory=list, description="Orchestration plan steps")
    parallel_groups: List[List[int]] = Field(default_factory=list, description="Parallel execution groups")
    plan_reasoning: str = Field("", description="Plan reasoning")
    agent_results: Dict[str, Any] = Field(default_factory=dict, description="Results from each agent")
    final_output: Dict[str, Any] = Field(default_factory=dict, description="Aggregated final output")
    execution_summary: str = Field("", description="Summary of execution")
    metadata: Dict[str, Any] = Field(default_factory=dict, description="Side-channel structured metadata")

class OrchestratorDeps(AgentDeps):
    """Dependencies for orchestrator agent (tenant-agnostic at startup)."""
    pass

class OrchestratorAgent(
    MemoryAwareMixin,
    A2AAgent[OrchestratorInput, OrchestratorOutput, OrchestratorDeps],
):
    """
    Multi-agent workflow orchestrator with planning and action phases.

    Coordinates agents via AgentRegistry discovery and A2A protocol.
    tenant_id and session_id arrive per-request in task payload.
    """

    def __init__(
        self,
        deps: OrchestratorDeps,
        registry: "AgentRegistry",
        config_manager: "ConfigManager" = None,
        port: int = 8013,
    ):
        """
        Initialize orchestrator agent.

        Args:
            deps: OrchestratorDeps (infrastructure only, no tenant_id)
            registry: AgentRegistry for dynamic agent discovery (REQUIRED)
            config_manager: ConfigManager for tenant memory initialization
            port: A2A HTTP server port
        """
        self.registry = registry

        # Initialize DSPy module
        orchestration_module = OrchestrationModule()

        config = A2AAgentConfig(
            agent_name="orchestrator_agent",
            agent_description="Type-safe orchestration with planning and action phases",
            capabilities=["orchestration", "planning", "multi_agent_coordination",
                          "parallel_execution", "result_aggregation"],
            port=port,
        )
        super().__init__(deps=deps, config=config, dspy_module=orchestration_module)

        # Memory initialized lazily per-tenant on first request
        self._memory_initialized_tenants: set = set()

    async def _process_impl(self, input: OrchestratorInput) -> OrchestratorOutput:
        """Type-safe orchestration — tenant_id/session_id from input."""
        # Lazily initialize memory for this tenant
        if hasattr(input, "tenant_id") and input.tenant_id:
            self._ensure_memory_for_tenant(input.tenant_id)

        result = await self.orchestrate(input.query)
        return OrchestratorOutput(
            query=result.query,
            plan_steps=result.plan_steps,
            parallel_groups=result.parallel_groups,
            plan_reasoning=result.plan_reasoning,
            agent_results=result.agent_results,
            final_output=result.final_output,
            execution_summary=result.execution_summary,
        )

Key Methods

_process_impl(input: Union[OrchestratorInput, Dict]) -> OrchestratorOutput

Main processing method (required by AgentBase).

async def _process_impl(
    self,
    input: Union[OrchestratorInput, Dict[str, Any]]
) -> OrchestratorOutput:
    """
    Process orchestration request with typed input/output.

    Args:
        input: Typed input with query field (or dict)

    Returns:
        OrchestratorOutput with plan, agent results, and final output
    """
    # Handle dict input for backward compatibility with tests
    if isinstance(input, dict):
        input = self.validate_input(input)

    query = input.query

    # Planning, with the agents that can serve this tenant
    plan = await self._create_plan(
        query, available_agents=await self._planning_agents(tenant_id)
    )

    # Action
    agent_results = await self._execute_plan(plan)

    # Aggregate results
    final_output = self._aggregate_results(query, agent_results)

    # Generate summary
    execution_summary = self._generate_summary(plan, agent_results)

    return OrchestratorOutput(
        query=query,
        plan_steps=[...],
        parallel_groups=plan.parallel_groups,
        plan_reasoning=plan.reasoning,
        agent_results=agent_results,
        final_output=final_output,
        execution_summary=execution_summary
    )

_planning_agents(tenant_id) -> List[str]

The registered agents a plan for the tenant may use. A retrieval agent — one whose configured modalities and capabilities make it the retrieval agent for a modality (search_agent for video, image_search_agent, audio_analysis_agent, document_agent) — is left out when none of the modalities it retrieves is served by the tenant's servable profiles (servable_tenant_profiles, mapped through MODALITY_PROFILE_TYPES). A video-only tenant is therefore never planned an image or audio search. A store or registry outage raises RuntimeError naming the tenant instead of planning against a guess.

_create_plan(query, conversation_context, gateway_context, *, available_agents) -> OrchestrationPlan

Planning Phase: the DSPy planner proposes an agent sequence from available_agents (the _planning_agents list).

async def _create_plan(
    self,
    query: str,
    conversation_context: str = "",
    gateway_context: str = "",
    *,
    available_agents: List[str],
) -> OrchestrationPlan:
    registered_agents = list(available_agents)

    result = await self.call_dspy(
        self.dspy_module,
        output_field="agent_sequence",
        query=query,
        available_agents=", ".join(registered_agents),
        conversation_context=conversation_context,
        gateway_context=gateway_context,
    )

    # Normalize the LM output into an executable plan:
    #   - "name" / "name_agent" aliases resolve to the offered name
    #   - an agent not offered is dropped and listed in unavailable_agents
    #   - a repeated agent keeps its first step only (results are keyed
    #     by agent name)
    #   - parallel_steps / dependency indices are remapped from the raw
    #     sequence to the surviving step positions
    #   - an empty result falls back to a single search_agent step when
    #     search_agent was offered
    ...
    return OrchestrationPlan(
        query=query,
        steps=steps,
        parallel_groups=parallel_groups,
        reasoning=result.reasoning,
        unavailable_agents=unavailable_agents,
    )

_execute_plan(plan: OrchestrationPlan) -> Dict[str, Any]

Action Phase: Execute orchestration plan with parallel execution support. Emits streaming progress events per step and checks for cancellation between steps.

async def _execute_plan(
    self,
    plan: OrchestrationPlan,
    tenant_id: str,
    workflow_id: str = "",
    session_id: Optional[str] = None,
) -> Dict[str, Any]:
    agent_results = {}
    executed = [False] * len(plan.steps)

    while not all(executed):
        # Check for cancellation
        if workflow_id and workflow_id in self._cancelled_workflows:
            break

        ready_steps = [
            (i, step) for i, step in enumerate(plan.steps)
            if not executed[i]
            and all(executed[dep] for dep in step.depends_on)
        ]

        async def execute_step(step_index: int, step: AgentStep):
            self.emit_progress("executing", f"Step {step_index}: {step.agent_name}")
            agent_endpoint = self.registry.get_agent(step.agent_name)
            if not agent_endpoint:
                return step.agent_name, {"status": "error", "message": "not found"}

            agent_input = step.input_data.copy()
            agent_input["tenant_id"] = tenant_id
            async with httpx.AsyncClient(timeout=60.0) as http_client:
                response = await http_client.post(
                    f"{agent_endpoint.url}/process",
                    json={"query": agent_input.pop("query", ""), **agent_input},
                )
                response.raise_for_status()
            self.emit_progress("step_complete", f"Step {step_index} complete")
            return step.agent_name, response.json()

        results = await asyncio.gather(
            *[execute_step(idx, step) for idx, step in ready_steps]
        )
        for (step_idx, _), (name, result) in zip(ready_steps, results):
            agent_results[name] = result
            executed[step_idx] = True

    return agent_results

_aggregate_results(query: str, agent_results: Dict) -> Dict[str, Any]

Cross-modal fusion of results from all agents. aggregated_content — the answer a reader sees — is built from each step's own answer text: the answer the serving runtime stamps on every completed dispatch (a search step's includes its hits, each by title and time range), else the step's message; never a step's payload. A search step's line naming the query it searched (the plan's rewrite) is replaced by one naming query, the question as asked. Enrichment steps (query_enhancement_agent, entity_extraction_agent, profile_selection_agent) feed later steps and are fused only when the plan produced nothing else. The fusion strategy (SCORE_BASED, TEMPORAL, HIERARCHICAL, or SIMPLE) is selected from the query and the fused steps' modalities: SIMPLE leads with the synthesized answers (a summary, a report) and follows with the search steps' hits, each group in execution order; SCORE_BASED puts each under **<Agent>** (<modality>, confidence <share>), highest first; HIERARCHICAL groups them under ## <Modality> results headings.

status is the orchestration's outcome and the one every consumer reads: failed when no step produced an answer, partial when a step failed or reported a partial answer of its own, success otherwise. Failed steps keep their entry in results but are never fused into aggregated_content. The dispatch envelope and _dspy_to_a2a_output carry this status through, and harness_turn treats failed as terminal (no answer text) while partial renders the answer the completed steps produced. Success memory is written only for success, in the background after the response returns; a write still pending when the process is killed is lost.

def _aggregate_results(
    self, query: str, agent_results: Dict[str, Any]
) -> Dict[str, Any]:
    # Fuse only the steps that answered; keep every step in ``results``
    answered = {
        name: result for name, result in agent_results.items()
        if result.get("status") not in FAILED_STEP_STATUSES
    }
    if not answered:
        return {
            "query": query, "status": FAILED_STATUS,
            "message": "No orchestration step completed successfully",
            "results": ..., "aggregated_content": "", "fusion_quality": ...,
        }
    fusion_strategy = self._select_fusion_strategy(query, agent_modalities)

    # Dispatch to fusion method
    if fusion_strategy == FusionStrategy.SCORE_BASED:
        fused = self._fuse_by_score(task_results)
    elif fusion_strategy == FusionStrategy.HIERARCHICAL:
        fused = self._fuse_hierarchically(task_results, agent_modalities)
    else:
        fused = self._fuse_simple(task_results)

    return {
        "query": query,
        "status": PARTIAL_STATUS if degraded else "success",
        "results": ...,
        "fusion_strategy": fusion_strategy.value,
        "fusion_quality": ...,
        "aggregated_content": fused["content"],
    }

_generate_summary(plan: OrchestrationPlan, agent_results: Dict) -> str

Generate execution summary.

def _generate_summary(
    self,
    plan: OrchestrationPlan,
    agent_results: Dict[str, Any]
) -> str:
    """Generate execution summary"""
    executed_steps = len(agent_results)
    successful_steps = sum(
        1 for result in agent_results.values()
        if not (isinstance(result, dict) and result.get("status") == "error")
    )
    total_steps = len(plan.steps)

    return (
        f"Executed {executed_steps}/{total_steps} steps "
        f"({successful_steps} successful). "
        f"Plan: {plan.reasoning}"
    )

Each plan-then-act request exports one checked cogniverse.orchestration span. agent_sequence is the planned sequence, while execution_order records the order in which concurrent agent calls actually completed. tasks_completed counts only results that did not return status="error"; success is true only when every planned agent completed without an error. Empty queries return their validation response without emitting a malformed training span. Plans with a duplicate agent name are rejected before execution because agent results are keyed by agent name. Planning, dispatch, aggregation, and deep-synthesis failures export the partial observed order with success=false, OpenTelemetry ERROR status, and a non-empty error_summary capped at 512 characters. A deep-synthesis result counts as successful only when it submitted a non-empty answer. Unexpected failures are re-raised after the checked export. The runtime must inject a telemetry manager; an unavailable exporter is a request failure because these outcomes are the training source for orchestration evaluation.

Workflow Types

The orchestrator supports different workflow patterns:

1. Simple Query Flow (Single profile, no planning):

User Query → EntityExtraction → Search (single profile) → Results

2. Complex Query Flow (Ensemble with planning):

User Query → [ProfileSelection + EntityExtraction] (parallel) → Search (ensemble) → RRF Fusion → Results

3. Sequential Dependencies:

User Query → Planning Phase → Action Phase → Post-processing → Results

Configuration

# Orchestrator agent configuration
{
    "orchestrator_agent": {
        "port": 8013,
        "agent_registry_url": "http://localhost:8000",
        "planning_timeout_seconds": 10.0,
        "action_timeout_seconds": 30.0,
        "enable_parallel_planning": True,
        "enable_result_caching": False
    }
}

API Usage

A2A Task Message:

{
  "id": "task_003",
  "messages": [
    {
      "role": "user",
      "parts": [
        {
          "type": "text",
          "text": "show me robots playing soccer in tournaments"
        },
        {
          "type": "data",
          "data": {
            "modality": "video",
            "top_k": 20
          }
        }
      ]
    }
  ]
}

Response:

{
  "id": "task_003",
  "messages": [
    {
      "role": "assistant",
      "parts": [
        {
          "type": "data",
          "data": {
            "results": [
              {"document_id": "video_001", "relevance_score": 0.92, ...},
              {"document_id": "video_002", "relevance_score": 0.87, ...}
            ],
            "metadata": {
              "planning_time_ms": 185,
              "action_time_ms": 620,
              "total_time_ms": 805,
              "selected_profile": "video_colpali_smol500_mv_frame",
              "entities": [...],
              "dominant_types": ["CONCEPT"],
              "confidence": 0.85
            }
          }
        }
      ]
    }
  ]
}

6. SearchAgent (Ensemble Mode)

Location: libs/agents/cogniverse_agents/search_agent.py Enhancement: Added ensemble search with RRF fusion

Overview

SearchAgent was enhanced to support ensemble mode, allowing it to query multiple backend profiles in parallel and fuse results using Reciprocal Rank Fusion (RRF).

New Capabilities:

  • Parallel profile execution (2-3 profiles)
  • RRF score calculation and fusion
  • Ensemble metadata tracking
  • Supports both single-profile and multi-profile (ensemble) mode

_search_ensemble returns an EnsembleOutcome: the fused hits, searched (the profiles whose search ran) and degraded ((profile, reason) for a leg that could not encode its query or whose search raised, with reason in encode_failed / search_failed; a leg is encode_failed when its search raised a typed encoder fault). Profiles sharing an embedding model share one SharedQueryEncoder, so one encode, and a leg encodes only when its ranking strategy needs embeddings. SearchOutput.profiles names the legs that ran and SearchOutput.degraded_profiles the ones that did not, so a caller reports a partial ensemble as partial. Every leg failing raises. The legs run on a pool of their own that is shut down without waiting, so a caller whose budget cancels the ensemble gets control back at once while a leg blocked on its encoder or backend finishes in its worker thread; the event loop never waits on a leg.

Both search paths rewrite the query once. _rewrite_query_for_search runs the DSPy rewrite (SearchOptimizationSignature) before the mode branch, so a one-profile search and an ensemble of the same question reach the backend with the same query, and an ensemble fans that single rewrite out to every leg instead of rewriting per leg. A rewrite the orchestrator already made (SearchInput.enhanced_query) is used as it stands. SearchInput.query_rewrite_timeout_s bounds the whole rewrite step and is the deadline its LM call carries, so no request reaches the router after the search has moved on; a rewrite that fails or overruns it searches the original query and names itself on SearchOutput.degraded_query_rewrite as query_rewrite_failed or query_rewrite_timed_out rather than raising. A rewrite whose LM endpoint answered 404 (nothing deployed) is query_rewrite_lm_not_serving: until the endpoint is rechecked (lm_endpoint_availability), every later search skips the rewrite step, context injection included, without a round trip, so an undeployed LLM costs one 404 per recheck window instead of one per search. The search span carries llm.endpoint.state, llm.endpoint.failed_fast and llm.endpoint.recheck_in_s, and the warning names the endpoint. A dispatched search sets the bound to dispatched_query_rewrite_budget_s of the configured grounding search budget: the smaller of QUERY_REWRITE_BUDGET_S (2.2 times the measured p95 of a rewrite through the semantic router's classification entry on a fresh upstream connection) and what that budget leaves after the retrieval reserve. The current span carries enhancement.path: lm when the rewrite served, heuristic_fallback when the search ran on the original query, and the failure log names the router model the call was bound to. A dispatched search reports both under the envelope's query_rewrite block (see docs/modules/runtime.md).

Ensemble Architecture

flowchart TB
    Input["<span style='color:#000'>Query</span>"] --> Encode{"<span style='color:#000'>For Each Profile</span>"}

    Encode -->|Profile 1| Enc1["<span style='color:#000'>ColPali Encoder</span>"]
    Encode -->|Profile 2| Enc2["<span style='color:#000'>X-CLIP Encoder</span>"]
    Encode -->|Profile 3| Enc3["<span style='color:#000'>Qwen Encoder</span>"]

    Enc1 --> Search1["<span style='color:#000'>Vespa Search<br/>Profile: colpali</span>"]
    Enc2 --> Search2["<span style='color:#000'>Vespa Search<br/>Profile: xclip</span>"]
    Enc3 --> Search3["<span style='color:#000'>Vespa Search<br/>Profile: qwen</span>"]

    Search1 --> Results1["<span style='color:#000'>Results 1<br/>Ranked by ColPali</span>"]
    Search2 --> Results2["<span style='color:#000'>Results 2<br/>Ranked by X-CLIP</span>"]
    Search3 --> Results3["<span style='color:#000'>Results 3<br/>Ranked by Qwen</span>"]

    Results1 --> RRF["<span style='color:#000'>RRF Fusion<br/>score = Σ 1/(k+rank)</span>"]
    Results2 --> RRF
    Results3 --> RRF

    RRF --> Sort["<span style='color:#000'>Sort by RRF Score</span>"]
    Sort --> TopN["<span style='color:#000'>Select Top N</span>"]
    TopN --> Rerank["<span style='color:#000'>MultiModalReranker</span>"]
    Rerank --> Final["<span style='color:#000'>Fused Results</span>"]

    style Input fill:#90caf9,stroke:#1565c0,color:#000
    style Encode fill:#ffcc80,stroke:#ef6c00,color:#000
    style Enc1 fill:#81d4fa,stroke:#0288d1,color:#000
    style Enc2 fill:#81d4fa,stroke:#0288d1,color:#000
    style Enc3 fill:#81d4fa,stroke:#0288d1,color:#000
    style Search1 fill:#90caf9,stroke:#1565c0,color:#000
    style Search2 fill:#90caf9,stroke:#1565c0,color:#000
    style Search3 fill:#90caf9,stroke:#1565c0,color:#000
    style Results1 fill:#a5d6a7,stroke:#388e3c,color:#000
    style Results2 fill:#a5d6a7,stroke:#388e3c,color:#000
    style Results3 fill:#a5d6a7,stroke:#388e3c,color:#000
    style RRF fill:#ffcc80,stroke:#ef6c00,color:#000
    style Sort fill:#ffcc80,stroke:#ef6c00,color:#000
    style TopN fill:#ffcc80,stroke:#ef6c00,color:#000
    style Rerank fill:#ffcc80,stroke:#ef6c00,color:#000
    style Final fill:#a5d6a7,stroke:#388e3c,color:#000

Ensemble Result Shape

SearchOutput.results carries public-shaped dicts. Identity and ranking sit at the top level; everything else is payload under metadata:

{
    "id": "v_-6dz6tBH77I_seg_0",
    "document_id": "id:content:video_colpali_smol500_mv_frame_acme_acme::v_-6dz6tBH77I_seg_0",
    "score": 3.04,               # backend relevance for the profile that matched
    "metadata": {...},           # remaining backend fields
    "temporal_info": {"start_time": 0.0, "end_time": 5.0},  # when both are present
    # ensemble mode only:
    "rrf_score": 0.0333,         # the value the result set is ORDERED by
    "profile_ranks": {"video_colpali_smol500_mv_frame": 0, "video_colqwen_omni_mv_chunk_30s": 1},
    "profile_scores": {"video_colpali_smol500_mv_frame": 3.04, "video_colqwen_omni_mv_chunk_30s": 1.1},
    "num_profiles": 2,
}

score and rrf_score are different quantities on different scales: score is one profile's backend relevance, rrf_score is the fused rank. Sort order follows rrf_score, so clients that re-sort should use it. A single-profile search carries none of the ensemble keys.

Key Methods (New/Enhanced)

_search_ensemble(query, profiles, modality, limit) -> List[SearchResult] (Internal)

Execute ensemble search across multiple profiles.

async def _search_ensemble(
    self,
    query: str,
    profiles: List[str],
    modality: str = "video",
    top_k: int = 10,
    rrf_k: int = 60
) -> List[Dict[str, Any]]:
    """
    Execute ensemble search with RRF fusion.

    Args:
        query: Search query
        profiles: List of profile names (2-3 profiles)
        modality: Content modality
        top_k: Number of results to return
        rrf_k: RRF constant (default: 60)

    Returns:
        Fused and reranked search results
    """
    # Execute searches in parallel
    profile_results = await self._execute_parallel_searches(
        query=query,
        profiles=profiles,
        modality=modality,
        top_k=top_k * 2  # Get more results for fusion
    )

    # Fuse results with RRF
    fused_results = self._fuse_results_rrf(
        profile_results=profile_results,
        k=rrf_k,
        top_k=top_k
    )

    # Optional: Rerank with MultiModalReranker
    if self.enable_reranking:
        fused_results = await self._rerank(fused_results, query)

    return fused_results

_execute_parallel_searches(query, profiles) -> Dict[str, List[SearchResult]]

Execute searches across all profiles concurrently.

async def _execute_parallel_searches(
    self,
    query: str,
    profiles: List[str],
    modality: str,
    limit: int
) -> Dict[str, List[SearchResult]]:
    """
    Execute searches in parallel with connection pooling.
    """
    tasks = [
        self._search_single(query, profile, modality, limit)
        for profile in profiles
    ]

    results = await asyncio.gather(*tasks, return_exceptions=True)

    # Build profile -> results mapping
    profile_results = {}
    for profile, result in zip(profiles, results):
        if isinstance(result, Exception):
            logger.warning(f"Profile {profile} failed: {result}")
            profile_results[profile] = []
        else:
            profile_results[profile] = result

    return profile_results

_fuse_results_rrf(profile_results, k, limit) -> List[SearchResult]

Fuse results using Reciprocal Rank Fusion.

def _fuse_results_rrf(
    self,
    profile_results: Dict[str, List[SearchResult]],
    k: int = 60,
    limit: int = 20
) -> List[SearchResult]:
    """
    Fuse results from multiple profiles using RRF.

    Algorithm:
    For each document across all profiles:
        RRF_score = Σ_profiles (1 / (k + rank_in_profile))

    Complexity: O(n_profiles × n_results) ~ 5-10ms typical
    """
    rrf_scores = {}
    doc_objects = {}

    # Calculate RRF scores
    for profile, results in profile_results.items():
        for rank, result in enumerate(results, start=1):
            doc_id = result.document_id

            # Accumulate RRF score
            rrf_scores[doc_id] = rrf_scores.get(doc_id, 0) + 1 / (k + rank)

            # Store document object (first occurrence)
            if doc_id not in doc_objects:
                doc_objects[doc_id] = result

    # Sort by RRF score (descending)
    sorted_docs = sorted(
        rrf_scores.items(),
        key=lambda x: x[1],
        reverse=True
    )

    # Return top results with updated scores
    return [
        doc_objects[doc_id]._replace(
            relevance_score=rrf_score,
            metadata={
                **doc_objects[doc_id].metadata,
                "rrf_score": rrf_score,
                "fusion_method": "rrf",
                "profiles_used": list(profile_results.keys())
            }
        )
        for doc_id, rrf_score in sorted_docs[:limit]
    ]

Configuration

# Search agent ensemble configuration
{
    "search_agent": {
        "ensemble_config": {
            "rrf_k": 60,
            "max_profiles": 3,
            "parallel_timeout": 5.0,
            "enable_reranking": True,
            "min_overlap": 0.1
        }
    }
}

Performance Characteristics

Configuration Latency Quality (NDCG@10) Notes
Single profile 400-600ms 0.72 Baseline
Ensemble (2 profiles) 500-700ms 0.78 +100-150ms overhead
Ensemble (3 profiles) 550-750ms 0.83 +150-200ms overhead
RRF fusion 5-10ms N/A Negligible overhead

Key Insight: Parallel execution keeps ensemble latency close to single-profile (not 2x or 3x).

API Usage

Single Profile Mode:

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps

# Note: schema_loader is REQUIRED, config_manager is optional
deps = SearchAgentDeps()  # tenant_id is per-request, not in deps
search_agent = SearchAgent(deps=deps, schema_loader=schema_loader, config_manager=config_manager)

# Use search_by_text for text queries (tenant_id required)
results = search_agent.search_by_text(
    query="robots playing soccer",
    tenant_id="acme",
    modality="video",
    top_k=20,
)

Ensemble Mode (Internal):

Ensemble search is handled internally via _search_ensemble(). The SearchAgent automatically uses ensemble when multiple profiles are configured:

# Ensemble is triggered internally based on configuration
# The _search_ensemble method uses RRF fusion across profiles
# See search_agent.py:691 for implementation details

7. DetailedReportAgent

Location: libs/agents/cogniverse_agents/detailed_report_agent.py

Generates comprehensive detailed reports with visual and technical analysis. Includes VLM integration and a "thinking phase" for complex queries.

flowchart TD
    Query["<span style='color:#000'>Query + Search Results</span>"] --> Think["<span style='color:#000'>Thinking Phase</span>"]
    Think --> VLM["<span style='color:#000'>VLM Visual Analysis</span>"]
    Think --> Tech["<span style='color:#000'>Technical Analysis</span>"]

    VLM --> Merge["<span style='color:#000'>Merge Insights</span>"]
    Tech --> Merge

    Merge --> Summary["<span style='color:#000'>Executive Summary</span>"]
    Merge --> Findings["<span style='color:#000'>Detailed Findings</span>"]
    Merge --> Recs["<span style='color:#000'>Recommendations</span>"]

    Summary --> Report["<span style='color:#000'>Final Report</span>"]
    Findings --> Report
    Recs --> Report

    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style Think fill:#ffcc80,stroke:#ef6c00,color:#000
    style VLM fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Tech fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Merge fill:#ffcc80,stroke:#ef6c00,color:#000
    style Summary fill:#a5d6a7,stroke:#388e3c,color:#000
    style Findings fill:#a5d6a7,stroke:#388e3c,color:#000
    style Recs fill:#a5d6a7,stroke:#388e3c,color:#000
    style Report fill:#a5d6a7,stroke:#388e3c,color:#000

Type Signature:

class DetailedReportAgent(
    MemoryAwareMixin,
    A2AAgent[DetailedReportInput, DetailedReportOutput, DetailedReportDeps],
)

Input Fields: | Field | Type | Description | |-------|------|-------------| | query | str | Query for report generation | | search_results | List[Dict] | Results to analyze | | report_type | str | Type: comprehensive, technical, analytical | | include_visual_analysis | bool | Include VLM visual analysis | | include_technical_details | bool | Include technical breakdown | | include_recommendations | bool | Include actionable recommendations | | max_results_to_analyze | int | Maximum results to process (default: 20) |

Output Fields: | Field | Type | Description | |-------|------|-------------| | executive_summary | str | High-level summary | | detailed_findings | List[Dict] | Detailed analysis results | | visual_analysis | List[Dict] | VLM visual insights | | technical_details | List[Dict] | Technical breakdown | | recommendations | List[str] | Actionable recommendations: the report LM writes one per line and each line, less a bullet or number, is one item, commas and parentheses included | | confidence_assessment | Dict[str, float] | Per-dimension confidence scores (keys: overall, data_quality, completeness, visual_analysis, technical_analysis) | | thinking_process | Dict | Thinking phase details | | metadata | Dict | Additional metadata |

Usage:

config_manager has a default of None in the signature but is required — the constructor raises ValueError if it is not supplied:

from cogniverse_agents.detailed_report_agent import (
    DetailedReportAgent,
    DetailedReportDeps,
    DetailedReportInput,
)
from cogniverse_foundation.config.utils import create_default_config_manager

config_manager = create_default_config_manager()
deps = DetailedReportDeps()
agent = DetailedReportAgent(deps=deps, config_manager=config_manager, port=8004)
result = await agent.process(DetailedReportInput(
    query="Analyze video content about machine learning",
    search_results=search_results,
    report_type="comprehensive",
    include_visual_analysis=True,
))
print(result.executive_summary)
# Access overall confidence:
print(result.confidence_assessment.get("overall", 0.0))

Multimodal generation (keyframe injection):

By default the report LLM sees only text (titles, scores). When multimodal_generation_enabled is set, the agent also attaches the top-K retrieved keyframes to the LLM so it grounds the report in the frames themselves rather than their filenames. Each keyframe is located through the shared keyframe-key contract (cogniverse_core.common.media.keyframe_uri) — the same function the ingestion write side uses, so the read and write paths cannot drift — fetched via MediaLocator, and passed as a keyframes: list[dspy.Image] input on ReportGenerationSignature.

flowchart LR
    Hit["<span style='color:#000'>Search hit<br/>source_url + video_id + segment_id</span>"] --> URI["<span style='color:#000'>keyframe_uri()</span>"]
    URI --> Loc["<span style='color:#000'>MediaLocator.localize (s3://)</span>"]
    Loc --> Img["<span style='color:#000'>dspy.Image<br/>bounded encode cache</span>"]
    Img --> LM["<span style='color:#000'>ReportGenerationSignature<br/>keyframes: list[dspy.Image]</span>"]

    style Hit fill:#90caf9,stroke:#1565c0,color:#000
    style URI fill:#ffcc80,stroke:#ef6c00,color:#000
    style Loc fill:#81d4fa,stroke:#0288d1,color:#000
    style Img fill:#ce93d8,stroke:#7b1fa2,color:#000
    style LM fill:#a5d6a7,stroke:#388e3c,color:#000

Dependencies (DetailedReportDeps): | Dep | Type | Default | Description | |-----|------|---------|-------------| | multimodal_generation_enabled | bool | True | Attach available retrieved keyframes and request images to the report LLM. | | max_keyframes_to_llm | int | 4 | Cap on keyframes attached per report. |

A keyframe not yet in object storage (or one that fails to fetch) is silently skipped — generation degrades to the text-only path, never errors. Encoded frames are cached in the agent by s3:// key so repeated reports over the same clips reuse the encoding.

max_keyframes_to_llm caps the count; fit_answer_images caps the bytes. One 768 px frame is a few hundred KB of base64, and the semantic-router Envoy in front of every agent LM call buffers the whole request body for its ext_proc routing decision, answering 413 above per_connection_buffer_limit_bytes. The agents keep frames in rank order while they fit multimodal.MAX_IMAGE_PAYLOAD_BYTES (LLM_REQUEST_BODY_LIMIT_BYTES less the non-image reserve) and shed whole frames from the tail; keyframes_attached and keyframes_shed report the outcome in the answer metadata. The chart renders the same number into the Envoy listener (semanticRouter.envoy.maxRequestBytes), so raising one without the other only moves where the request is refused.


8. DocumentAgent

Location: libs/agents/cogniverse_agents/document_agent.py

Document analysis and search with dual strategy support (visual + text). Uses ColPali for visual document understanding and traditional text extraction for semantic search.

flowchart LR
    Query["<span style='color:#000'>Query</span>"] --> Strategy{"<span style='color:#000'>Strategy<br/>Selection</span>"}

    Strategy -->|Visual| ColPali["<span style='color:#000'>ColPali<br/>Page-as-Image</span>"]
    Strategy -->|Text| TextEmbed["<span style='color:#000'>Text Extraction<br/>+ Embeddings</span>"]
    Strategy -->|Hybrid| Both["<span style='color:#000'>Both Strategies</span>"]
    Strategy -->|Auto| AutoSelect["<span style='color:#000'>Auto-detect<br/>Best Strategy</span>"]

    ColPali --> Vespa["<span style='color:#000'>Vespa Search</span>"]
    TextEmbed --> Vespa
    Both --> Vespa
    AutoSelect --> Vespa

    Vespa --> Fusion["<span style='color:#000'>Result Fusion</span>"]
    Fusion --> Results["<span style='color:#000'>Document Results</span>"]

    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style Strategy fill:#ffcc80,stroke:#ef6c00,color:#000
    style ColPali fill:#ce93d8,stroke:#7b1fa2,color:#000
    style TextEmbed fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Both fill:#ce93d8,stroke:#7b1fa2,color:#000
    style AutoSelect fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Vespa fill:#90caf9,stroke:#1565c0,color:#000
    style Fusion fill:#ffcc80,stroke:#ef6c00,color:#000
    style Results fill:#a5d6a7,stroke:#388e3c,color:#000

Type Signature:

class DocumentAgent(MemoryAwareMixin, A2AAgent[DocumentSearchInput, DocumentSearchOutput, DocumentAgentDeps])

Strategies: - Visual (ColPali): Treats document pages as images for visual understanding - Text: Traditional text extraction + semantic embeddings - Hybrid: Combines both strategies with fusion - Auto: Automatically selects best strategy based on query type

Input Fields: | Field | Type | Description | |-------|------|-------------| | query | str | Search query | | tenant_id | Optional[str] | Tenant whose project the process span is recorded in | | strategy | str | Strategy: visual, text, hybrid, auto | | limit | int | Number of results (default: 20) |

Output Fields: | Field | Type | Description | |-------|------|-------------| | results | List[DocumentResult] | Search results with page info | | count | int | Total result count | | span_id | Optional[str] | Id of the DocumentAgent.process span, which records the query, modality document and one row per hit (document_id, score, content, title); None when telemetry is off |

The runtime dispatches a document search through process, with its telemetry manager attached, and its envelope carries span_id, so a client rates the hits against that span (POST /ag-ui/results/relevance).

Usage:

Unlike SearchAgent/OrchestratorAgent, DocumentAgent reads deps.tenant_id in __init__ — tenant_id must be passed to DocumentAgentDeps at construction time, not per-request:

from cogniverse_agents.document_agent import DocumentAgent, DocumentAgentDeps, DocumentSearchInput

deps = DocumentAgentDeps(tenant_id="acme")
agent = DocumentAgent(deps=deps)
result = await agent._process_impl(DocumentSearchInput(
    query="quarterly financial report",
    strategy="hybrid",
    limit=10
))
for doc in result.results:
    print(f"{doc.title} - Page {doc.page_number} ({doc.strategy_used})")

9. ImageSearchAgent

Location: libs/agents/cogniverse_agents/image_search_agent.py

Image similarity search using ColPali multi-vector embeddings. Uses the same approach as video frame search.

flowchart LR
    Query["<span style='color:#000'>Text Query</span>"] --> Encode["<span style='color:#000'>ColPali Encoder</span>"]
    Encode --> Search["<span style='color:#000'>Vespa Image Search</span>"]
    Search --> Objects["<span style='color:#000'>Object Detection</span>"]
    Objects --> Filter["<span style='color:#000'>Visual Filters</span>"]
    Filter --> Results["<span style='color:#000'>Image Results</span>"]

    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style Encode fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Search fill:#90caf9,stroke:#1565c0,color:#000
    style Objects fill:#ffcc80,stroke:#ef6c00,color:#000
    style Filter fill:#ffcc80,stroke:#ef6c00,color:#000
    style Results fill:#a5d6a7,stroke:#388e3c,color:#000

Type Signature:

class ImageSearchAgent(A2AAgent[ImageSearchInput, ImageSearchOutput, ImageSearchDeps])

Capabilities:

  • Image similarity search using ColPali embeddings
  • Hybrid search (BM25 text + ColPali semantic)
  • Object and scene detection in results
  • Visual filtering support

Input Fields:

Field Type Description
query str Search query
search_mode str Mode: semantic, hybrid
limit int Number of results (default: 20)
visual_filters Dict Optional visual filters

Output Fields:

Field Type Description
results List[ImageResult] Image search results
count int Total result count

Search by image content:

Method Description
search_images(query, search_mode="semantic", limit=20, visual_filters=None) Text query encoded via query_encoder.encode
search_by_image(image_bytes, limit=20) Decodes raw bytes to a PIL image and runs image-to-image search; raises ValueError on empty or undecodable bytes rather than returning an empty result
find_similar_images(reference_image, limit=20) Same path from an already-decoded PIL image

_encode_image delegates to query_encoder.encode_image, so the query image is embedded by the deployed encoder (the image profile's sidecar via QueryEncoderFactory) — the same model that embedded the stored images, which is what makes MaxSim between them meaningful. Query patches are not padded; MaxSim ranks over the real patch count, as on the text side.

A chat client that sends a photo reaches this through the runtime: the gateway base64-encodes the photo into the dispatch context as media_content_b64, and AgentDispatcher._execute_image_search_task(query, tenant_id, top_k, image_b64=...) decodes it and calls search_by_image. Caption text is not mixed into an image query — the image profile's ranking expresses one query tensor, so a caption would silently override the visual match.


10. SummarizerAgent

Location: libs/agents/cogniverse_agents/summarizer_agent.py

Intelligent summarization of search results with VLM visual content analysis and a "think phase" for complex queries.

flowchart TD
    Results["<span style='color:#000'>Search Results</span>"] --> Think["<span style='color:#000'>Think Phase</span>"]
    Query["<span style='color:#000'>Query Context</span>"] --> Think

    Think --> VLM["<span style='color:#000'>VLM Analysis</span>"]
    Think --> Text["<span style='color:#000'>Text Processing</span>"]

    VLM --> Extract["<span style='color:#000'>Extract Key Points</span>"]
    Text --> Extract

    Extract --> Gen["<span style='color:#000'>Generate Summary</span>"]
    Gen --> Output["<span style='color:#000'>Summary Output</span>"]

    style Results fill:#90caf9,stroke:#1565c0,color:#000
    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style Think fill:#ffcc80,stroke:#ef6c00,color:#000
    style VLM fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Text fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Extract fill:#ffcc80,stroke:#ef6c00,color:#000
    style Gen fill:#ffcc80,stroke:#ef6c00,color:#000
    style Output fill:#a5d6a7,stroke:#388e3c,color:#000

Type Signature:

class SummarizerAgent(
    MemoryAwareMixin,
    A2AAgent[SummarizerInput, SummarizerOutput, SummarizerDeps],
)

Summary Types:

  • brief - Short, concise summary
  • comprehensive - Detailed summary with context
  • bullet_points - Key points as bullet list

Input Fields:

Field Type Description
query str Query for context
search_results List[Dict] Results to summarize
summary_type str Type: brief, comprehensive, bullet_points
include_visual_analysis bool Include VLM insights
max_results_to_analyze int Max results to process (default: 10)

Output Fields:

Field Type Description
summary str Generated summary
key_points List[str] Extracted key points
visual_insights List[str] VLM visual insights
confidence_score float Summary confidence
thinking_process Dict Thinking phase details
metadata Dict Additional metadata

Multimodal generation (keyframe injection):

Like the DetailedReportAgent, the summarizer can attach the top-K retrieved keyframes to its answer LLM via the shared KeyframeImageResolver and a keyframes: list[dspy.Image] input on SummaryGenerationSignature (see DetailedReportAgent for the shared mechanism and the keyframe-key contract). Gated behind SummarizerDeps.multimodal_generation_enabled (bool, default True) with max_keyframes_to_llm (int, default 4); a keyframe not yet in object storage is silently skipped (text-only fallback). Frames over the request-body allowance are shed and counted in metadata.keyframes_shed, alongside metadata.keyframes_attached.


Summary, detailed-report, and deep-research inputs accept attachments as image URIs. Worker threads prepare them with attachments_to_images, retaining input order before retrieved keyframes and applying max_keyframes_to_llm and the request-body allowance to the combined list via fit_answer_images. validate_attachments(input) rejects disabled visual input before streaming with attachments_disabled: visual inputs are disabled for this request. Summaries answer supplied text or images even when retrieval returns no hits. Attachment failures appear in summary/research degradation metadata and in the report's report_degraded_reason. Summary key points and report recommendations remain paired with their generated text.

11. AudioAnalysisAgent

Location: libs/agents/cogniverse_agents/audio_analysis_agent.py

Audio search and analysis using Whisper transcription. Supports transcript, semantic, and acoustic search.

flowchart LR
    Audio["<span style='color:#000'>Audio File</span>"] --> Whisper["<span style='color:#000'>Whisper<br/>Transcription</span>"]
    Whisper --> Transcript["<span style='color:#000'>Transcript Text</span>"]
    Whisper --> Acoustic["<span style='color:#000'>Acoustic Features</span>"]

    Query["<span style='color:#000'>Text Query</span>"] --> Mode{"<span style='color:#000'>Search Mode</span>"}

    Mode -->|transcript| TranscriptSearch["<span style='color:#000'>Transcript Search</span>"]
    Mode -->|semantic| SemanticSearch["<span style='color:#000'>Semantic Search</span>"]
    Mode -->|hybrid| HybridSearch["<span style='color:#000'>Hybrid Search</span>"]
    Mode -->|acoustic| AcousticSearch["<span style='color:#000'>Acoustic Search</span>"]

    Transcript --> TranscriptSearch
    Transcript --> SemanticSearch
    Transcript --> HybridSearch
    Acoustic --> AcousticSearch

    TranscriptSearch --> Results["<span style='color:#000'>Audio Results</span>"]
    SemanticSearch --> Results
    HybridSearch --> Results
    AcousticSearch --> Results

    style Audio fill:#90caf9,stroke:#1565c0,color:#000
    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style Whisper fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Transcript fill:#ffcc80,stroke:#ef6c00,color:#000
    style Acoustic fill:#ffcc80,stroke:#ef6c00,color:#000
    style Mode fill:#ffcc80,stroke:#ef6c00,color:#000
    style TranscriptSearch fill:#90caf9,stroke:#1565c0,color:#000
    style SemanticSearch fill:#90caf9,stroke:#1565c0,color:#000
    style HybridSearch fill:#90caf9,stroke:#1565c0,color:#000
    style AcousticSearch fill:#90caf9,stroke:#1565c0,color:#000
    style Results fill:#a5d6a7,stroke:#388e3c,color:#000

Type Signature:

class AudioAnalysisAgent(A2AAgent[AudioSearchInput, AudioSearchOutput, AudioAnalysisDeps])

Search Modes:

  • semantic - Default. ColBERT (phased_semantic) over transcript text
  • transcript - Lexical BM25 (transcript_search) over audio_title and audio_transcript
  • hybrid - Every clip ranked by ColBERT binary MaxSim plus nativeRank of audio_title and audio_transcript (hybrid_semantic_bm25)
  • acoustic - CLAP text-to-audio similarity over acoustic_embedding

Text modes go through the shared search backend with type="audio"; the backend selects the tenant's audio profile and encodes the query on demand. An unknown mode raises ValueError.

find_similar_audio(..., similarity_type="semantic") transcribes the reference clip and reuses transcript search. Acoustic similarity encodes the reference clip with CLAP and searches acoustic_embedding.

Input Fields:

Field Type Description
query str Search query
search_mode str Mode: transcript, semantic, acoustic, hybrid
limit int Number of results (default: 20)

Output Fields:

Field Type Description
results List[AudioResult] Audio search results
count int Total result count

AudioResult Fields:

  • audio_id, audio_url, title
  • transcript - Full transcript text
  • duration - Duration in seconds
  • speaker_labels - Detected speakers
  • detected_events - Audio events
  • language - Detected language
  • metadata - Raw backend fields preserved on the result

Transcription (transcribe_audio):

transcribe_audio(audio_url) resolves the URL through MediaLocator to a local path, decodes it to 16 kHz mono and POSTs it multipart to {whisper_endpoint}/v1/audio/transcriptions (OpenAI-compatible vLLM Whisper), each chunk of at most 30 s with timestamps and, unless the timed segments run to the end of the chunk, without, through cogniverse_core.common.models.whisper_transcription: the text comes from the untimed answer, timed by the timed answer's segments; an empty or looping answer is asked again at the next sampling temperature, and a chunk whose untimed text never comes back usable raises EmptyTranscriptError or GarbledTranscriptError. The merged answer is mapped to a TranscriptionResult(text, segments, language, confidence). The pinned vLLM response represents duration as a non-negative decimal string; the agent validates that exact wire type and converts it to seconds before checking segment bounds. A segment time past duration (Whisper times text into the padding after short audio) is clamped to it and logged at DEBUG with the original value.

whisper_endpoint and whisper_model are fields on AudioAnalysisDeps. The runtime populates whisper_endpoint from system_config.inference_service_urls['vllm_asr']. How the pod is deployed and which model it serves are chart concerns — see docs/operations/setup-installation.md and the inference.vllm_asr block in charts/cogniverse/values.yaml.

Modal Whisper URLs obtain their bearer credential only from COGNIVERSE_INFERENCE_API_KEY and reject whisper_headers; non-Modal endpoints may accept one canonical Authorization: Bearer ... mapping. AudioAnalysisDeps validates HTTPS and authentication before any request, and the agent snapshots the resolved headers as one immutable mapping for every multipart request. An authentication error, timeout, or serving failure is raised to the caller without exposing the credential or producing an empty transcript.

Remote acoustic-query encoding uses AudioAnalysisDeps.clap_endpoint. Modal URLs obtain their bearer credential from COGNIVERSE_INFERENCE_API_KEY and reject clap_headers; non-Modal endpoints may accept one canonical Authorization: Bearer ... mapping. AudioEmbeddingGenerator resolves the credential before constructing its pooled HTTP client. Per-keyframe face extraction obtains Modal credentials from COGNIVERSE_INFERENCE_API_KEY; non-Modal endpoints may pass a canonical bearer mapping through extract_faces_per_keyframe(..., headers=...). The internally owned HTTP client uses the resolved immutable headers for every concurrent keyframe request.

extract_faces_per_keyframe reads the keyframes KeyframeProcessor wrote, processing_results["keyframes"]["keyframes"], each {frame_number, timestamp, filename, path}. A keyframe's segment id is its index in that list, the id of its content document and of the transcript segment aligned to it. Each request reads the image at path and sends it as image_b64, with at most max_concurrency requests in flight (default FACE_EMBED_MAX_CONCURRENCY, which the chart sets to the sidecar's CPU count) and a 120 s budget per request. It returns a FaceExtraction: mentions, the faces sorted by (segment_id, bbox), and failed, a FailedKeyframe (segment_id, cause) for each keyframe whose image could not be read or whose request failed on its retry; the other keyframes' faces are kept. A keyframe without a path or timestamp raises ValueError, and RuntimeError names every keyframe's cause when all of them failed.


12. TextAnalysisAgent

Location: libs/agents/cogniverse_agents/text_analysis_agent.py

Text analysis agent with runtime-configurable DSPy modules. Supports dynamic reconfiguration of modules and optimizers via REST API.

Mixins Used: - DynamicDSPyMixin - Runtime DSPy module switching - ConfigAPIMixin - REST API for configuration - HealthCheckMixin - Health monitoring - TenantAwareAgentMixin - Multi-tenancy

Analysis Types: - sentiment - Sentiment analysis - summary - Text summarization - entities - Entity extraction

Usage:

from cogniverse_agents.text_analysis_agent import TextAnalysisAgent
from cogniverse_foundation.config.utils import create_default_config_manager

config_manager = create_default_config_manager()
agent = TextAnalysisAgent(tenant_id="acme", config_manager=config_manager)

result = agent.analyze_text(text="Your text here", analysis_type="sentiment")
print(f"Result: {result}")


13. SearchAgent (Refactored)

Refactored video search agent using the unified search service architecture. Provides a simpler interface compared to the original SearchAgent.

Constructor:

SearchAgent(
    deps: SearchAgentDeps,
    schema_loader=None,     # REQUIRED
    config_manager=None,
    port: int = 8002,
)

Methods:

def search_by_text(
    query: str,
    *,
    tenant_id: str,             # Required per-request
    modality: str = "video",
    top_k: int = 10,
    **kwargs,
) -> List[Dict[str, Any]]

Usage:

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_foundation.config.utils import create_default_config_manager
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from pathlib import Path

config_manager = create_default_config_manager()
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))
deps = SearchAgentDeps()
agent = SearchAgent(deps=deps, config_manager=config_manager, schema_loader=schema_loader)

results = agent.search_by_text("cooking tutorial", tenant_id="acme", top_k=10)
for result in results:
    print(result)


14. QueryEnhancementAgent

Location: libs/agents/cogniverse_agents/query_enhancement_agent.py

Enhances user queries by adding synonyms, context, and related terms to improve search recall. It takes the sampled source text alongside the query, and the DSPy prompt requires expansion_terms to be token-grounded in that text while synonyms remain free-form. Every non-stopword alphanumeric token in an expansion term must appear in the sampled source text; multi-word phrases are allowed when each substantive token is grounded. Runs a DSPy QueryEnhancementModule and returns the enhanced query alongside the expansion terms it generated.

The cogniverse.query_enhancement span always includes enhancement.path in the declared span contract. lm marks a genuine enhancement; heuristic_fallback marks the heuristic used when the LM call fails, echoes the query or leaves enhanced_query blank (blank expansion_terms is a valid LM answer). The heuristic only spells out an acronym the query contains and otherwise returns the query unchanged, so it never adds words that change what is searched.

Constructor: QueryEnhancementAgent(deps: QueryEnhancementDeps, port: int = 8012) (standalone A2A server default; runs in-process on port 8000 in cogniverse_runtime)

Usage:

from cogniverse_agents.query_enhancement_agent import (
    QueryEnhancementAgent,
    QueryEnhancementDeps,
    QueryEnhancementInput,
)

agent = QueryEnhancementAgent(deps=QueryEnhancementDeps())
output = await agent.process(
    QueryEnhancementInput(
        query="find cooking videos",
        source_text="cooking videos source text",
        tenant_id="acme",
    )
)
print(output.enhanced_query, output.expansion_terms)


15. CodingAgent

Location: libs/agents/cogniverse_agents/coding_agent.py

Iterative code generation agent: searches code semantically via the code_lateon_mv Vespa profile (LateOn-Code-edge multi-vector embeddings with tree-sitter AST chunking), plans an implementation with DSPy, generates code, executes it in a sandbox, evaluates the result, and iterates up to max_iterations times.

Advertised client tools select workspace mode. The agent suspends with pending_tool_calls and continuation_state; replayed observations match calls by ID. WORKSPACE_MAX_ROUNDS is the shared workspace default and limit (8); an explicit max_iterations can lower it. WORKSPACE_ACTION_MAX_ATTEMPTS bounds malformed-action retries (3). Each step either calls one advertised tool (its JSON arguments; the step's summary may be blank) or finishes with a summary of the completed work; a finish with a blank summary is a malformed action and is retried. Failed steps expose success=False and error with an empty summary. Context reads run in workers; sandbox staging directories are removed on success, failure, and cancellation.

Constructor:

CodingAgent(
    deps: CodingDeps,
    config: A2AAgentConfig | None = None,
    search_fn: Any = None,
    sandbox_manager: Any = None,
    config_manager=None,
)

Usage:

from cogniverse_agents.coding_agent import CodingAgent, CodingDeps, CodingInput

agent = CodingAgent(deps=CodingDeps(tenant_id="acme"))
output = await agent.process(
    CodingInput(task="add input validation to the login form", tenant_id="acme")
)
print(output.plan, output.summary)


16. DeepResearchAgent

Location: libs/agents/cogniverse_agents/deep_research_agent.py

Multi-step research agent: decomposes a complex query into sub-questions, dispatches parallel searches for each, evaluates evidence sufficiency, refines and re-searches if gaps remain, then synthesizes a cited report.

Constructor:

DeepResearchAgent(
    deps: DeepResearchDeps,
    config: A2AAgentConfig | None = None,
    search_fn: Any = None,
    config_manager=None,
)

Usage:

from cogniverse_agents.deep_research_agent import (
    DeepResearchAgent,
    DeepResearchDeps,
    DeepResearchInput,
)

agent = DeepResearchAgent(deps=DeepResearchDeps(tenant_id="acme"))
output = await agent.process(
    DeepResearchInput(query="how did the outage affect checkout latency?", tenant_id="acme")
)
print(output.summary, output.citations)

Through the runtime, POST /agents/deep_research_agent/process forwards context.max_iterations into DeepResearchInput.max_iterations; a request without it runs on the field default (3).

Progress and cancellation: through the runtime a run is a workflow task (context.workflow_id, or a new id the result's workflow_id names), streamed from /events/workflows/{workflow_id}. The agent reports decompose, then search and evaluate each iteration, synthesize and rlm_synthesis (AgentBase.report_phase); its RLM synthesis reports on the same task. A cancellation stops the run at its next phase boundary, or after the RLM synthesis, and the dispatch answers {"status": "cancelled", "workflow_id", "message"}.

Multimodal generation (keyframe injection):

_synthesize flattens the nested evidence hits (evidence[i]["results"]) and, via the shared KeyframeImageResolver, attaches the top-K retrieved keyframes to a keyframes: list[dspy.Image] input on SynthesisSignature (see DetailedReportAgent for the shared mechanism). Gated behind DeepResearchDeps.multimodal_generation_enabled (bool, default True) with max_keyframes_to_llm (int, default 4); an evidence hit lacking the source_url/video_id/segment_id fields (or whose keyframe isn't in object storage yet) is silently skipped.


Agent Architecture

ConfigManager Injection

ConfigManagerAware (libs/core/cogniverse_core/agents/base.py) is the base of both AgentBase and MemoryAwareMixin, so every agent holds the injected ConfigManager in one slot: bind_config_manager(cm) writes it, the config_manager property reads it. Agents whose constructor takes config_manager bind it there; agents the runtime builds from deps alone are bound by the dispatcher and the knowledge router right after construction. bind_config_manager(None) and a read before any bind both raise AgentConfigurationError naming the agent — a missing manager is a construction bug, never a fall back to the process singleton, which would serve tenant instructions and per-tenant LM routing from a different config store. A build site that reads config before the agent exists refuses the same way through require_config_manager(config_manager, owner=<agent class name>), which bind_config_manager itself applies.

Type-Safe A2AAgent Base Class with Generics

All agents extend A2AAgent[InputT, OutputT, DepsT] from cogniverse_core, providing compile-time type safety:

# libs/core/cogniverse_core/agents/a2a_agent.py

from cogniverse_core.agents.base import AgentBase, AgentDeps, AgentInput, AgentOutput
from typing import Generic, TypeVar
import dspy

InputT = TypeVar("InputT", bound=AgentInput)
OutputT = TypeVar("OutputT", bound=AgentOutput)
DepsT = TypeVar("DepsT", bound=AgentDeps)
class A2AAgentConfig(BaseModel):
    """Configuration for A2A agents."""
    agent_name: str
    agent_description: str
    capabilities: list[str] = []
    port: int = 8000
    version: str = "1.0.0"
class A2AAgent(AgentBase[InputT, OutputT, DepsT]):
    """
    Type-safe base class that bridges DSPy 3.0 modules with A2A protocol.

    Architecture:
    - Generic Types: InputT, OutputT, DepsT for compile-time type checking
    - DSPy 3.0 Core: Advanced AI capabilities and optimization
    - Pydantic Validation: Automatic input/output validation

    Features:
    - Type-safe process(input: InputT) -> OutputT method
    - IDE autocomplete and type checking support
    - Multi-tenant support via tenant_id arriving per-request in A2A task payload
    - Multi-modal support (text, images, video, audio)

    Note: HTTP server concerns (A2A protocol endpoints) are handled by the
    runtime's A2AStarletteApplication (official a2a-sdk). This class holds
    configuration metadata and the optional DSPy module only.
    """

    def __init__(
        self,
        deps: DepsT,
        config: A2AAgentConfig,
        dspy_module: Optional[dspy.Module] = None,
    ):
        """
        Initialize type-safe A2A agent.

        Args:
            deps: Agent dependencies (tenant-agnostic at startup)
            config: A2AAgentConfig with name, description, etc.
            dspy_module: Optional DSPy 3.0 module
        """
        super().__init__(deps=deps)

        # A2A Protocol Configuration
        self.config = config
        self.dspy_module = dspy_module

    @abstractmethod
    async def _process_impl(self, input: InputT) -> OutputT:
        """
        Type-safe processing method.

        Must be implemented by subclass. IDE provides autocomplete
        for both input fields and return type.
        """
        pass

Key Benefits of Type-Safe Architecture:

  • Generic Types: A2AAgent[InputT, OutputT, DepsT] enables IDE autocomplete
  • Pydantic Validation: Input/output automatically validated at runtime
  • Tenant-Agnostic Startup: AgentDeps has no tenant_id — agents start without tenant context; tenant_id arrives per-request in the A2A task payload
  • Abstract _process_impl(): Clear contract with type-safe signature
  • A2A Protocol: HTTP server endpoints handled by the runtime's A2AStarletteApplication (official a2a-sdk), not embedded in each agent

MemoryAwareMixin

Location: libs/agents/cogniverse_agents/memory_aware_mixin.py

Provides memory integration for all agents:

class MemoryAwareMixin:
    """
    Mixin for agent memory integration via Mem0.

    Provides:
    - Memory initialization per tenant
    - Context retrieval
    - Success/failure recording
    """

    def initialize_memory(
        self,
        agent_name: str,
        tenant_id: str,
        embedder_base_url: str,      # Required — OpenAI-compatible /v1/embeddings endpoint
        *,
        llm_model: str,              # Required — no 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,
        schema_loader=None,
        backend_config_port: Optional[int] = None,
        auto_create_schema: bool = True,
    ) -> bool:
        """
        Initialize memory for agent.

        Creates tenant-specific Mem0MemoryManager instance.
        Raises ValueError if tenant_id is empty or None. A failed init returns
        False and leaves the instance's memory state as it was.
        """
        ...  # Implementation in cogniverse_agents.memory_aware_mixin

    def get_relevant_context(self, query: str, top_k: Optional[int] = 5) -> Optional[str]:
        """Retrieve relevant memories for query. Returns None if memory not initialized."""
        ...

    def inject_context_into_prompt(self, prompt: str, query: str) -> str:
        """Enrich a prompt with tenant instructions, strategies, and memory context.

        The tenant-instruction fetch status ("loaded" / "none" / "unavailable")
        is recorded on last_tenant_instructions_status and as the span attribute
        enrichment.tenant_instructions, so a ConfigStore outage is
        distinguishable from a tenant with no instructions.
        """
        ...

    def update_memory(self, content: str, metadata: Optional[Dict[str, Any]] = None) -> bool:
        """Add content to agent's memory. Returns success status."""
        ...

    def write_memory_in_background(self, write: Callable[..., Any], /, *args: Any, **kwargs: Any) -> bool:
        """Queue write(*args, **kwargs) on the shared background memory writer
        (cogniverse_agents.background_memory_writes). False when its queue is
        full and the write was dropped."""
        ...

    def remember_success(self, query: str, result: Any, metadata: Optional[Dict[str, Any]] = None) -> bool:
        """Remember a successful interaction."""
        ...

    def remember_failure(self, query: str, error: str, metadata: Optional[Dict[str, Any]] = None) -> bool:
        """Remember a failed interaction to avoid repeating mistakes."""
        ...

    def is_memory_enabled(self) -> bool:
        """Check if memory is enabled and initialized."""
        ...

    def get_memory_summary(self) -> Dict[str, Any]:
        """Get summary of memory status (enabled, agent_name, tenant_id, initialized)."""
        ...

Wiki and graph backend resolution

WikiManager(backend_resolver=...) and GraphManager(backend_resolver=...) take a zero-arg callable resolving the tenant's backend through BackendRegistry and lease it per operation; the runtime's factories hand down that resolver, so a manager cached for the process never holds an instance the registry may close.

Document graph extraction

Location: libs/agents/cogniverse_agents/graph/doc_extractor.py

DocExtractor uses an explicitly injected remote GLiNER endpoint when one is provided: GraphManager injects the already validated SystemConfig.inference_service_urls["gliner"] value alongside the ColBERT endpoint, and both production graph factories reject a missing endpoint before building a manager. Built without an explicit URL, the extractor resolves the endpoint from validated system configuration at first model load; only when no endpoint is configured anywhere does it load the model in-process, which requires an environment that installs the local gliner stack. No extractor parses environment JSON or recognizes a separate GLINER_INFERENCE_URL value. Entity extraction is all-or-error for each source: if GLiNER fails on any chunk, the call raises with the one-based chunk number and source document ID, chained from the inference error. The claim pass runs only for transcript, document, and code modalities; VLM and OCR segments still get entity extraction and back-refs, but they do not trigger claim extraction. It never reports a partial knowledge graph as a successful extraction.

GraphBindableMixin

Location: libs/agents/cogniverse_agents/graph_bindable.py

Binds a single GraphManager to a KG-aware agent. The seven single-graph KG agents — CitationTracingAgent, TemporalReasoningAgent, AuditExplanationAgent, KnowledgeGraphTraversalAgent, KnowledgeSummarizationAgent, MultiDocumentSynthesisAgent, ContradictionReconciliationAgent — mix this in to share one setter and one guard instead of redeclaring them:

class GraphBindableMixin:
    """Bind a single GraphManager to a KG-aware agent."""

    _graph_manager: Optional["GraphManager"] = None

    def set_graph_manager(self, graph_manager: "GraphManager") -> None:
        """Bind the GraphManager this agent reads Node/Edge rows from."""
        self._graph_manager = graph_manager

    def _require_graph_manager(self, method: str) -> "GraphManager":
        """Return the bound GraphManager or raise, naming the calling method."""
        ...

Usage in an agent (mixed in ahead of MemoryAwareMixin so the binding API sits at the front of the MRO):

class CitationTracingAgent(
    GraphBindableMixin,
    MemoryAwareMixin,
    A2AAgent[CitationTracingInput, CitationTracingOutput, CitationTracingDeps],
):
    def trace(self, claim_id: str) -> Dict[str, Any]:
        graph_manager = self._require_graph_manager("trace")
        edge_fields = graph_manager.get_edge_by_id(claim_id)
        ...

Agents that bind multiple managers (FederatedQueryAgent, CrossTenantComparisonAgent) keep their own plural set_graph_managers and do not use this mixin.

Binding at dispatch (complementary to Mem0). set_graph_manager is called on every request path, not only in tests: agent_dispatcher._bind_graph_manager binds it on the orchestrator-routing path, and routers/knowledge.py::_bind_graph binds it on the /admin/.../knowledge/... routes. With the graph bound, the agent's _process_impl walks its own Mem0 memory and consults the shared, provenance-rich Vespa KG, merging the KG result into a dedicated typed kg_* output field — the two stores are complementary, not competing:

Agent kg_* field Graph method called
KnowledgeGraphTraversalAgent merged nodes/edges traverse(seed)
TemporalReasoningAgent kg_timeline compare_over_time(subject)
MultiDocumentSynthesisAgent kg_claim_groups synthesize()
ContradictionReconciliationAgent kg_conflict_entries detect(subject, predicate)
KnowledgeSummarizationAgent kg_video_summaries summarize(video) per video
CitationTracingAgent kg_primary_sources trace(claim_id)

The complement is fail-safe: when no graph is bound (or the backend is unconfigured), the bind is a logged no-op, the kg_* fields stay empty, and the agent returns its Mem0-only answer.

TenantAwareAgentMixin

Location: libs/core/cogniverse_core/agents/tenant_aware_mixin.py

Provides standardized multi-tenant support for all agents, eliminating ~10 lines of duplicated validation code per agent:

class TenantAwareAgentMixin:
    """
    Mixin class that adds multi-tenant capabilities to agents.

    Design Philosophy:
    - REQUIRED tenant_id: No defaults, explicit tenant identification
    - Fail-fast validation: Raises ValueError immediately on invalid tenant_id
    - Context helpers: Provides utilities for tenant-scoped operations
    - Config integration: loads tenant-specific configuration through the
      injected ConfigManager

    Key Benefits:
    - Eliminates ~10 lines of duplicated validation code per agent
    - Consistent error messages across all agents
    - Standardized tenant context API
    - Easy to extend with additional tenant utilities
    """

    def __init__(
        self,
        tenant_id: str,
        config: Optional[SystemConfig] = None,
        config_manager: Optional["ConfigManager"] = None,
        **kwargs
    ):
        """
        Initialize tenant-aware agent mixin.

        Args:
            tenant_id: Tenant identifier (REQUIRED - no default)
            config: Optional system configuration
            config_manager: Injected ConfigManager (REQUIRED - bound via
                bind_config_manager)
            **kwargs: Passed to other base classes in MRO chain

        Raises:
            ValueError: If tenant_id is empty, None, or invalid format
            AgentConfigurationError: If config_manager is None
        """
        # Validate tenant_id (fail fast)
        if not tenant_id:
            raise ValueError(
                "tenant_id is required - no default tenant. "
                "Agents must be explicitly initialized with a valid tenant identifier."
            )

        # Strip whitespace and validate again
        tenant_id = tenant_id.strip()
        if not tenant_id:
            raise ValueError(
                "tenant_id cannot be empty or whitespace only. "
                "Provide a valid tenant identifier (e.g., 'customer_a', 'acme:production')."
            )

        # Store tenant_id
        self.tenant_id = tenant_id

        # Store or load configuration
        self.config = config
        if config is None:
            try:
                self.config = get_config()
            except Exception as e:
                logger.warning(f"Failed to load system config for tenant {tenant_id}: {e}")
                self.config = None

        # Initialize tenant-aware flag
        self._tenant_initialized = True

        logger.debug(f"Tenant context initialized: {tenant_id}")

        # Call super for MRO chain (if needed)
        if hasattr(super(), '__init__'):
            super().__init__(**kwargs)

    def get_tenant_context(self) -> Dict[str, Any]:
        """
        Get tenant context for operations.

        Returns a dictionary with tenant information useful for:
        - Logging and debugging
        - Telemetry span attributes
        - Database query filtering
        - Cache key prefixes

        Returns:
            Dictionary with tenant context information
        """
        context = {"tenant_id": self.tenant_id}

        # Add environment if available from config
        if self.config:
            if hasattr(self.config, 'environment'):
                context["environment"] = self.config.environment
            elif hasattr(self.config, 'get') and callable(self.config.get):
                env = self.config.get('environment')
                if env:
                    context["environment"] = env

        # Add agent type and name if available
        if hasattr(self, '__class__'):
            context["agent_type"] = self.__class__.__name__
        if hasattr(self, 'agent_name'):
            context["agent_name"] = self.agent_name

        return context

    def validate_tenant_access(self, resource_tenant_id: str) -> bool:
        """
        Validate that this agent can access a resource owned by a tenant.

        Used for:
        - Cross-tenant data access checks
        - Security validation
        - Resource authorization

        Returns:
            True if agent's tenant matches resource tenant, False otherwise
        """
        if not resource_tenant_id:
            logger.warning(
                f"Attempted to validate access to resource with no tenant_id "
                f"(agent tenant: {self.tenant_id})"
            )
            return False

        return self.tenant_id == resource_tenant_id

    def get_tenant_scoped_key(self, key: str) -> str:
        """
        Generate a tenant-scoped key for caching, storage, etc.

        Example:
            agent.get_tenant_scoped_key("embeddings/video_123")
            # Returns: "customer_a:embeddings/video_123"
        """
        return f"{self.tenant_id}:{key}"

    def log_tenant_operation(
        self,
        operation: str,
        details: Optional[Dict[str, Any]] = None,
        level: str = "info"
    ):
        """
        Log an operation with tenant context.

        Example:
            agent.log_tenant_operation(
                "search_completed",
                {"query": "machine learning", "results": 10}
            )
            # Logs: [customer_a] [OrchestratorAgent] search_completed: {'query': 'machine learning', 'results': 10}
        """
        log_func = getattr(logger, level, logger.info)

        agent_info = f"[{self.tenant_id}]"
        if hasattr(self, '__class__'):
            agent_info += f" [{self.__class__.__name__}]"

        message = f"{agent_info} {operation}"
        if details:
            message += f": {details}"

        log_func(message)

Usage in Agents

With Type-Safe A2AAgent (Recommended):

# libs/agents/cogniverse_agents/orchestrator_agent.py

from cogniverse_core.agents.a2a_agent import A2AAgent, A2AAgentConfig
from cogniverse_core.agents.base import AgentDeps
from cogniverse_agents.memory_aware_mixin import MemoryAwareMixin

class OrchestratorDeps(AgentDeps):
    """Infrastructure dependencies — no tenant_id (it's per-request)"""
    ...

class OrchestratorAgent(A2AAgent[OrchestratorInput, OrchestratorOutput, OrchestratorDeps], MemoryAwareMixin):
    """Orchestrator agent with type-safe deps and memory support"""

    def __init__(self, deps: OrchestratorDeps, registry: AgentRegistry, ...):
        config = A2AAgentConfig(agent_name="orchestrator_agent", ...)
        super().__init__(deps=deps, config=config, dspy_module=None)

        logger.info("OrchestratorAgent initialized (tenant-agnostic)")

Key Methods

Note: All examples below assume agent is initialized with deps (tenant-agnostic):

from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps
from cogniverse_core.registries.agent_registry import AgentRegistry

registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
agent = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)
# tenant_id arrives per-request in A2A task payload

get_tenant_context() -> Dict[str, Any]

Returns tenant context for logging, telemetry, and debugging:

context = agent.get_tenant_context()
# {
#     "tenant_id": "acme",
#     "agent_type": "OrchestratorAgent",
#     "agent_name": "orchestrator_agent"
# }

validate_tenant_access(resource_tenant_id: str) -> bool

Validates cross-tenant access attempts:

# Same tenant - allow
assert agent.validate_tenant_access("acme") is True

# Different tenant - deny
assert agent.validate_tenant_access("startup") is False

get_tenant_scoped_key(key: str) -> str

Generates tenant-scoped keys for caching/storage:

cache_key = agent.get_tenant_scoped_key("embeddings/video_123")
# "acme:embeddings/video_123"

log_tenant_operation(operation: str, details: Dict, level: str)

Logs operations with full tenant context:

agent.log_tenant_operation(
    "search_completed",
    {"query": "machine learning", "results": 10},
    level="info"
)
# Logs: [acme] [OrchestratorAgent] search_completed: {'query': 'machine learning', 'results': 10}

Routing Execution Paths

AgentDispatcher processes queries through two paths based on complexity classification:

Non-Orchestration Path (Single Agent)

When GatewayAgent classifies the query as simple:

  1. GatewayAgent detects clear entities/modality via GLiNER (<100ms, no LLM)
  2. _execute_downstream_agent dispatches directly to the execution agent
  3. conversation_history is threaded through, enabling query rewrite on multi-turn conversations
  4. Response includes the agent output with plan_steps, agent_results, and execution_summary from OrchestratorOutput

Orchestration Path (Multi-Agent)

When GatewayAgent classifies the query as complex:

  1. GatewayAgent classifies the query using GLiNER entity detection (no LLM, <100ms)
  2. If no entities, low confidence, or multi-modal → complexity="complex", forward to OrchestratorAgent
  3. OrchestratorAgent plans a workflow using DSPy, executes agents via A2A HTTP, and aggregates results
  4. A cogniverse.orchestration telemetry span is persisted through a checked synchronous exporter while ordinary observability spans retain the configured batch exporter. The blocking export runs off the request event loop. A rejected export raises with tenant and workflow context, so a completed response cannot hide the loss of the execution record used for learning.
Component Role Location
GatewayAgent Entry point, classifies queries as simple/complex via GLiNER cogniverse_agents/gateway_agent.py
OrchestratorAgent Autonomous A2A orchestrator: planning, execution, fusion cogniverse_agents/orchestrator_agent.py
OrchestrationEvaluator Reads persisted executions in deterministic, lossless batches for workflow learning cogniverse_agents/routing/orchestration_evaluator.py
AgentTask HTTP wire schema carrying enrichment fields to execution agents cogniverse_runtime/routers/agents.py

OrchestrationEvaluator orders Phoenix rows by aware UTC start_time and context.span_id, then advances its (start_time, span_id) cursor only after WorkflowIntelligence.record_execution() succeeds. Repeated calls therefore continue through every row when the query returns more than batch_size; equal timestamps use the span ID as a stable tie-breaker. Each Phoenix read has a 30-second timeout and at most three attempts. It uses the public lossless TraceStore.get_all_spans() cursor walk rather than placing a fixed ceiling on the history window. The workflow optimization command reuses one fixed UTC end time while draining every batch, so new spans cannot move the upper bound during a run. A Phoenix read failure or a workflow-store failure raises and leaves the failed row eligible for the next call. Every selected row is validated before the first workflow is recorded, so a malformed row cannot produce a partially learned batch. Failed rows must carry both an ERROR status and the bounded error_summary. task_count is the planned agent_sequence length; the observed successful tasks_completed value remains ordinary execution metadata rather than part of the required-field semantic map.

After the drain, WorkflowIntelligence.derive_learning_artifacts(executions) derives exact per-agent counts, success rates, observed workflow timing, confidence averages, and preferred query types. Those profile values come only from the agent_observations recorded around each real A2A dispatch. An observation identifies the dispatched agent and its own duration and success; confidence is included only when that response reported a valid confidence value. Agents with no confidence samples still get a profile, with average_confidence=None. Workflow-level duration, success, and confidence are never copied onto every participating agent, and the performance score leaves missing confidence out instead of treating it as zero. Deep-synthesis observations likewise record only their measured duration and success. Successful execution shapes become deterministic templates whose dependency lists can be replayed by _apply_template. Executions, profiles, query patterns, and templates are persisted through the configured workflow store; a fresh WorkflowIntelligence instance loads those artifacts before serving and can match the learned query patterns without relying on the batch process's memory.

Content-generated workflow plans are persisted as reusable templates with the generated query and exact ordered agent sequence. Their expected_execution_time and success_rate are None: generation proposes a plan but does not claim an observed outcome. Serving can match these templates without a fabricated success-rate boost. A generated response must provide the requested number of unique, non-empty plans. Generated-template batches for one tenant share a renewable Redis lease across application replicas before they write Phoenix. Redis must be configured; there is no process-local persistence fallback. If a batch fails, the store restores or removes only IDs whose current content exactly matches a template successfully written by that batch. Existing unrelated templates and a later replica's successful batch remain untouched.


Multi-Tenant Integration

Tenant Context Flow

sequenceDiagram
    participant API as FastAPI Router
    participant Agent as OrchestratorAgent
    participant SchemaManager as VespaSchemaManager
    participant Memory as Mem0MemoryManager
    participant Vespa as Vespa Backend

    API->>API: require_tenant_id(body.tenant_id) — no header, no middleware
    API->>SchemaManager: get_tenant_schema_name("acme", "video_frames")
    Note over SchemaManager: bare "acme" canonicalizes to "acme:acme"
    SchemaManager-->>API: "video_frames_acme_acme"

    API->>Agent: A2A task with tenant_id="acme" in payload
    Agent->>Memory: initialize_memory("orchestrator_agent", "acme", ...config)
    Memory-->>Agent: Memory ready (agent_memories_acme_acme)

    API->>Agent: _process_impl(OrchestratorInput("cooking videos"))
    Agent->>Memory: get_relevant_context("cooking videos")
    Memory->>Vespa: search(schema="agent_memories_acme_acme")
    Vespa-->>Memory: Relevant memories
    Memory-->>Agent: Context
    Agent->>Agent: Process query
    Agent-->>API: OrchestratorOutput

Tenant Isolation

Key Points:

  • Each agent instance is tenant-scoped
  • Vespa schemas are tenant-specific (video_frames_acme)
  • Memory managers are per-tenant singletons
  • Telemetry projects are per-tenant (acme_orchestrator_agent)

Example:

from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps
from cogniverse_core.registries.agent_registry import AgentRegistry

# ONE orchestrator serves ALL tenants — tenant_id arrives per-request
registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
agent = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

# Memory is initialized lazily per-tenant on first request via MemoryAwareMixin:
# - Tenant "acme" → agent_memories_acme schema (first request initializes)
# - Tenant "startup" → agent_memories_startup schema (first request initializes)
# Memory namespaced by (tenant_id, agent_name) — no cross-tenant leakage

Usage Examples

Example 1: Basic Routing via Orchestrator

from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps, OrchestratorInput
from cogniverse_core.registries.agent_registry import AgentRegistry

registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
orchestrator = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

result = await orchestrator._process_impl(
    OrchestratorInput(
        query="Show me videos about machine learning",
        tenant_id="acme",
    )
)

print(f"Plan steps: {len(result.plan_steps)}")
print(f"Agents executed: {list(result.agent_results.keys())}")
print(f"Summary: {result.execution_summary}")
from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_foundation.config.utils import create_default_config_manager

config_manager = create_default_config_manager()
deps = SearchAgentDeps()
agent = SearchAgent(deps=deps, config_manager=config_manager, schema_loader=schema_loader)

# Search (synchronous) — tenant_id required per-request
results = agent.search_by_text(
    query="Python programming tutorial",
    tenant_id="acme",
    modality="video",
    top_k=5,
)

for result in results:
    print(result)

Example 3: Multi-Agent Orchestration

from cogniverse_agents.orchestrator_agent import (
    OrchestratorAgent, OrchestratorDeps, OrchestratorInput,
)
from cogniverse_core.registries.agent_registry import AgentRegistry

# Create orchestrator with agent registry (discovers agents from config.json)
registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
orchestrator = OrchestratorAgent(
    deps=OrchestratorDeps(), registry=registry, config_manager=config_manager
)

# Execute orchestration — tenant_id and session_id per-request
result = await orchestrator._process_impl(
    OrchestratorInput(
        query="Find and summarize AI research videos from 2024",
        tenant_id="acme_corp",
        session_id="sess-uuid-123",
    )
)

print(f"Plan steps: {len(result.plan_steps)}")
print(f"Agents executed: {list(result.agent_results.keys())}")
print(f"Summary: {result.execution_summary}")
from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_core.memory.manager import Mem0MemoryManager
from cogniverse_foundation.config.utils import create_default_config_manager

config_manager = create_default_config_manager()
deps = SearchAgentDeps()
agent = SearchAgent(deps=deps, config_manager=config_manager, schema_loader=schema_loader)

# Initialize memory manager separately — embedder_base_url is required
memory = Mem0MemoryManager(tenant_id="acme")
memory.initialize(
    backend_host="localhost",
    backend_port=8080,
    llm_model="google/gemma-4-e4b-it",
    embedding_model="lightonai/DenseOn",
    llm_base_url="http://localhost:11434",
    embedder_base_url="http://localhost:29010/v1",
    config_manager=config_manager,
    schema_loader=schema_loader,
    base_schema_name="agent_memories",
)

# First search — tenant_id is required
results1 = agent.search_by_text(query="cooking tutorials", tenant_id="acme", top_k=5)

# Store in memory for future context
# add_memory takes: content (str), tenant_id, agent_name, optional metadata
memory.add_memory(
    content=f"User searched for cooking tutorials, found {len(results1)} results",
    tenant_id="acme",
    agent_name="video_search_agent",
    metadata={"preference": "high_relevance"},
)

# Second search (memory context retrieved separately)
results2 = agent.search_by_text(query="advanced cooking techniques", tenant_id="acme", top_k=5)

Example 5: Streaming Results

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from pathlib import Path

# Initialize agent (schema_loader is REQUIRED)
deps = SearchAgentDeps(
    backend_url="http://localhost",
    backend_port=8080,
)
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))
agent = SearchAgent(deps=deps, schema_loader=schema_loader)

# Non-streaming call — tenant_id per-request in task payload
result = await agent.process({"query": "machine learning", "top_k": 10, "tenant_id": "acme"})
print(f"Found {result.total_results} results")

# Streaming call (returns AsyncGenerator)
async for event in agent.process({"query": "machine learning", "top_k": 10}, stream=True):
    if event["type"] == "status":
        print(f"Status: {event['message']}")
    elif event["type"] == "partial":
        print(f"Partial results: {event['data']}")
    elif event["type"] == "final":
        print(f"Final: {event['data']}")

Event Types:

  • status - Progress updates (e.g., "Searching...", "Encoding query...")
  • partial - Intermediate results (e.g., results from first profile in ensemble)
  • token - DSPy token streaming for reasoning fields
  • task_complete - Workflow task completion (orchestrator)
  • final - Complete result
  • error - Error information

Streaming API

All agents support OpenAI-style streaming via the stream=True parameter on the process() method.

Architecture

flowchart LR
    subgraph Process["<span style='color:#000'>Agent.process()</span>"]
        NonStream["<span style='color:#000'>stream=False (default)<br/>→ _process_impl()<br/>→ Returns OutputT</span>"]
        Stream["<span style='color:#000'>stream=True<br/>→ _stream_with_progress()<br/>→ emit_progress() events + final</span>"]
    end

    style Process fill:#ffcc80,stroke:#ef6c00,color:#000
    style NonStream fill:#a5d6a7,stroke:#388e3c,color:#000
    style Stream fill:#90caf9,stroke:#1565c0,color:#000

Method Pattern

When creating agents, override _process_impl() for core logic:

class MyAgent(A2AAgent[MyInput, MyOutput, MyDeps]):

    async def _process_impl(self, input: MyInput) -> MyOutput:
        """Core processing logic (required).
        Call self.emit_progress() for streaming events."""
        self.emit_progress("processing", "Working on it...")
        result = do_work(input)
        self.emit_progress("processing", "Done", data={"partial": result})
        return MyOutput(...)

HTTP SSE Integration

The A2A endpoint /tasks/send supports streaming via stream field:

# Non-streaming request
POST /tasks/send
{"query": "...", "context": "..."}
→ Returns JSON

# Streaming request
POST /tasks/send
{"query": "...", "context": "...", "stream": true}
→ Returns Server-Sent Events (SSE)

Event Format

All streaming events are plain dicts with a type field:

{"type": "status", "phase": "encoding", "message": "Encoding query..."}
{"type": "partial", "data": {"results_so_far": 5}}
{"type": "token", "field": "reasoning", "text": "The query..."}
{"type": "task_complete", "task": "entity_extraction", "success": True}
{"type": "final", "data": {"results": [...], "total": 10}}
{"type": "error", "agent": "SearchAgent", "error_type": "VespaSearchDegraded",
 "message": "SearchAgent streaming failed with VespaSearchDegraded. See server logs for detail."}

An unhandled exception in _process_impl yields the error event above as the terminal event: agent is the agent class (the A2A executor replaces it with the dispatch name), error_type is the leaf exception type — ExceptionGroups from TaskGroups are flattened to their leaves — and the exception text stays server-side (logged with the full traceback), since it can carry credentialed backend URLs.


RLM Inference (Recursive Language Models)

RLM (Recursive Language Models) enables agents to handle near-infinite context by programmatically examining, decomposing, and recursively calling LLMs. This is useful for processing large result sets, long transcripts, or multi-document analysis.

Reference: RLM Paper (arXiv:2512.24601)

Architecture

RLM is implemented using DSPy's built-in dspy.RLM module (requires dspy-ai>=3.1.0):

libs/agents/cogniverse_agents/
├── inference/
│   ├── __init__.py
│   ├── ab_harness.py             # RLMABRunner: with-RLM vs without-RLM comparison
│   ├── deno_check.py             # Boot probe: fail-fast if Deno missing
│   ├── instrumented_rlm.py       # InstrumentedRLM with EventQueue + fallback marker
│   ├── rlm_inference.py          # RLMInference wrapper, RLMResult
│   └── tolerant_interpreter.py   # TolerantPythonInterpreter/TolerantRLM: skip
│                                 # stale id-null messages on the Deno channel;
│                                 # RLMTimeoutError and the iteration deadline;
│                                 # unparseable model turns as failed iterations
├── mixins/
│   ├── __init__.py
│   └── rlm_aware_mixin.py        # RLMAwareMixin for agents

libs/core/cogniverse_core/agents/
└── rlm_options.py                # RLMOptions schema for query-level config

Runtime requirement — Deno. DSPy's RLM module spawns a Deno subprocess to execute LLM-generated code in a sandboxed JavaScript runtime. The Deno binary must be reachable on PATH (the runtime container installs it to /usr/local in libs/runtime/Dockerfile). RLMInference.__init__ probes for Deno and raises DenoNotInstalledError at construction if missing — failing fast at boot rather than on the first call. Set COGNIVERSE_RLM_SKIP_DENO_CHECK=1 to bypass the probe (only when you are certain RLM will not be invoked); the runtime entrypoint resolves it once at process start via configure_deno_check.

Key Components

RLMOptions (Query-Level Configuration)

from cogniverse_core.agents.rlm_options import RLMOptions

# Configuration for A/B testing
rlm_opts = RLMOptions(
    enabled=True,              # Explicitly enable RLM
    auto_detect=False,         # Or auto-enable based on context size
    context_threshold=50_000,  # Threshold for auto_detect (chars)
    max_iterations=3,          # Maximum REPL iterations (1-10)
    max_llm_calls=30,          # Maximum LLM sub-calls (1-100)
    timeout_seconds=300,       # Timeout for RLM processing (1-1800s)
    cache=True,                # False forces fresh inference for live timing
    backend="openai",          # LLM backend (openai, anthropic, litellm)
    model="gpt-4o",            # Model override
)

RLMInference (Core Wrapper)

from cogniverse_agents.inference.rlm_inference import RLMInference, RLMResult
from cogniverse_foundation.config.unified_config import LLMEndpointConfig

rlm = RLMInference(
    llm_config=LLMEndpointConfig(model="openai/gpt-4o"),
    max_iterations=10,      # Maximum REPL iterations
    max_llm_calls=30,
    timeout_seconds=300,
)

result: RLMResult = rlm.process(
    query="Summarize the main findings",
    context=large_context_string,  # Can be 100K+ chars
)

print(f"Answer: {result.answer}")
print(f"Depth: {result.depth_reached}, Calls: {result.total_calls}")
print(f"Latency: {result.latency_ms}ms")

build_rlm_from_options(llm_config, rlm_options, *, config_manager=None, tenant_id="") (same module) is the shared constructor used by the KG summariser agents (FederatedQueryAgent, TemporalReasoningAgent, KnowledgeGraphTraversalAgent, KnowledgeSummarizationAgent, MultiDocumentSynthesisAgent). It resolves the endpoint config through rlm_endpoint — the request's rlm.model (with its api_base and api_key) when it names one, else the agent's endpoint, else the tenant's configured endpoint for rlm_inference (the primary LM unless llm_config.overrides names one), and RLMEndpointNotConfiguredError when none is available — and applies the option's iteration / call / timeout caps. RLMAwareMixin.process_with_rlm resolves its endpoint the same way, so an RLM run that names no model never defaults to a provider's public API. When config_manager and tenant_id are supplied and gateway routing is enabled for that tenant, the resolved endpoint is routed through the gateway (task rlm_inference) before the RLM's LM is built; tenant_id is also threaded onto the RLMInference for event scoping. Omitting them keeps the direct-to-backend path:

from cogniverse_agents.inference.rlm_inference import build_rlm_from_options

rlm = build_rlm_from_options(
    self._llm_config,
    rlm_options,
    config_manager=self._config_manager,
    tenant_id=self._memory_tenant_id or "",
)
result = rlm.process(query=query, context=block)

route_rlm_endpoint(endpoint, config_manager, tenant_id) (same module) is the underlying primitive build_rlm_from_options uses to route an endpoint through the gateway (task rlm_inference); it returns the endpoint unchanged when config_manager/tenant_id are absent, routing is disabled, or the config is unreachable. Call sites that construct RLMInference directly rather than via build_rlm_from_options — e.g. the orchestrator's deep-synthesis workflow — route their resolved endpoint through it before construction:

from cogniverse_agents.inference.rlm_inference import (
    RLMInference,
    route_rlm_endpoint,
)

llm_primary = route_rlm_endpoint(llm_primary, self._config_manager, tenant_id)
rlm = RLMInference(llm_config=llm_primary, tenant_id=tenant_id)

RLMAwareMixin (Agent Integration)

Agents inherit from RLMAwareMixin to gain RLM capabilities:

from cogniverse_agents.mixins.rlm_aware_mixin import RLMAwareMixin

class SearchAgent(RLMAwareMixin, MemoryAwareMixin, A2AAgent[...]):
    async def _process_impl(self, input: SearchInput) -> SearchOutput:
        # ... perform search ...

        # Check if RLM should be used for this query
        if self.should_use_rlm_for_query(input.rlm, results_context):
            rlm_result = self.process_with_rlm(
                query=input.query,
                context=results_context,
                rlm_options=input.rlm,
                tenant_id=input.tenant_id,
            )
            return SearchOutput(
                results=results,
                rlm_synthesis=rlm_result.answer,
                rlm_telemetry=self.get_rlm_telemetry(rlm_result, len(results_context)),
            )

process_with_rlm and get_rlm take the request's tenant_id as a required keyword argument (an empty one raises ValueError): one agent instance serves every tenant, so the tenant travels with the call, never on the instance. get_rlm routes the RLM endpoint through the gateway (task rlm_inference) for that tenant before building the LM, reading the host's config_manager (SearchAgent exposes config_manager; CodingAgent, DeepResearchAgent, and DetailedReportAgent store _config_manager); with routing disabled the endpoint is used unchanged. The non-event instance is reused only while the tenant, the routed identity (model + api_base + headers) and the caps match, and each call returns the instance built for its own key even when a concurrent request replaces the cached one. WikiManager._merge_with_rlm routes the same way via its own config_manager.

A/B Testing with RLM

RLM is query-level configurable to enable A/B testing:

from cogniverse_agents.search_agent import SearchInput
from cogniverse_core.agents.rlm_options import RLMOptions

# Group A: Standard search (no RLM)
input_a = SearchInput(query="machine learning tutorials", rlm=None)

# Group B: RLM-enabled search
input_b = SearchInput(
    query="machine learning tutorials",
    rlm=RLMOptions(enabled=True, max_iterations=3),
)

# Auto-detect mode: Enable RLM only for large context
input_c = SearchInput(
    query="machine learning tutorials",
    rlm=RLMOptions(auto_detect=True, context_threshold=50_000),
)

Telemetry Metrics

RLM results include telemetry for comparison in Phoenix dashboard:

Metric Description
rlm_enabled Boolean flag indicating RLM was used
rlm_depth_reached Actual recursion depth achieved
rlm_total_calls Number of LLM sub-calls made
rlm_tokens_used Total tokens consumed across the recursive run (sum across all sub-LMs, populated via DSPy track_usage)
rlm_latency_ms End-to-end RLM processing time
rlm_was_fallback True when answer came from fallback extraction (max iterations exhausted without SUBMIT)
rlm_trajectory_length Number of REPL iterations captured in RLMResult.trajectory (0 unless RLMOptions.include_trajectory=True)
context_size_chars Input context size

Orchestrator RLM promotion

Before dispatching to a sub-agent, the orchestrator now estimates the projected payload size (sum of stringified field lengths). When the projection exceeds 75% of RLMOptions.context_threshold (default 50_000 chars × 0.75 = 37_500), the orchestrator stamps agent_input["rlm"] = {"enabled": True, "auto_detect": True, "context_threshold": 50_000} so the sub-agent runs through RLMInference instead of stuffing everything into a single prompt.

Eligible agents (their input schema accepts an rlm field): search_agent, deep_research_agent, detailed_report_agent, coding_agent.

Env var Default Effect
COGNIVERSE_ORCH_RLM_PROMOTION unset Set to disabled to skip promotion entirely.
COGNIVERSE_ORCH_RLM_PROMOTION_FRACTION 0.75 Override the threshold fraction.

Both are resolved once at process start by the runtime entrypoint and applied via configure_rlm_promotion; the promotion module itself reads no env.

The promotion is idempotent: if the caller already supplied an rlm field (any value, including None for explicit opt-out), the orchestrator does not touch it.

The orchestrator's iterative-retrieval sufficiency gate has a separate, smaller promotion threshold. When the JSON-serialized accumulated evidence exceeds 6,000 characters, the gate runs through InstrumentedRLM with a hard cap of three iterations; smaller evidence sets use dspy.ChainOfThought(SufficientContextSignature). A False gate decision does not mean the accumulated evidence is discarded: its rationale must cite concrete support already found and name only the facets still missing. The next retrieval round can therefore reformulate for the gaps without losing supported facts. The integration golden in tests/agents/integration/test_iterative_retrieval_loop.py pins both the multi-round behavior and the three-iteration RLM cap.

RLM A/B harness

RLMABRunner runs the same query through both arms — once without RLM (single LM call on the raw context) and once with RLM (recursive REPL) — and returns a typed ABResult with both arms plus a side-by-side comparison. Both arms share an ab_id stamped in their result metadata so Phoenix spans correlate across the pair.

from cogniverse_agents.inference.ab_harness import RLMABRunner

runner = RLMABRunner(
    llm_config=llm_config,
    judge=lambda q, ctx, ans: 1.0 if "Paris" in ans else 0.0,
    rlm_max_iterations=4,
)
result = runner.run(query="What is the capital of France?", context=context)

# Both arms ran with the same llm_config; ab_id correlates them in Phoenix.
print(result.ab_id)
print(result.comparison.latency_delta_ms, result.comparison.tokens_delta)

# Flattened payload for direct span attribute set:
phoenix_attrs = result.to_telemetry_dict()

The judge parameter is optional. Without it the comparison still tracks latency / tokens / was_fallback, just not quality. With it, every arm's answer gets scored and comparison.judge_delta reports with_rlm − without.

When constructed with tenant_id and config_manager, both arms route their endpoint through the gateway (task rlm_inference). Routing is resolved once and shared by both arms, so the gateway returns the same model for each — the comparison still isolates the RLM machinery, now measured against the production (routed) path. optimization_cli.run_ab_compare passes both, so the web client's RLM A/B view reflects what production actually runs.

Deep synthesis workflow

Opt-in entry point for queries that need recursive multi-agent orchestration over a knowledge subgraph that doesn't fit in a normal plan ("compare these 50 documents across 5 tenants and produce a unified report"). Distinct from the default OrchestratorAgent plan- then-act path: this workflow runs Orchestrator inside an RLM-style trajectory so partial results inform the next round of fan-out.

This is not the default execution path. Default Orchestrator stays plan-then-act with parallel sub-agent fan-out. DeepSynthesisWorkflow is a separate class so its cost is local and explicit, the default trace shape stays clean, and the RLM A/B harness can compare the two paths on a curated benchmark before promotion.

from cogniverse_agents.deep_synthesis_workflow import (
    DeepSynthesisConfig,
    DeepSynthesisWorkflow,
)
from cogniverse_agents.inference.rlm_inference import RLMInference

rlm = RLMInference(llm_config=llm_config, max_iterations=8)

async def dispatch(query: str, sub_agent_name: str) -> str:
    """Caller-supplied: route the sub-query to the right A2A sub-agent."""
    out = await registry.send(sub_agent_name, query)
    return out.answer or ""

workflow = DeepSynthesisWorkflow(
    rlm=rlm,
    sub_agent_dispatcher=dispatch,
    config=DeepSynthesisConfig(
        rate_limit_per_hour=5,            # per-tenant ceiling
        hard_call_cap=200,                # total LLM + sub-agent calls
        max_iterations=8,
        max_subagent_calls_per_round=6,
    ),
)
result = await workflow.run(
    query="Compare refund policies across all subsidiaries",
    tenant_id="acme:production",
    seed_subagents=["search_agent", "document_agent", "kg_traversal_agent"],
)
print(result.answer, result.was_submitted, result.was_capped, result.was_rate_limited)
print(result.iterations_used, result.subagent_calls_made, result.llm_calls_used)

The RLM step controls each iteration. Two protocol tokens:

Token Meaning
SUBMIT() anywhere in the answer Workflow returns the answer (token stripped).
ASK(<subagent>: <subquery>) Workflow dispatches that sub-agent next round.

Anything else terminates the trajectory (was_capped=True, kind stalled_no_asks in the trajectory log) — a runaway never silently spends budget.

Bound Default Behaviour at limit
Per-tenant rate limit / hour 5 was_rate_limited=True, RLM never consulted.
Total LLM + sub-agent calls 200 was_capped=True, returns gathered evidence.
Iterations 8 was_capped=True, returns gathered evidence.
Sub-agent calls per round 6 Excess ASK()s for this round are dropped.

Sub-agent failures are non-fatal — the dispatcher's exception is logged at debug and that name is dropped from the round's evidence. The trajectory records every step (subagent, rlm_step, cap_reached, stalled_no_asks, iteration_cap_exhausted) so callers can audit exactly how the answer was produced.


Knowledge Agents

Nine A2A agents that operate on the Knowledge Management Layer (schema-driven memory, provenance, trust, contradiction, federation). All live at top level in libs/agents/cogniverse_agents/. Each is a full A2AAgent[InputT, OutputT, DepsT] implementation; Deps carries a tenant_id and an optional memory_manager_factory for override in tests.

AuditExplanationAgent

Read-only A2A agent that explains why a system answer was produced. Walks the provenance chain to surface every source memory the answer was derived from, attaches the decayed trust score per source, and flags any contradictions touching those sources' subjects. Compliance deployments need this surface for every answer.

from cogniverse_agents.audit_explanation_agent import (
    AuditExplanationAgent,
    AuditExplanationDeps,
    AuditExplanationInput,
)

agent = AuditExplanationAgent(deps=AuditExplanationDeps(tenant_id="acme"))
out = await agent._process_impl(AuditExplanationInput(
    tenant_id="acme",
    answer_memory_id="answer:42",
    include_trust=True,
    include_contradictions=True,
    max_chain_depth=10,
    max_chain_nodes=100,
))
print(out.explanation)         # ready-to-render structured text
for src in out.sources:
    print(src.memory_id, src.depth, src.trust_score, src.derivation_kind)
for c in out.contradictions_touched:
    print(c.subject_key, c.conflicting_memory_ids)
Property Behaviour
Answer with no provenance Returns one source (the answer itself, depth 0).
Chain exceeds max_chain_depth / max_chain_nodes truncated_chain=True; remaining refs surfaced as primary sources.
include_trust=False Skips per-source extract_trust / apply_decay.
include_contradictions=False Skips ContradictionDetector pass.
Source memory missing trust metadata trust_score is None.
Source memory absent The source row remains present with no trust metadata.
Memory backend read fails Raises RuntimeError with the memory and tenant identifiers instead of returning an incomplete explanation.

The explanation field is human-readable structured text, intended to be rendered as-is in audit UIs without further LLM post-processing. Capability strings: audit_explanation, audit, provenance_consumer. Default port=8027.

KnowledgeSummarizationAgent

A2A agent that distills a slice of the knowledge layer (a subject area, a kind, optionally a time window) into a structured summary with citations. Distinct from SummarizerAgent (which summarises retrieval results in-flight): KnowledgeSummarizationAgent summarises the knowledge layer itself and can promote the result into the org trunk via federation.

from cogniverse_agents.knowledge_summarization_agent import (
    KnowledgeSummarizationAgent,
    KnowledgeSummarizationDeps,
    KnowledgeSummarizationInput,
)

agent = KnowledgeSummarizationAgent(
    deps=KnowledgeSummarizationDeps(tenant_id="acme:production"),
)
out = await agent._process_impl(KnowledgeSummarizationInput(
    tenant_id="acme:production",
    subject_keys=["policy:refunds", "policy:returns"],
    title="Q1 buyer-protection summary",
    since="2026-01-01T00:00:00Z",
    until="2026-04-01T00:00:00Z",
    promote=True,
    actor_role="tenant_admin",
    actor_id="tadm",
))
print(out.source_count, out.summary)
print(out.promoted_to_org_trunk, out.promoted_memory_id)

The agent auto-registers a knowledge_summary schema (permanent, org_shared, tenant_admin pin floor) into the supplied KnowledgeRegistry if missing — so promotion works out of the box.

Behaviour Detail
No matching memories Empty summary; metadata.reason = no_matching_memories; promotion skipped.
actor_role < tenant_admin + promote=True Promotion refused (logged warning); summary still returned.
Both subject_keys and kinds empty Pulls every memory in the agent_name namespace (cap by max_memories).
since or until set Bounds must be timezone-aware ISO-8601 values; a matching memory without a valid timezone-aware written_at is excluded.
RLM enabled AND context > threshold RLM summariser fires; used_rlm=True.
Org-trunk promotion Storage write runs outside the event-loop thread; a write failure propagates and no successful promotion is reported.

Capability strings: knowledge_summarization, audit, federation_promoter. Default port=8026.

TemporalReasoningAgent

A2A agent that compares knowledge about one subject across multiple time windows. Useful for "how has our refund policy evolved" or "what did we know about X in Q2 vs Q4." Read-only.

from cogniverse_agents.temporal_reasoning_agent import (
    TemporalReasoningAgent,
    TemporalReasoningDeps,
    TemporalReasoningInput,
    TimeWindow,
)

agent = TemporalReasoningAgent(deps=TemporalReasoningDeps(tenant_id="acme"))
out = await agent._process_impl(TemporalReasoningInput(
    tenant_id="acme",
    subject_key="policy:refunds",
    windows=[
        TimeWindow(label="Q1", start="2026-01-01T00:00:00Z",
                   end="2026-04-01T00:00:00Z"),
        TimeWindow(label="Q2", start="2026-04-01T00:00:00Z",
                   end="2026-07-01T00:00:00Z"),
        TimeWindow(label="from_q3", start="2026-07-01T00:00:00Z"),  # open-ended
    ],
))
print(out.distinct_signatures_count)  # 1 = unchanged, >1 = evolved
print(out.undated_count)  # memories on the subject lacking written_at

The agent uses the written_at field that provenance stamps onto every metadata block — no Vespa-side time-version index is required. Each window emits a stable content-hash signature so a caller can detect whether knowledge actually changed without invoking an LLM. The optional RLM summariser narrates the deltas when the per-window context exceeds the RLM threshold.

Property Behaviour
Memories without written_at Counted in undated_count, excluded from windows.
Window bounds start and end must be timezone-aware ISO-8601 values; naive timestamps are rejected.
windows[i].end is None Open-ended (matches everything ≥ start).
Same content in two windows One distinct signature (knowledge unchanged).
Subject filter Strict metadata.subject_key == subject_key match.

Capability strings: temporal_reasoning, audit. Default port=8025.

FederatedQueryAgent

A2A agent that answers a free-text query by aggregating federated reads across multiple tenants in the same org. Distinct from CrossTenantComparisonAgent: that one compares tenant views of a single known subject; this one answers a query by finding any matching memories. Both share the federation read path so the org trunk is included automatically and cross-org reads are structurally prevented.

from cogniverse_agents.federated_query_agent import (
    FederatedQueryAgent,
    FederatedQueryDeps,
    FederatedQueryInput,
)
from cogniverse_core.agents.rlm_options import RLMOptions

agent = FederatedQueryAgent(
    deps=FederatedQueryDeps(tenant_id="acme:production"),
)
out = await agent._process_impl(FederatedQueryInput(
    tenant_id="acme:production",
    query="What is our refund policy in the EU?",
    tenant_ids=["acme:alpha", "acme:beta", "acme:gamma"],
    actor_role="org_admin",
    actor_id="oadm",
    top_k_per_tenant=20,
    rlm=RLMOptions(enabled=True),  # optional: summarise large result sets
))
for hit in out.hits:
    print(hit.tenant_id, hit.origin, hit.memory_id, hit.excerpt)
print(out.summary)  # populated when RLM ran
ACL rule Behaviour
actor_role < tenant_admin Rejected.
actor_role = tenant_admin / org_admin Allowed.
Any tenant_ids[i] outside caller's org Rejected — never reads cross-org.
tenant_id=None (admin CLI) Cross-org check skipped.

Default agent_name_filter="_promoted" so the agent reads from the canonical promoted-knowledge namespace; pass an explicit value to scope to a different agent's namespace. top_k_per_tenant (default 20) caps per-tenant fan-in. Federated reads run outside the event-loop thread, and memory-manager construction or read failures propagate instead of being returned as an empty hit list. The optional RLM summariser only fires when both rlm.enabled and the merged context exceeds the RLM threshold. Capability strings: federated_query, audit, federation_consumer. Default port=8024.

CrossTenantComparisonAgent

A2A agent that compares per-tenant views of a subject across multiple tenants in the same org. Built on federation: each per-tenant fetch goes through FederationService.federated_get_all, so the org trunk is included automatically and cross-org reads are structurally prevented.

from cogniverse_agents.cross_tenant_comparison_agent import (
    CrossTenantComparisonAgent,
    CrossTenantComparisonDeps,
    CrossTenantComparisonInput,
)

agent = CrossTenantComparisonAgent(
    deps=CrossTenantComparisonDeps(tenant_id="acme:production"),
)
out = await agent._process_impl(CrossTenantComparisonInput(
    tenant_id="acme:production",
    subject_key="france:capital",
    tenant_ids=["acme:alpha", "acme:beta", "acme:gamma"],
    actor_role="org_admin",
    actor_id="oadm",
))
print(out.distinct_signatures_count)  # 1 = all tenants agree
for view in out.tenant_views:
    print(view.tenant_id, view.matching_memory_ids, view.origin_tags)
ACL rule Behaviour
actor_role < tenant_admin Rejected.
actor_role = tenant_admin / org_admin Allowed.
Any tenant_ids[i] outside the caller's org Rejected — never reads cross-org.

distinct_signatures_count collapses tenants by content signature so callers can quickly tell agreement from disagreement. Per-tenant federated reads run outside the event-loop thread, and backend failures propagate instead of producing an apparently empty tenant view. Capability strings: cross_tenant_comparison, audit, federation_consumer. Default port=8023.

KnowledgeGraphTraversalAgent

A2A agent that walks the entity/edge memories of the knowledge graph from a starting subject_key (or memory id resolving to one) and returns a structured (nodes, edges) view plus an optional RLM-summarised narrative for large traversals.

from cogniverse_agents.kg_traversal_agent import (
    KGTraversalDeps,
    KGTraversalInput,
    KnowledgeGraphTraversalAgent,
)
from cogniverse_core.agents.rlm_options import RLMOptions

agent = KnowledgeGraphTraversalAgent(
    deps=KGTraversalDeps(tenant_id="acme"),
    llm_config=llm_config,
)
out = await agent._process_impl(KGTraversalInput(
    tenant_id="acme",
    start_subject_key="company:acme",
    max_depth=3,
    max_edges=200,
    relation_allowlist=["depends_on", "owns"],
    rlm=RLMOptions(auto_detect=True, context_threshold=20_000),
))
for node in out.nodes:
    print(node.subject_key, node.label, node.excerpt)
for edge in out.edges:
    print(edge.from_subject_key, edge.relation, edge.to_subject_key)
print(out.summary)  # populated when RLM ran
Memory shape required Metadata keys
Node kind=kg_node (or entity_fact), subject_key, optional label
Edge kind=kg_edge, from_subject_key, to_subject_key, relation

The walker is deterministic (BFS over edges with depth + edge cap) and honours an optional relation_allowlist so callers can isolate a sub-graph. Capability strings: kg_traversal, graph_walk. Default port=8022.

MultiDocumentSynthesisAgent

A2A agent that produces a coherent answer across N source documents (10–500) and persists it as a new memory of kind synthesis_fact whose provenance carries derivation_kind=synthesis and derived_from referencing every input document. RLM-capable for large document sets.

from cogniverse_agents.multi_document_synthesis_agent import (
    DocumentRef,
    MultiDocSynthesisDeps,
    MultiDocSynthesisInput,
    MultiDocumentSynthesisAgent,
)
from cogniverse_core.agents.rlm_options import RLMOptions

agent = MultiDocumentSynthesisAgent(
    deps=MultiDocSynthesisDeps(tenant_id="acme"),
    llm_config=llm_config,
)
out = await agent._process_impl(MultiDocSynthesisInput(
    tenant_id="acme",
    query="What does the literature say about X?",
    documents=[
        DocumentRef(memory_id="m_paper_1", label="Paper #1"),
        DocumentRef(memory_id="m_paper_2", label="Paper #2"),
        DocumentRef(content="Inline excerpt from a third source", label="Inline"),
    ],
    rlm=RLMOptions(auto_detect=True, context_threshold=50_000),
    persist=True,
))
print(out.answer, out.persisted_memory_id, out.used_rlm)
Behaviour Outcome
documents referenced by memory_id Fetched via Mem0.memory.get() outside the event-loop thread; rendered in the prompt with the supplied label. A read failure propagates so a partial synthesis is never persisted as complete.
documents supplied as inline content Rendered directly; cited as an external ref using the label.
rlm.enabled=True or auto-detect over threshold Synthesis runs through RLMInference; used_rlm=True.
persist=True Output is returned only after the synthesis_fact primary and indexed provenance both persist. A provenance write or compensation failure propagates as a non-success request; it never becomes persisted_memory_id=None.
persist=False Read-only / audit run; nothing written.

Capability strings: multi_document_synthesis, citation_preservation. Default port=8021.

ContradictionReconciliationAgent

Read-only A2A agent that consumes a ConflictSet and resolves it per the target schema's contradiction_policy (or an explicit per-call override). Returns a survivor list plus per-member outcomes for the UI.

from cogniverse_agents.contradiction_reconciliation_agent import (
    ContradictionReconciliationAgent,
    ContradictionReconciliationDeps,
    ContradictionReconciliationInput,
)

agent = ContradictionReconciliationAgent(
    deps=ContradictionReconciliationDeps(tenant_id="acme"),
)
out = await agent._process_impl(ContradictionReconciliationInput(
    tenant_id="acme",
    target_kind="entity_fact",
    conflict_member_ids=["m_a", "m_b", "m_c"],
    policy_override=None,  # use schema's default
))
print(out.policy_used, out.survivors)
for member in out.resolved:
    print(member.memory_id, "survived" if member.survived else "dropped")

Requested conflict members are fetched outside the event-loop thread. A missing memory is reported in metadata.missing, while a backend read failure propagates instead of being treated as a deleted member.

V1 is deterministic: it applies reconcile() from cogniverse_core.memory.contradiction directly. An RLM trajectory for fetching extra evidence per side is not yet implemented; that enrichment is left as a follow-up. Capability strings: contradiction_reconciliation, audit. Default port=8020.

CitationTracingAgent

Read-only A2A agent that wraps ProvenanceWalker. Given a memory id, returns the BFS-walked citation chain plus the structured primary-source list — useful for "show me the sources" UX, compliance audit, and debugging why a synthesised claim was promoted.

from cogniverse_agents.citation_tracing_agent import (
    CitationTracingAgent,
    CitationTracingDeps,
    CitationTracingInput,
)

agent = CitationTracingAgent(deps=CitationTracingDeps(tenant_id="acme"))
out = await agent._process_impl(CitationTracingInput(
    memory_id="m_synthesis",
    tenant_id="acme",
    max_depth=10,
    max_nodes=100,
))
for node in out.nodes:
    print(node.depth, node.memory_id, node.derivation_kind, node.confidence)
for source in out.primary_sources:
    print("primary source:", source.ref_kind, source.ref_id)

The agent does not call the LLM and does not write to memory; it is deterministic and cheap. It rejects a primary/index provenance mismatch instead of rendering the memory as a provenance-free leaf; operators use Mem0MemoryManager.repair_provenance() explicitly before retrying the read. Defaults: max_depth=10, max_nodes=100, port=8019. Capability strings: citation_tracing, provenance_walk, audit.

RLMResult.metadata always carries trajectory_length and a bounded trajectory_summary (first 8 entries) regardless of the opt-in, so Phoenix spans can debug recursion behaviour without forcing the full trajectory back to the caller.

SearchOutput with RLM

When RLM is enabled, SearchOutput includes:

class SearchOutput(AgentOutput):
    # ... standard fields ...
    results: List[Dict[str, Any]]
    total_results: int

    # Id of the SearchAgent.process span this result set was recorded on, so a
    # client can attach result_relevance annotations for embedding-triplet mining.
    span_id: Optional[str]

    # RLM fields (populated when RLM enabled)
    rlm_synthesis: Optional[str]           # Synthesized answer
    rlm_telemetry: Optional[Dict[str, Any]] # Telemetry metrics

Timeout and Error Handling

timeout_seconds is enforced inside the REPL loop: TolerantRLM checks the deadline at each iteration boundary and in the in-REPL llm_query / llm_query_batched tools, raising RLMTimeoutError at the first of them reached after expiry, so the computation stops rather than continuing unobserved. The overrun is therefore bounded by the model call in flight when the deadline passes. The same check guards aforward/acall. process is synchronous; async callers run it through asyncio.to_thread.

A model reply that does not parse into reasoning and code (dspy's AdapterParseError, typically a reply cut off at max_tokens) is recorded as a failed iteration: its trajectory entry has empty reasoning and code and an output telling the model which fields to send, and the next iteration sees it. Each such turn spends one of max_iterations; when all of them fail the run ends in fallback extraction (was_fallback=True), and an unparseable extraction reply raises AdapterParseError. Endpoint failures (HTTP errors, timeouts, connection errors) propagate from process unchanged.

from cogniverse_agents.inference import RLMTimeoutError

try:
    result = rlm.process(query=query, context=large_context)
except RLMTimeoutError as e:
    # Handle timeout (default: 300 seconds)
    logger.error(f"RLM timed out: {e}")
except Exception as e:
    # Handle other errors
    logger.error(f"RLM failed: {e}")

Real-Time Progress with EventQueue

RLM operations can be long-running (up to 5 minutes). Use InstrumentedRLM with EventQueue for real-time progress tracking and cancellation support.

InstrumentedRLM

InstrumentedRLM subclasses dspy.RLM (via TolerantRLM) to emit events at each iteration:

from cogniverse_agents.inference import InstrumentedRLM, RLMCancelledError
from cogniverse_core.events import get_queue_manager

# Create event queue for real-time progress
manager = get_queue_manager()
event_queue = await manager.create_queue("task_123", "tenant_1")

# InstrumentedRLM emits events automatically
rlm = InstrumentedRLM(
    "context, query -> answer",
    max_iterations=10,
    event_queue=event_queue,
    task_id="task_123",
    tenant_id="tenant_1",
)

# Events emitted during processing:
# - StatusEvent(WORKING, phase="rlm_start")
# - ProgressEvent(current=0, total=10, step="iteration_1")
# - ProgressEvent(current=1, total=10, step="iteration_2")
# - ...
# - StatusEvent(COMPLETED, phase="rlm_complete")

result = rlm(context=large_context, query="Summarize this")

Events require a queue and a task id, so _emit_sync takes a builder and calls it only when both are present.

rlm_run_span

Every promoted recursive-LM call is wrapped in rlm_run_span, the one seam that emits the InstrumentedRLM.run span (RLM_RUN_SPAN_NAME). The orchestrator's sufficiency gate and the ingest path's claim extraction both use it, so both produce the same span name and the same max_iterations / rlm_iterations attributes (RLM_RUN_SPAN_ATTRIBUTES); each call site adds its own on top.

from cogniverse_agents.inference import rlm_run_span

with rlm_run_span(
    telemetry_manager,
    tenant_id=tenant_id,
    max_iterations=module.max_iterations,
    attributes={"segment_id": segment_id},
) as run_span:
    prediction = module(**inputs)
    run_span.record(prediction)

The span opens in the tenant's project and nests under whatever span is current, so it hangs off the caller's trace rather than starting a new one. record stamps rlm_iterations from Prediction.trajectory; the attribute is written at open time too, so its presence does not depend on the call reaching that point. Telemetry never gates the call: with no manager, or when the span cannot be opened or closed, the body still runs and only the recorder's writes are dropped. Exceptions from the body propagate.

Cancellation Support

Users can cancel RLM operations mid-execution via the CancellationToken:

# Cancel from another coroutine (cancel() is synchronous, not a coroutine)
event_queue.cancellation_token.cancel(reason="User requested")

# RLM raises RLMCancelledError when cancelled
try:
    result = rlm(context=context, query=query)
except RLMCancelledError as e:
    logger.info(f"RLM cancelled: {e.reason}")

Integration with RLMInference

RLMInference automatically uses InstrumentedRLM when event_queue is provided:

from cogniverse_agents.inference import RLMInference
from cogniverse_core.events import get_queue_manager
from cogniverse_foundation.config.unified_config import LLMEndpointConfig

manager = get_queue_manager()
event_queue = await manager.create_queue("task_123", "tenant_1")

rlm = RLMInference(
    llm_config=LLMEndpointConfig(model="openai/gpt-4o"),
    max_iterations=10,
    event_queue=event_queue,  # Enables InstrumentedRLM
    task_id="task_123",
    tenant_id="tenant_1",
)

# Progress events are emitted automatically
result = rlm.process(query="Summarize", context=large_context)

High-Value Use Cases

Use Case Description
Video Analysis Process large frame counts recursively
Multi-Document Search Aggregate results from many sources
Transcript Analysis Process long video/audio transcripts
Cross-Modal Fusion Combine results from multiple modalities

Testing

Unit Tests

Location: tests/agents/unit/

# tests/agents/unit/test_orchestrator_agent.py

import pytest
from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps, OrchestratorInput
from cogniverse_core.registries.agent_registry import AgentRegistry

class TestOrchestratorAgent:
    def test_initialization(self, config_manager):
        """Test agent initialization with deps (no tenant_id)"""
        registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
        agent = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

        assert agent.deps is not None

    async def test_process_query(self, config_manager):
        """Test query processing — tenant_id in request"""
        registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
        agent = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

        result = await agent._process_impl(
            OrchestratorInput(
                query="machine learning videos",
                tenant_id="test:unit",
            )
        )

        assert result is not None

    def test_tenant_agnostic_construction(self, config_manager):
        """Agent serves all tenants — one instance, per-request tenant_id"""
        registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
        agent = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

        # Same agent handles different tenants at request time
        # Memory is namespaced by (tenant_id, agent_name) at request time

Integration Tests

Location: tests/agents/integration/

# Example integration test (actual tests exist in tests/agents/integration/)

import pytest

@pytest.mark.integration
class TestVideoSearchAgentIntegration:
    @pytest.fixture
    def tenant_id(self):
        return "test_tenant_integration"

    @pytest.fixture
    def agent(self, tenant_id):
        """Create agent with real Vespa connection"""
        from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
        from cogniverse_foundation.config.utils import create_default_config_manager

        config_manager = create_default_config_manager()

        # Create agent — schema_loader is required
        agent = SearchAgent(
            deps=SearchAgentDeps(),
            config_manager=config_manager,
            schema_loader=schema_loader,
        )
        return agent

    def test_search_end_to_end(self, agent):
        """Test complete search flow — tenant_id required per-request"""
        results = agent.search_by_text(  # synchronous
            query="test query",
            tenant_id="test_tenant_integration",
            top_k=5,
        )

        assert isinstance(results, list)
        # Results depend on ingested data

Test Utilities

# tests/conftest.py

import pytest
from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager

@pytest.fixture
def test_tenant_id():
    """Unique tenant ID for tests"""
    import uuid
    return f"test_tenant_{uuid.uuid4().hex[:8]}"

@pytest.fixture
def cleanup_tenant_schemas(test_tenant_id):
    """Cleanup tenant schemas after test"""
    yield

    # Cleanup
    schema_manager = VespaSchemaManager(
        backend_endpoint="http://localhost",
        backend_port=8080
    )
    schema_manager.delete_tenant_schemas(test_tenant_id)

Audit Checklist

The criteria below are what every audit pass must check before accepting a planned feature as "done." The list combines the user-stated criteria and the failure modes prior audits surfaced; expand it as new failure patterns are found.

A. Plan completeness

  • Every plan item is either WIRED with a real-service integration test, or SCOPED OUT with an entry in the plan's "Scoping decisions" section.
  • No silent skips. If a planned feature was deferred, it is documented in the plan, not just in commit messages.
  • Per-item verification (round-trip / integration test / A/B / adversarial) described in the plan was run and recorded.

B. Test integrity (real-service flushes plumbing breaks)

  • Tests claiming "real Vespa" / "real Mem0" / "real DSPy" actually construct against the live fixture, not a MagicMock.
  • No _dspy_module = MagicMock(...) followed by an assertion on the stub's return value (assert out.answer == "STUB-...").
  • No constructor-acceptance tests masquerading as integration tests (the test must call the method whose wire it claims to verify).
  • No substring-only assertions that pass on garbage output. LM tests assert content + length bounds + sentence structure, not just "non-empty string."
  • No silent fallbacks accepted by the test (e.g., if the agent catches DSPy exceptions and returns [FALLBACK: ...], the test must reject that prefix).
  • Tests assert the exact primary outcome of the method (the row in Vespa, the response body, the captured HTTP payload), not just that the call returned a non-None value.

C. Production code quality

  • No TODO, FIXME, NotImplementedError (except in deliberate "not supported" rejections), STUB, or hard-coded placeholder return values in production paths.
  • No silent try/except blocks returning a sentinel that masks a real failure.
  • No backwards-compat shims left after the migration they cover has been completed.
  • No dead code: classes / files / methods with no production callers. (Prior audits caught one dead Mem0 vector-store adapter; recheck periodically.)
  • No half-finished implementations gated behind a flag that is never flipped.

D. Comment / label hygiene

  • No plan-internal labels in any form — A.1, B.5, C.6, C3.4, D.1, F2.x, F3.x, F4.x, F5.x, F7.x, H5–H12, M13–M17, L18–L22 — anywhere in source tree: docstrings, comments, class names, method names, variable names, constants, fixture data, file names, URLs, doc cross- references, commit messages.
  • No multi-paragraph comment blocks (3+ consecutive comment lines beyond a one-line WHY annotation).
  • No comments that explain WHAT instead of WHY ("# delete the row" above delete(row)).
  • No phase markers (Phase 1, Step N, for now, temporary, placeholder, this commit).
  • No silly variable names: data, result, temp, foo, helper, single-letter outside narrow loops.
  • Test failure messages preserved (don't strip the assertion's f"...; got {actual!r}" payload during a comment cleanup).

E. Code duplication

  • Vespa client construction is in one factory, not repeated across files.
  • Memory write paths route through one adapter (the SDK Backend interface), not bypassed via raw HTTP.
  • Knowledge-agent boilerplate (Deps + Input + Output Pydantic classes, memory_manager_factory constructor parameter, default-lambda _resolve_mm) lives in a shared base, not copy-pasted across each agent.
  • Test fixtures shared across files via re-export from a common conftest.py, not re-implemented per file.

F. Auth + session contract (added when auth lands)

  • Every public-facing request type (AgentTask, SearchRequest, future IngestionRequest variants that touch memory) takes session_id as a top-level optional field, AND the handler merges it into dispatch_context["session_id"] before dispatcher.dispatch.
  • The auth middleware / dependency populates session_id from the auth context (token, cookie, JWT claim) on the path the plan picks (trust-server vs trust-client).
  • One end-to-end test posts a real HTTP request through the middleware and asserts the resulting Mem0 row in Vespa carries the auth-derived session_id. Today this hop is exercised only with client-supplied session_id; the auth integration point is the gap that will move this from "session_id propagation works" to "auth produces session_id."

G. User-flow integrity

  • The chain HTTP entry → router → dispatcher → agent → memory write/read → cleanup has no missing link and no inconsistent contract between routers (e.g., routers/search.py and routers/agents.py must agree on whether session_id is top-level or context-bag).
  • Cleanup endpoints (per-tenant DELETE, cross-tenant POST close) each have a real-Vespa test that proves the right rows are deleted and unrelated rows survive.

H. Documentation

  • Public docs (under docs/) do not reference plan-internal labels. Cross-references use class/method names, not plan IDs.
  • Architecture decisions (S1, S2, future S-numbers) are recorded in the plan's "Scoping decisions" section AND mirrored in the relevant operations doc, so an operator sees the rationale without reading the plan.

Best Practices

1. Always Use Deps with Tenant ID

from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps
from cogniverse_core.registries.agent_registry import AgentRegistry

# ✅ Good: Explicit deps (no tenant_id — it arrives per-request)
registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
agent = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

# ❌ Bad: No deps or registry
agent = OrchestratorAgent()  # TypeError: missing deps

2. Initialize Config Manager for SearchAgent

from cogniverse_agents.search_agent import SearchAgent, SearchAgentDeps
from cogniverse_foundation.config.utils import create_default_config_manager

# SearchAgent requires schema_loader; config_manager is optional
config_manager = create_default_config_manager()
deps = SearchAgentDeps()
agent = SearchAgent(deps=deps, config_manager=config_manager, schema_loader=schema_loader)

# Search is synchronous — tenant_id required per-request
results = agent.search_by_text("query", tenant_id="acme", top_k=10)

3. Use Telemetry for Observability

# OrchestratorAgent automatically initializes telemetry per-request
from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps, OrchestratorInput
from cogniverse_core.registries.agent_registry import AgentRegistry

registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
orchestrator = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

# Operations traced per-request with tenant_id from A2A payload
result = await orchestrator._process_impl(
    OrchestratorInput(query="query", tenant_id="acme")
)

4. Test Tenant Isolation

# Always verify tenants are isolated
def test_tenant_isolation(config_manager):
    # ONE orchestrator serves all tenants — tenant_id is per-request
    from cogniverse_agents.orchestrator_agent import OrchestratorAgent, OrchestratorDeps
    from cogniverse_core.registries.agent_registry import AgentRegistry

    registry = AgentRegistry(tenant_id=tenant_id, config_manager=config_manager)
    orchestrator = OrchestratorAgent(deps=OrchestratorDeps(), registry=registry, config_manager=config_manager)

    # Tenant isolation verified at request time, not construction
    # Telemetry projects isolated: cogniverse-{tenant_id}-orchestrator

Real-Time Event Notifications

Overview

An OrchestratorAgent run reports its progress as a workflow task on the runtime's shared task event store, so any runtime process streams, lists and cancels it:

  • Multiple Subscribers: Web client + CLI can watch the same workflow simultaneously
  • Phase Events: Each phase boundary is a StatusEvent; the sufficiency-gate InstrumentedRLM adds Status/Progress events per REPL iteration
  • Graceful Cancellation: A cancelled workflow stops at its next phase boundary
  • Reconnection with Replay: Clients can resume from a specific event offset

Binding

One cached OrchestratorAgent serves every request, so its queue is bound per request, never held on the agent. The runtime's dispatcher binds the run's queue (AgentDispatcher.workflow_run); a caller outside the runtime binds its own:

from cogniverse_core.events import InMemoryEventQueue, bind_event_queue

queue = InMemoryEventQueue(task_id="workflow_123", tenant_id="acme:acme")
with bind_event_queue(queue):
    output = await orchestrator.process(OrchestratorInput(query="find cats", tenant_id="acme:acme"))
assert output.workflow_id == "workflow_123"

Event Flow

With a bound queue the orchestrator reports (AgentBase.report_phase, which also streams each phase to a process(stream=True) caller):

  1. memory_context, planning, execution — then retrieval_iteration at each iteration of the retrieval loop and executing at each plan step — aggregating and complete (deep_synthesis for a deep-synthesis run)
  2. The sufficiency-gate InstrumentedRLM (evidence too large for a single ChainOfThought call) emits its per-iteration StatusEvent/ProgressEvent sequence on the same task

The run takes the bound queue's task id as its workflow_id. The runtime ends the task with a CompleteEvent, a cancelled StatusEvent, or an ErrorEvent.

Cancellation

A cancellation of the task (POST /events/workflows/{workflow_id}/cancel, on any runtime process) reaches the process running it; the orchestrator raises TaskCancelled at its next phase boundary (every boundary but executing and complete, and before each group of plan steps), and the dispatch answers {"status": "cancelled", "workflow_id", "message"}. An inbound stop message (POST /agents/{name}/message) is separate: it ends the retrieval loop with the evidence gathered so far and the run completes.

See Events Module for complete EventQueue documentation.


Approval Workflow System

The Approval Workflow System provides human-in-the-loop (HITL) approval for AI-generated outputs.

Location: libs/agents/cogniverse_agents/approval/

Full Documentation: Approval Workflow Module

flowchart TB
    subgraph "Data Generation"
        SyntheticGen["<span style='color:#000'>Synthetic Data Generator</span>"]
        Extractor["<span style='color:#000'>Confidence Extractor</span>"]
        SyntheticGen --> Extractor
    end

    subgraph "Approval Workflow"
        ApprovalAgent["<span style='color:#000'>HumanApprovalAgent</span>"]
        Storage["<span style='color:#000'>ApprovalStorageImpl</span>"]

        Extractor --> ApprovalAgent
        ApprovalAgent --> Storage
    end

    subgraph "Review Interface"
        WebClient["<span style='color:#000'>Web Client</span>"]
        WebClient --> ApprovalAgent
    end

    subgraph "Training Pipeline"
        Optimizer["<span style='color:#000'>DSPy Optimizer</span>"]
        Storage --> Optimizer
    end

    style SyntheticGen fill:#ffcc80,stroke:#ef6c00,color:#000
    style Extractor fill:#ffcc80,stroke:#ef6c00,color:#000
    style ApprovalAgent fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Storage fill:#90caf9,stroke:#1565c0,color:#000
    style WebClient fill:#b0bec5,stroke:#546e7a,color:#000
    style Optimizer fill:#ffcc80,stroke:#ef6c00,color:#000

Key Components

Component Description
HumanApprovalAgent Orchestrates HITL approval with confidence-based auto-approval
DecisionOrchestrator Combines WorkflowStateMachine with approval checkpoints
ApprovalStorageImpl Phoenix approval spans and Redis-backed canonical replacement selection
ConfidenceExtractor Domain-specific confidence scoring (abstract)
FeedbackHandler Domain-specific rejection handling (abstract)

Quick Example

from cogniverse_agents.approval import (
    ApprovalStorageImpl,
    HumanApprovalAgent,
    ReviewDecision,
)
from cogniverse_synthetic.approval import (
    SyntheticDataConfidenceExtractor,
    SyntheticDataFeedbackHandler,
)
from cogniverse_synthetic.dspy_modules import ValidatedSyntheticExampleRegenerator

# Initialize
storage = ApprovalStorageImpl(
    grpc_endpoint="http://localhost:4317",
    http_endpoint="http://localhost:6006",
    tenant_id="acme:production",
    redis_url="redis://redis:6379/0",
)
regenerator = ValidatedSyntheticExampleRegenerator(max_retries=3)
regenerator.lm = configured_dspy_lm

agent = HumanApprovalAgent(
    storage=storage,
    confidence_extractor=SyntheticDataConfidenceExtractor(),
    feedback_handler=SyntheticDataFeedbackHandler(
        generator=regenerator,
        generation_timeout_seconds=primary_lm_config.request_timeout,
    ),
    confidence_threshold=0.8  # Auto-approve >= 0.8
)

# Process a batch of raw items — splits by confidence, auto-approves >= threshold,
# persists the batch to storage so the rest surface for human review
batch = await agent.process_batch(
    items=synthetic_data,
    batch_id="batch_001",
    context={"tenant_id": "acme:production", "agent_type": "routing"},
)

# Apply a human decision to one of the pending items. When approved and
# storage is an ApprovalStorageImpl, the item is appended to the training
# dataset named in batch.context["dataset_name"] automatically — there is
# no separate export call.
await agent.apply_decision(
    batch.batch_id,
    ReviewDecision(item_id=f"{batch.batch_id}_0", approved=True),
)

# Inspect approval stats for the batch (counts + rates, no LLM call)
print(agent.get_approval_stats(batch))

Workflow States

stateDiagram-v2
    [*] --> Generated: Synthetic data created
    Generated --> AutoApproved: confidence >= threshold
    Generated --> PendingReview: confidence < threshold
    PendingReview --> Approved: Human approves
    PendingReview --> Rejected: Human rejects
    AutoApproved --> TrainingDataset: Export
    Approved --> TrainingDataset: Export
    Rejected --> Regenerated: Feedback handler creates replacement
    Regenerated --> [*]: Persist canonical replacement

    classDef orange fill:#ffcc80,stroke:#ef6c00,color:#000
    classDef green fill:#a5d6a7,stroke:#388e3c,color:#000
    classDef blue fill:#90caf9,stroke:#1565c0,color:#000
    classDef purple fill:#ce93d8,stroke:#7b1fa2,color:#000

    class Generated orange
    class AutoApproved,Approved green
    class PendingReview blue
    class Rejected,Regenerated purple
    class TrainingDataset green

See Approval Workflow Module for complete documentation including:

  • ApprovalStorageImpl with Phoenix integration
  • ConfidenceExtractor implementations
  • Web client integration
  • Testing patterns

Inference System

The Inference subsystem provides RLM (Recursive Language Model) inference for handling large contexts that exceed model limits.

Location: libs/agents/cogniverse_agents/inference/

flowchart TD
    Query["<span style='color:#000'>Query + Large Context</span>"] --> RLM["<span style='color:#000'>RLMInference</span>"]

    RLM --> Check{"<span style='color:#000'>EventQueue<br/>provided?</span>"}

    Check -->|Yes| Instrumented["<span style='color:#000'>InstrumentedRLM</span>"]
    Check -->|No| Standard["<span style='color:#000'>TolerantRLM</span>"]

    Instrumented --> Events["<span style='color:#000'>Progress Events</span>"]
    Instrumented --> Cancel["<span style='color:#000'>Cancellation Check</span>"]

    Standard --> Process["<span style='color:#000'>REPL Iterations</span>"]
    Instrumented --> Process

    Process --> Result["<span style='color:#000'>RLMResult</span>"]

    style Query fill:#90caf9,stroke:#1565c0,color:#000
    style RLM fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Check fill:#ffcc80,stroke:#ef6c00,color:#000
    style Instrumented fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Standard fill:#b0bec5,stroke:#546e7a,color:#000
    style Events fill:#a5d6a7,stroke:#388e3c,color:#000
    style Cancel fill:#ffcc80,stroke:#ef6c00,color:#000
    style Process fill:#ffcc80,stroke:#ef6c00,color:#000
    style Result fill:#a5d6a7,stroke:#388e3c,color:#000

RLM Overview

RLM (Recursive Language Model) enables LLMs to handle near-infinite context by:

  1. Storing context as Python variables (not in LLM prompt)
  2. LLM generates code to inspect, filter, and partition context
  3. Spawns sub-LLMs recursively to process partitions
  4. Aggregates results via SUBMIT({fields})

Use Cases:

  • Large video frame analysis (100s of frames)
  • Multi-document aggregation
  • Long transcript processing
  • Search result synthesis

RLMInference

Location: libs/agents/cogniverse_agents/inference/rlm_inference.py

Wrapper around DSPy's RLM module with timeout handling and optional EventQueue integration.

from cogniverse_agents.inference import RLMInference, RLMResult, RLMTimeoutError
from cogniverse_foundation.config.unified_config import LLMEndpointConfig

# Basic usage
rlm = RLMInference(
    llm_config=LLMEndpointConfig(model="openai/gpt-4o"),
    max_iterations=10,
    max_llm_calls=30,
    timeout_seconds=300
)

result = rlm.process(
    query="Summarize the main findings",
    context=large_context_string,  # Can be 100K+ chars
    system_prompt="Focus on key insights"
)

print(f"Answer: {result.answer}")
print(f"Depth: {result.depth_reached}, Calls: {result.total_calls}")
print(f"Latency: {result.latency_ms:.0f}ms")

With EventQueue for Progress Tracking:

from cogniverse_core.events import get_queue_manager
from cogniverse_foundation.config.unified_config import LLMEndpointConfig

manager = get_queue_manager()
event_queue = await manager.create_queue("rlm_task_001", "acme")

rlm = RLMInference(
    llm_config=LLMEndpointConfig(model="openai/gpt-4o"),
    event_queue=event_queue,
    task_id="rlm_task_001",
    tenant_id="acme"
)

# Client receives progress events as RLM iterates
result = rlm.process(query="Analyze documents", context=docs)

RLMResult

Dataclass with telemetry data for A/B testing and monitoring.

@dataclass
class RLMResult:
    answer: str                    # Final answer from RLM
    depth_reached: int             # Actual recursion depth used
    total_calls: int               # Number of LLM sub-calls
    tokens_used: int               # Total tokens (if available)
    latency_ms: float              # End-to-end latency
    was_fallback: bool             # True when answer came from fallback extraction
                                   # because max_iterations was reached without SUBMIT().
    trajectory: list[dict]         # Bounded structured trajectory snapshots,
                                   # populated only when RLMOptions.include_trajectory=True
                                   # (cap via trajectory_max_entries).
    metadata: dict                 # Includes trajectory_length + trajectory_summary
                                   # for server-side debug regardless of opt-in.

    def to_telemetry_dict(self) -> Dict:
        """Export for telemetry/Phoenix."""
        return {
            "rlm_enabled": True,
            "rlm_depth_reached": self.depth_reached,
            "rlm_total_calls": self.total_calls,
            "rlm_was_fallback": self.was_fallback,
            "rlm_trajectory_length": len(self.trajectory),
            # Additional telemetry fields are derived from metadata.
        }

InstrumentedRLM

Location: libs/agents/cogniverse_agents/inference/instrumented_rlm.py

Subclass of dspy.RLM (via TolerantRLM, which hardens the Deno JSON-RPC channel against stale id: null messages) with real-time progress events and cancellation support.

from cogniverse_agents.inference import InstrumentedRLM, RLMCancelledError

rlm = InstrumentedRLM(
    "context, query -> answer",
    event_queue=queue,
    task_id="task_123",
    tenant_id="acme",
    max_iterations=10,
)

try:
    result = rlm(context=large_context, query="Summarize this")
except RLMCancelledError as e:
    print(f"Cancelled: {e.reason}")

Events Emitted:

Event Type Phase Description
StatusEvent rlm_start RLM processing started
ProgressEvent iteration_N Per-iteration progress (current/total)
StatusEvent rlm_extracting Max iterations, extracting fallback
StatusEvent rlm_complete Processing completed

Cancellation Support:

# Client can cancel via CancellationToken (cancel() is synchronous)
event_queue.cancel("User requested cancellation")

# InstrumentedRLM checks token at each iteration
# Raises RLMCancelledError if cancelled

Model String Format

RLMInference receives an LLMEndpointConfig whose model field is a litellm-prefixed model id. The Cogniverse chart always emits openai/<bare> for every engine (vLLM, Ollama, external); api_base selects the actual destination. External SaaS providers use their own litellm prefix:

Provider Model format Example
In-cluster vLLM / Ollama (chart) openai/{model} openai/google/gemma-4-e4b-it
OpenAI (SaaS) openai/{model} openai/gpt-4o
Anthropic anthropic/{model} anthropic/claude-3-5-sonnet-20241022
Together, Anyscale, etc. Provider-prefixed together/meta-llama/Llama-3-70b

Integration-test generation endpoints

Agent integration tests use the same pinned production models as the deployed runtime. INFERENCE_SERVICE_URLS may name the authenticated Modal endpoints for vllm_llm_student and vllm_asr, while COGNIVERSE_INFERENCE_API_KEY supplies their bearer credential. The fixture validates /v1/models through the canonical endpoint resolver before exposing either URL. A configured endpoint with the wrong model, revision, or credentials fails immediately and never falls through to a local model.

The Gemma fixture injects openai/google/gemma-4-e4b-it, the endpoint /v1 URL, and the resolved bearer key into LLMEndpointConfig; with no Modal URL configured it reads the Gemma deployment through the Modal lifecycle. The Whisper fixture prefers a configured Modal openai/whisper-large-v3-turbo and otherwise uses the cogniverse-e2e cluster's vllm_asr service. No model is started on the test host.


Architecture Documentation

Module Documentation

Integration Documentation


Summary: The Agents package provides tenant-aware agent implementations that integrate with the core SDK. All agents require tenant_id, use tenant-specific schemas, and support memory, telemetry, and health checks. The package includes intelligent profile selection (ProfileSelectionAgent), entity extraction (EntityExtractionAgent), multi-agent orchestration (OrchestratorAgent), ensemble search with RRF fusion (SearchAgent), real-time event notifications, human-in-the-loop approval workflows, A2A protocol tools, video playback tools, and RLM inference for large-context processing in fault-tolerant, observable workflows.