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¶
- Overview
- Architecture
- Key Features
- Module Structure
- API Reference
- Usage Examples
- Dependencies
- Testing
- 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;
Backendhas minimal concrete scaffolding (idempotentinitialize()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
ContentTypeenum (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
ConfigStoreinterface (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, orstable. - 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 processingAUDIO: Audio files and speech contentIMAGE: Images and visual contentTEXT: Natural language text and documentsDOCUMENT: PDF, DOCX, and structured documentsDATAFRAME: Tabular data (CSV, Excel, Pandas DataFrames)ProcessingStatus: Enum for processing status (PENDING, PROCESSING, COMPLETED, FAILED, SKIPPED)Document: Main document class with metadata, embeddings, and statusSearchResult: 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 interfaceSearchBackend: Search-only operationsIngestionBackend: Ingestion-only operations
Methods (SearchBackend):
initialize(config): Initialize search backendsearch(query_dict): Search documents; returns a list ofSearchResultget_document(document_id): Retrieve a document by IDbatch_get_documents(document_ids): Retrieve multiple documents by IDget_statistics(): Get search backend statisticshealth_check(): Check backend healthget_embedding_requirements(schema_name): Get embedding requirements for schemaexport_embeddings(schema=None, max_documents=None, filters=None, include_embeddings=True): Read up tomax_documentsdocuments of a deployed schema, each as itsidand stored fields (embeddings included unlessinclude_embeddingsis False); raises when the backend cannot be read
Methods (IngestionBackend):
initialize(config): Initialize ingestion backendingest_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 ingestionupdate_document(document_id, document, schema_name=None): Update an existing documentdelete_document(document_id): Delete a documentget_schema_info(): Get backend schema informationvalidate_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 retryabledeployment_lease(): Context manager holding the lease that serialises schema deployment across processes; reentrant on the calling thread, sodeploy_schemascalled inside it runs under it and the caller keeps it through the registration that followsdeploy_schemas(schema_definitions): Deploy multiple schemas togetherdelete_schema(schema_name, tenant_id): Delete tenant schema(s); returnsList[str]of deleted namesschema_exists(schema_name, tenant_id): Check if schema existsget_tenant_schema_name(tenant_id, base_schema_name): Get tenant-specific schema namecreate_metadata_document(schema, doc_id, fields): Create or update metadata documentget_metadata_document(schema, doc_id): Get metadata document by IDquery_metadata_documents(schema, query, yql, **kwargs): Query metadata documentsdelete_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 historyConfigScope: Enum for config scope (SYSTEM, AGENT, ROUTING, TELEMETRY, SCHEMA, BACKEND)ConfigEntry: Dataclass for configuration entry with versionConfigStoreUnavailableError: a read the store could not answer within its retry budgetConfigWriteConflictError: anupdate_configwhose every compare-and-set lost to a concurrent writer (config_id,attempts); nothing was written
Methods:
initialize(): Initialize the configuration storeset_config(tenant_id, scope, service, config_key, config_value): Set config with versioningcompare_and_set_config(tenant_id, scope, service, config_key, config_value, *, expected_version): Append exactly the next version, orNoneon contentionupdate_config(tenant_id, scope, service, config_key, update, *, max_attempts=CONFIG_UPDATE_MAX_ATTEMPTS): Read-modify-write throughcompare_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 historylist_configs(tenant_id, scope, service): List configs with filterslist_all_configs(scope, service): List all configurations across all tenantsdelete_config(tenant_id, scope, service, config_key): Delete all versionsexport_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 undertenant_id; a payload carrying a schema-scope row is refused withValueErrorbefore any write; a row that cannot be written raisesRuntimeErrorafter every version the import already wrote is removed, so an import lands whole or not at allget_stats(): Get storage statisticshealth_check(): Check storage health
Lines of Code: ~290
interfaces/schema_loader.py¶
Purpose: Schema loading and ranking strategies
Key Classes:
SchemaLoader: Abstract schema loaderSchemaNotFoundException: Raised when schema not foundSchemaLoadError: Raised when schema fails to load/parse
Methods:
load_schema(schema_name): Load schema definition by namelist_available_schemas(): List all available schema namesschema_exists(schema_name): Check if schema existsload_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 executionAgentPerformance: Dataclass for an agent performance profile;average_confidenceisNoneuntil sampledWorkflowTemplate: Dataclass for a reusable workflow templateWorkflowLearningState: 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 anExceptionGroup.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 writeload_learning_state(tenant_id): load one coherent four-channel generation under the same distributed tenant lock used by replacementsave_template(tenant_id, template): Create or update a templatesave_generated_templates(tenant_id, templates): Persist one generated batch atomically and return its template IDs in input orderload_templates(tenant_id): Load all templates for tenantdelete_template(tenant_id, template_id): Delete a template by idhealth_check(): Check storage healthget_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 storesave_adapter(metadata): Save adapter metadataget_adapter(adapter_id): Get adapter by IDlist_adapters(tenant_id, agent_type, model_type, status, limit): List adapters with filtersget_active_adapter(tenant_id, agent_type, model_type): Get active adapterset_active(adapter_id, tenant_id, agent_type): Set adapter as activedeactivate_adapter(adapter_id): Deactivate adapterdeprecate_adapter(adapter_id): Deprecate adapterdelete_adapter(adapter_id): Delete adapterhealth_check(): Check storage healthget_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_pathmap a generic field to a schema field name. created_at/updated_atwithcreated_at_format:"epoch"(int seconds),"epoch_ms"(int milliseconds — for fields likecreation_timestamp), or"iso"(UTC string).metadata_fields:{metadata_key: schema_field}renames a value carried inDocument.metadatato 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 exactadd_embedding()wrapper:data,metadata, and an integer-secondcreated_at. Backends that hex/binary-encode embeddings (the ingestion path) override these with their processed vectors.include_metadata: whenfalse, 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:
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¶
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¶
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¶
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.pypins document, config, and workflow-record serialization boundaries.tests/backends/unit/test_backend_interface_contract.pyverifies concrete backend method signatures against SDK contracts.tests/backends/unit/test_backend_bool_contracts.pyverifies boolean return contracts.tests/backends/integration/test_document_mapping_roundtrip.pyexercises document mapping against a real Vespa service.- Concrete ConfigStore, WorkflowStore, and AdapterStore tests live alongside their implementations under
tests/backends/andtests/agents/.
Development¶
Adding New Interfaces¶
To add a new interface:
-
Create interface file in
cogniverse_sdk/interfaces/: -
Export in
__init__.py: -
Add tests:
Design Guidelines¶
- Keep It Pure: No implementation, only interfaces
- Type Hints: Always use type hints for clarity
- Documentation: Every method needs docstring
- Flexibility: Design for extensibility
- 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-vespafor Backend implementation - Usage: See
cogniverse-corefor how SDK is used - Foundation: See
cogniverse-foundationfor config base classes