Skip to content

Cogniverse SDK Module Documentation

Package: cogniverse-sdk Import Name: cogniverse_sdk Layer: Foundation Layer Version: 0.1.0 Last Updated: 2026-07-05


Table of Contents

  1. Overview
  2. Architecture
  3. Key Features
  4. Module Structure
  5. API Reference
  6. Usage Examples
  7. Dependencies
  8. Testing
  9. Development

Overview

Purpose and Responsibilities

The cogniverse-sdk package is the pure foundation of the Cogniverse system, providing core interfaces and data models with zero internal Cogniverse dependencies. It defines the contracts that all backend implementations must follow.

Key Responsibilities:

  • Backend Interface: Abstract base classes for search and ingestion backends
  • Document Model: Universal document representation across all content types
  • Configuration Interface: Config storage abstraction for multi-tenancy
  • Schema Loading: Template loading interface for schema management

Design Philosophy:

  • Zero Dependencies: Only depends on standard library and numpy
  • Pure Interfaces: Primarily abstract base classes; Backend has minimal concrete scaffolding (idempotent initialize() and __init__) but all domain logic is abstract
  • Content-Agnostic: Generic models that work for video, audio, text, images
  • Extensible: Easy to add new backend implementations

Position in Architecture

flowchart TB
    subgraph AppLayer["<span style='color:#000'>Application Layer</span>"]
        Runtime["<span style='color:#000'>cogniverse-runtime</span>"]
    end

    subgraph ImplLayer["<span style='color:#000'>Implementation Layer</span>"]
        Agents["<span style='color:#000'>cogniverse-agents</span>"]
        Vespa["<span style='color:#000'>cogniverse-vespa</span>"]
        Synthetic["<span style='color:#000'>cogniverse-synthetic</span>"]
        Finetuning["<span style='color:#000'>cogniverse-finetuning</span>"]
    end

    subgraph CoreLayer["<span style='color:#000'>Core Layer</span>"]
        Core["<span style='color:#000'>cogniverse-core</span>"]
        Evaluation["<span style='color:#000'>cogniverse-evaluation</span>"]
    end

    subgraph FoundationLayer["<span style='color:#000'>Foundation Layer</span>"]
        Foundation["<span style='color:#000'>cogniverse-foundation</span>"]
        SDK["<span style='color:#000'>cogniverse-sdk ◄─ YOU ARE HERE<br/>Pure Interfaces, Zero Dependencies</span>"]
    end

    AppLayer --> ImplLayer
    ImplLayer --> CoreLayer
    CoreLayer --> FoundationLayer

    style AppLayer fill:#90caf9,stroke:#1565c0,color:#000
    style Runtime fill:#90caf9,stroke:#1565c0,color:#000
    style ImplLayer fill:#ffcc80,stroke:#ef6c00,color:#000
    style Agents fill:#ffcc80,stroke:#ef6c00,color:#000
    style Vespa fill:#ffcc80,stroke:#ef6c00,color:#000
    style Synthetic fill:#ffcc80,stroke:#ef6c00,color:#000
    style Finetuning fill:#ffcc80,stroke:#ef6c00,color:#000
    style CoreLayer fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Core fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Evaluation fill:#ce93d8,stroke:#7b1fa2,color:#000
    style FoundationLayer fill:#a5d6a7,stroke:#388e3c,color:#000
    style Foundation fill:#a5d6a7,stroke:#388e3c,color:#000
    style SDK fill:#a5d6a7,stroke:#388e3c,color:#000

SDK is the foundation - every other workspace package that touches backends, config, schemas, or workflow/adapter storage depends on it directly (agents, core, evaluation, finetuning, foundation, runtime, synthetic, vespa), but it depends on nothing except numpy. (cli, messaging, and telemetry-phoenix have no direct dependency on cogniverse-sdk.)


Architecture

Design Patterns

1. Interface Segregation

Separate interfaces for different concerns:

  • Backend: General backend interface (combined search + ingestion)

  • SearchBackend: Search-specific operations

  • IngestionBackend: Ingestion-specific operations

  • ConfigStore: Configuration storage with versioning and scoping

  • SchemaLoader: Schema loading and ranking strategies

  • WorkflowStore: Workflow execution records and agent performance

  • AdapterStore: Adapter registry and activation management

2. Abstract Base Classes

All interfaces use ABC (Abstract Base Class) pattern:

from abc import ABC, abstractmethod

class SearchBackend(ABC):
    @abstractmethod
    def search(self, query_dict: Dict[str, Any]) -> List[SearchResult]:
        """Search implementation must be provided by subclass"""
        pass

3. Generic Document Model

Single Document class for all content types:

  • Uses ContentType enum (VIDEO, AUDIO, IMAGE, TEXT, DOCUMENT)

  • Flexible metadata dictionary

  • Extensible embeddings storage

  • No content-specific fields


Key Features

1. Backend Interface

Defines the contract for all backend implementations (Vespa, Qdrant, etc.):

from threading import Lock
from typing import Any, Dict

from cogniverse_sdk.interfaces.backend import IngestionBackend, SearchBackend

class Backend(IngestionBackend, SearchBackend):
    """Combined search and ingestion backend interface"""

    def __init__(self, name: str):
        """Initialize backend with a name for identification"""
        self.name = name
        self._initialized = False
        self._initialization_lock = Lock()

    def initialize(self, config: Dict[str, Any]) -> None:
        """Initialize once across concurrent callers."""
        if self._initialized:
            return
        with self._initialization_lock:
            if self._initialized:
                return
            self._initialize_backend(config)
            self._initialized = True

    # _initialize_backend(config) -> None: abstract - subclasses implement
    # backend-specific connection setup.

    # Inherited from SearchBackend (abstract - must implement):
    # search(query_dict) -> List[SearchResult]
    # get_document(document_id) -> Optional[Document]
    # batch_get_documents(document_ids) -> List[Optional[Document]]
    # get_statistics() -> Dict[str, Any]
    # health_check() -> bool
    # get_embedding_requirements(schema_name) -> Dict[str, Any]
    # export_embeddings(schema=None, max_documents=None, filters=None,
    #                   include_embeddings=True) -> List[Dict[str, Any]]

    # Inherited from IngestionBackend (abstract - must implement):
    # ingest_documents(documents, schema_name, operation_type="feed") -> Dict[str, Any]
    # ingest_stream(documents, schema_name) -> Iterator[Dict[str, Any]]
    # update_document(document_id, document, schema_name=None) -> bool
    # delete_document(document_id) -> bool
    # get_schema_info() -> Dict[str, Any]
    # validate_schema(schema_name) -> bool

    # Schema management (abstract - must implement):
    # deployment_lease() -> ContextManager
    # deploy_schemas(schema_definitions) -> bool
    # delete_schema(schema_name, tenant_id) -> List[str]
    # schema_exists(schema_name, tenant_id) -> bool
    # get_tenant_schema_name(tenant_id, base_schema_name) -> str

    # Metadata document operations (abstract - must implement):
    # create_metadata_document(schema, doc_id, fields) -> bool
    # get_metadata_document(schema, doc_id) -> Optional[Dict]
    # query_metadata_documents(schema, query, yql, **kwargs) -> List[Dict]
    # delete_metadata_document(schema, doc_id) -> bool

Benefits:

  • Swappable Backends: Easy to switch from Vespa to Qdrant or other backends
  • Consistent API: All backends expose the same interface
  • Type Safety: Clear type hints for all methods
  • Testable: Easy to create mock backends for testing

2. Universal Document Model

Single document representation for all content:

from cogniverse_sdk.document import Document, ContentType

@dataclass
class Document:
    # Core identification
    id: str = field(default_factory=lambda: str(uuid.uuid4()))
    content_type: ContentType = ContentType.DOCUMENT

    # Content information
    content_path: Optional[Path] = None
    content_id: Optional[str] = None
    title: Optional[str] = None

    # Generic content data
    text_content: Optional[str] = None
    description: Optional[str] = None

    # Embeddings - flexible storage for any embedding type
    embeddings: Dict[str, Any] = field(default_factory=dict)

    # Processing metadata
    status: ProcessingStatus = ProcessingStatus.PENDING
    processing_time: Optional[float] = None
    error_message: Optional[str] = None

    # Flexible metadata for any additional fields
    metadata: Dict[str, Any] = field(default_factory=dict)

    # System metadata (Unix timestamps as int)
    created_at: int = field(default_factory=lambda: int(time.time()))
    updated_at: int = field(default_factory=lambda: int(time.time()))

Features:

  • Content-Agnostic: Works for video, audio, image, text, dataframes
  • Flexible Metadata: Store any metadata as dict
  • Multiple Embeddings: Store different embedding types (ColPali, X-CLIP, etc.)
  • Processing Status: Track document lifecycle with error messages
  • Timestamps: Automatic creation and update tracking (Unix int timestamps)
  • Auto-Detection: Content type auto-detected from file extension

3. Configuration Interface

Abstract config storage for multi-tenancy with versioning and scoping:

from cogniverse_sdk.interfaces.config_store import ConfigStore, ConfigScope, ConfigEntry

class ConfigStore(ABC):
    """Interface for configuration storage with versioning"""

    @property
    @abstractmethod
    def source(self) -> Hashable:
        """Where this store keeps its entries; equal sources read and write
        the same configuration."""

    @abstractmethod
    def initialize(self) -> None:
        """Initialize the configuration store."""
        pass

    @abstractmethod
    def set_config(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        config_value: Dict[str, Any],
    ) -> ConfigEntry:
        """Store or update a configuration entry. Returns ConfigEntry with version."""
        pass

    @abstractmethod
    def get_config(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        version: Optional[int] = None,
    ) -> Optional[ConfigEntry]:
        """Retrieve configuration, optionally by specific version."""
        pass

    @abstractmethod
    def get_config_history(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        limit: int = 10,
    ) -> List[ConfigEntry]:
        """Get configuration version history."""
        pass

    @abstractmethod
    def list_configs(
        self,
        tenant_id: str,
        scope: Optional[ConfigScope] = None,
        service: Optional[str] = None,
    ) -> List[ConfigEntry]:
        """List all configurations matching criteria."""
        pass

    @abstractmethod
    def delete_config(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
    ) -> bool:
        """Delete all versions of a configuration entry."""
        pass

    @abstractmethod
    def export_configs(self, tenant_id: str, include_history: bool = False) -> Dict[str, Any]:
        """Export all configurations for a tenant."""
        pass

    @abstractmethod
    def import_configs(self, tenant_id: str, configs: Dict[str, Any]) -> int:
        """Import configurations for a tenant. Returns count imported."""
        pass

    @abstractmethod
    def list_all_configs(
        self,
        scope: Optional[ConfigScope] = None,
        service: Optional[str] = None,
    ) -> List[ConfigEntry]:
        """List all configurations across all tenants."""
        pass

    @abstractmethod
    def get_stats(self) -> Dict[str, Any]:
        """Get storage statistics."""
        pass

    @abstractmethod
    def health_check(self) -> bool:
        """Check if storage backend is healthy."""
        pass

ConfigScope Enum:

class ConfigScope(Enum):
    SYSTEM = "system"
    AGENT = "agent"
    ROUTING = "routing"
    TELEMETRY = "telemetry"
    SCHEMA = "schema"
    BACKEND = "backend"
    DURABLE = "durable"

ConfigEntry Dataclass:

@dataclass
class ConfigEntry:
    tenant_id: str
    scope: ConfigScope
    service: str
    config_key: str
    config_value: Dict[str, Any]
    version: int
    created_at: datetime
    updated_at: datetime

    def get_config_id(self) -> str:
        """Generate unique config ID: tenant_id:scope:service:config_key"""
        return f"{self.tenant_id}:{self.scope.value}:{self.service}:{self.config_key}"

created_at and updated_at must be timezone-aware datetimes. Construction normalizes them to UTC. from_dict() accepts exactly the fields emitted by to_dict() and requires canonical UTC ISO-8601 strings; missing fields, unknown fields, alternate offsets, and obsolete payload shapes raise ValueError. Identifiers are strings, scope is a ConfigScope, config_value is a dictionary, and version is a positive Python integer.

Benefits:

  • Storage-Agnostic: Pluggable backend via ConfigStore interface (e.g., Vespa)
  • Versioned: All updates are versioned with history retrieval
  • Scoped: Configuration organized by scope (system, agent, routing, etc.)
  • Service-Aware: Configuration tied to specific services

4. Schema Loading Interface

Schema loading for backend schema definitions:

from cogniverse_sdk.interfaces.schema_loader import SchemaLoader
from pathlib import Path

class SchemaLoader(ABC):
    """Interface for loading backend schema definitions"""

    @property
    @abstractmethod
    def source(self) -> Hashable:
        """Where this loader reads schemas from; equal sources load the
        same definitions."""

    @abstractmethod
    def load_schema(self, schema_name: str) -> Dict[str, Any]:
        """Load a schema definition by name."""
        pass

    @abstractmethod
    def list_available_schemas(self) -> List[str]:
        """List all available schema names."""
        pass

    @abstractmethod
    def schema_exists(self, schema_name: str) -> bool:
        """Check if a schema exists."""
        pass

    @abstractmethod
    def load_ranking_strategies(self) -> Dict[str, Dict[str, Any]]:
        """Load ranking strategies configuration."""
        pass

Use Case:

# Load schema definition
schema = loader.load_schema("video_frames")

# List available schemas
schemas = loader.list_available_schemas()  # ["video_frames", "audio_chunks", ...]

# Check if schema exists
if loader.schema_exists("video_frames"):
    schema = loader.load_schema("video_frames")

# Load ranking strategies
strategies = loader.load_ranking_strategies()
# Returns: {"strategy_name": {"ranking_profile": "...", "parameters": {...}}}

5. Workflow Store Interface

Storage for workflow execution records and agent performance:

The interface is typed to the workflow domain. The reader (WorkflowIntelligence, at orchestrator startup) and the writer (the batch optimizer) share this contract. Save methods replace the stored set for the tenant; save_template upserts one template. Data methods are async (the telemetry backend and both callers are async). Implementations register against the cogniverse.workflow.stores entry-point group and are resolved via WorkflowStoreRegistry.

from cogniverse_sdk.interfaces.workflow_store import (
    WorkflowStore,
    WorkflowExecution,
    AgentPerformance,
    WorkflowLearningState,
    WorkflowTemplate,
)

class WorkflowStore(ABC):
    """Typed persistence for workflow intelligence."""

    def initialize(self) -> None:
        """Provision backing storage. Default: no-op (lazy creation)."""

    @abstractmethod
    async def save_executions(
        self, tenant_id: str, executions: List[WorkflowExecution]
    ) -> None: ...
    @abstractmethod
    async def load_executions(self, tenant_id: str) -> List[WorkflowExecution]: ...

    @abstractmethod
    async def save_agent_profiles(
        self, tenant_id: str, profiles: List[AgentPerformance]
    ) -> None: ...
    @abstractmethod
    async def load_agent_profiles(self, tenant_id: str) -> List[AgentPerformance]: ...

    @abstractmethod
    async def save_query_patterns(
        self, tenant_id: str, patterns: Dict[str, List[str]]
    ) -> None: ...
    @abstractmethod
    async def load_query_patterns(self, tenant_id: str) -> Dict[str, List[str]]: ...

    # Concrete template method (not abstract): serializes writes per tenant,
    # writes executions last, and restores every prior channel after a failure.
    async def save_learning_corpus(
        self,
        tenant_id: str,
        executions: List[WorkflowExecution],
        profiles: List[AgentPerformance],
        patterns: Dict[str, List[str]],
    ) -> None: ...

    @abstractmethod
    async def replace_learning_state(
        self,
        tenant_id: str,
        executions: List[WorkflowExecution],
        profiles: List[AgentPerformance],
        patterns: Dict[str, List[str]],
        templates: List[WorkflowTemplate],
    ) -> None: ...

    @abstractmethod
    async def load_learning_state(
        self, tenant_id: str
    ) -> WorkflowLearningState: ...

    @abstractmethod
    async def save_template(self, tenant_id: str, template: WorkflowTemplate) -> str: ...
    @abstractmethod
    async def save_generated_templates(
        self, tenant_id: str, templates: List[WorkflowTemplate]
    ) -> List[str]: ...
    @abstractmethod
    async def load_templates(self, tenant_id: str) -> List[WorkflowTemplate]: ...
    @abstractmethod
    async def delete_template(self, tenant_id: str, template_id: str) -> bool: ...

    @abstractmethod
    def health_check(self) -> bool: ...
    @abstractmethod
    def get_stats(self) -> Dict[str, Any]: ...

save_learning_corpus() replaces every channel, so an empty patterns mapping clears stale patterns. Saves for the same tenant cannot interleave, while different tenants retain independent locks. If a forward write fails, every restore step is attempted. A successful restore re-raises the forward error; if any restore also fails, an ExceptionGroup contains the forward error followed by each restore error.

The live optimizer and serving loader use replace_learning_state() and load_learning_state(). The telemetry implementation holds one renewable per-tenant Redis lease across all four Phoenix channels, including replacement compensation and complete-state reads. A load therefore returns one generation; it cannot combine executions or profiles from one replacement with patterns or templates from another.

Data Classes:

from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Dict, List, Optional

@dataclass
class WorkflowExecution:
    workflow_id: str
    query: str
    query_type: str
    execution_time: float
    success: bool
    agent_sequence: List[str]
    task_count: int
    parallel_efficiency: float
    confidence_score: float
    user_satisfaction: Optional[float] = None
    error_details: Optional[str] = None
    timestamp: datetime = field(
        default_factory=lambda: datetime.now(timezone.utc)
    )
    metadata: Dict[str, Any] = field(default_factory=dict)

@dataclass
class AgentPerformance:
    agent_name: str
    total_executions: int = 0
    successful_executions: int = 0
    average_execution_time: float = 0.0
    average_confidence: Optional[float] = None
    error_rate: float = 0.0
    preferred_query_types: List[str] = field(default_factory=list)
    performance_trend: str = "stable"
    last_updated: datetime = field(
        default_factory=lambda: datetime.now(timezone.utc)
    )

@dataclass
class WorkflowTemplate:
    template_id: str
    name: str
    description: str
    query_patterns: List[str]
    task_sequence: List[Dict[str, Any]]
    expected_execution_time: float
    success_rate: float
    usage_count: int = 0
    created_at: datetime = field(
        default_factory=lambda: datetime.now(timezone.utc)
    )
    last_used: Optional[datetime] = None

Workflow record datetimes follow the same canonical form: defaults are UTC, aware offsets are normalized to UTC, naive datetimes are rejected, and stored payloads use timezone-bearing ISO-8601 strings. Each record accepts exactly its declared fields and validates values during both direct construction and deserialization:

  • Execution durations are finite, non-negative Python floats; efficiencies, confidence scores, and optional satisfaction scores are floats between zero and one. Success is a boolean, task counts are non-negative Python integers, agent sequences contain only strings, and metadata is a dictionary.
  • Performance counts are non-negative Python integers, successful executions cannot exceed total executions, timing is a finite non-negative float, and confidence and error rate are floats between zero and one. Preferred query types contain only strings, and trend is improving, degrading, or stable.
  • Template query patterns contain only strings and task sequences contain only dictionaries. Expected timing is a finite non-negative float, success rate is a float between zero and one, and usage count is a non-negative Python integer.

6. Adapter Store Interface

Storage for adapter metadata and activation management. Implementations register against the cogniverse.adapter.stores entry-point group and are resolved via AdapterStoreRegistry (e.g., VespaAdapterStore registered by cogniverse-vespa):

from cogniverse_sdk.interfaces.adapter_store import AdapterStore

class AdapterStore(ABC):
    """Interface for adapter metadata storage"""

    @abstractmethod
    def initialize(self) -> None:
        """Initialize the adapter store."""
        pass

    @abstractmethod
    def save_adapter(self, metadata: Dict[str, Any]) -> str:
        """Save adapter metadata. Returns adapter_id."""
        pass

    @abstractmethod
    def get_adapter(self, adapter_id: str) -> Optional[Dict[str, Any]]:
        """Get adapter metadata by ID."""
        pass

    @abstractmethod
    def list_adapters(
        self,
        tenant_id: Optional[str] = None,
        agent_type: Optional[str] = None,
        model_type: Optional[str] = None,
        status: Optional[str] = None,
        limit: int = 100,
    ) -> List[Dict[str, Any]]:
        """List adapters with optional filters."""
        pass

    @abstractmethod
    def get_active_adapter(
        self,
        tenant_id: str,
        agent_type: str,
        model_type: str = "llm",
    ) -> Optional[Dict[str, Any]]:
        """Get the active adapter for a tenant/agent/model combination."""
        pass

    @abstractmethod
    def set_active(self, adapter_id: str, tenant_id: str, agent_type: str) -> None:
        """Set an adapter as active for a tenant/agent combination."""
        pass

    @abstractmethod
    def deactivate_adapter(self, adapter_id: str) -> None:
        """Deactivate an adapter."""
        pass

    @abstractmethod
    def deprecate_adapter(self, adapter_id: str) -> None:
        """Deprecate an adapter (marks as deprecated and deactivates)."""
        pass

    @abstractmethod
    def delete_adapter(self, adapter_id: str) -> bool:
        """Delete an adapter. Returns True if deleted."""
        pass

    @abstractmethod
    def health_check(self) -> bool:
        """Check if the store is healthy."""
        pass

    @abstractmethod
    def get_stats(self) -> Dict[str, Any]:
        """Get storage statistics."""
        pass

Module Structure

Directory Layout

cogniverse_sdk/
├── __init__.py                    # Package exports
├── document.py                    # Universal Document model
└── interfaces/
    ├── __init__.py                # Interface exports
    ├── backend.py                 # Backend interfaces (Search, Ingestion, Combined)
    ├── config_store.py            # Config storage interface
    ├── schema_loader.py           # Schema loading interface
    ├── workflow_store.py          # Workflow execution and agent performance
    └── adapter_store.py           # Adapter registry interface

File Descriptions

document.py

Purpose: Universal document model for all content types

Key Classes:

  • ContentType: Enum for content types supporting multi-modal content:
  • VIDEO: Video files with frame-based or chunk-based processing
  • AUDIO: Audio files and speech content
  • IMAGE: Images and visual content
  • TEXT: Natural language text and documents
  • DOCUMENT: PDF, DOCX, and structured documents
  • DATAFRAME: Tabular data (CSV, Excel, Pandas DataFrames)
  • ProcessingStatus: Enum for processing status (PENDING, PROCESSING, COMPLETED, FAILED, SKIPPED)
  • Document: Main document class with metadata, embeddings, and status
  • SearchResult: Represents a search result with document and score

Lines of Code: ~267

interfaces/backend.py

Purpose: Backend interface definitions

Key Classes:

  • Backend: Combined search + ingestion interface
  • SearchBackend: Search-only operations
  • IngestionBackend: Ingestion-only operations

Methods (SearchBackend):

  • initialize(config): Initialize search backend
  • search(query_dict): Search documents; returns a list of SearchResult
  • get_document(document_id): Retrieve a document by ID
  • batch_get_documents(document_ids): Retrieve multiple documents by ID
  • get_statistics(): Get search backend statistics
  • health_check(): Check backend health
  • get_embedding_requirements(schema_name): Get embedding requirements for schema
  • export_embeddings(schema=None, max_documents=None, filters=None, include_embeddings=True): Read up to max_documents documents of a deployed schema, each as its id and stored fields (embeddings included unless include_embeddings is False); raises when the backend cannot be read

Methods (IngestionBackend):

  • initialize(config): Initialize ingestion backend
  • ingest_documents(documents, schema_name, operation_type="feed"): Ingest a batch of documents (operation_type="update" for partial field assignment)
  • ingest_stream(documents, schema_name): Stream documents for ingestion
  • update_document(document_id, document, schema_name=None): Update an existing document
  • delete_document(document_id): Delete a document
  • get_schema_info(): Get backend schema information
  • validate_schema(schema_name): Validate schema exists and is configured

Methods (Backend — schema management and metadata ops):

  • _initialize_backend(config): Abstract backend-specific connection/client setup; concurrent calls on one instance invoke it once after a successful initialization, while a raised exception leaves initialization retryable
  • deployment_lease(): Context manager holding the lease that serialises schema deployment across processes; reentrant on the calling thread, so deploy_schemas called inside it runs under it and the caller keeps it through the registration that follows
  • deploy_schemas(schema_definitions): Deploy multiple schemas together
  • delete_schema(schema_name, tenant_id): Delete tenant schema(s); returns List[str] of deleted names
  • schema_exists(schema_name, tenant_id): Check if schema exists
  • get_tenant_schema_name(tenant_id, base_schema_name): Get tenant-specific schema name
  • create_metadata_document(schema, doc_id, fields): Create or update metadata document
  • get_metadata_document(schema, doc_id): Get metadata document by ID
  • query_metadata_documents(schema, query, yql, **kwargs): Query metadata documents
  • delete_metadata_document(schema, doc_id): Delete metadata document

Lines of Code: ~523

interfaces/config_store.py

Purpose: Configuration storage interface with versioning

Key Classes:

  • ConfigStore: Abstract config storage with version history
  • ConfigScope: Enum for config scope (SYSTEM, AGENT, ROUTING, TELEMETRY, SCHEMA, BACKEND)
  • ConfigEntry: Dataclass for configuration entry with version
  • ConfigStoreUnavailableError: a read the store could not answer within its retry budget
  • ConfigWriteConflictError: an update_config whose every compare-and-set lost to a concurrent writer (config_id, attempts); nothing was written

Methods:

  • initialize(): Initialize the configuration store
  • set_config(tenant_id, scope, service, config_key, config_value): Set config with versioning
  • compare_and_set_config(tenant_id, scope, service, config_key, config_value, *, expected_version): Append exactly the next version, or None on contention
  • update_config(tenant_id, scope, service, config_key, update, *, max_attempts=CONFIG_UPDATE_MAX_ATTEMPTS): Read-modify-write through compare_and_set_config (concrete on the ABC)
  • get_config(tenant_id, scope, service, config_key, version): Get config (optionally by version)
  • get_config_history(tenant_id, scope, service, config_key, limit): Get version history
  • list_configs(tenant_id, scope, service): List configs with filters
  • list_all_configs(scope, service): List all configurations across all tenants
  • delete_config(tenant_id, scope, service, config_key): Delete all versions
  • export_configs(tenant_id, include_history): Export tenant configs, without schema-scope rows (the schema registry's deployment records)
  • import_configs(tenant_id, configs): Import configs under tenant_id; a payload carrying a schema-scope row is refused with ValueError before any write; a row that cannot be written raises RuntimeError after every version the import already wrote is removed, so an import lands whole or not at all
  • get_stats(): Get storage statistics
  • health_check(): Check storage health

Lines of Code: ~290

interfaces/schema_loader.py

Purpose: Schema loading and ranking strategies

Key Classes:

  • SchemaLoader: Abstract schema loader
  • SchemaNotFoundException: Raised when schema not found
  • SchemaLoadError: Raised when schema fails to load/parse

Methods:

  • load_schema(schema_name): Load schema definition by name
  • list_available_schemas(): List all available schema names
  • schema_exists(schema_name): Check if schema exists
  • load_ranking_strategies(): Load ranking strategies configuration

Lines of Code: ~103

interfaces/workflow_store.py

Purpose: Workflow execution and agent performance storage

Key Classes:

  • WorkflowStore: Abstract workflow storage (typed to the workflow domain)
  • WorkflowExecution: Dataclass for a historical workflow execution
  • AgentPerformance: Dataclass for an agent performance profile; average_confidence is None until sampled
  • WorkflowTemplate: Dataclass for a reusable workflow template
  • WorkflowLearningState: Dataclass containing one complete four-channel generation

Methods (data methods are async; save methods replace the tenant's set):

  • initialize(): Provision backing storage (default no-op)
  • save_executions(tenant_id, executions) / load_executions(tenant_id)
  • save_agent_profiles(tenant_id, profiles) / load_agent_profiles(tenant_id)
  • save_query_patterns(tenant_id, patterns) / load_query_patterns(tenant_id)
  • save_learning_corpus(tenant_id, executions, profiles, patterns): concrete template method — replaces all three corpora, including an empty patterns mapping; same-tenant calls are serialized, executions are written last, and every prior channel is restored after a failure. Restore failures are collected with the forward error in an ExceptionGroup.
  • replace_learning_state(tenant_id, executions, profiles, patterns, templates): replace all four channels under the implementation's distributed tenant lock, restoring the complete prior generation before propagating a failed write
  • load_learning_state(tenant_id): load one coherent four-channel generation under the same distributed tenant lock used by replacement
  • save_template(tenant_id, template): Create or update a template
  • save_generated_templates(tenant_id, templates): Persist one generated batch atomically and return its template IDs in input order
  • load_templates(tenant_id): Load all templates for tenant
  • delete_template(tenant_id, template_id): Delete a template by id
  • health_check(): Check storage health
  • get_stats(): Get storage statistics

Lines of Code: ~470

interfaces/adapter_store.py

Purpose: Adapter registry and activation management

Key Classes:

  • AdapterStore: Abstract adapter storage

Methods:

  • initialize(): Initialize the adapter store
  • save_adapter(metadata): Save adapter metadata
  • get_adapter(adapter_id): Get adapter by ID
  • list_adapters(tenant_id, agent_type, model_type, status, limit): List adapters with filters
  • get_active_adapter(tenant_id, agent_type, model_type): Get active adapter
  • set_active(adapter_id, tenant_id, agent_type): Set adapter as active
  • deactivate_adapter(adapter_id): Deactivate adapter
  • deprecate_adapter(adapter_id): Deprecate adapter
  • delete_adapter(adapter_id): Delete adapter
  • health_check(): Check storage health
  • get_stats(): Get storage statistics

Lines of Code: ~159


API Reference

Document Class

Constructor

doc = Document(
    id="doc_123",                              # Optional: auto-generated UUID
    content_type=ContentType.VIDEO,
    content_path=Path("video.mp4"),
    metadata={
        "title": "My Video",
        "duration": 120.5,
        "tags": ["tutorial", "python"]
    },
    embeddings={
        "colpali": np.array([...]),
        "xclip": np.array([...])
    }
)

Properties

doc.id                  # str: Unique document ID
doc.content_type        # ContentType: Type of content
doc.content_path        # Optional[Path]: Path to content file
doc.content_id          # Optional[str]: Content identifier
doc.title               # Optional[str]: Document title
doc.text_content        # Optional[str]: Text content
doc.description         # Optional[str]: Document description
doc.metadata            # Dict[str, Any]: Flexible metadata
doc.embeddings          # Dict[str, Any]: Multiple embeddings
doc.status              # ProcessingStatus: Current status
doc.processing_time     # Optional[float]: Processing duration
doc.error_message       # Optional[str]: Error message if failed
doc.created_at          # int: Unix timestamp
doc.updated_at          # int: Unix timestamp

Methods

# Update processing status
doc.set_processing_status(ProcessingStatus.COMPLETED)
doc.set_processing_status(ProcessingStatus.FAILED, error_message="Error details")

# Mark as completed
doc.mark_completed(processing_time=1.5)

# Mark as failed
doc.mark_failed("Error message")

# Add metadata
doc.add_metadata("author", "John Doe")

# Get metadata
author = doc.get_metadata("author", default="Unknown")

# Add embedding with optional metadata
doc.add_embedding("colpali", embedding_array, metadata={"model": "colpali-v1"})

# Get embedding data
embedding = doc.get_embedding("colpali")

# Get embedding metadata
emb_meta = doc.get_embedding_metadata("colpali")

# To dict
doc_dict = doc.to_dict()

# From dict
doc = Document.from_dict(doc_dict)

# Serialize into one schema's field names via its declared mapping.
# Schemas opt in with a "document_mapping" block in their schema JSON
# (see configs/schemas/document_text_schema.json); unmapped generic
# fields are omitted.
from cogniverse_sdk.document import DocumentFieldMapping

mapping = DocumentFieldMapping.from_dict(schema_json["document_mapping"])
fields = doc.to_schema_fields(mapping)   # {"document_id": ..., "full_text": ...}

A document_mapping block maps generic Document fields to one schema's field names:

  • Core fields — id, title, text_content, description, content_type, content_id, content_path map a generic field to a schema field name.
  • created_at / updated_at with created_at_format: "epoch" (int seconds), "epoch_ms" (int milliseconds — for fields like creation_timestamp), or "iso" (UTC string).
  • metadata_fields: {metadata_key: schema_field} renames a value carried in Document.metadata to a schema field name (e.g. {"segment_index": "segment_id"}). When the source and destination names differ, the source key is consumed rather than also being passed through, so feeds contain only the schema's declared destination field.
  • embeddings: {embedding_name: schema_field} maps a stored embedding to its field. Stored embeddings use the exact add_embedding() wrapper: data, metadata, and an integer-second created_at. Backends that hex/binary-encode embeddings (the ingestion path) override these with their processed vectors.
  • include_metadata: when false, only the explicitly mapped/renamed fields are fed (no blanket passthrough of every metadata key) — used by schemas whose values all live in metadata under non-matching names.

Mapping construction, loading, and serialization reject unknown keys, mistyped field names, booleans, and mapping dictionaries. from_schema_json likewise requires a dictionary (or None), so malformed schema data fails with the boundary field named instead of surfacing later during a feed.

Backends apply this automatically: VespaBackend.put_document(document, schema_name=..., base_schema_name=...) loads the base schema's document_mapping, serializes, and feeds — raising ValueError when the schema declares no mapping rather than guessing field names.

Both write paths load the block through one helper — DocumentFieldMapping.from_schema_json(schema_json, schema_name=..., required=...) — so they cannot drift. The ingestion serializer (VespaPyClient.process, used by ingest_documents) applies the metadata_fields renames to the per-segment metadata a video Document carries (e.g. segment_index → segment_id); editing a schema's metadata_fields changes what ingestion feeds. It does not require a mapping — a schema without one (memory, graph) still feeds by passing metadata keys straight through to matching schema fields.

SearchResult Class

Pairs a Document with a relevance score for search-response construction:

from cogniverse_sdk.document import Document, SearchResult

result = SearchResult(
    document=doc,               # Document: the matched document
    score=0.87,                 # float: relevance score
    highlights={"text": "..."},  # Optional[Dict[str, Any]]: highlighted snippets
)

# Convert to a dict for an API response
result_dict = result.to_dict()
# {
#     "document_id": doc.id,
#     "score": 0.87,
#     "metadata": doc.metadata,
#     "highlights": {"text": "..."},
#     # "source_id" included if present in doc.metadata
#     # "source_title" included if present in doc.metadata
#     # "temporal_info" (start_time/end_time/duration) included if both
#     # start_time and end_time are present in doc.metadata
# }

source_title_key(title) returns the key golden evaluation sets name a source by: the title's basename without its file extension (v_-uJnucdW6DY.mp4 → v_-uJnucdW6DY; a non-alphabetic suffix such as clip.2024 is kept). result_source_title_key(row) applies it to a result row's top-level source_title (to_dict() output or a search span result row) and raises ValueError naming the row's document id when the row carries none.

document must be a Document, score must be a finite Python float, and highlights must be a dictionary or None. Invalid values raise at construction instead of reaching API serialization.

Backend Interface

Search Method

results = backend.search(
    {
        "query": "machine learning",  # str: text query
        "tenant_id": "acme_corp",     # str: tenant for schema scoping
        "type": "video",              # content type
        "top_k": 10,                  # int: number of results
        "filters": {"duration_min": 60},  # optional field filters
    }
)

Returns: List[SearchResult] (each carries document, score, highlights)

Ingest Method

result = backend.ingest_documents(
    documents=[doc1, doc2, doc3],
    schema_name="video_frames",
    operation_type="feed",  # or "update" for partial field assignment
)

Returns: Dict[str, Any] with ingestion stats:

{
    "success_count": 3,
    "failed_count": 0,
    "failed_documents": [],
    "total_documents": 3
}

ConfigStore Interface

Set Config

from cogniverse_sdk.interfaces.config_store import ConfigScope

entry = config_store.set_config(
    tenant_id="acme_corp",
    scope=ConfigScope.AGENT,
    service="video_search_agent",
    config_key="embedding_model",
    config_value={"model": "TomoroAI/tomoro-colqwen3-embed-4b", "dimension": 320}
)
# Returns ConfigEntry with version number

Get Config

entry = config_store.get_config(
    tenant_id="acme_corp",
    scope=ConfigScope.AGENT,
    service="video_search_agent",
    config_key="embedding_model",
    version=None  # None = latest, or specific version number
)
# Returns ConfigEntry or None

Get Config History

history = config_store.get_config_history(
    tenant_id="acme_corp",
    scope=ConfigScope.AGENT,
    service="video_search_agent",
    config_key="embedding_model",
    limit=10
)
# Returns List[ConfigEntry] sorted by version (newest first)

List Configs

configs = config_store.list_configs(
    tenant_id="acme_corp",
    scope=ConfigScope.AGENT,  # Optional filter
    service="video_search_agent"  # Optional filter
)
# Returns List[ConfigEntry] (latest versions only)

SchemaLoader Interface

Load Schema

schema = schema_loader.load_schema("video_frames")

Returns: Dict[str, Any] containing the complete schema definition

Raises: SchemaNotFoundException if schema doesn't exist, SchemaLoadError on parse error

List Available Schemas

schemas = schema_loader.list_available_schemas()
# Returns: ["video_frames", "audio_chunks", "text_documents", ...]

Check Schema Exists

if schema_loader.schema_exists("video_frames"):
    schema = schema_loader.load_schema("video_frames")

Load Ranking Strategies

strategies = schema_loader.load_ranking_strategies()
# Returns: {
#     "hybrid_search": {"ranking_profile": "hybrid", "parameters": {...}},
#     "semantic_only": {"ranking_profile": "semantic", "parameters": {...}}
# }

Usage Examples

Example 1: Creating a Custom Backend

from cogniverse_sdk.interfaces.backend import SearchBackend
from cogniverse_sdk.document import Document, ContentType, SearchResult
from typing import List, Dict, Any, Optional

class MyCustomSearchBackend(SearchBackend):
    """Custom search backend implementation (simplified example)"""

    def __init__(self):
        self.documents = {}
        self.endpoint = None

    def initialize(self, config: Dict[str, Any]) -> None:
        """Initialize with config"""
        self.endpoint = config.get("endpoint")
        self.api_key = config.get("api_key")

    def search(self, query_dict: Dict[str, Any]) -> List[SearchResult]:
        """Implement search"""
        query_text = query_dict.get("query")
        top_k = query_dict.get("top_k", 10)
        results = []
        if query_text:
            for doc in list(self.documents.values())[:top_k]:
                results.append(SearchResult(document=doc, score=0.5))
        return results

    def get_document(self, document_id: str) -> Optional[Document]:
        """Get document by ID"""
        return self.documents.get(document_id)

    def batch_get_documents(self, document_ids: List[str]) -> List[Optional[Document]]:
        """Get multiple documents by IDs"""
        return [self.documents.get(doc_id) for doc_id in document_ids]

    def get_statistics(self) -> Dict[str, Any]:
        """Get backend statistics"""
        return {"document_count": len(self.documents)}

    def health_check(self) -> bool:
        """Check backend health"""
        return self.endpoint is not None

    def get_embedding_requirements(self, schema_name: str) -> Dict[str, Any]:
        """Get embedding requirements for schema"""
        return {"dimension": 768, "type": "float"}

    def export_embeddings(
        self, schema=None, max_documents=None, filters=None, include_embeddings=True
    ) -> List[Dict[str, Any]]:
        """Read documents with their stored fields"""
        return [{"id": doc_id, **doc} for doc_id, doc in self.documents.items()][
            :max_documents
        ]

# Usage
backend = MyCustomSearchBackend()
backend.initialize({"endpoint": "http://localhost:9200"})

# Add a document to the in-memory store
doc = Document(content_type=ContentType.VIDEO, title="Tutorial Video")
backend.documents[doc.id] = doc

# Search
results = backend.search({"query": "tutorial", "top_k": 10})

Example 2: Working with Multi-Modal Documents

from cogniverse_sdk.document import Document, ContentType, ProcessingStatus
from pathlib import Path
import numpy as np

# ===== VIDEO CONTENT =====
video_doc = Document(
    id="video_123",
    content_type=ContentType.VIDEO,
    content_path=Path("videos/tutorial.mp4"),
    metadata={
        "title": "Python Tutorial",
        "duration": 600.0,
        "resolution": [1920, 1080],
        "fps": 30.0,
        "tags": ["python", "tutorial", "beginner"]
    }
)

# Add one frame's ColPali/ColQwen3 patch embeddings
colpali_embeddings = np.random.randn(1024, 320)  # 1024 patches, 320 dims
video_doc.add_embedding("colpali_patch", colpali_embeddings)

# Add global embeddings (X-CLIP)
xclip_embedding = np.random.randn(768)  # Global video embedding
video_doc.add_embedding("xclip_global", xclip_embedding)

# ===== AUDIO CONTENT =====
audio_doc = Document(
    id="audio_456",
    content_type=ContentType.AUDIO,
    content_path=Path("audio/podcast.mp3"),
    metadata={
        "title": "Tech Podcast Episode 42",
        "duration": 3600.0,
        "sample_rate": 44100
    }
)

# ===== IMAGE CONTENT =====
image_doc = Document(
    id="image_789",
    content_type=ContentType.IMAGE,
    content_path=Path("images/diagram.png"),
    metadata={
        "title": "System Architecture Diagram",
        "width": 1920,
        "height": 1080,
        "format": "PNG"
    }
)

# Add image embeddings (ColQwen)
colqwen_embedding = np.random.randn(1024, 320)
image_doc.add_embedding("colqwen", colqwen_embedding)

# ===== DOCUMENT CONTENT =====
document_doc = Document(
    id="doc_101",
    content_type=ContentType.DOCUMENT,
    content_path=Path("documents/research_paper.pdf"),
    metadata={
        "title": "Multi-Modal RAG Research",
        "authors": ["Alice", "Bob"],
        "pages": 25,
        "format": "PDF"
    }
)

# ===== TEXT CONTENT =====
text_doc = Document(
    id="text_202",
    content_type=ContentType.TEXT,
    text_content="This is a text document about machine learning.",
    metadata={
        "title": "ML Introduction",
        "word_count": 1500,
        "language": "en"
    }
)

# ===== DATAFRAME CONTENT =====
# Tabular data (CSV, Excel, Pandas DataFrames)
dataframe_doc = Document(
    id="df_303",
    content_type=ContentType.DATAFRAME,
    content_path=Path("data/sales_data.csv"),
    metadata={
        "title": "Q4 Sales Data",
        "rows": 10000,
        "columns": 15,
        "format": "CSV",
        "schema": {
            "date": "datetime",
            "product": "string",
            "revenue": "float",
            "units_sold": "int"
        }
    }
)

# Add dataframe embeddings (text-based representation)
df_text_embedding = np.random.randn(768)
dataframe_doc.add_embedding("text", df_text_embedding)

# ===== COMMON OPERATIONS =====
# Update status
video_doc.set_processing_status(ProcessingStatus.COMPLETED)

# Serialize to dict
doc_dict = video_doc.to_dict()

# Save to JSON. Convert numpy arrays before serialization; Document.to_dict()
# preserves embedding values and does not silently change their types.
import json
for stored_embedding in doc_dict["embeddings"].values():
    data = stored_embedding.get("data")
    if hasattr(data, "tolist"):
        stored_embedding["data"] = data.tolist()
with open("document.json", "w") as f:
    json.dump(doc_dict, f)

# Load from dict
loaded_doc = Document.from_dict(doc_dict)

# Check content type
if loaded_doc.content_type == ContentType.VIDEO:
    print(f"Processing video: {loaded_doc.metadata.get('title')}")
elif loaded_doc.content_type == ContentType.DATAFRAME:
    print(f"Processing dataframe with {loaded_doc.metadata.get('rows')} rows")

Document.from_dict() consumes only the exact to_dict() field set. It rejects missing or unknown fields, null metadata/embedding mappings, and alternate timestamp encodings; created_at and updated_at must be Python integers containing epoch seconds. Construction and serialization enforce the same field types: content_type and status are enum values, mappings are dictionaries, and timestamps are never coerced from strings, floats, NumPy scalars, or milliseconds.

Example 3: Config Store Interface

The ConfigStore interface defines the contract for configuration persistence. The default implementation uses Vespa via VespaConfigStore.

compare_and_set_config(..., expected_version=n) conditionally appends revision n + 1; zero requires an absent key. Concurrent writers cannot both append the same revision. It returns the new ConfigEntry when confirmed current, or None on contention, including a write superseded before confirmation. Callers reread before retrying. Successful writes follow the store's history retention policy. Negative expected versions raise ValueError; storage failures propagate.

update_config(tenant_id, scope, service, config_key, update) is the read-modify-write built on it, concrete on the ABC so every store has it. update(entry) receives the latest entry (None when absent) and returns the value to write, or None to leave the config as it is. When the compare-and-set loses, the entry is re-read and update runs again on it, with full-jitter backoff between attempts (CONFIG_UPDATE_BACKOFF_BASE_S 0.05 s doubling, capped at CONFIG_UPDATE_BACKOFF_CAP_S 1 s). It returns the written entry or the one update declined to change, and raises ConfigWriteConflictError once CONFIG_UPDATE_MAX_ATTEMPTS (10) attempts have all lost; storage failures and whatever update raises propagate with nothing written. Writers on different processes or replicas each keep their change.

Selected interface methods (cogniverse_sdk):

from cogniverse_sdk.interfaces.config_store import ConfigStore, ConfigScope, ConfigEntry
from abc import ABC, abstractmethod
from typing import Optional, List, Dict, Any

class ConfigStore(ABC):
    """Abstract interface for configuration storage."""

    @property
    @abstractmethod
    def source(self) -> Hashable:
        """Where this store keeps its entries."""

    @abstractmethod
    def initialize(self) -> None:
        """Initialize the config store."""
        pass

    @abstractmethod
    def set_config(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        config_value: Dict[str, Any],
    ) -> ConfigEntry:
        """Store configuration with versioning."""
        pass

    @abstractmethod
    def compare_and_set_config(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        config_value: Dict[str, Any],
        *,
        expected_version: int,
    ) -> Optional[ConfigEntry]:
        """Append the next revision, returning None on contention."""
        pass

    @abstractmethod
    def get_config(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        version: Optional[int] = None,
    ) -> Optional[ConfigEntry]:
        """Retrieve configuration, optionally by version."""
        pass

    @abstractmethod
    def get_config_history(
        self,
        tenant_id: str,
        scope: ConfigScope,
        service: str,
        config_key: str,
        limit: int = 10,
    ) -> List[ConfigEntry]:
        """Get version history for a configuration."""
        pass

Example Implementation (VespaConfigStore):

from cogniverse_vespa.config.config_store import VespaConfigStore

# VespaConfigStore implements ConfigStore using Vespa's config_metadata schema.
# keep_versions bounds per-config history through best-effort pruning
# after set_config and compare_and_set_config writes (default 10).
store = VespaConfigStore(
    backend_url="http://localhost",
    backend_port=8080,
    schema_name="config_metadata",
    keep_versions=10,
)
store.initialize()

# For typical usage, use the factory function:
from cogniverse_foundation.config.utils import create_default_config_manager
config_manager = create_default_config_manager()

Example 4: Schema Loader Implementation

The real production implementation is cogniverse_core.schemas.filesystem_loader.FilesystemSchemaLoader (constructor: FilesystemSchemaLoader(base_path: Path); raises ValueError if base_path is missing or not a directory, and SchemaLoadError if ranking_strategies.json is missing). The example below implements the same interface from scratch, with intentionally simplified behavior, to show what a minimal SchemaLoader implementation looks like:

from cogniverse_sdk.interfaces.schema_loader import (
    SchemaLoader,
    SchemaNotFoundException,
    SchemaLoadError
)
from pathlib import Path
import json
from typing import Any, Dict, List

class CustomSchemaLoader(SchemaLoader):
    """Minimal filesystem-based schema loader (example)"""

    def __init__(self, schema_dir: Path):
        self.schema_dir = schema_dir

    def load_schema(self, schema_name: str) -> Dict[str, Any]:
        """Load schema definition by name"""
        schema_path = self.schema_dir / f"{schema_name}_schema.json"
        if not schema_path.exists():
            raise SchemaNotFoundException(f"Schema not found: {schema_name}")
        try:
            with open(schema_path) as f:
                return json.load(f)
        except json.JSONDecodeError as e:
            raise SchemaLoadError(f"Failed to parse schema: {e}")

    def list_available_schemas(self) -> List[str]:
        """List all available schema names"""
        schemas = []
        for path in self.schema_dir.glob("*_schema.json"):
            name = path.stem.replace("_schema", "")
            schemas.append(name)
        return schemas

    def schema_exists(self, schema_name: str) -> bool:
        """Check if schema exists"""
        schema_path = self.schema_dir / f"{schema_name}_schema.json"
        return schema_path.exists()

    def load_ranking_strategies(self) -> Dict[str, Dict[str, Any]]:
        """Load ranking strategies configuration"""
        strategies_path = self.schema_dir / "ranking_strategies.json"
        if not strategies_path.exists():
            return {}
        with open(strategies_path) as f:
            return json.load(f)

# Usage
loader = CustomSchemaLoader(Path("schemas/"))

# List available schemas
available = loader.list_available_schemas()
print(f"Available schemas: {available}")

# Check and load schema
if loader.schema_exists("video_frames"):
    schema = loader.load_schema("video_frames")
    print(f"Loaded schema with {len(schema.get('fields', []))} fields")

# Load ranking strategies
strategies = loader.load_ranking_strategies()
for name, config in strategies.items():
    print(f"Strategy '{name}': profile={config['ranking_profile']}")

Dependencies

External Dependencies

dependencies = [
    "numpy==2.4.4",  # For embedding arrays in backend interface
]

Why numpy? - Embeddings are typically stored as numpy arrays - Standard format across all ML libraries - Efficient storage and computation

Internal Dependencies

None - SDK has zero internal Cogniverse dependencies by design.


Testing

Running Tests

SDK data-record and interface contracts, concrete implementations, and real-service round trips are tested in the project-level tests/ directory:

uv run pytest \
  tests/backends/unit/test_sdk_document_contracts.py \
  tests/backends/unit/test_backend_interface_contract.py \
  tests/backends/unit/test_backend_bool_contracts.py \
  -v --tb=long > /tmp/sdk-tests.log 2>&1

Test Structure

  • tests/backends/unit/test_sdk_document_contracts.py pins document, config, and workflow-record serialization boundaries.
  • tests/backends/unit/test_backend_interface_contract.py verifies concrete backend method signatures against SDK contracts.
  • tests/backends/unit/test_backend_bool_contracts.py verifies boolean return contracts.
  • tests/backends/integration/test_document_mapping_roundtrip.py exercises document mapping against a real Vespa service.
  • Concrete ConfigStore, WorkflowStore, and AdapterStore tests live alongside their implementations under tests/backends/ and tests/agents/.

Development

Adding New Interfaces

To add a new interface:

  1. Create interface file in cogniverse_sdk/interfaces/:

    # cogniverse_sdk/interfaces/my_interface.py
    from abc import ABC, abstractmethod
    
    class MyInterface(ABC):
        """Description of interface"""
    
        @abstractmethod
        def my_method(self, arg: str) -> str:
            """Method description"""
            pass
    

  2. Export in __init__.py:

    # cogniverse_sdk/interfaces/__init__.py
    from .my_interface import MyInterface
    
    __all__ = ["MyInterface", ...]
    

  3. Add tests:

    # tests/test_my_interface.py
    def test_my_interface():
        # Test implementation
        pass
    

Design Guidelines

  1. Keep It Pure: No implementation, only interfaces
  2. Type Hints: Always use type hints for clarity
  3. Documentation: Every method needs docstring
  4. Flexibility: Design for extensibility
  5. Simplicity: Prefer simple interfaces over complex ones

Building and Publishing

# Build package
cd libs/sdk
uv build

# Outputs:
# dist/cogniverse_sdk-0.1.0-py3-none-any.whl
# dist/cogniverse_sdk-0.1.0.tar.gz

# Publish to PyPI (if needed)
uv publish --token $PYPI_TOKEN

Summary

Key Points

  • Foundation Package: Zero internal dependencies, pure interfaces
  • Backend Abstraction: Swap backends (Vespa, Qdrant, etc.) easily
  • Universal Document: Single model for all content types
  • Config Interface: Multi-tenant configuration storage
  • Schema Loading: Template-based schema deployment

Package Stats

  • Total Lines: ~1,577
  • Files: 8 Python files
  • Interfaces: 6 main interfaces (Backend, ConfigStore, SchemaLoader, WorkflowStore, AdapterStore, Document)
  • Dependencies: 1 external (numpy)

Next Steps

  • Implementation: See cogniverse-vespa for Backend implementation
  • Usage: See cogniverse-core for how SDK is used
  • Foundation: See cogniverse-foundation for config base classes