Skip to content

Backends Module (Vespa Integration)

Package: cogniverse_vespa (Implementation Layer) Location: libs/vespa/cogniverse_vespa/


Table of Contents

  1. Module Overview
  2. Package Structure
  3. Backend Configuration Architecture
  4. Backend Abstraction Layer
  5. Profile-Based Architecture
  6. Connection Pool Management
  7. Multi-Tenant Schema Management
  8. Search Backend
  9. Ingestion Client
  10. Schema Deployment
  11. Metadata Schema Management
  12. Supporting Modules
  13. Usage Examples
  14. Testing
  15. Best Practices
  16. VespaConfigStore API
  17. VespaEmbeddingProcessor
  18. BackendVectorStore (Mem0 Backend)
  19. RankingStrategyExtractor
  20. Tenant-Scoped Search via VespaSearchBackend

Module Overview

The Vespa package (cogniverse-vespa) provides backend integration for vector and hybrid search with multi-tenant support.

Key Features

  1. Multi-Tenant Schema Management: Physical isolation via schema-per-tenant
  2. Search Backend: Video search with ColPali and X-CLIP embeddings
  3. Ingestion: Batch document feeding with retry logic
  4. Schema Deployment: JSON-based schema parsing and deployment
  5. Tenant Isolation: Dedicated schemas for each tenant

Design Principles

  • Tenant-Aware: All clients require tenant-specific schema names
  • Schema-Per-Tenant: Physical data isolation via dedicated Vespa schemas
  • Core Integration: Depends on cogniverse_sdk and cogniverse_core packages
  • Production-Ready: Retry logic, health checks, batch processing

Package Dependencies

# Vespa package depends on:
from cogniverse_sdk.document import Document
from cogniverse_sdk.interfaces.backend import Backend, SearchBackend
from cogniverse_core.common.utils.retry import RetryConfig, retry_with_backoff
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from cogniverse_foundation.config.utils import get_config  # Lazy import

External Dependencies:

  • pyvespa==1.1.2: Official Vespa Python client
  • numpy==2.4.4: Array operations
  • pydantic==2.13.5: Configuration models
  • requests==2.33.1: Document API and deployment requests

Package Structure

libs/vespa/cogniverse_vespa/
├── __init__.py
├── _vespa_factory.py               # Internal: Vespa app factory (make_vespa_app, PersistentVespaOps, make_persistent_vespa_ops)
├── _yql.py                         # Internal: YQL escaping utilities (yql_quote)
├── backend.py                      # Backend abstraction
├── config/
│   ├── __init__.py
│   └── config_store.py             # Vespa-based config storage
├── config_utils.py                 # Configuration utilities
├── embedding_processor.py          # Embedding processing
├── ingestion_client.py             # Ingestion client (VespaPyClient class)
├── json_schema_parser.py           # JSON schema parsing
├── memory_config.py                # Memory configuration
├── metadata_schemas.py             # Metadata schema definitions
├── ranking_strategy_extractor.py   # Ranking strategy extraction
├── registry/
│   ├── __init__.py
│   └── adapter_store.py            # Adapter registry storage
├── search_backend.py               # Tenant-scoped search backend (VespaSearchBackend)
├── strategy_aware_processor.py     # Strategy-aware processing
└── vespa_schema_manager.py         # Multi-tenant schema management

Total Files: 15 Python modules (excluding __init__.py), including 2 subdirectories: config/, registry/ (plus 2 private modules: _vespa_factory.py, _yql.py)

Key Files:

  • vespa_schema_manager.py: Core tenant management
  • json_schema_parser.py: Schema parsing
  • ingestion_client.py: PyVespa wrapper for ingestion
  • search_backend.py: Search backend with connection pooling
  • backend.py: Unified backend abstraction

Note: Schema templates are JSON files located in configs/schemas/ at project root


Backend Configuration Architecture

Overview

Cogniverse uses a profile-based backend configuration system with multi-tenant support. Configuration is loaded from config.json with auto-discovery and supports deep merging of system base config with tenant-specific overlays.

Key Features:

  • Auto-Discovery: Automatic config.json discovery from standard locations
  • Profile-Based: Multiple processing profiles per backend (ColPali, X-CLIP, ColQwen-Omni, etc.)
  • Tenant Overlays: Tenant-specific config merges with system base
  • Deep Merge: System profiles + Tenant overrides = Merged configuration
  • Type-Safe: BackendConfig and BackendProfileConfig dataclasses

Configuration Auto-Discovery

Search Order (defined in cogniverse_foundation/config/utils.py:_discover_config_file()):

  1. COGNIVERSE_CONFIG environment variable (if set)
  2. configs/config.json (from current directory)
  3. ../configs/config.json (one level up)
  4. ../../configs/config.json (two levels up)
from cogniverse_foundation.config.utils import ConfigUtils, create_default_config_manager

config_manager = create_default_config_manager()
config_utils = ConfigUtils(tenant_id="acme", config_manager=config_manager)
backend_config = config_utils.get("backend")  # Auto-discovered and merged

The production factory requires BACKEND_URL in the process environment. BACKEND_PORT is optional and defaults to 8080. The backend type still comes from configs/config.json.

Auto-Discovery Flow

flowchart TD
    Start["<span style='color:#000'>ConfigUtils initialized<br/>tenant_id: acme</span>"] --> CheckEnv{"<span style='color:#000'>COGNIVERSE_CONFIG<br/>env var set?</span>"}

    CheckEnv -->|Yes| LoadEnv["<span style='color:#000'>Load from env var path</span>"]
    CheckEnv -->|No| Check1["<span style='color:#000'>Check: configs/config.json</span>"]

    Check1 --> Exists1{"<span style='color:#000'>File exists?</span>"}
    Exists1 -->|Yes| Load1["<span style='color:#000'>Load configs/config.json</span>"]
    Exists1 -->|No| Check2["<span style='color:#000'>Check: ../configs/config.json</span>"]

    Check2 --> Exists2{"<span style='color:#000'>File exists?</span>"}
    Exists2 -->|Yes| Load2["<span style='color:#000'>Load ../configs/config.json</span>"]
    Exists2 -->|No| Check3["<span style='color:#000'>Check: ../../configs/config.json</span>"]

    Check3 --> Exists3{"<span style='color:#000'>File exists?</span>"}
    Exists3 -->|Yes| Load3["<span style='color:#000'>Load ../../configs/config.json</span>"]
    Exists3 -->|No| UseDefaults["<span style='color:#000'>Use empty JSON base</span>"]

    LoadEnv --> Parse["<span style='color:#000'>Parse JSON</span>"]
    Load1 --> Parse
    Load2 --> Parse
    Load3 --> Parse
    UseDefaults --> GetTenantOverride

    Parse --> SystemBase["<span style='color:#000'>System Base Config</span>"]
    SystemBase --> GetTenantOverride["<span style='color:#000'>Check for tenant override<br/>ConfigScope.BACKEND<br/>tenant_id: acme</span>"]

    GetTenantOverride --> HasOverride{"<span style='color:#000'>Tenant override<br/>exists?</span>"}

    HasOverride -->|Yes| LoadOverride["<span style='color:#000'>Load tenant config from<br/>ConfigManager</span>"]
    HasOverride -->|No| Merge{"<span style='color:#000'>Deep Merge</span>"}

    LoadOverride --> Merge

    Merge --> FinalConfig["<span style='color:#000'>Final Backend Config<br/>System base + Tenant overlays</span>"]

    style Start fill:#90caf9,stroke:#1565c0,color:#000
    style CheckEnv fill:#ffcc80,stroke:#ef6c00,color:#000
    style LoadEnv fill:#ffcc80,stroke:#ef6c00,color:#000
    style Check1 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Check2 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Check3 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Exists1 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Exists2 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Exists3 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Load1 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Load2 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Load3 fill:#ffcc80,stroke:#ef6c00,color:#000
    style UseDefaults fill:#b0bec5,stroke:#546e7a,color:#000
    style Parse fill:#ffcc80,stroke:#ef6c00,color:#000
    style SystemBase fill:#b0bec5,stroke:#546e7a,color:#000
    style GetTenantOverride fill:#ffcc80,stroke:#ef6c00,color:#000
    style HasOverride fill:#ffcc80,stroke:#ef6c00,color:#000
    style LoadOverride fill:#ffcc80,stroke:#ef6c00,color:#000
    style Merge fill:#ce93d8,stroke:#7b1fa2,color:#000
    style FinalConfig fill:#a5d6a7,stroke:#388e3c,color:#000

Backend Configuration Structure

config.json Structure

{
  "backend": {
    "type": "vespa",
    "url": "http://localhost",
    "port": 8080,
    "profiles": {
      "video_colpali_smol500_mv_frame": {
        "type": "video",
        "description": "Frame-based ColPali for patch-level visual search",
        "schema_name": "video_colpali_smol500_mv_frame",
        "embedding_model": "TomoroAI/tomoro-colqwen3-embed-4b",
        "pipeline_config": {
          "extract_keyframes": true,
          "transcribe_audio": true,
          "keyframe_fps": 0.5
        },
        "strategies": {
          "segmentation": {"class": "FrameSegmentationStrategy", "params": {}},
          "embedding": {"class": "MultiVectorEmbeddingStrategy", "params": {}}
        },
        "embedding_type": "multi_vector",
        "schema_config": {
          "num_patches": 1024,
          "embedding_dim": 320,
          "binary_dim": 40
        }
      },
      "video_xclip_sv_chunk_6s": {
        "type": "video",
        "description": "X-CLIP for 30-second chunk embeddings",
        "schema_name": "video_xclip_sv_chunk_6s",
        "embedding_model": "microsoft/xclip-large-patch14",
        "embedding_type": "multi_vector",
        "schema_config": {
          "embedding_dim": 768,
          "binary_dim": 96
        }
      }
    }
  }
}

BackendProfileConfig Dataclass

from cogniverse_foundation.config.unified_config import BackendProfileConfig

profile = BackendProfileConfig(
    profile_name="video_colpali_smol500_mv_frame",
    type="video",
    description="Frame-based ColPali processing",
    schema_name="video_colpali_smol500_mv_frame",  # Vespa schema name
    embedding_model="TomoroAI/tomoro-colqwen3-embed-4b",
    model_loader="colpali",
    pipeline_config={
        "extract_keyframes": True,
        "transcribe_audio": True,
        "keyframe_fps": 0.5
    },
    strategies={
        "segmentation": {"class": "FrameSegmentationStrategy"},
        "embedding": {"class": "MultiVectorEmbeddingStrategy"}
    },
    embedding_type="multi_vector",
    schema_config={
        "num_patches": 1024,
        "embedding_dim": 320,
        "binary_dim": 40
    }
)

Profile Fields:

  • profile_name: Unique identifier for the profile
  • schema_name: Vespa schema name (without tenant suffix)
  • embedding_model: HuggingFace model ID or local path
  • model_loader: Loader class key (colpali, colqwen, xclip, colbert)
  • pipeline_config: Video processing pipeline settings
  • strategies: Processing strategy classes and params
  • embedding_type: Type of embeddings (multi_vector or single_vector)
  • schema_config: Schema-specific metadata (dimensions, patches, etc.)

BackendConfig Dataclass

from cogniverse_foundation.config.unified_config import BackendConfig, BackendProfileConfig

config = BackendConfig(
    tenant_id="acme",
    backend_type="vespa",
    url="http://localhost",
    port=8080,
    profiles={
        "video_colpali_smol500_mv_frame": profile1,
        "video_xclip_sv_chunk_6s": profile2
    },
    default_profiles={
        "video": {"profile": "video_colpali_smol500_mv_frame"}
    },
)

# Get specific profile
profile = config.get_profile("video_colpali_smol500_mv_frame")

# Add new profile
config.add_profile(new_profile)

Tenant Configuration Overlay

Deep Merge Algorithm

System base config + Tenant-specific overrides = Merged configuration

flowchart TB
    SystemConfig["<span style='color:#000'>System Base Config<br/>config.json backend section</span>"] --> Merge["<span style='color:#000'>Deep Merge Algorithm</span>"]
    TenantConfig["<span style='color:#000'>Tenant Override Config<br/>ConfigManager.get_backend_config</span>"] --> Merge
    Merge --> MergedConfig["<span style='color:#000'>Merged BackendConfig<br/>System profiles + Tenant profiles</span>"]

    SystemConfig -.->|"type: vespa<br/>url: localhost<br/>profiles: [colpali, xclip]"| Merge
    TenantConfig -.->|"url: custom-vespa.acme.com<br/>profiles: [acme_custom_profile]"| Merge
    MergedConfig -.->|"All system profiles +<br/>Tenant custom profiles +<br/>Tenant URL override"| Result["<span style='color:#000'>Available to Application</span>"]

    style SystemConfig fill:#b0bec5,stroke:#546e7a,color:#000
    style TenantConfig fill:#ffcc80,stroke:#ef6c00,color:#000
    style Merge fill:#ce93d8,stroke:#7b1fa2,color:#000
    style MergedConfig fill:#a5d6a7,stroke:#388e3c,color:#000
    style Result fill:#90caf9,stroke:#1565c0,color:#000

Merge Rules (from config/utils.py:_ensure_backend_config()):

  1. Profiles: Dict merge - tenant profiles override system profiles with same name
  2. Default Profiles: Modality merge - tenant selections override stored-system and shipped selections
  3. Backend Type: Tenant value OR system value (tenant takes precedence)
  4. URL: Tenant value if not default, otherwise system value
  5. Port: Tenant value if not default, otherwise system value
  6. Metadata: Dict merge - tenant metadata extends system metadata
# System config.json
{
  "backend": {
    "url": "http://localhost",
    "port": 8080,
    "profiles": {
      "video_colpali": {...},
      "video_xclip": {...}
    }
  }
}

# Tenant "acme" override (stored in ConfigManager)
tenant_config = BackendConfig(
    tenant_id="acme",
    url="http://vespa.acme.com",
    port=8080,
    profiles={
        "acme_custom_profile": {...}
    },
    default_profiles={
        "video": {"profile": "acme_custom_profile"}
    },
)

# Merged result for tenant "acme"
# → url: http://vespa.acme.com (tenant override)
# → profiles: {video_colpali, video_xclip, acme_custom_profile} (merged)
# → default_profiles.video.profile: acme_custom_profile (tenant override)

Partial Profile Updates

from cogniverse_foundation.config.unified_config import BackendConfig

# Merge overrides into existing profile
modified_profile = config.merge_profile(
    profile_name="video_colpali_smol500_mv_frame",
    overrides={
        "pipeline_config": {"keyframe_fps": 2.0},  # Only override FPS
        "embedding_model": "TomoroAI/tomoro-colqwen3-embed-4b"  # Update model
    }
)

# Original profile unchanged, returns new profile with merged values

Using Backend Configuration

Example 1: Load Merged Config for Tenant

from cogniverse_foundation.config.utils import ConfigUtils, create_default_config_manager

# Auto-discovers config.json and merges with tenant overrides
config_manager = create_default_config_manager()
config_utils = ConfigUtils(tenant_id="acme", config_manager=config_manager)

# Get merged backend config
backend_dict = config_utils.get("backend")

# Access profile
profiles = backend_dict["profiles"]
colpali_profile = profiles["video_colpali_smol500_mv_frame"]

Example 2: Get BackendConfig Object

from cogniverse_foundation.config.utils import create_default_config_manager
from cogniverse_foundation.config.unified_config import BackendConfig

manager = create_default_config_manager()

# Get tenant backend config (includes system base + tenant overlay merge)
backend_config: BackendConfig = manager.get_backend_config(tenant_id="acme")

# Get specific profile
profile = backend_config.get_profile("video_colpali_smol500_mv_frame")

print(f"Schema: {profile.schema_name}")
print(f"Model: {profile.embedding_model}")
print(f"Strategies: {profile.strategies.keys()}")

Example 3: Set Tenant-Specific Backend Config

from cogniverse_foundation.config.utils import create_default_config_manager
from cogniverse_foundation.config.unified_config import BackendConfig, BackendProfileConfig

manager = create_default_config_manager()

# Create tenant-specific profile
tenant_profile = BackendProfileConfig(
    profile_name="acme_high_fps",
    schema_name="video_colpali_smol500_mv_frame",
    embedding_model="TomoroAI/tomoro-colqwen3-embed-4b",
    model_loader="colpali",
    pipeline_config={"keyframe_fps": 5.0},  # 5 FPS instead of 1 FPS
    embedding_type="multi_vector"
)

# Set tenant backend config
tenant_backend = BackendConfig(
    tenant_id="acme",
    url="http://vespa.acme.com",
    profiles={"acme_high_fps": tenant_profile}
)

manager.set_backend_config(tenant_backend)

Architecture Diagram

flowchart TB
    App["<span style='color:#000'>Application</span>"] --> ConfigUtils["<span style='color:#000'>ConfigUtils<br/>tenant_id='acme'</span>"]

    ConfigUtils --> AutoDiscover["<span style='color:#000'>Auto-Discover<br/>config.json</span>"]
    AutoDiscover --> SearchPath1["<span style='color:#000'>1. COGNIVERSE_CONFIG env</span>"]
    AutoDiscover --> SearchPath2["<span style='color:#000'>2. configs/config.json</span>"]
    AutoDiscover --> SearchPath3["<span style='color:#000'>3. ../configs/config.json</span>"]

    ConfigUtils --> LoadSystem["<span style='color:#000'>Load System Base<br/>backend section</span>"]
    ConfigUtils --> LoadTenant["<span style='color:#000'>Load Tenant Override<br/>ConfigManager</span>"]

    LoadSystem --> Merge["<span style='color:#000'>Deep Merge</span>"]
    LoadTenant --> Merge

    Merge --> MergedBackend["<span style='color:#000'>Merged BackendConfig</span>"]
    MergedBackend --> Profiles{"<span style='color:#000'>Available Profiles</span>"}

    Profiles --> ColPali["<span style='color:#000'>video_colpali<br/>System</span>"]
    Profiles --> X-CLIP["<span style='color:#000'>video_xclip<br/>System</span>"]
    Profiles --> AcmeCustom["<span style='color:#000'>acme_custom<br/>Tenant Override</span>"]

    MergedBackend --> VespaBackend["<span style='color:#000'>VespaBackend</span>"]
    VespaBackend --> Application["<span style='color:#000'>Application Logic</span>"]

    style App fill:#90caf9,stroke:#1565c0,color:#000
    style ConfigUtils fill:#ffcc80,stroke:#ef6c00,color:#000
    style AutoDiscover fill:#ffcc80,stroke:#ef6c00,color:#000
    style SearchPath1 fill:#b0bec5,stroke:#546e7a,color:#000
    style SearchPath2 fill:#b0bec5,stroke:#546e7a,color:#000
    style SearchPath3 fill:#b0bec5,stroke:#546e7a,color:#000
    style LoadSystem fill:#ffcc80,stroke:#ef6c00,color:#000
    style LoadTenant fill:#ffcc80,stroke:#ef6c00,color:#000
    style Merge fill:#ce93d8,stroke:#7b1fa2,color:#000
    style MergedBackend fill:#a5d6a7,stroke:#388e3c,color:#000
    style Profiles fill:#b0bec5,stroke:#546e7a,color:#000
    style ColPali fill:#90caf9,stroke:#1565c0,color:#000
    style X-CLIP fill:#90caf9,stroke:#1565c0,color:#000
    style AcmeCustom fill:#ffcc80,stroke:#ef6c00,color:#000
    style VespaBackend fill:#90caf9,stroke:#1565c0,color:#000
    style Application fill:#a5d6a7,stroke:#388e3c,color:#000

Backend Abstraction Layer

VespaBackend Class

Location: libs/vespa/cogniverse_vespa/backend.py

Purpose: Unified backend interface that wraps VespaSearchBackend and VespaPyClient, providing a single abstraction for both search and ingestion operations.

Self-registration: backend.py calls its own module-level register() function on import, which registers VespaBackend under the name "vespa" via cogniverse_core.registries.backend_registry.register_backend — this is what makes BackendRegistry.get_search_backend(name="vespa", ...) and get_ingestion_backend(name="vespa", ...) resolve below.

Recommended Pattern: Use the BackendRegistry to obtain backend instances:

from cogniverse_core.registries.backend_registry import BackendRegistry
from cogniverse_foundation.config.utils import create_default_config_manager
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from pathlib import Path

# Create required dependencies
config_manager = create_default_config_manager()
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))

# Get shared search backend from registry (handles instantiation and caching)
backend = BackendRegistry.get_search_backend(
    name="vespa",
    config_manager=config_manager,
    schema_loader=schema_loader
)

# Search — tenant_id is passed in query_dict for schema name derivation
results = backend.search({
    "query": "cooking video",
    "type": "video",
    "tenant_id": "acme",
    "top_k": 10
})

# For ingestion, use get_ingestion_backend
ingestion_backend = BackendRegistry.get_ingestion_backend(
    name="vespa",
    tenant_id="acme",
    config_manager=config_manager,
    schema_loader=schema_loader
)
ingestion_backend.ingest_documents(documents, schema_name="video_colpali_smol500_mv_frame")

Key Features:

  • Unified Interface: Single class for search + ingestion
  • Profile-Aware: Automatically uses profile config from BackendConfig
  • Shared Search Backend: One search backend instance per backend endpoint serves all tenants; tenant_id passed in query_dict
  • Tenant-Isolated Ingestion: Ingestion backends are per (tenant, endpoint) for schema isolation
  • Lazy Initialization: Components created on-demand per operation

Architecture Diagram

flowchart TB
    App["<span style='color:#000'>Application Code</span>"] --> VespaBackend["<span style='color:#000'>VespaBackend<br/>Unified Interface</span>"]

    VespaBackend --> SearchBackend["<span style='color:#000'>VespaSearchBackend</span>"]
    VespaBackend --> IngestionClient["<span style='color:#000'>VespaPyClient</span>"]

    VespaBackend --> SchemaManager["<span style='color:#000'>VespaSchemaManager</span>"]
    VespaBackend --> TenantManager["<span style='color:#000'>VespaSchemaManager</span>"]

    SearchBackend --> VespaInst["<span style='color:#000'>Vespa Instance</span>"]

    IngestionClient --> PyVespa["<span style='color:#000'>PyVespa feed_iterable</span>"]
    PyVespa --> VespaInst

    SchemaManager --> VespaInst
    TenantManager --> VespaInst

    style App fill:#90caf9,stroke:#1565c0,color:#000
    style VespaBackend fill:#ce93d8,stroke:#7b1fa2,color:#000
    style SearchBackend fill:#ffcc80,stroke:#ef6c00,color:#000
    style IngestionClient fill:#ffcc80,stroke:#ef6c00,color:#000
    style SchemaManager fill:#ffcc80,stroke:#ef6c00,color:#000
    style TenantManager fill:#ffcc80,stroke:#ef6c00,color:#000
    style VespaInst fill:#a5d6a7,stroke:#388e3c,color:#000
    style PyVespa fill:#b0bec5,stroke:#546e7a,color:#000

Why VespaBackend? - Eliminates Vespa-specific imports: Application code doesn't import VespaSearchBackend or VespaPyClient directly - Simplified API: One class instead of multiple clients - Consistent interface: Same initialization and method signatures - Future-proof: Can swap Vespa with other backends without changing application code

Key Methods

VespaBackend implements both IngestionBackend and SearchBackend (from cogniverse_sdk.interfaces.backend) plus schema-lifecycle and metadata-document operations:

Method Purpose
search(query_dict) Tenant-scoped search (delegates to VespaSearchBackend)
ingest_documents(documents, schema_name) Batch ingest via a per-tenant VespaPyClient
put_document(document, schema_name, namespace=None, base_schema_name=None) Full-put a generic cogniverse_sdk Document, serialized through the base schema's declared document_mapping block; raises ValueError when the schema declares no mapping
conditional_put_document(document, *, condition, schema_name=None, namespace=None, base_schema_name=None, create=True) Test-and-set variant of put_document: issues a Document v1 conditional partial update so a read-modify-write only lands while condition still holds. Returns True when applied, False on a 412 condition mismatch (a racing writer advanced the doc), and raises on transport/other errors. create=True inserts a missing doc (Vespa ignores the condition when the target is absent)
put_document_fields / get_document_fields / update_document_fields / delete_document_fields Namespace-aware raw-fields document CRUD (the primitives wiki/graph/ingestion use)
feed(document, schema_name) Feed a single document
ingest_stream(documents, schema_name) Stream ingestion for large datasets
update_document(document_id, document, schema_name) Partial or full document update; raises on backend failure or id mismatch (False means the write was rejected, never that the backend was unreachable)
delete_document(document_id, schema_name) Delete a single document. A genuine 404 is an idempotent success; connection failures and every other rejected status raise with the document route.
get_document(document_id, schema_name) / batch_get_documents(document_ids, schema_name) Point lookups that reconstruct stored Vespa tensors through Document.add_embedding, so Document.get_embedding(name) returns the embedding data rather than a storage envelope. When a shared search backend handles a batch read, the unified backend resolves the matching Document v1 namespace and passes both schema and namespace explicitly.
deploy_schemas(schema_definitions, allow_schema_removal=False) Low-level deploy of one or more schema definitions in a single Vespa application package. Registry or config-server enumeration failures abort before the package is sent. Every live schema the package does not carry is rebuilt from its registry row or, for an activation another process has not registered yet, from its deployment intent (VespaSchemaManager.reconstruct_unknown_schemas); one neither can rebuild refuses the deploy. The registry and config-server enumeration, the merge and the activation all run under the cross-process deployment lease, so the package cannot be built from a snapshot another process has already moved past; the convergence wait below runs outside the lease. A package the config server refuses raises BackendDeploymentError carrying the server's reason; a 400 INVALID_APPLICATION_PACKAGE, which a change needing a validation override (a field type or indexing change that requires a refeed) gets, raises its subclass SchemaChangeRefusedError, since posting the same package again is refused the same way. Returns only once the activated config generation runs on every Vespa service (serviceconverge) and each schema new to the cluster has accepted a probe feed over document/v1, within SCHEMA_CONVERGENCE_TIMEOUT_S (120s, sized to outlast a configproxy restart); a service that never reaches the generation or a schema whose feed is refused raises SchemaConvergenceError carrying the activated generation.
deployment_lease() Context manager holding the cross-process deployment lease on the calling thread (VespaSchemaManager.deployment_lease). A deploy_schemas made while it is held reuses it, so a caller keeps every other deployer out from its deploy decision through the registration after the convergence wait.
delete_schema(schema_name, tenant_id=None) / schema_exists(schema_name, tenant_id=None) Schema lifecycle. Tenant deletion uses the canonical tenant suffix only. Registry tombstone failures surface after Vespa removal so a retry can finish durable cleanup. schema_exists reads the tenant's stored registry row on every call, never what this process's registry holds, so a schema another process dropped or registered is seen at once; it (and validate_schema) raises on an enumeration/registry outage rather than returning False.
get_tenant_schema_name(tenant_id, base_schema_name) Delegates to self.schema_manager
create_metadata_document / get_metadata_document / query_metadata_documents / delete_metadata_document Organization/tenant/config metadata CRUD; writes raise on a backend outage (a bool False is a rejected write, not an unreachable backend). Passing tenant_id to query_metadata_documents resolves the base schema to the canonical tenant schema and rewrites a direct YQL source only when it names that base schema exactly. Whether that tenant schema is deployed comes from the backend's own DeployedSchemaNames, with no config-store read per query: an undeployed schema answers no rows without a query, and another process's drop is seen within DEPLOYED_SCHEMAS_MAX_STALENESS_S (30 s). A query sent meanwhile that Vespa refuses because it cannot resolve the schema answers no rows when the stored registry row confirms the drop, and raises otherwise.
health_check() / close() Lifecycle management

Instance Identity and Caching

BackendRegistry caches instances in a bounded LRU (TenantLRUCache, capacity 16, configure_tenant_cache_capacity to change it). Keys carry the endpoint the instance binds at initialize():

Backend kind Cache key
search search_<name>@<url>:<port>
ingestion ingestion_<name>_<tenant_id>@<url>:<port>
full (search + ingestion) backend_<name>_<tenant_id>@<url>:<port>

url/port come from SystemConfig, overridden by config["backend"] and then by top-level config. Profiles and default_profiles are not part of the key: search() merges the query tenant's profiles per request, read through get_config (the shipped catalog, the system tenant's stored profiles and the tenant's own), into a local copy of the profiles the instance was built with. A profile written at runtime reaches searches through that read; no profile write changes a cached instance.

Concurrent cold starts of one key build in parallel and resolve through set_if_absent; the losing builds are closed. config_manager and schema_loader are not part of the key. Every registry-built backend exposes the ones it was constructed with (VespaBackend.config_manager, VespaBackend.schema_loader), and a cache hit compares them by the source they read — config_manager.store.source and schema_loader.source — not by object: a fresh ConfigManager over the same store and a fresh FilesystemSchemaLoader over the same directory share the cached instance, while a requester reading another store or schema directory gets BackendBindingConflictError naming both sources. Sources are canonical: VespaConfigStore.source spells its endpoint through canonical_endpoint in cogniverse_vespa/_vespa_factory.py (scheme and host lowercased, explicit port, every loopback spelling written localhost, no name resolution), and FilesystemSchemaLoader.source is the resolved directory, so http://127.0.0.1:8080 and http://localhost:8080 are one store.

BackendRegistry.lease_instance(instance) checks a cached backend out for the length of a request: eviction skips it and takes the least-recently-used free entry instead, and a close aimed at it waits for the request to finish. VespaBackend.search() holds a checkout for the whole query, because other requesters insert into the same cache the backend running the query lives in. A backend the cache does not hold — built directly, or already evicted — checks out nothing.

Eviction closes the instance. VespaBackend.close() releases the search backend, the ingestion clients and the metadata client, all of which are otherwise rebuilt lazily on next use, so a closed instance refuses work: search(), the ingestion-client accessor and the metadata accessor raise BackendClosedError naming the endpoint. A caller holding an evicted backend is refused rather than quietly building a second set of connections outside the cache.

Undeployed Tenant Schemas

Schemas are deployed per tenant, so a search whose profile has no schema for the requesting tenant has nothing to read. VespaSearchBackend.search raises SchemaNotDeployedError (cogniverse_sdk.interfaces.backend) naming the tenant, the tenant-scoped schema, the base schema and the profile, before any query is issued. Returning an empty result made a never-deployed schema indistinguishable from a corpus with no match.

Query Encoder Faults

When a strategy needs query embeddings and the caller supplies none, the search backend builds the encoder the profile declares. Two distinct failures surface as distinct types (cogniverse_core.query.encoders):

  • EncoderNotConfiguredError (a ValueError) — the profile names no model, names an inference service with no configured URL, or omits a dimension the encoder needs.
  • EncoderUnavailableError (a RuntimeError) — the encoder is configured but its inference service did not serve the request. Carries profile, service, endpoint, and chains the underlying error.

Both encoder branches classify the same way. Dense profiles (embedding_type: "dense" or encoder: "denseon") resolve their embedder from the profile's own inference_services.embedding, so a profile's queries are answered from the space that indexed its documents; profiles naming no service use the runtime's configured embedder. An InferenceServiceUnavailableError that carries module means the service has no URL and no in-process backend, which is a configuration gap, and is reported as EncoderNotConfiguredError.

Profile-Based Architecture

What is a Profile?

A profile is a complete content processing configuration that defines: 1. Model Loader: Which loader class to use (colpali, colqwen, xclip, colbert) — the model_loader config key 2. Embedding Model: Which model to use (ColPali, X-CLIP, ColQwen, ColBERT) 3. Embedding Type: Processing mode (multi_vector or single_vector) 4. Processing Pipeline: Keyframe extraction, transcription, description generation 5. Segmentation Strategy: Frame-based, chunk-based, direct video, document segments, or audio segments 6. Vespa Schema: Which schema structure to use (document_text, audio_content, or video schemas) 7. Ranking Strategies: How to score and rank results

Profile Types

Multi-Profile Architecture

flowchart TB
    subgraph Profiles["<span style='color:#000'>Backend Profiles</span>"]
        ColPali["<span style='color:#000'>video_colpali_smol500_mv_frame<br/>Frame-Based<br/>1024 patches × 320-dim<br/>Binary embeddings</span>"]
        X-CLIP["<span style='color:#000'>video_xclip_sv_chunk_6s<br/>Direct Video<br/>768-dim global<br/>30s chunks</span>"]
        ColQwen["<span style='color:#000'>video_colqwen_omni_mv_chunk_30s<br/>Chunk-Based<br/>Multi-modal<br/>Audio + Visual</span>"]
    end

    subgraph QueryTime["<span style='color:#000'>Query-Time Selection</span>"]
        Query["<span style='color:#000'>User Query</span>"] --> AutoSelect{"<span style='color:#000'>Auto-Select Profile</span>"}
        AutoSelect -->|has_video| SelectX-CLIP
        AutoSelect -->|Fine-grained search| SelectColPali
        AutoSelect -->|Multimodal| SelectColQwen
    end

    subgraph Strategies["<span style='color:#000'>Processing Strategies</span>"]
        SelectColPali["<span style='color:#000'>ColPali</span>"] --> FrameSeg["<span style='color:#000'>FrameSegmentationStrategy<br/>1 FPS keyframe extraction</span>"]
        SelectX-CLIP["<span style='color:#000'>X-CLIP</span>"] --> DirectVideo["<span style='color:#000'>DirectVideoStrategy<br/>No frame extraction</span>"]
        SelectColQwen["<span style='color:#000'>ColQwen</span>"] --> ChunkSeg["<span style='color:#000'>ChunkSegmentationStrategy<br/>30s audio+visual chunks</span>"]
    end

    subgraph VespaSchemas["<span style='color:#000'>Vespa Schemas per Tenant (bare tenant_id 'acme' canonicalizes to 'acme:acme')</span>"]
        FrameSeg --> ColPaliSchema["<span style='color:#000'>video_colpali_smol500_mv_frame_acme_acme<br/>Multi-vector binary</span>"]
        DirectVideo --> X-CLIPSchema["<span style='color:#000'>video_xclip_sv_chunk_6s_acme_acme<br/>Global float vectors</span>"]
        ChunkSeg --> ColQwenSchema["<span style='color:#000'>video_colqwen_omni_mv_chunk_30s_acme_acme<br/>Multi-modal vectors</span>"]
    end

    style Profiles fill:#90caf9,stroke:#1565c0,color:#000
    style QueryTime fill:#ffcc80,stroke:#ef6c00,color:#000
    style Strategies fill:#ce93d8,stroke:#7b1fa2,color:#000
    style VespaSchemas fill:#a5d6a7,stroke:#388e3c,color:#000
    style ColPali fill:#90caf9,stroke:#1565c0,color:#000
    style X-CLIP fill:#90caf9,stroke:#1565c0,color:#000
    style ColQwen fill:#90caf9,stroke:#1565c0,color:#000
    style Query fill:#ffcc80,stroke:#ef6c00,color:#000
    style AutoSelect fill:#ffcc80,stroke:#ef6c00,color:#000
    style SelectColPali fill:#ce93d8,stroke:#7b1fa2,color:#000
    style SelectX-CLIP fill:#ce93d8,stroke:#7b1fa2,color:#000
    style SelectColQwen fill:#ce93d8,stroke:#7b1fa2,color:#000
    style FrameSeg fill:#ce93d8,stroke:#7b1fa2,color:#000
    style DirectVideo fill:#ce93d8,stroke:#7b1fa2,color:#000
    style ChunkSeg fill:#ce93d8,stroke:#7b1fa2,color:#000
    style ColPaliSchema fill:#a5d6a7,stroke:#388e3c,color:#000
    style X-CLIPSchema fill:#a5d6a7,stroke:#388e3c,color:#000
    style ColQwenSchema fill:#a5d6a7,stroke:#388e3c,color:#000

Frame-Based Profiles

Example: video_colpali_smol500_mv_frame - Extracts keyframes at fixed FPS (1-5 FPS) - Generates patch-level embeddings per frame - Schema: Multi-vector with 1024 patches × 320 dimensions - Best for: Fine-grained visual search, specific objects/text in frames

Chunk-Based Profiles

Example: video_colqwen_omni_mv_chunk_30s - Segments video into 30-second chunks - Processes audio + visual together - Schema: Multi-vector with multimodal understanding - Best for: Semantic content search, audio+visual comprehension

Direct Video Profiles

Example: video_xclip_sv_chunk_6s - Native video understanding without keyframes - Global 768-dim or 1024-dim embeddings - Schema: High-dimensional global vectors - Best for: Video-level semantic similarity, scene understanding

Profile Selection at Query Time

from cogniverse_foundation.config.utils import create_default_config_manager, get_config
from cogniverse_vespa.backend import VespaBackend
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from pathlib import Path

# Get configuration
config_manager = create_default_config_manager()
config = get_config(tenant_id="acme", config_manager=config_manager)

# List available profiles from backend config
backend_config = config_manager.get_backend_config("acme")
profiles = list(backend_config.profiles.keys())
# → ['video_colpali_smol500_mv_frame', 'video_xclip_sv_chunk_6s', ...]

# Select profile dynamically
profile_name = "video_colpali_smol500_mv_frame"

# Get shared search backend from registry with profile configuration
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))
backend = BackendRegistry.get_search_backend(
    name="vespa",
    config={"profile": profile_name},
    config_manager=config_manager,
    schema_loader=schema_loader
)

Creating Custom Profiles

# Add new profile to tenant config
custom_profile = BackendProfileConfig(
    profile_name="acme_ultra_high_quality",
    schema_name="video_colpali_smol500_mv_frame",  # Reuse existing schema
    embedding_model="TomoroAI/tomoro-colqwen3-embed-4b",
    model_loader="colpali",
    pipeline_config={
        "extract_keyframes": True,
        "keyframe_fps": 10.0,  # 10 FPS for ultra-high temporal resolution
        "transcribe_audio": True,
        "generate_descriptions": True
    },
    strategies={
        "segmentation": {
            "class": "FrameSegmentationStrategy",
            "params": {"fps": 10.0, "max_frames": 10000}
        },
        "embedding": {"class": "MultiVectorEmbeddingStrategy"}
    },
    embedding_type="multi_vector"
)

# Save to tenant config
backend_config.add_profile(custom_profile)
manager.set_backend_config(backend_config)

Advanced Query-Time Resolution

The VespaSearchBackend implements query-time resolution for profiles and strategies directly in the search() method with a 4-step fallback approach:

Resolution Flow

flowchart TD
    Start["<span style='color:#000'>Query Request<br/>with type + query</span>"] --> ProfileCheck{"<span style='color:#000'>Has explicit<br/>profile param?</span>"}

    ProfileCheck -->|Yes| ValidateProfile["<span style='color:#000'>Validate profile exists</span>"]
    ProfileCheck -->|No| VideoCheck{"<span style='color:#000'>type is<br/>video?</span>"}
    VideoCheck -->|Yes| SelectedDefault{"<span style='color:#000'>Selected default<br/>video profile?</span>"}
    VideoCheck -->|No| CountTypeProfiles["<span style='color:#000'>Count profiles<br/>matching type</span>"]
    SelectedDefault -->|Yes| UseDefault
    SelectedDefault -->|No| ErrorNoDefault["<span style='color:#000'>Error:<br/>No default video profile</span>"]

    ValidateProfile --> UseExplicit["<span style='color:#000'>Use Explicit Profile</span>"]

    CountTypeProfiles --> CheckCount{"<span style='color:#000'>How many<br/>profiles?</span>"}

    CheckCount -->|"1"| AutoSelect["<span style='color:#000'>Auto-select<br/>single profile</span>"]
    CheckCount -->|"> 1"| CheckDefault{"<span style='color:#000'>Has default<br/>for type?</span>"}
    CheckCount -->|"0"| ErrorNoProfile["<span style='color:#000'>Error:<br/>No profiles for type</span>"]

    CheckDefault -->|Yes| UseDefault["<span style='color:#000'>Use Default Profile</span>"]
    CheckDefault -->|No| ErrorMultiple["<span style='color:#000'>Error:<br/>Multiple profiles,<br/>no default</span>"]

    UseExplicit --> StrategyCheck
    AutoSelect --> StrategyCheck
    UseDefault --> StrategyCheck

    StrategyCheck{"<span style='color:#000'>Has explicit<br/>strategy param?</span>"} -->|Yes| UseExplicitStrategy["<span style='color:#000'>Use Explicit Strategy</span>"]
    StrategyCheck -->|No| StrategySimilar["<span style='color:#000'>Similar logic:<br/>Count → Auto-select → Default</span>"]

    UseExplicitStrategy --> TenantScoping
    StrategySimilar --> TenantScoping

    TenantScoping["<span style='color:#000'>Tenant Schema Scoping<br/>base_schema + tenant_id</span>"] --> ExecuteSearch["<span style='color:#000'>Execute Search</span>"]

    style Start fill:#90caf9,stroke:#1565c0,color:#000
    style ExecuteSearch fill:#a5d6a7,stroke:#388e3c,color:#000
    style ErrorNoProfile fill:#ffcccc,stroke:#c62828,color:#000
    style ErrorMultiple fill:#ffcccc,stroke:#c62828,color:#000
    style ErrorNoDefault fill:#ffcccc,stroke:#c62828,color:#000
    style VideoCheck fill:#ffcc80,stroke:#ef6c00,color:#000
    style SelectedDefault fill:#ffcc80,stroke:#ef6c00,color:#000
    style ProfileCheck fill:#ffcc80,stroke:#ef6c00,color:#000
    style StrategyCheck fill:#ffcc80,stroke:#ef6c00,color:#000
    style CheckCount fill:#ffcc80,stroke:#ef6c00,color:#000
    style CheckDefault fill:#ffcc80,stroke:#ef6c00,color:#000
    style CountTypeProfiles fill:#b0bec5,stroke:#546e7a,color:#000
    style ValidateProfile fill:#b0bec5,stroke:#546e7a,color:#000
    style UseExplicit fill:#ce93d8,stroke:#7b1fa2,color:#000
    style AutoSelect fill:#ce93d8,stroke:#7b1fa2,color:#000
    style UseDefault fill:#ce93d8,stroke:#7b1fa2,color:#000
    style UseExplicitStrategy fill:#ce93d8,stroke:#7b1fa2,color:#000
    style StrategySimilar fill:#ce93d8,stroke:#7b1fa2,color:#000
    style TenantScoping fill:#b0bec5,stroke:#546e7a,color:#000

Implementation Details

Location: libs/vespa/cogniverse_vespa/search_backend.py - search() method (starting line 592)

Profile Resolution Logic (inline in search() method):

# Priority order:
# 1. Explicit 'profile' parameter in query_dict
requested_profile = query_dict.get("profile")
if requested_profile:
    if requested_profile not in self.profiles:
        raise ValueError(f"Requested profile '{requested_profile}' not found")
    profile_name = requested_profile
elif content_type == "video":
    # Video never auto-selects: resolve_default_profile over the merged
    # default_profiles and the tenant's active_video_profile, or refuse.
    profile_name = self._selected_default_profile(tenant_id, default_profiles, "video")
    if not profile_name:
        raise ValueError(
            "No profile specified on the request and tenant ... has no "
            "configured default video profile."
        )
else:
    # 2. Auto-select if only one profile for content type
    type_profiles = {
        name: config
        for name, config in self.profiles.items()
        if config.get("type") == content_type
    }

    if len(type_profiles) == 1:
        profile_name = list(type_profiles.keys())[0]
    elif len(type_profiles) > 1:
        # 3. Use default profile for type
        default_config = self.default_profiles.get(content_type, {})
        profile_name = default_config.get("profile")
        if not profile_name:
            raise ValueError(
                f"Multiple profiles for '{content_type}' but no default configured"
            )
    else:
        # 4. No profiles for type - error
        raise ValueError(f"No profiles found for type '{content_type}'")

Strategy Resolution Logic (similar fallback approach):

# Same 4-step fallback:
# 1. Explicit strategy in query_dict
# 2. Auto-select if single strategy for profile
# 3. Use default strategy for profile/type
# 4. Error if no strategies found
requested_strategy = query_dict.get("strategy")
# ... similar logic pattern ...

Tenant Schema Scoping (inline construction):

# tenant_id is extracted from query_dict (REQUIRED), then canonicalized.
# canonical_tenant_id() maps a BARE tenant_id to "org:tenant" form
# ("acme" -> "acme:acme") before the colon is replaced with "_" — so a
# bare tenant_id's suffix is doubled. Only an already-canonical
# "org:tenant" input (e.g. "acme:prod") produces a single suffix.
tenant_id = query_dict.get("tenant_id")  # raises ValueError if missing
safe_tenant_id = canonical_tenant_id(tenant_id).replace(":", "_")
base_schema_name = profile_config.get("schema_name", profile_name)
schema_name = f"{base_schema_name}_{safe_tenant_id}"
# tenant_id="acme"      -> "video_colpali_smol500_mv_frame_acme_acme"
# tenant_id="acme:prod" -> "video_colpali_smol500_mv_frame_acme_prod"

Usage Example

Request with Auto-Resolution:

# Client request without explicit profile/strategy (REQUIRES 'type' key)
query_dict = {
    "query": "machine learning tutorial",
    "type": "video",  # REQUIRED for profile resolution
    "tenant_id": "acme:prod",  # REQUIRED for every search() call
    "top_k": 10,
    # No 'profile' or 'strategy' specified
}

# Backend resolves:
# 1. type="video" uses the tenant's selected default video profile
#    (default_profiles["video"]["profile"], then active_video_profile) and
#    raises ValueError when there is none; other types auto-select their
#    only profile, or use default_profiles[type]["profile"] when several
# 3. Similar logic for strategy
# 4. Schema: base_schema_name + "_" + canonicalized tenant_id
#    ("acme:prod" canonicalizes to itself -> "..._acme_prod")

results = backend.search(query_dict)

Once a profile is resolved (requested or auto-selected), its declared type owns the content type of every returned hit; the query's type is only the profile-selection key. A profile with no type is rejected with ValueError.

Request with Explicit Parameters:

query_dict = {
    "query": "cooking videos",
    "type": "video",  # REQUIRED
    "tenant_id": "acme:prod",  # REQUIRED
    "profile": "video_xclip_sv_chunk_6s",  # Explicit
    "strategy": "float_float",  # Explicit
    "top_k": 20
}

# Backend uses explicit values:
# 1. Profile: "video_xclip_sv_chunk_6s" (explicit)
# 2. Strategy: "float_float" (explicit, validated against profile)
# 3. Schema: "video_xclip_sv_chunk_6s_acme_prod" (tenant-scoped)

results = backend.search(query_dict)

Benefits

  1. Flexibility: Clients can control or let backend auto-select
  2. Sensible Defaults: Automatic selection based on query characteristics
  3. Tenant Isolation: Automatic schema scoping per tenant
  4. Performance: Strategy selection optimized for embedding type
  5. Simplicity: Clients don't need to know all configuration details

Connection Pool Management

Overview

The VespaSearchBackend implements connection pooling for efficient Vespa client management with health monitoring and automatic recovery.

Key Features:

  • Connection Reuse: A single, URL-scoped pool of persistent Vespa HTTP clients (not per-schema — the same pool serves every tenant schema at that backend URL)
  • Health Monitoring: A background thread probes every connection on health_check_interval with a real Vespa query
  • Automatic Recovery: Unhealthy or over-idle connections are closed and dropped; a fresh connection is created on the next demand up to max_connections
  • Bounded Growth: Pool grows from min_connections up to max_connections, blocking (with timeout) once the ceiling is reached
  • Metrics Tracking: Query-level latency and success/failure counts via SearchMetrics (separate from connection health)
  • Ordered Shutdown: VespaBackend.close() serializes concurrent callers, stops the search pool's health thread, and closes every search, ingestion, and metadata client before a service is removed

Architecture

flowchart TB
    Backend["<span style='color:#000'>VespaSearchBackend</span>"] --> GetConn["<span style='color:#000'>pool.get_connection()</span>"]

    GetConn --> HasAvailable{"<span style='color:#000'>_available<br/>non-empty?</span>"}
    HasAvailable -->|Yes| PopConn["<span style='color:#000'>Pop connection<br/>from _available</span>"]
    HasAvailable -->|No| UnderMax{"<span style='color:#000'>len(_connections)<br/>&lt; max_connections?</span>"}

    UnderMax -->|Yes| CreateConn["<span style='color:#000'>Create new VespaConnection</span>"]
    UnderMax -->|No| WaitCond["<span style='color:#000'>Wait on condition variable<br/>until connection_timeout</span>"]

    WaitCond --> Returned{"<span style='color:#000'>Connection returned<br/>before timeout?</span>"}
    Returned -->|Yes| PopConn
    Returned -->|No| TimeoutErr["<span style='color:#000'>Raise TimeoutError:<br/>No connections available</span>"]

    PopConn --> ExecuteQuery["<span style='color:#000'>conn.query() over persistent<br/>VespaSync HTTP client</span>"]
    CreateConn --> ExecuteQuery

    ExecuteQuery --> ReturnPool["<span style='color:#000'>finally: append to _available,<br/>notify waiters</span>"]

    subgraph HealthLoop["<span style='color:#000'>Background thread: _health_check_loop (every health_check_interval)</span>"]
        Snapshot["<span style='color:#000'>Snapshot _connections under lock</span>"] --> Probe["<span style='color:#000'>conn.health_check():<br/>real Vespa query</span>"]
        Probe --> Healthy{"<span style='color:#000'>Query<br/>succeeded?</span>"}
        Healthy -->|No| Unhealthy["<span style='color:#000'>is_healthy = False</span>"]
        Healthy -->|Yes| IdleCheck{"<span style='color:#000'>idle_time &gt; idle_timeout<br/>and above min_connections?</span>"}
        IdleCheck -->|Yes| Unhealthy
        IdleCheck -->|No| KeepConn["<span style='color:#000'>Keep connection in pool</span>"]
        Unhealthy --> RemoveConn["<span style='color:#000'>_remove_connection:<br/>close + drop from pool</span>"]
    end

    style Backend fill:#90caf9,stroke:#1565c0,color:#000
    style GetConn fill:#90caf9,stroke:#1565c0,color:#000
    style HasAvailable fill:#ffcc80,stroke:#ef6c00,color:#000
    style UnderMax fill:#ffcc80,stroke:#ef6c00,color:#000
    style Returned fill:#ffcc80,stroke:#ef6c00,color:#000
    style PopConn fill:#ce93d8,stroke:#7b1fa2,color:#000
    style CreateConn fill:#ce93d8,stroke:#7b1fa2,color:#000
    style WaitCond fill:#b0bec5,stroke:#546e7a,color:#000
    style ExecuteQuery fill:#a5d6a7,stroke:#388e3c,color:#000
    style ReturnPool fill:#a5d6a7,stroke:#388e3c,color:#000
    style TimeoutErr fill:#ffcccc,stroke:#c62828,color:#000
    style HealthLoop fill:#b0bec5,stroke:#546e7a,color:#000
    style Snapshot fill:#b0bec5,stroke:#546e7a,color:#000
    style Probe fill:#ffcc80,stroke:#ef6c00,color:#000
    style Healthy fill:#ffcc80,stroke:#ef6c00,color:#000
    style IdleCheck fill:#ffcc80,stroke:#ef6c00,color:#000
    style KeepConn fill:#a5d6a7,stroke:#388e3c,color:#000
    style Unhealthy fill:#ffcccc,stroke:#c62828,color:#000
    style RemoveConn fill:#ffcccc,stroke:#c62828,color:#000

Connection Pool Implementation

Location: libs/vespa/cogniverse_vespa/search_backend.py (lines 89-340)

ConnectionPoolConfig Class

@dataclass
class ConnectionPoolConfig:
    """Configuration for connection pool."""

    max_connections: int = 10           # Maximum connections in pool
    min_connections: int = 2            # Minimum connections to maintain
    connection_timeout: float = 30.0    # Timeout waiting for connection (seconds)
    idle_timeout: float = 300.0         # Remove idle connections after (seconds)
    health_check_interval: float = 60.0 # Health check frequency (seconds)

VespaConnection Class

class VespaConnection:
    """
    Managed Vespa connection with health checking.

    Queries run over a persistent VespaSync HTTP client (self._sync) so
    TCP connections are reused across searches — plain Vespa.query()
    builds and tears down a fresh client per call. health_check() uses
    the same fail-fast persistent session from the pool's background
    thread, including while the connection is checked out.

    Attributes:
        url: Vespa endpoint URL
        connection_id: Unique connection identifier
        vespa: Vespa client instance (created internally)
        created_at: Connection creation timestamp
        last_used: Last query execution timestamp
        is_healthy: Current health status
    """

    def __init__(self, url: str, connection_id: str):
        self.url = url
        self.connection_id = connection_id
        self.vespa = make_vespa_app(url=url)  # Created internally via make_vespa_app, not passed in
        self._sync = self.vespa.syncio(connections=4)
        self._sync._open_http_client()
        self.created_at = time.time()
        self.last_used = time.time()
        self.is_healthy = True
        self._lock = threading.Lock()

    def query(self, *args, **kwargs):
        """Execute query over the persistent client and update last used time."""
        with self._lock:
            self.last_used = time.time()
        return self._sync.query(*args, **kwargs)

    def close(self) -> None:
        """Release the persistent HTTP client."""
        self._sync._close_http_client()

    def health_check(self) -> bool:
        """
        Check connection health with simple query.

        Returns:
            True if connection is healthy
        """
        try:
            result = self.vespa.query(yql="select * from sources * where true limit 1")
            self.is_healthy = result is not None
            return self.is_healthy
        except Exception as e:
            logger.warning(f"Health check failed for {self.connection_id}: {e}")
            self.is_healthy = False
            return False

    @property
    def idle_time(self) -> float:
        """Time since last use in seconds."""
        return time.time() - self.last_used

ConnectionPool Class

class ConnectionPool:
    """
    Thread-safe connection pool with health monitoring.

    Features:
    - Connection reuse for performance
    - Automatic health checks in background thread
    - Dynamic connection creation up to max limit
    - Idle connection cleanup
    - Context manager pattern for safe connection handling
    """

    def __init__(self, url: str, config: ConnectionPoolConfig):
        self.url = url
        self.config = config
        self._connections: List[VespaConnection] = []
        self._available: List[VespaConnection] = []
        self._lock = threading.Lock()
        # Signalled whenever a connection returns to the pool, so waiters
        # wake immediately instead of polling on a sleep loop.
        self._returned = threading.Condition(self._lock)
        self._stop_health_check = threading.Event()

        # Initialize minimum connections
        self._initialize_connections()

        # Start background health check thread
        self._health_check_thread = threading.Thread(
            target=self._health_check_loop, daemon=True
        )
        self._health_check_thread.start()

    @contextmanager
    def get_connection(self):
        """
        Get a connection from the pool (context manager).

        Usage:
            with pool.get_connection() as conn:
                result = conn.query(yql="...")

        Yields:
            VespaConnection: A healthy connection

        Raises:
            TimeoutError: If no connection available within timeout
        """
        conn = None
        deadline = time.monotonic() + self.config.connection_timeout

        try:
            with self._returned:
                while conn is None:
                    if self._available:
                        conn = self._available.pop()
                    elif len(self._connections) < self.config.max_connections:
                        conn = VespaConnection(self.url, f"conn-{uuid.uuid4().hex[:8]}")
                        self._connections.append(conn)
                    else:
                        # Block until a connection returns, instead of
                        # sleep-polling
                        remaining = deadline - time.monotonic()
                        if remaining <= 0 or not self._returned.wait(remaining):
                            if not self._available:
                                raise TimeoutError("No connections available")

            yield conn

        finally:
            # Return connection to pool
            if conn is not None:
                with self._returned:
                    self._available.append(conn)
                    self._returned.notify()

    def close(self):
        """Close all connections (releasing their persistent HTTP
        clients) and stop health checks."""
        self._stop_health_check.set()
        if self._health_check_thread.is_alive():
            self._health_check_thread.join(timeout=5)

        with self._lock:
            for conn in self._connections:
                conn.close()
            self._connections.clear()
            self._available.clear()

Usage in VespaSearchBackend

Simplified illustration (the real profile/strategy resolution is inline in search() itself — see Profile Resolution Logic above, not delegated to separate helper methods):

class VespaSearchBackend:
    def __init__(self, config: Dict[str, Any], **kwargs):
        self.backend_url = config.get("url", "http://localhost")
        self.backend_port = config.get("port", 8080)
        full_url = f"{self.backend_url}:{self.backend_port}"
        pool_config = ConnectionPoolConfig()
        # Single connection pool (URL-based, not schema-based)
        self.pool = ConnectionPool(full_url, pool_config)

    def search(
        self,
        query_dict: Dict[str, Any],
    ) -> List[SearchResult]:
        """Execute search using pooled connection.

        tenant_id is REQUIRED in query_dict and is used for schema scoping
        per-call. Profile and strategy are resolved from query_dict at call time.
        """
        # tenant_id is REQUIRED in query_dict; raises ValueError if missing
        tenant_id = query_dict.get("tenant_id")
        # Profile and strategy resolution happens inline (see above)
        profile_name = ...
        strategy = ...

        # Get connection from pool (context manager pattern)
        with self.pool.get_connection() as conn:
            results = conn.query(
                yql=self._build_query(query_dict, strategy),
                ranking=strategy,
                hits=query_dict.get("top_k", 10)
            )
            return results

Health Metrics

SearchMetrics Integration:

class SearchMetrics:
    """Track search backend metrics for latency and success rates."""

    def record_search(
        self,
        success: bool,
        latency_ms: float,
        strategy: str,
        error: Optional[Exception] = None,
    ):
        """Record a search operation with latency and success/failure."""
        ...

    @property
    def success_rate(self) -> float:
        """Percentage of successful searches (0.0–100.0)."""
        ...

    @property
    def avg_latency_ms(self) -> float:
        """Average search latency in milliseconds."""
        ...

    @property
    def p95_latency_ms(self) -> float:
        """95th-percentile search latency in milliseconds."""
        ...

Connection health is tracked per-connection via VespaConnection.is_healthy (set by the background _health_check_loop in ConnectionPool). The pool itself maintains healthy/unhealthy state; SearchMetrics tracks query-level latency and success rates only.

Benefits

  1. Performance: Connection reuse eliminates connection overhead per query
  2. Reliability: Automatic removal and replacement of connections that fail their periodic health check
  3. Observability: Health metrics for monitoring connection status
  4. Bounded resource use: One pool per backend URL grows from min_connections to max_connections, shared across every tenant schema at that URL
  5. Deterministic shutdown: Closing the unified backend stops the health thread before releasing its HTTP clients; repeated or concurrent closes execute once, and every client is attempted before any close failure is reported

Configuration

ConnectionPoolConfig is passed as the pool_config constructor argument — it is not read from the config dict used for url/port/profiles:

from cogniverse_vespa.search_backend import ConnectionPoolConfig, VespaSearchBackend

pool_config = ConnectionPoolConfig(
    max_connections=10,          # Maximum connections in pool
    min_connections=2,           # Minimum connections to maintain
    connection_timeout=30.0,     # Seconds to wait for a free connection
    idle_timeout=300.0,          # Close idle connections above min_connections after this
    health_check_interval=60.0,  # Background health-check frequency (seconds)
)

backend = VespaSearchBackend(
    config={"url": "http://localhost", "port": 8080, "profiles": profiles},
    pool_config=pool_config,
    config_manager=config_manager,
    schema_loader=schema_loader,
    is_schema_deployed=DeployedSchemaNames(config_manager),
)

Multi-Tenant Schema Management

VespaSchemaManager

Location: libs/vespa/cogniverse_vespa/vespa_schema_manager.py Purpose: Manage tenant-specific Vespa schemas with physical isolation

See Multi-Tenant Architecture for comprehensive details.

Architecture

flowchart TB
    API["<span style='color:#000'>API Request<br/>body.tenant_id: acme</span>"] --> Router["<span style='color:#000'>FastAPI Router<br/>require_tenant_id()</span>"]
    Router --> SchemaManager["<span style='color:#000'>VespaSchemaManager</span>"]

    SchemaManager --> CheckCache{"<span style='color:#000'>Schema in cache?</span>"}
    CheckCache -->|Yes| UseSchema["<span style='color:#000'>Use schema: video_frames_acme_acme</span>"]
    CheckCache -->|No| LoadTemplate["<span style='color:#000'>Load base template</span>"]

    LoadTemplate --> Transform["<span style='color:#000'>Transform for tenant:<br/>video_frames → video_frames_acme_acme<br/>(bare 'acme' canonicalizes to 'acme:acme')</span>"]
    Transform --> Deploy["<span style='color:#000'>Deploy to Vespa</span>"]
    Deploy --> Cache["<span style='color:#000'>Cache deployment</span>"]
    Cache --> UseSchema

    UseSchema --> VespaClient["<span style='color:#000'>VespaSearchBackend<br/>schema=video_frames_acme_acme</span>"]
    VespaClient --> Search["<span style='color:#000'>Search tenant data</span>"]

    style API fill:#90caf9,stroke:#1565c0,color:#000
    style Router fill:#90caf9,stroke:#1565c0,color:#000
    style SchemaManager fill:#ffcc80,stroke:#ef6c00,color:#000
    style CheckCache fill:#ffcc80,stroke:#ef6c00,color:#000
    style UseSchema fill:#ce93d8,stroke:#7b1fa2,color:#000
    style LoadTemplate fill:#b0bec5,stroke:#546e7a,color:#000
    style Transform fill:#b0bec5,stroke:#546e7a,color:#000
    style Deploy fill:#b0bec5,stroke:#546e7a,color:#000
    style Cache fill:#b0bec5,stroke:#546e7a,color:#000
    style VespaClient fill:#a5d6a7,stroke:#388e3c,color:#000
    style Search fill:#a5d6a7,stroke:#388e3c,color:#000

Constructor

from pathlib import Path
from cogniverse_foundation.config.utils import create_default_config_manager
from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from cogniverse_core.registries.backend_registry import get_backend_registry

# Basic initialization (for get_tenant_schema_name and JSON schema uploads only)
schema_manager = VespaSchemaManager(
    backend_endpoint="http://localhost",  # REQUIRED
    backend_port=8080                      # REQUIRED
)

# Full initialization (for tenant schema operations like delete_tenant_schemas, tenant_schema_exists)
# Use BackendRegistry — the returned backend already has a fully-configured schema_manager
config_manager = create_default_config_manager()
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))
registry = get_backend_registry()
backend = registry.get_ingestion_backend(
    "vespa",
    tenant_id="your_org:production",
    config_manager=config_manager,
    schema_loader=schema_loader,
)
schema_manager = backend.schema_manager  # Already has schema_registry, schema_loader injected

Key Methods

# Deploy metadata schemas (organization/tenant) for multi-tenant management.
# Schema-aware: the package carries every schema live when it is built, inside
# the deployment lease, and is rebuilt from a fresh enumeration on every
# conflict retry. allow_schema_removal defaults to False — Vespa refuses a
# deploy that would drop schemas instead of executing it. Dropping the schemas
# of deleted tenants belongs to POST /admin/reconcile-orphans.
schema_manager.upload_metadata_schemas(app_name="cogniverse")

# Deploy content-type schemas together in one application package
schema_manager.upload_content_type_schemas(
    app_name="contenttypes",
    schemas=["image_content", "audio_content", "document_visual", "document_text"],
)
# Schema definitions live in configs/schemas/ as JSON (single source of truth).

# Get tenant-specific schema name. tenant_id is canonicalized to "org:tenant"
# first (a bare "acme" becomes "acme:acme"), then the colon is converted to
# an underscore, so deploy and search paths always agree on the schema name.
schema_name = schema_manager.get_tenant_schema_name(
    tenant_id="acme",
    base_schema_name="video_colpali_smol500_mv_frame"
)
# Returns: "video_colpali_smol500_mv_frame_acme_acme"
# Example: "acme:production" -> "video_colpali_smol500_mv_frame_acme_production"

# Check if tenant schema exists
# REQUIRES: schema_registry in constructor, raises ValueError if not provided
exists = schema_manager.tenant_schema_exists(
    tenant_id="acme",
    base_schema_name="video_colpali_smol500_mv_frame"
)
# Returns: True/False

# Delete tenant schemas (cleanup) — immediately redeploys to Vespa
# REQUIRES: schema_registry in constructor, raises ValueError if not provided
# Internally: redeploys the application package, then unregisters each schema
# with allow_schema_removal=True (Vespa validation override for content type removal)
# Drops the tenant's registered schemas and its suffix-matched Vespa
# orphans in one redeploy. Registered peers are excluded from suffix matching.
# Serialized across processes on the deployment lease
# (SchemaRegistry.deployment_lease), with the deployed snapshot taken inside
# it and retaken before every conflict retry, so no deploy or delete can
# activate a package built from a stale survivor set. The lease covers target
# selection and registry removal for single and bulk deletes. A redeploy that
# removes a schema keeps the lease until the config server's serviceconverge
# reports every service, the content nodes among them, at the generation that
# removed it (REMOVAL_CONVERGENCE_TIMEOUT_S, 120 s, else RuntimeError): a
# content node that skipped straight to a later generation re-adding the
# schema would keep serving the removed documents under it. A pending
# activation of a deleted tenant's schema (a deployment intent) is never
# carried into the package as a survivor.
deleted = schema_manager.delete_tenant_schemas(tenant_id="old_tenant")
# Returns: List of deleted schema names (schemas removed from Vespa via redeployment)

Schema Naming Convention

Pattern: {base_schema}_{canonical_tenant_id with ":" replaced by "_"} — a bare tenant id (no org prefix) is canonicalized to {tenant_id}:{tenant_id} before the suffix is built.

Examples:

Base Schema Tenant ID Tenant Schema
video_colpali_smol500_mv_frame acme video_colpali_smol500_mv_frame_acme_acme
video_xclip_sv_chunk_6s startup video_xclip_sv_chunk_6s_startup_startup
agent_memories acme:production agent_memories_acme_production

Schema Lifecycle

  1. Load Template: Base schema from configs/schemas/{base_schema}_schema.json
  2. Transform: Rename schema and document to include tenant suffix
  3. Deploy: Create Vespa application package and deploy
  4. Cache: Store deployment in memory for fast lookups

Application services.xml

build_services_config(app_package) renders the services.xml deployed with every application package. Both deploy funnels (VespaSchemaManager._deploy_package and VespaBackend._deploy_package) assign its result to app_package.services_config before zipping. A package without a services tree gets pyvespa's default layout (one container cluster, one cogniverse_content content cluster with every schema in index mode); a package that already carries one keeps it. Either way the content cluster's engine/proton/tuning/searchnode/flushstrategy/native/component/maxage is FLUSH_COMPONENT_MAXAGE_S (1800 s), which bounds how long a document-less DocumentDB retains config operations in its transaction log (see Vespa Restart Cost).

Fenced Activation

VespaSchemaManager._post_package creates, prepares and activates one config server session (POST /application/v2/tenant/default/session, then PUT .../session/<id>/prepared and PUT .../session/<id>/active) rather than posting prepareandactivate. The config server activates a session only while the generation it was created from is still the active one and answers 409 ACTIVATION_CONFLICT otherwise, so a deployer stalled between preparing its package and activating it cannot replace an application a successor activated meanwhile; the 409 sends the deploy back through a fresh package build. The deployment lease is renewed once more immediately before the activate, so a deployer whose lease has already been taken over abandons its prepared session instead. VespaBackend._deploy_package posts prepareandactivate under the same lease and never passes the removal override, so a package that went stale there is refused by Vespa's schema-removal validation instead of executed.


Search Backend

VespaSearchBackend

Location: libs/vespa/cogniverse_vespa/search_backend.py Purpose: Production search backend with tenant-scoped schema routing, connection pooling, retries, and metrics

Construction

VespaSearchBackend is dependency-injected with a config dict, a config_manager, and a schema_loader. Profiles come from the same config the search router uses, so a query's profile resolves to the correct tenant schema and its strategy is validated against that schema's rank profiles.

from cogniverse_core.registries.schema_registry import DeployedSchemaNames
from cogniverse_vespa.search_backend import VespaSearchBackend

backend = VespaSearchBackend(
    config={
        "url": "http://localhost",
        "port": 8080,
        "profiles": backend_section["profiles"],
        "default_profiles": backend_section["default_profiles"],
    },
    config_manager=config_manager,
    schema_loader=schema_loader,
    is_schema_deployed=DeployedSchemaNames(config_manager),
)

is_schema_deployed is required: (tenant_id, base_schema_name) -> bool, asked before every query. False raises SchemaNotDeployedError; a read failure propagates (RegistryStorageError chained from the store's error), so an outage never reads as "not deployed". VespaBackend passes a DeployedSchemaNames over its own config_manager — the same schema-registry rows plus pending deployment intents that decide whether a profile is servable — so a search through the registry builds no other backend to answer it. The reader caches deployed names only: a warm search of a deployed schema reads nothing on the request thread, a schema missing from the cache is re-read before the search is refused (a deployment by any process is visible to the next search), and a deletion by another process is seen within DEPLOYED_SCHEMAS_MAX_STALENESS_S (see the core module).

A lower-level create_vespa_search_backend(schema_name, backend_url="http://localhost:8080", *, is_schema_deployed, **kwargs) factory function is also available; it builds a VespaSearchBackend from the non-config (single fixed schema_name/profile) constructor path rather than the multi-profile config dict shown above.

search(query_dict)

search(query_dict: Dict[str, Any]) -> List[SearchResult]. The query_dict accepts these keys:

Key Type Required Notes
query str yes* Text query (*or query_embeddings)
type str yes Content type, e.g. "video"
tenant_id str yes Tenant scope, e.g. "acme:prod" — routes to the tenant schema
profile str no Profile name, e.g. "test_colpali". Omitted: video uses the tenant's selected default video profile or raises; other types auto-select their only profile
strategy str no Rank-profile name (auto-selected if only one available)
top_k int no Result count (defaults to 10)
query_embeddings numpy array no Pre-computed embeddings for visual/hybrid strategies
filters dict no Optional metadata filters
result_granularity str no "source" or "segment"; defaults from the profile

Without query_embeddings, a strategy that reads query tensors has the backend encode the query text, one encoder per input (cogniverse_vespa.query_inputs.query_input_encodings). An input scored against a field that the profile embeds with a service of its own, named under the field's name in inference_services, takes that service's text encoder: acoustic_query, scored against acoustic_embedding, takes the clap_embed CLAP text vector on audio_clap_semantic. Every other input takes the profile's query encoder, the int8 inputs packing its float output. Each encoder encodes the query once. A caller's query_embeddings array binds to every input.

# Text search (tenant-scoped)
results = backend.search({
    "query": "cooking videos",
    "type": "video",
    "profile": "test_colpali",
    "strategy": "hybrid_bm25_binary",
    "top_k": 10,
    "tenant_id": "acme:prod",
})
# Searches ONLY the acme:prod tenant schema
# Physical isolation - no access to other tenants' data

# SearchResult shape (from cogniverse_sdk.document)
for result in results:
    print(f"Score: {result.score}")                          # float
    print(f"Source video: {result.document.metadata['source_id']}")

Every hit's metadata carries source_id, the value of the schema's document_mapping.id field (the document id when the schema declares none), and, when the hit stores one, source_title, the value of its document_mapping.title field. Ingestion fills the title with the source's original upload basename, so source_title stays the same whether source_id is a filename stem or a content hash; golden evaluation matches on it.

SearchResult.to_dict() carries matched_segments and segments_in_window on source-granularity results. Each matched segment row includes document_id, score, and any temporal keys present in the schema; segment granularity omits both fields.

Source granularity sends one Vespa query that groups the matches by the schema's document_mapping.id field, which must be an attribute. It returns the best top_k sources ordered by their best segment's score, equal scores ordered by source id, including at the top_k cut-off. Each source carries its source_collapse_oversample (a profile key, default 4) best segments in matched_segments, highest score first, equal scores by document id; which of several segments tied at that window edge are kept is unspecified. segments_in_window is the source's number of matched segments. Rank profiles without nearestNeighbor group every match. A nearestNeighbor rank profile uses the candidate budget max(top_k, min(top_k * source_collapse_oversample, 256)) as targetHits, so only that many segments reach grouping. When fewer than top_k sources come back and the match count reached the budget, the returned SearchResultBatch has source_search_incomplete=True: sources beyond the budget may exist. False does not make the search exhaustive: approximate nearestNeighbor retrieval can still miss segments. targetHits applies per content node, so on a multi-node content cluster more segments than the budget can reach grouping and the flag can be True although no node's budget was full. total_count is the number of matched segments. source_collapse_oversample must be between 1 and 64, and a source search whose top_k * (1 + source_collapse_oversample) exceeds 10000 grouping rows raises ValueError without querying Vespa. A response with errors, degraded coverage or missing grouping raises VespaError.

export_embeddings() defaults to the backend's configured schema; an explicit schema argument overrides it for that call. It walks Vespa's Document v1 continuation pages. Every page must return HTTP 200; an initial or continuation failure raises with the visit route instead of returning an empty or partial export.

Search retries use the RetryConfig supplied to the constructor, including after initialize() is called. Reconstructed results retain the query content type; memory, wiki, and code profiles produce ContentType.DOCUMENT, not video documents.

Point reads

The public search-backend point-read methods require both physical routing values as keyword-only arguments:

document = backend.get_document(
    "memory-42",
    schema_name="agent_memories_acme_production",
    namespace="memory_content",
)
documents = backend.batch_get_documents(
    ["memory-42", "memory-43"],
    schema_name="agent_memories_acme_production",
    namespace="memory_content",
)

get_document returns None only for a genuine HTTP 404. batch_get_documents preserves input order and uses None for each genuine 404. Transport failures and every other non-success status raise with the physical schema and namespace instead of being reported as missing documents.

Multi-Tenant Search Example

One backend instance serves every tenant; the tenant_id in each query_dict selects the tenant schema, so there is no per-tenant client.

# Tenant A: acme:prod
results_acme = backend.search({
    "query": "cooking videos",
    "type": "video",
    "profile": "test_colpali",
    "tenant_id": "acme:prod",
})

# Tenant B: startup:prod
results_startup = backend.search({
    "query": "cooking videos",
    "type": "video",
    "profile": "test_colpali",
    "tenant_id": "startup:prod",
})

# Complete physical isolation via tenant-specific schemas

Ranking Strategies

strategy is a plain string equal to a rank-profile name defined in the profile's schema. The backend validates each query's strategy against the schema's available rank profiles (there is no enum). Valid names:

Strategy Type Notes
bm25_only, bm25_no_description Text Pure text search, no embeddings
float_float, binary_binary, float_binary, phased Visual Require query_embeddings
hybrid_float_bm25, hybrid_binary_bm25, hybrid_bm25_binary, hybrid_bm25_float Hybrid Visual + text; embeddings required
hybrid_*_no_description variants Hybrid Same as above, ignoring the description field

In the video_colpali_smol500_mv_frame, video_colqwen_omni_mv_chunk_30s and video_xclip_sv_chunk_6s schemas, default inherits phased. On the two ColQwen3 video schemas its first phase, max_sim_hamming, scores every segment by binary MaxSim: for each query token, 1 - 2h/320, where h is the Hamming distance to the token's nearest patch, summed over the query tokens. The second phase reranks the top 100 segments by float MaxSim (max_sim). The ColQwen3 image_colpali_mv and document_visual schemas rank default by max_sim_hamming alone.

The visual-first hybrids (hybrid_float_bm25, hybrid_binary_bm25, their _no_description variants, and hybrid_semantic_bm25 on audio_content) rank every document in one phase by a visual score plus a text score:

  • visual: MaxSim averaged over the query tokens, so query length does not scale it. Each query token contributes its dot product with its best stored vector (float), or 1 - 2h/bits with h the Hamming distance to its nearest stored vector (binary);
  • text: nativeRank over the profile's text fields, between 0 and 1.

For them VespaSearchBackend matches every document with rank(true, {grammar: "any"}userInput(@userQuery)): the query terms only rank, a document without a text match keeps its visual score, and each document's text match counts in full. On the single-vector schemas the visual score is the dense similarity: hybrid_float_bm25 on video_xclip_sv_chunk_6s adds the angular closeness to nativeRank, hybrid_binary_bm25 the same closeness estimated from the Hamming distance h of the 768-bit codes, 1/(1 + πh/768), and hybrid_acoustic_bm25 on audio_content the acoustic closeness; the hybrid profile of agent_memories adds closeness to nativeRank(text). Every hybrid that retrieves through nearestNeighbor matches ({grammar: "any"}userInput(@userQuery)) OR nearestNeighbor(...): the nearest neighbours and every document holding a query term are candidates, each with its full text features.

The text-first hybrid_bm25_* profiles rank the matches weakAnd keeps for the query text, and no others, by the same visual score plus nativeRank in one phase. Their rank profile declares "candidates": "text_matches", which the RankingStrategyExtractor reports as text_candidates_only: the backend then matches userInput(@userQuery) although the first phase scores an embedding, and never adds a nearestNeighbor term. On video_xclip_sv_chunk_6s they compute the visual score from the stored vector, as their query carries no nearestNeighbor term.

bm25_only, bm25_no_description and the text-first hybrids match userInput(@userQuery), Vespa's weakAnd: it walks the matching documents in the order they were first indexed and skips those that cannot beat its running threshold, so they rank the matches weakAnd keeps, not every document holding a query term. That set depends on the order the documents were fed and on the request's hit count: a source-grouped search sends hits=0 and keeps weakAnd's default target, while a request for more hits than there are matches keeps them all.

Where a strategy's phase order lives. The ranking phases (first_phase / second_phase) that define a strategy's actual behavior are authoritative in the schema's rank_profiles (the schema JSON). By naming convention hybrid_binary_bm25* ranks every document by binary visual similarity plus text and hybrid_bm25_binary* ranks the text matches alone by the same sum. configs/schemas/ranking_strategies.json is a generated artifact holding phase-agnostic metadata (which embeddings/tensors each strategy needs) — StrategyAwareProcessor writes it at ingestion via extract_all_ranking_strategies, and query time re-extracts from the schema, not from that file. To change ranking behavior (e.g. phase order), edit the schema rank_profiles; do not hand-edit ranking_strategies.json (it is regenerated). See docs/architecture/schema_driven_flow.md.

# Pure text search (fast, no embeddings)
results = backend.search({
    "query": "machine learning tutorial",
    "type": "video",
    "profile": "test_colpali",
    "strategy": "bm25_only",
    "tenant_id": "acme:prod",
})

# Visual + text hybrid (requires query_embeddings)
results = backend.search({
    "query": "robot arm demonstration",
    "type": "video",
    "profile": "test_colpali",
    "strategy": "hybrid_float_bm25",
    "tenant_id": "acme:prod",
    "query_embeddings": query_embeddings,  # numpy array
})

Real-Vespa integration coverage for every ranking strategy lives in tests/runtime/integration/test_ranking_strategies_real.py, which drives VespaSearchBackend.search against real Vespa and real ColPali.


Ingestion Client

VespaPyClient

Location: libs/vespa/cogniverse_vespa/ingestion_client.py Purpose: PyVespa wrapper for document ingestion with automatic format conversion

Architecture

flowchart TB
    Documents["<span style='color:#000'>Documents<br/>cogniverse_sdk.Document</span>"] --> Client["<span style='color:#000'>VespaPyClient</span>"]
    Client --> Process["<span style='color:#000'>process(doc)<br/>Convert to Vespa format</span>"]

    Process --> Embeddings["<span style='color:#000'>VespaEmbeddingProcessor<br/>Float + Binary + Hex</span>"]
    Process --> Fields["<span style='color:#000'>Map to schema fields</span>"]

    Embeddings --> VespaDoc["<span style='color:#000'>Vespa Document</span>"]
    Fields --> VespaDoc

    VespaDoc --> Feed["<span style='color:#000'>app.feed_iterable()<br/>PyVespa batch feed</span>"]
    Feed --> Retry["<span style='color:#000'>Automatic Retry<br/>pyvespa handles retries</span>"]
    Retry --> Success["<span style='color:#000'>Track Success/Failure</span>"]

    style Documents fill:#90caf9,stroke:#1565c0,color:#000
    style Client fill:#ffcc80,stroke:#ef6c00,color:#000
    style Process fill:#ffcc80,stroke:#ef6c00,color:#000
    style Embeddings fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Fields fill:#ce93d8,stroke:#7b1fa2,color:#000
    style VespaDoc fill:#b0bec5,stroke:#546e7a,color:#000
    style Feed fill:#b0bec5,stroke:#546e7a,color:#000
    style Retry fill:#b0bec5,stroke:#546e7a,color:#000
    style Success fill:#a5d6a7,stroke:#388e3c,color:#000

Tenant-Aware Ingestion

from cogniverse_vespa.ingestion_client import VespaPyClient
from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager
from cogniverse_sdk.document import Document
from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from pathlib import Path
import numpy as np

# 1. Create schema loader (required for VespaPyClient)
schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))

# 2. VespaSchemaManager for schema parsing (deployment uses pyvespa)
schema_manager = VespaSchemaManager(
    backend_endpoint="http://localhost",
    backend_port=8080
)

# 3. Get tenant-specific schema name using the canonical naming convention
# (never string-format the suffix by hand — get_tenant_schema_name()
# canonicalizes tenant_id first, so a bare "acme" doubles to "acme_acme")
base_schema_name = "video_colpali_smol500_mv_frame"
tenant_schema = schema_manager.get_tenant_schema_name(
    tenant_id="acme", base_schema_name=base_schema_name
)
# -> "video_colpali_smol500_mv_frame_acme_acme"

# 4. Create sample documents
documents = [
    Document(
        id="video123_segment_0",
        content="Cooking demonstration",
        metadata={"start_time": 0.0, "end_time": 1.0},
        embeddings={"embedding": np.random.randn(1024, 320)}
    )
]

# 5. Initialize client with configuration
config = {
    "schema_name": tenant_schema,  # video_colpali_smol500_mv_frame_acme_acme
    "base_schema_name": "video_colpali_smol500_mv_frame",
    "url": "http://localhost",
    "port": 8080,
    "schema_loader": schema_loader,  # Required: SchemaLoader instance
    "feed_max_queue_size": 500,
    "feed_max_workers": 4,
    "feed_max_connections": 8
}

client = VespaPyClient(config=config)

# 6. Connect to Vespa
client.connect()

# 7. Process documents and feed
processed_docs = [client.process(doc) for doc in documents]
success_count, failed_ids = client._feed_prepared_batch(processed_docs, batch_size=100)
print(f"Ingested {success_count}/{len(documents)} documents to {tenant_schema}")

Point-read and delete contracts

check_document_exists(document_id) returns True for HTTP 200 and False only for HTTP 404. get_document_data(document_id) returns the stored fields for HTTP 200 and None only for HTTP 404. delete_document(document_id) returns True for HTTP 200 and treats HTTP 404 as idempotent success. Connection failures and every other non-success status raise; returned non-success responses include the complete Document v1 route and are never converted into absent-document results.

Document Processing

from cogniverse_sdk.document import Document
import numpy as np

# Create Document (universal format)
doc = Document(
    id="video123_segment_0",
    content="Chopping vegetables",
    metadata={
        "start_time": 2.5,
        "end_time": 3.0,
        "segment_index": 0,
        "total_segments": 10,
        "audio_transcript": "First, we chop the vegetables",
        "description": "Cooking tutorial scene"
    },
    embeddings={
        "embedding": np.random.randn(1024, 320)  # ColPali embeddings
    }
)

# Process converts to Vespa format automatically:
# 1. Extracts embeddings and converts to hex/binary (VespaEmbeddingProcessor)
# 2. Passes metadata keys through to matching schema fields, and applies the
#    schema's document_mapping.metadata_fields renames for values carried under
#    a different name (segment_index -> segment_id, description ->
#    segment_description); a rename whose target the schema lacks is dropped
# 3. Adds a millisecond creation timestamp (feed only; omitted on update)
# 4. Creates proper Vespa document structure
vespa_doc = client.process(doc)

# vespa_doc structure:
# {
#     "put": "id:video:video_colpali_smol500_mv_frame_acme_acme::video123_segment_0",
#     "fields": {
#         "creation_timestamp": 1729350000000,
#         "embedding": "0x4142...",  # Hex-encoded float embeddings
#         "embedding_binary": [1, 0, 1, ...],  # Binary embeddings
#         "start_time": 2.5,
#         "end_time": 3.0,
#         "segment_id": 0,
#         "total_segments": 10,
#         "audio_transcript": "First, we chop the vegetables",
#         "segment_description": "Cooking tutorial scene"
#     }
# }

Batch Feed Configuration

# Production-ready configuration (via config dict or environment variables)
config = {
    "schema_name": tenant_schema,
    "base_schema_name": "video_colpali_smol500_mv_frame",
    "url": "http://localhost",
    "port": 8080,
    "schema_loader": schema_loader,  # Required: SchemaLoader instance

    # Feed configuration (set via config dict — no env var fallbacks)
    "feed_max_queue_size": 500,
    "feed_max_workers": 4,
    "feed_max_connections": 8,
    "feed_compress": "auto"
}

client = VespaPyClient(config=config)

# Feed uses pyvespa's feed_iterable with these settings automatically

Schema Deployment

JSON Schema Parser

Location: libs/vespa/cogniverse_vespa/json_schema_parser.py Purpose: Parse JSON schema definitions to PyVespa objects

Schema Template Structure

Base schemas are stored in configs/schemas/:

{
  "name": "video_colpali_smol500_mv_frame",
  "document": {
    "name": "video_colpali_smol500_mv_frame",
    "fields": [
      {
        "name": "video_id",
        "type": "string",
        "indexing": ["summary", "attribute"],
        "attribute": ["fast-search"]
      },
      {
        "name": "embedding",
        "type": "tensor<float>(patch{}, v[320])",
        "indexing": ["attribute"]
      }
    ]
  },
  "rank_profiles": [
    {
      "name": "colpali",
      "inputs": [
        {"name": "query(qt)", "type": "tensor<float>(querytoken{}, v[320])"}
      ],
      "first_phase": {
        "expression": "sum(reduce(sum(query(qt) * attribute(embedding), v), max, patch), querytoken)"
      }
    }
  ]
}

Parsing and Deployment

from cogniverse_vespa.json_schema_parser import JsonSchemaParser
from cogniverse_core.registries.schema_registry import SchemaRegistry

# Parse JSON schema (pyvespa Schema object) directly when needed
parser = JsonSchemaParser()
schema = parser.load_schema_from_json_file(
    "configs/schemas/video_colpali_smol500_mv_frame_schema.json"
)

# Deploy a tenant-scoped schema — primary entry point.
# deploy_schema() loads the base JSON definition, transforms it to a
# tenant-specific schema, and deploys it via the backend.
registry = SchemaRegistry(
    config_manager=config_manager, backend=backend, schema_loader=schema_loader
)
tenant_schema_name = registry.deploy_schema(
    tenant_id="acme:production",
    base_schema_name="video_colpali_smol500_mv_frame",
)
# Returns the deployed tenant schema name, e.g.
# "video_colpali_smol500_mv_frame_acme_production"

For new-schema deployments, SchemaRegistry persists the exact registration payload before activation. Vespa package construction probes live document types before calling SchemaRegistry.reconcile_deployment_intents(live_names). After a 90-second grace period, recovery conditionally registers schemas confirmed live using that stored payload, with at most three recovery attempts per intent. Recovered definitions are included in the package, as is every schema still reserved by a pending intent (SchemaRegistry.reserved_schemas), rebuilt from the intent's definition. Recovery never removes schemas, deletes documents, or enables schema-removal overrides. Live schemas without a valid reconstruction still trigger the deployment refusal guard. upload_metadata_schemas (the runtime's startup deploy) builds its package the same way when a registry is injected: registry rows plus reserved intents, refusing on any live schema it cannot rebuild.


Metadata Schema Management

Cogniverse uses JSON-based metadata schemas for multi-tenant management data stored in Vespa. These schemas are the single source of truth and are loaded dynamically at runtime.

Overview

Metadata schemas store operational data (not video content):

  • Organization/tenant hierarchy for multi-tenancy

  • Configuration key-value pairs for VespaConfigStore

  • Adapter registry for model management

configs/schemas/
├── organization_metadata_schema.json   # Organization-level data
├── tenant_metadata_schema.json         # Tenant-level data
├── config_metadata_schema.json         # Configuration storage (VespaConfigStore)
├── adapter_registry_schema.json        # Trained adapter metadata
├── agent_memories_schema.json          # Agent memory storage
└── video_*_schema.json                 # Video content schemas (profiles)

Metadata Schema Types

Schema Purpose Key Fields
organization_metadata Multi-tenant org hierarchy org_id, org_name, status, tenant_count
tenant_metadata Tenant information tenant_full_id, org_id, status, schemas_deployed
config_metadata VespaConfigStore backend config_id, tenant_id, scope, config_key, config_value
adapter_registry Trained LoRA adapters adapter_id, tenant_id, base_model, status, is_active

Loading Schemas from JSON

All metadata schemas are loaded via metadata_schemas.py:

from cogniverse_vespa.metadata_schemas import (
    create_organization_metadata_schema,
    create_tenant_metadata_schema,
    create_config_metadata_schema,
    create_adapter_registry_schema,
    add_metadata_schemas_to_package,
)

# Load individual schema
org_schema = create_organization_metadata_schema()

# Or add all metadata schemas to an ApplicationPackage
from vespa.package import ApplicationPackage
app_package = ApplicationPackage(name="cogniverse")
add_metadata_schemas_to_package(app_package)
# Adds: organization_metadata, tenant_metadata, config_metadata, adapter_registry

Schema File Location

Schemas are auto-discovered from configs/schemas/:

from cogniverse_vespa.metadata_schemas import get_schemas_dir, set_schemas_dir

# Get current schemas directory
schemas_path = get_schemas_dir()
print(schemas_path)  # /path/to/cogniverse/configs/schemas

# Override for testing
set_schemas_dir(Path("/tmp/test_schemas"))

JSON Schema Format

Metadata schemas follow the same JSON format as video schemas:

{
  "name": "config_metadata",
  "document": {
    "fields": [
      {
        "name": "config_id",
        "type": "string",
        "indexing": ["attribute", "summary"],
        "attribute": ["fast-search"]
      },
      {
        "name": "tenant_id",
        "type": "string",
        "indexing": ["attribute", "summary"],
        "attribute": ["fast-search"]
      },
      {
        "name": "config_value",
        "type": "string",
        "indexing": ["summary"]
      }
    ]
  }
}

Field Attributes:

  • indexing: ["attribute", "summary"] - Stored and searchable
  • attribute: ["fast-search"] - Optimized for exact matching
  • indexing: ["summary"] - Stored but not indexed (for large values)

Adding New Metadata Schemas

  1. Create JSON schema file in configs/schemas/my_metadata_schema.json:
{
  "name": "my_metadata",
  "document": {
    "fields": [
      {"name": "id", "type": "string", "indexing": ["attribute", "summary"], "attribute": ["fast-search"]},
      {"name": "tenant_id", "type": "string", "indexing": ["attribute", "summary"], "attribute": ["fast-search"]},
      {"name": "data", "type": "string", "indexing": ["summary"]}
    ]
  }
}
  1. Add loader function in metadata_schemas.py:
def create_my_metadata_schema() -> Schema:
    """Create my_metadata schema. Loads from configs/schemas/my_metadata_schema.json."""
    return _load_schema("my_metadata")
  1. Include in package deployment (if needed globally):
def add_metadata_schemas_to_package(app_package) -> None:
    # ... existing schemas ...
    app_package.add_schema(create_my_metadata_schema())

Schema Deployment

Metadata schemas are deployed automatically via VespaSchemaManager:

from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager
from cogniverse_vespa.metadata_schemas import add_metadata_schemas_to_package
from vespa.package import ApplicationPackage

# Create application package
app_package = ApplicationPackage(name="cogniverse")

# Add all metadata schemas
add_metadata_schemas_to_package(app_package)

# Deploy to Vespa using internal _deploy_package method. It takes a builder,
# not a package: the builder runs again inside the deployment lease before
# every attempt, so the package carries the survivor set Vespa holds then.
schema_manager = VespaSchemaManager(
    backend_endpoint="http://localhost",
    backend_port=19071
)
schema_manager._deploy_package(lambda: app_package)

Best Practices

  1. Single Source of Truth: Always define schemas in JSON files, never duplicate in Python code
  2. Tenant Isolation: Include tenant_id field with fast-search for multi-tenant queries
  3. Versioning: Use version field for optimistic locking when needed
  4. Timestamps: Include created_at and updated_at for auditing
  5. JSON Storage: Store complex objects as JSON strings in string fields

Supporting Modules

Smaller modules in libs/vespa/cogniverse_vespa/ that back the classes above.

Vespa app factory (_vespa_factory.py)

Single source of truth for constructing pyvespa Vespa clients — every other module in the package (search backend, config store, adapter store, schema manager) builds its Vespa instance through this module instead of calling Vespa(url=...) directly. It exposes two construction paths that share the same underlying Vespa(url=...) construction:

  • make_vespa_app(*, url, port=None) -> Vespa — a plain pyvespa client. Each data-plane call (query(), feed_data_point(), get_data(), delete_data()) opens its own VespaSync(self, pool_maxsize=1) under the hood — a fresh connection pool and TCP(+TLS) handshake per call. This is fine for callers that are themselves already pooling connections or that call rarely: the search backend's per-connection VespaConnection (search_backend.py, see Connection Pool Management) and the ingestion client VespaPyClient (ingestion_client.py) both still construct their Vespa instance this way — unchanged.
  • make_persistent_vespa_ops(*, url, port=None, connections=4) -> PersistentVespaOps — calls make_vespa_app and wraps the result in a PersistentVespaOps that opens ONE syncio(connections=connections) session at construction time (_open_http_client() called immediately) and routes query(), feed_data_point(), get_data(), and delete_data() through that single session instead of paying a handshake per call. .url proxies the wrapped app's URL (needed for Document v1 visit-URL construction); any other attribute access (get_application_status, deploy helpers, …) falls through to the wrapped Vespa app via __getattr__. close() releases the session.

Callers that issue many operations over the process lifetime construct their vespa_app this way whenever one isn't injected: - VespaConfigStore (config/config_store.py) — frequent config reads/writes for the life of the store. - VespaAdapterStore (registry/adapter_store.py) — set_active alone issues several sequential reads/writes. - VespaBackend's cached metadata client (backend.py, _metadata_vespa_app()) — metadata CRUD runs on every ingest/deploy. The cache key is (url, port); when it changes, the old PersistentVespaOps is close()d before a new one is built. VespaBackend.close() also releases it, together with the search connection pool and every schema-specific ingestion client.

from cogniverse_vespa._vespa_factory import make_persistent_vespa_ops, make_vespa_app

# Plain client: url + port combined into "host:port" before handing to pyvespa
app = make_vespa_app(url="http://localhost", port=8080)

# url already fully-formed (e.g. a connection-pool entry)
app = make_vespa_app(url="http://localhost:8080")

# Persistent client: same construction, plus one long-lived sync session
# reused across query()/feed_data_point()/get_data()/delete_data()
ops = make_persistent_vespa_ops(url="http://localhost", port=8080, connections=4)
ops.query(yql="select * from sources * where true limit 1")
ops.close()  # releases the session

yql_quote (_yql.py)

Escapes a value for safe interpolation into a YQL string literal (field contains "value"). Every module that builds YQL by hand (config_store.py, vespa_schema_manager.py) uses this instead of ad hoc .replace('"', ...) calls, since an unescaped " or \ both breaks the query and is a YQL injection vector.

from cogniverse_vespa._yql import yql_quote

yql = f"select * from config_metadata where tenant_id contains {yql_quote(tenant_id)}"

Port utilities (config_utils.py)

from cogniverse_vespa.config_utils import (
    VESPA_DEFAULT_DATA_PORT,      # 8080
    VESPA_DEFAULT_CONFIG_PORT,    # 19071
    calculate_config_port,
)

calculate_config_port(8080)   # -> 19071 (standard)
calculate_config_port(8100)   # -> 19091 (custom data port + 10991 offset)

VespaConfig (memory_config.py)

Pydantic model for Mem0's Vespa vector-store config block; raises on any field not in {collection_name, embedding_model_dims, host, port}.

from cogniverse_vespa.memory_config import VespaConfig

config = VespaConfig(collection_name="agent_memories", host="localhost", port=8080)

StrategyAwareProcessor (strategy_aware_processor.py)

Determines which embedding formats (float, binary) a schema's ranking strategies actually read, so ingestion only computes and stores the formats that will be queried.

from cogniverse_vespa.strategy_aware_processor import StrategyAwareProcessor

processor = StrategyAwareProcessor(schema_loader)  # schema_loader is REQUIRED

needs = processor.get_required_embeddings("video_colpali_smol500_mv_frame")
# {"needs_float": True, "needs_binary": True}

fields = processor.get_embedding_field_names("video_colpali_smol500_mv_frame")
# {"binary_field": "embedding_binary", "float_field": "embedding"}

VespaAdapterStore (registry/adapter_store.py)

Vespa-backed storage for trained LoRA adapter metadata, implementing cogniverse_sdk.interfaces.adapter_store.AdapterStore. Stores documents in the adapter_registry schema (see Metadata Schema Types).

from cogniverse_vespa.registry.adapter_store import VespaAdapterStore

store = VespaAdapterStore(
    backend_url="http://localhost",
    backend_port=8080,
    schema_name="adapter_registry",
)
store.initialize()

get_adapter, list_adapters, get_active_adapter and get_stats raise VespaQueryDegraded on a degraded answer: root.errors, a coverage.degraded flag, or a hit without adapter_id (tenant_id for get_stats), the shape Vespa gives a row deleted between match and summary fill. get_stats counts up to 1000 adapters and raises RuntimeError when more match.

set_active(adapter_id, tenant_id, agent_type) verifies that the target adapter belongs to that exact tenant and agent type before changing either the current or target adapter.

delete_adapter(adapter_id) returns True for HTTP 200 and False only for a genuine HTTP 404. Transport failures raise, and returned non-success responses raise with adapter context rather than being reported as "not found."


Usage Examples

Example 1: Tenant Onboarding

from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager

# New tenant "acme" starts using the system
schema_manager = VespaSchemaManager(
    backend_endpoint="http://localhost",
    backend_port=8080
)

# Deploy all required schemas for tenant
schemas_to_deploy = [
    "video_colpali_smol500_mv_frame",
    "video_xclip_sv_chunk_6s",
    "agent_memories"
]

for base_schema in schemas_to_deploy:
    tenant_schema = schema_manager.get_tenant_schema_name(
        tenant_id="acme",
        base_schema_name=base_schema
    )
    print(f"Tenant schema name: {tenant_schema}")
    # Deploy schema via Vespa CLI: vespa deploy

# Expected tenant schemas follow naming convention (a bare tenant_id
# canonicalizes to "acme:acme" before the suffix is built, doubling it):
# ['video_colpali_smol500_mv_frame_acme_acme',
#  'video_xclip_sv_chunk_6s_acme_acme',
#  'agent_memories_acme_acme']
from cogniverse_vespa.search_backend import VespaSearchBackend

def search_for_tenant(backend: VespaSearchBackend, tenant_id: str, query: str) -> list:
    """
    Search videos for specific tenant.

    Args:
        backend: Shared VespaSearchBackend instance
        tenant_id: Tenant identifier (e.g. "acme:prod")
        query: Search query

    Returns:
        List[SearchResult] from the tenant-specific schema
    """
    # tenant_id in the query_dict routes to the tenant schema
    return backend.search({
        "query": query,
        "type": "video",
        "profile": "test_colpali",
        "strategy": "hybrid_float_bm25",
        "top_k": 10,
        "tenant_id": tenant_id,
    })

# One backend instance serves every tenant
acme_results = search_for_tenant(backend, "acme:prod", "cooking videos")
startup_results = search_for_tenant(backend, "startup:prod", "cooking videos")

# Completely isolated - different data sets

Example 3: Tenant-Scoped Ingestion

from pathlib import Path

import numpy as np

from cogniverse_core.schemas.filesystem_loader import FilesystemSchemaLoader
from cogniverse_vespa.ingestion_client import VespaPyClient
from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager

def ingest_videos_for_tenant(
    tenant_id: str,
    video_frames: list,
    schema_loader: FilesystemSchemaLoader,
) -> tuple[int, list]:
    """
    Ingest video frames for specific tenant.

    Args:
        tenant_id: Tenant identifier
        video_frames: List of frame documents
        schema_loader: SchemaLoader instance (required by VespaPyClient)

    Returns:
        (success_count, failed_ids)
    """
    # Get tenant schema name
    schema_manager = VespaSchemaManager(
        backend_endpoint="http://localhost",
        backend_port=8080
    )

    tenant_schema = schema_manager.get_tenant_schema_name(
        tenant_id=tenant_id,
        base_schema_name="video_colpali_smol500_mv_frame"
    )

    # Initialize VespaPyClient
    config = {
        "schema_name": tenant_schema,
        "base_schema_name": "video_colpali_smol500_mv_frame",
        "url": "http://localhost",
        "port": 8080,
        "schema_loader": schema_loader,  # Required
    }
    client = VespaPyClient(config=config)
    client.connect()

    # Process and ingest
    processed_docs = [client.process(doc) for doc in video_frames]
    success_count, failed_ids = client._feed_prepared_batch(processed_docs, batch_size=100)

    print(f"Ingested {success_count}/{len(video_frames)} frames")
    print(f"Schema: {tenant_schema}")

    return success_count, failed_ids
# Ingest for tenant "acme"
frames_acme = [
    {
        "id": f"acme_video1_frame_{i}",
        "fields": {
            "video_id": "video1",
            "frame_id": i,
            "embedding": np.random.randn(1024, 320),
            "video_title": "Cooking Tutorial"
        }
    }
    for i in range(100)
]

schema_loader = FilesystemSchemaLoader(Path("configs/schemas"))
success, failed = ingest_videos_for_tenant("acme", frames_acme, schema_loader)
# Ingests to video_colpali_smol500_mv_frame_acme_acme (bare "acme" doubles)

Example 4: Agent Integration

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"))

# Agent is tenant-agnostic at construction; profile set via deps
agent = SearchAgent(
    deps=SearchAgentDeps(profile="video_colpali_smol500_mv_frame"),
    config_manager=config_manager,
    schema_loader=schema_loader,
)

# Agent internally:
# 1. Uses ConfigManager to get backend settings
# 2. Gets tenant-specific schema name per request
# 3. Initializes search with tenant schema
# 4. All searches automatically scoped to tenant

results = agent.search_by_text(query="cooking videos", tenant_id="acme", top_k=10)  # synchronous
# Searches video_colpali_smol500_mv_frame_acme_acme (bare "acme" doubles)

Testing

Unit Tests

Location: tests/backends/unit/ — real files include test_schema_registry.py (TestSchemaRegistryValidation, TestSchemaRegistryDeployment, TestSchemaRegistryTracking, TestSchemaRegistryInitialization), test_backend_config.py, test_backend_registry_tenant.py, test_schema_name_matching.py, test_embedding_binarization.py, test_ranking_strategy_extractor.py, and others listed in Package Structure. The pattern below illustrates the naming-convention contract those tests pin (it is not a copy of any single file):

import pytest
from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager

class TestTenantSchemaNaming:
    @pytest.fixture
    def manager(self):
        return VespaSchemaManager(
            backend_endpoint="http://localhost",
            backend_port=8080
        )

    def test_bare_tenant_id_canonicalizes_and_doubles(self, manager):
        """A bare tenant_id canonicalizes to 'org:tenant' (acme -> acme:acme)
        before the schema suffix is built, so the suffix is doubled."""
        schema = manager.get_tenant_schema_name("acme", "video_frames")
        assert schema == "video_frames_acme_acme"

    def test_org_tenant_format_is_not_doubled(self, manager):
        """An already-canonical 'org:tenant' id is not doubled."""
        schema = manager.get_tenant_schema_name("acme:production", "video_frames")
        assert schema == "video_frames_acme_production"

    def test_tenant_isolation(self, manager):
        """Verify tenants get distinct schema names."""
        schema_a = manager.get_tenant_schema_name("tenant_a", "video_frames")
        schema_b = manager.get_tenant_schema_name("tenant_b", "video_frames")

        assert schema_a != schema_b
        assert schema_a == "video_frames_tenant_a_tenant_a"
        assert schema_b == "video_frames_tenant_b_tenant_b"

Integration Tests

Location: tests/backends/integration/ — real files include test_tenant_schema_lifecycle.py (TestSchemaRegistryDeployment, TestSchemaRegistryDeletion, against a real Vespa Docker instance via the vespa_instance fixture), test_config_store.py, test_dynamic_profile_search_visibility.py, test_partial_update_roundtrip.py, and test_vespa_factory.py. Illustrative pattern (real fixtures wire config_manager/schema_loader through get_backend(tenant_id) factory fixtures, not self. attributes):

VespaPyClient._feed_prepared_batch tracks the terminal callback state of every document by document id. PyVespa completes concurrent feed callbacks out of input order, so a raised feed marks as failed every id that did not receive a successful callback (explicit failure or unresolved), while preserving ids that already succeeded. The returned (success_count, failed_documents) is therefore an id-set result, never a positional tail of the submitted batch.

Each failed document carries a state: "rejected" when the backend inspected the document and refused it (a schema/field validation failure), with the refusal's status_code; "unresolved" when the backend never gave the document a verdict — the feed ended before it answered, or it stopped answering mid-feed and the transport failed. PyVespa collapses both a schema rejection and a mid-feed connection loss into the same synthetic HTTP 599, so the two are told apart by the wrapped exception, not the status code: a "rejected" reason reads HTTP 599: {...backend message...} while an "unresolved" transport reason reads the backend stopped answering during this batch: .... When no document reaches a verdict because the backend was unreachable for the whole batch, the feed raises ConnectionError naming the endpoint rather than returning a (0, [...]) that reads as a completed feed with rejections. test_partial_update_roundtrip.py fault-injects a real mid-feed backend death (via vespa-sentinel-cmd stop) and checks both the reported states and the reported ids against Document v1 state.

import pytest

from cogniverse_core.registries.schema_registry import DeployedSchemaNames
from cogniverse_vespa.search_backend import VespaSearchBackend

@pytest.mark.integration
class TestTenantScopedSearch:
    @pytest.fixture
    def backend(self, config_manager, schema_loader):
        """Search backend wired to a real Vespa connection."""
        backend_section = config_manager.get_backend_config(tenant_id="__system__")
        return VespaSearchBackend(
            config={
                "url": "http://localhost",
                "port": 8080,
                "profiles": {n: p.to_dict() for n, p in backend_section.profiles.items()},
            },
            config_manager=config_manager,
            schema_loader=schema_loader,
            is_schema_deployed=DeployedSchemaNames(config_manager),
        )

    def test_search_with_tenant_schema(self, backend):
        """Search routes to the tenant-scoped schema."""
        results = backend.search({
            "query": "test query",
            "type": "video",
            "profile": "test_colpali",
            "top_k": 5,
            "tenant_id": "acme:prod",
        })

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

    def test_tenant_isolation(self, backend):
        """tenant_id in each query routes to a different physical schema."""
        results_a = backend.search({
            "query": "test",
            "type": "video",
            "profile": "test_colpali",
            "tenant_id": "acme:prod",
        })
        results_b = backend.search({
            "query": "test",
            "type": "video",
            "profile": "test_colpali",
            "tenant_id": "startup:prod",
        })

        # Results are from different schemas (different data)
        # Physical isolation ensures no cross-tenant access

Test Fixtures

# tests/backends/integration/conftest.py

import uuid

import pytest

from cogniverse_vespa.vespa_schema_manager import VespaSchemaManager

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

@pytest.fixture
def schema_manager():
    """VespaSchemaManager instance"""
    return VespaSchemaManager(
        backend_endpoint="http://localhost",
        backend_port=8080
    )

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

    # Cleanup
    schema_manager.delete_tenant_schemas(test_tenant_id)

Best Practices

1. Always Pass tenant_id in the Query

from cogniverse_vespa.search_backend import VespaSearchBackend

# ✅ Good: tenant_id is supplied with every query
results = backend.search({
    "query": "cooking videos",
    "type": "video",
    "profile": "test_colpali",
    "tenant_id": "acme:prod",  # REQUIRED
})

# ❌ Bad: Missing tenant_id (will raise ValueError)
# backend.search({"query": "cooking videos", "type": "video"})  # ValueError!

2. Construct the Backend with Injected Dependencies

from cogniverse_core.registries.schema_registry import DeployedSchemaNames
from cogniverse_vespa.search_backend import VespaSearchBackend

# config_manager and schema_loader are injected once at construction;
# tenant_id is then provided per query.
backend = VespaSearchBackend(
    config={
        "url": "http://localhost",
        "port": 8080,
        "profiles": backend_section["profiles"],
        "default_profiles": backend_section["default_profiles"],
    },
    config_manager=config_manager,
    schema_loader=schema_loader,
    is_schema_deployed=DeployedSchemaNames(config_manager),
)

3. Test Tenant Isolation

# Always verify tenants are isolated
def test_tenant_isolation():
    schema_a = schema_manager.get_tenant_schema_name("tenant_a", "video_frames")
    schema_b = schema_manager.get_tenant_schema_name("tenant_b", "video_frames")

    assert schema_a != schema_b

4. Use Batch Ingestion

# ✅ Good: Batch ingestion
config = {
    "schema_name": tenant_schema,
    "url": "http://localhost",
    "port": 8080,
    "schema_loader": schema_loader,
}
client = VespaPyClient(config=config)
client.connect()
processed = [client.process(doc) for doc in documents]
success, failed = client._feed_prepared_batch(processed, batch_size=100)

# ❌ Bad: One document per feed call defeats pyvespa's connection reuse
for doc in documents:
    client._feed_prepared_batch([client.process(doc)], batch_size=1)  # Slow!

VespaConfigStore API

Location: libs/vespa/cogniverse_vespa/config/config_store.py

Vespa-based configuration storage with multi-tenant support, implementing the ConfigStore interface. Reads retry transient Vespa failures (connection refused, timeouts, 5xx) with backoff and raise ConfigStoreUnavailableError (from cogniverse_sdk.interfaces.config_store) once the budget is spent; a clean absence returns None. The quality-monitor startup wait retries on that error rather than exiting the sidecar.

Document Structure

config_value is stored as a JSON-serialized string (json.dumps(config_value) on write, json.loads(...) on read) since Vespa's config_metadata schema types the field as string, not a nested object:

{
  "fields": {
    "config_id": "tenant_id:scope:service:config_key",
    "tenant_id": "default",
    "scope": "system",
    "service": "system",
    "config_key": "system_config",
    "config_value": "{\"model\": \"gemini-pro\", \"temperature\": 0.7}",
    "version": 1,
    "created_at": "2024-01-01T00:00:00+00:00",
    "updated_at": "2024-01-01T00:00:00+00:00"
  }
}

Every reader — the visit-based get_config, get_config_history, list_configs, list_all_configs, count_version_rows and prune_all_configs, and the single-document get_immutable_config — shares one bounded retry/backoff helper over the Document v1 API, so a backend that accepts a connection and never answers ends in ConfigStoreUnavailableError rather than holding its caller open. The per-attempt budget is 5s for a single document and for one bounded page, 30s and 60s for the wider scans. Connection errors, timeouts, and 5xx responses retry; 404 or a genuinely empty visit returns None or an empty collection; other failures raise. A get_immutable_config answer without fields raises ConfigStoreUnavailableError; a visited document without fields raises ValueError naming its id, and list_all_configs and get_stats skip it with a warning as they skip any malformed row. A completed set_config is therefore immediately visible to those readers without sleeps or search-index convergence retries.

A write takes its version from the key's version counter: one document per config under the config_version_counter namespace, which pruning never deletes. The writer reads the counter, moves it one version forward with a conditional update (version == <read>), and only then writes that version's document, so no two writers are handed the same version however long either stalls. A version document's absence cannot grant a version: pruning deletes old version documents, and Vespa applies a conditional put with create to a missing document without evaluating its condition. A missing counter is created, conditionally on its absence, at the latest stored version, read with the same visit get_config uses. A counter or version read the store does not answer raises ConfigStoreUnavailableError; nothing is written. Visits run in the config_metadata namespace and the prune query matches on config_id, which the counter does not carry, so no reader sees a counter. delete_config removes the counter after the versions, so a recreated key starts at version 1. The write's prune of old versions is the one step that queries; a listing Vespa answers degraded (root.errors, a coverage.degraded flag, or a hit without its version, which is how Vespa answers for a row another writer deleted between match and summary fill) raises ConfigStoreUnavailableError chained to VespaQueryDegraded (cogniverse_vespa._vespa_factory) inside it, and the prune deletes nothing until the next write.

compare_and_set_config(..., expected_version=n) conditionally writes revision n + 1 and returns None on contention. Its version checks and history retention are described in Configuration Scopes.

Key Methods

from cogniverse_vespa.config.config_store import VespaConfigStore
from cogniverse_sdk.interfaces.config_store import ConfigScope

store = VespaConfigStore(
    backend_url="http://localhost",
    backend_port=8080,
    schema_name="config_metadata"
)

# Store configuration (versioned)
entry = store.set_config(
    tenant_id="acme",
    scope=ConfigScope.ROUTING,
    service="routing",
    config_key="model_settings",
    config_value={"model": "gemini-pro", "temperature": 0.7}
)
# Creates new version on each update. Concurrent writers reserve versions
# on the key's version counter, so each receives a distinct version.

# Retrieve latest version
entry = store.get_config(
    tenant_id="acme",
    scope=ConfigScope.ROUTING,
    service="routing",
    config_key="model_settings"
)

# List all configs for tenant
entries = store.list_configs(tenant_id="acme")

# Delete config
store.delete_config(
    tenant_id="acme",
    scope=ConfigScope.ROUTING,
    service="routing",
    config_key="model_settings"
)

VespaEmbeddingProcessor

Location: libs/vespa/cogniverse_vespa/embedding_processor.py

Handles Vespa-specific embedding format conversions (numpy → hex/binary).

Format Conversions

Schema Type Float Format Binary Format
Single-vector (_sv_, _lvt_ tokens) Raw float list Hex-encoded int8
Patch-based Dict of hex-encoded bfloat16 Dict of hex-encoded int8

Key Methods

from cogniverse_vespa.embedding_processor import VespaEmbeddingProcessor
import numpy as np
import logging

# Create processor (logger is optional first parameter)
processor = VespaEmbeddingProcessor(
    logger=logging.getLogger(__name__),
    model_name="TomoroAI/tomoro-colqwen3-embed-4b",
    schema_name="video_colpali_smol500_mv_frame"
)

# Process raw embeddings
raw = np.random.randn(1024, 320)  # ColPali (Tomoro ColQwen3): 1024 patches × 320 dims
result = processor.process_embeddings(raw)
# Returns: {"embedding": {...}, "embedding_binary": {...}}

# Single-vector processing (X-CLIP)
raw = np.random.randn(768)  # Global embedding
result = processor.process_embeddings(raw)
# Returns: {"embedding": [float, float, ...], "embedding_binary": "hex..."}

Binarization

# Binarization: positive values → 1, negative/zero → 0
binarized = np.packbits(np.where(embeddings > 0, 1, 0), axis=1).astype(np.int8)
# Then hex-encoded for storage

BackendVectorStore (Mem0 Backend)

Location: libs/core/cogniverse_core/memory/backend_vector_store.py

Implements Mem0's VectorStoreBase interface for agent memory persistence, routing every operation through the SDK's Backend interface (cogniverse_sdk.interfaces.backend.Backend) instead of talking to Vespa directly. Registered by the module-level _register_backend_provider() (in cogniverse_core.memory.manager, invoked automatically on import) as the "backend" provider, so any Mem0 config with vector_store.provider: "backend" lands here.

Why through the Backend interface

Routing through Backend inherits the SDK's typed Document, the per-tenant schema scoping, the egress policy that the dispatcher's make_http_client(agent_type) enforces, and the telemetry spans backends emit on every call — all for free.

Capabilities

  • Multi-tenant isolation (user_id → tenant scoping on the backend)
  • Per-agent namespacing (agent_id)
  • Semantic search via embeddings
  • Canonical Document.get_embedding("embedding") reads on point lookups
  • Metadata filtering
  • Telemetry spans on every operation (inherited from the Backend impl)

Key methods

from cogniverse_core.memory.backend_vector_store import BackendVectorStore

# In production this is constructed by Mem0's VectorStoreFactory after
# the module-level _register_backend_provider() registers "backend".
# Direct construction is only used in tests.
store = BackendVectorStore(
    # Tenant-scoped schema name, e.g. from
    # schema_manager.get_tenant_schema_name("acme", "agent_memories")
    # -> "agent_memories_acme_acme" for the bare tenant_id "acme"
    collection_name="agent_memories_acme_acme",
    backend_client=vespa_backend,    # an SDK Backend implementation
    embedding_model_dims=768,
)

# Insert (routes to backend.ingest_documents)
store.insert(
    vectors=[[0.1, 0.2, ...]],
    payloads=[{"data": "...", "user_id": "alice", "agent_id": "search"}],
    ids=["mem_001"],
)

# Search (routes to backend.search)
hits = store.search(query="...", vectors=[0.1, 0.2, ...], limit=5)

# Delete (routes to backend.delete_document)
store.delete(vector_id="mem_001")

RankingStrategyExtractor

Location: libs/vespa/cogniverse_vespa/ranking_strategy_extractor.py

Extracts ranking profile configurations from schema JSON files.

Strategy Types

Type Description Use Case
PURE_VISUAL Embedding-only ranking Image/video similarity
PURE_TEXT BM25 text ranking Text search
HYBRID Embedding + BM25 Multi-modal search

RankingStrategyInfo

from cogniverse_vespa.ranking_strategy_extractor import (
    RankingStrategyExtractor,
    RankingStrategyInfo,
    SearchStrategyType
)

extractor = RankingStrategyExtractor()
strategies = extractor.extract_from_schema(
    Path("configs/schemas/video_colpali_smol500_mv_frame_schema.json")
)

# Strategy info fields
strategy = strategies["hybrid_float_bm25"]
print(strategy.name)                    # "hybrid_float_bm25"
print(strategy.strategy_type)           # SearchStrategyType.HYBRID
print(strategy.needs_float_embeddings)  # True
print(strategy.needs_binary_embeddings) # False
print(strategy.needs_text_query)        # True
print(strategy.use_nearestneighbor)     # False (patch-based schema; True for global schemas)
print(strategy.first_phase_embedding_field)  # "embedding"
print(strategy.inputs)                  # {"qt": "tensor<float>(...)"}
print(strategy.input_fields)            # {"qt": "embedding"}
print(strategy.query_tensors_needed)    # ["qt"]

Detection Logic

  • needs_text_query: Profile name contains "bm25" or "text" OR first-phase has "bm25(" OR "userInput"
  • needs_float_embeddings: Input types contain "float"
  • input_fields: the field each input is scored against: the field nearestNeighbor searches for its input, else the one attribute a phase reads beside the input, functions expanded
  • needs_binary_embeddings: Input types contain "int8"
  • use_nearestneighbor: Global schemas + visual strategies
  • first_phase_embedding_field: The tensor field a visual or hybrid strategy's first phase scores, resolved through profile functions; a text-seeking strategy that has one matches every document and ranks by the query terms

Tenant-Scoped Search via VespaSearchBackend

Location: libs/vespa/cogniverse_vespa/search_backend.py

VespaSearchBackend is the tenant-scoped search entry point. Tenant isolation is enforced per query: tenant_id is required in the query_dict passed to search(), and the backend resolves the tenant-specific schema name before issuing the Vespa query.

Key Features

  • Per-query tenant scoping: tenant_id is required in query_dict
  • Automatic schema name resolution: base_schema + canonicalized tenant_id → tenant_schema (the base schema name comes from the resolved profile's schema_name, not from a query_dict key — there is no "schema" key in query_dict)
  • Schema/profile resolved at query time when constructed with config
  • Thread-safe profile management

Usage

from cogniverse_core.registries.schema_registry import DeployedSchemaNames
from cogniverse_vespa.search_backend import VespaSearchBackend

backend = VespaSearchBackend(
    config=backend_config,        # preferred; carries url/port/profiles
    query_encoder=query_encoder,
    config_manager=config_manager,
    schema_loader=schema_loader,
    is_schema_deployed=DeployedSchemaNames(config_manager),
)

# tenant_id is REQUIRED in query_dict; search() raises if it is missing
results = backend.search({
    "query": "robots playing soccer",
    "tenant_id": "acme:prod",                           # REQUIRED
    "profile": "video_colpali_smol500_mv_frame",
    "strategy": "hybrid_float_bm25",
    "top_k": 10,
})
# Searches tenant-scoped schema: video_colpali_smol500_mv_frame_acme_prod

Schema Resolution

# Pattern: {base_schema}_{canonicalized tenant_id} (colon replaced with underscore).
# canonical_tenant_id() maps a bare tenant_id to "org:tenant" form first
# (e.g. "acme" -> "acme:acme"), so a bare tenant_id's suffix is doubled.

# Bare tenant_id
# Input: tenant_id="acme", base schema "video_colpali_smol500_mv_frame"
# Result: "video_colpali_smol500_mv_frame_acme_acme"

# Org:tenant format (already canonical — not doubled)
# Input: tenant_id="acme:production"
# Result: "video_colpali_smol500_mv_frame_acme_production"


Summary: The Vespa package provides tenant-aware backend integration with physical data isolation via schema-per-tenant. All clients are tenant-scoped, and VespaSchemaManager handles schema lifecycle management transparently.