Skip to content

Telemetry Module Documentation

Packages:

  • Foundation Interfaces: cogniverse-foundation (libs/foundation/cogniverse_foundation/telemetry/)

  • Phoenix Plugin: cogniverse-telemetry-phoenix (libs/telemetry-phoenix/) Layer: Foundation Layer (interfaces and infrastructure) + Plugin (Phoenix provider)

Architecture Note: Cogniverse uses a plugin-based telemetry architecture with two layers: 1. Foundation Layer (cogniverse-foundation): Telemetry provider interfaces, infrastructure, manager, and configuration. Zero knowledge of any specific backend. 2. Plugin Layer (cogniverse-telemetry-phoenix): Phoenix-specific implementation auto-discovered via Python entry points

There is no cogniverse_core.telemetry re-export layer — a compatibility shim existed briefly during the foundation-package migration but was removed; all telemetry imports go directly through cogniverse_foundation.

This design enables clean separation between telemetry infrastructure and provider-specific implementations, allowing easy swapping of telemetry backends (Phoenix, Datadog, New Relic, etc.).


Package Structure

1. Foundation Layer (cogniverse-foundation)

libs/foundation/cogniverse_foundation/telemetry/
├── __init__.py              # Package initialization
├── manager.py               # TelemetryManager singleton
├── config.py                # TelemetryConfig and BatchExportConfig
├── context.py               # Span context helpers
├── span_contract.py         # Canonical span I/O contract (record_span_io / read_span_io)
├── span_metrics.py          # Aggregates over span frames for the operations views
├── registry.py              # Provider registry for auto-discovery
└── providers/               # Provider interfaces
    ├── __init__.py
    └── base.py              # TelemetryProvider abstract base class

Purpose: Telemetry provider interfaces and infrastructure with multi-tenant management, configuration, and context helpers. Zero dependencies on specific telemetry backends (Phoenix, LangSmith, etc.).

Key Classes:

  • TelemetryProvider: Abstract base class for all telemetry providers
  • TelemetryConfig: Configuration for telemetry systems with BatchExportConfig
  • TelemetryManager: Singleton manager for multi-tenant tracer providers
  • Context helpers for common operations (search, encode, backend)
  • span_metrics: aggregates over a TraceStore span frame. trace_rows(spans) turns root spans into trace rows newest first, trace_statistics(rows) gives their counts, latency percentiles, outlier bounds (outlier_bounds(durations)) and per-operation figures. profile_selection_metrics(spans) gives per-modality count, p50/p95/p99 latency in ms and success rate, most-used first; aggregate_ab_compare(spans) gives the ABCompareAggregate of rlm.ab_compare spans (AB_COMPARE_SPAN_NAME), with averages, per_dataset and per_row (newest first) frames. span_succeeded(status) is false only for ERROR, so a span that finished without setting a status (UNSET) counts as a success.
  • span_contract: the one span I/O shape every operation uses, and the result annotation contract: RESULT_RELEVANCE, RELEVANCE_SCORES and persist_result_relevance(provider, project, span_id, result_id, label, readable_within_s=), which checks the span is in project and stores the rating under the result's own identifier. A span no project holds yet is looked up again with backoff for readable_within_s — span_readable_within_s(batch_config), the exporter's schedule_delay_millis plus SPAN_INGEST_MARGIN_S (5 s) — since a search hands out its span id before the span is exported; a span another project holds, or none by then, raises SpanNotInProjectError. persist_session_evaluation(provider, project, session_id, span_ids, outcome, score, readable_within_s=) stores a conversation's verdict (SESSION_OUTCOMES) as a SESSION_EVALUATION annotation on each of its spans, after finding every one in project the same way, keyed by the session so a new verdict replaces it. record_span_io(span, input_value=, output=, operation=, modality=) writes the input on input.value, the output as JSON on output.value, and the type on operation; read_span_io(row) reads {input, output, operation, modality} back and read_span_attributes(row) returns every attribute as a flat dotted-key dict from a Phoenix span row. Search (list output) and domain spans like query_enhancement (dict output) share the same writer and reader, so eval, dataset-building, experiments, and optimization read every operation uniformly. current_span_id() is the active span's 16-hex id (None with no span), and record_search_io_on_current_span(query, results, modality) records a search on the active span (one search_result_row per result); the search and document agents stamp both on their process span, whose id a client annotates. Query-enhancement spans also carry enhancement.path with lm for a genuine enhancement and heuristic_fallback for the heuristic expansion, so served rows stay machine-readable when the LM echoes the query or returns empty fields. Also holds the annotation constants (RESULT_RELEVANCE, RESULT_CLICK, RESULT_ID_META_KEY, RELEVANCE_POSITIVE_THRESHOLD, PREFERENCE_CHOSEN_THRESHOLD) that TripletExtractor, PreferencePairExtractor, the trace converter, and the relevance writer (persist_result_relevance) all use instead of hardcoding the annotation names, metadata key, and thresholds at each site.

2. Plugin Layer: Phoenix Telemetry Provider (cogniverse-telemetry-phoenix)

libs/telemetry-phoenix/cogniverse_telemetry_phoenix/
├── __init__.py              # Package initialization & exports
├── provider.py              # PhoenixProvider implementation
└── evaluation/              # Phoenix evaluation subsystem
    ├── __init__.py
    ├── analytics.py         # Analytics utilities (PhoenixAnalytics)
    ├── evaluation_provider.py   # PhoenixEvaluationProvider
    └── framework.py         # Evaluation framework

Purpose: Phoenix-specific implementation of telemetry interfaces. Auto-discovered via Python entry points.

Plugin Registration: The Phoenix provider is auto-discovered via entry points defined in pyproject.toml:

[project.entry-points."cogniverse.telemetry.providers"]
phoenix = "cogniverse_telemetry_phoenix:PhoenixProvider"

[project.entry-points."cogniverse.evaluation.providers"]
phoenix = "cogniverse_telemetry_phoenix.evaluation.evaluation_provider:PhoenixEvaluationProvider"

Benefits of Plugin Architecture:

  • Swappable Providers: Easy to add Datadog, New Relic, Jaeger, or custom providers

  • Zero Foundation Dependencies: cogniverse-foundation doesn't depend on the Phoenix SDK

  • Auto-discovery: Providers automatically registered via entry points

  • Clean Separation: Provider-specific code isolated in plugins


Table of Contents

  1. Module Overview
  2. Architecture Diagrams
  3. Provider Abstraction
  4. Phoenix Implementation
  5. Core Components
  6. Usage Examples
  7. Production Considerations
  8. Testing

Module Overview

Purpose and Responsibilities

The Telemetry Module provides multi-tenant observability infrastructure for the Cogniverse system with:

  • Multi-Tenant Isolation: Separate Phoenix projects per tenant for data isolation
  • Lazy Initialization: Tracer providers created on-demand with LRU caching
  • Batch Export: Configurable batch processing for high-throughput span export
  • Graceful Degradation: System continues functioning even when telemetry fails
  • Per-Modality Observability: Profile routing metrics aggregated from cogniverse.profile_selection spans; surfaced in the web client's "Profile metrics" view
  • OpenTelemetry Integration: Standards-based distributed tracing with Phoenix backend

Key Features

  1. Singleton Manager Pattern
  2. Thread-safe singleton TelemetryManager
  3. Global instance accessible throughout application
  4. Configuration from environment variables

  5. Tenant-Aware Tracing

  6. Project-based separation: cogniverse-{tenant_id}-{service}
  7. LRU cache for tenant providers (configurable max tenants)
  8. Automatic tenant context propagation

  9. Flexible Export Modes

  10. Production: Batch export with queue management (async)
  11. Testing: Synchronous export for immediate flush
  12. Configurable queue size, batch size, timeout

  13. Per-Modality Observability

  14. cogniverse.profile_selection spans emitted by ProfileSelectionAgent carry a profile_selection.modality attribute
  15. The runtime's GET /admin/tenant/{tenant_id}/telemetry/profile-selection route (routers/telemetry_metrics.py), shown in the web client's "Profile metrics" view, queries those spans from Phoenix and aggregates P50/P95/P99 latency, success rate, and request count per modality
  16. No separate metrics class required — all observability flows through standard OTel spans

  17. Context Helpers

  18. Pre-built span creators for search, encoding, backend operations
  19. Exact format matching with old instrumentation system
  20. OpenInference semantic conventions

Dependencies

Foundation Package (cogniverse-foundation):

  • opentelemetry-api: OpenTelemetry API
  • opentelemetry-sdk: OpenTelemetry Python SDK
  • pydantic: Data validation
  • sqlalchemy: Database ORM
  • pandas: Data structures

Phoenix Plugin (cogniverse-telemetry-phoenix):

  • arize-phoenix-client: Phoenix client
  • arize-phoenix-otel: Phoenix OpenTelemetry setup
  • pandas: Data structures

Architecture Diagrams

1. Multi-Tenant Telemetry Architecture

flowchart TB
    subgraph AppLayer["<span style='color:#000'>Application Layer</span>"]
        OrchestratorA["<span style='color:#000'>OrchestratorAgent<br/>Tenant A</span>"]
        VideoB["<span style='color:#000'>VideoAgent<br/>Tenant B</span>"]
        SummarizerA["<span style='color:#000'>SummarizerAgent<br/>Tenant A</span>"]
    end

    OrchestratorA --> Manager
    VideoB --> Manager
    SummarizerA --> Manager

    subgraph Manager["<span style='color:#000'>Telemetry Manager Singleton</span>"]
        Cache["<span style='color:#000'>Thread-Safe LRU Cache<br/>Tenant A Cache: Tracer Providers<br/>Tenant B Cache: Tracer Providers<br/>Configuration:<br/>• max_cached_tenants: 100 LRU eviction<br/>• tenant_project_template: cogniverse-tenant_id-service<br/>• batch_export_config: queue, batch size, timeout</span>"]
    end

    Manager --> ProviderA["<span style='color:#000'>TracerProvider<br/>Project:<br/>cogniverse-tenant-a-routing</span>"]
    Manager --> ProviderB["<span style='color:#000'>TracerProvider<br/>Project:<br/>cogniverse-tenant-b-video</span>"]
    Manager --> ProviderC["<span style='color:#000'>TracerProvider<br/>Project:<br/>cogniverse-tenant-a-summary</span>"]

    ProviderA --> Phoenix["<span style='color:#000'>Phoenix Backend<br/>OpenTelemetry<br/>• Project-based separation<br/>• Span storage<br/>• Analytics UI</span>"]
    ProviderB --> Phoenix
    ProviderC --> Phoenix

    style AppLayer fill:#90caf9,stroke:#1565c0,color:#000
    style Manager fill:#ffcc80,stroke:#ef6c00,color:#000
    style Cache fill:#b0bec5,stroke:#546e7a,color:#000
    style OrchestratorA fill:#90caf9,stroke:#1565c0,color:#000
    style VideoB fill:#90caf9,stroke:#1565c0,color:#000
    style SummarizerA fill:#90caf9,stroke:#1565c0,color:#000
    style ProviderA fill:#ce93d8,stroke:#7b1fa2,color:#000
    style ProviderB fill:#ce93d8,stroke:#7b1fa2,color:#000
    style ProviderC fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Phoenix fill:#a5d6a7,stroke:#388e3c,color:#000

Key Points:

  • Single TelemetryManager instance per application

  • Lazy creation of TracerProviders per tenant

  • LRU cache evicts old providers when max_cached_tenants exceeded

  • Each tenant gets isolated Phoenix project


2. Span Lifecycle and Export Flow

flowchart TB
    AppCode["<span style='color:#000'>Application Code<br/>with telemetry.span search, tenant_id=acme as span:<br/>span.set_attribute query, test<br/>business logic<br/>span.set_status Status.OK</span>"]

    AppCode --> TelManager["<span style='color:#000'>Telemetry Manager<br/>1. Get or create tracer for tenant:<br/>• Check LRU cache<br/>• Create TracerProvider if needed<br/>• Add tenant context attributes<br/><br/>2. Start span with tracer:<br/>• Add default attributes tenant.id, service.name, env<br/>• Add user attributes openinference.project.name, etc<br/>• Start as current context span</span>"]

    TelManager --> ProdProcessor{"<span style='color:#000'>Mode?</span>"}

    ProdProcessor -->|Production| BatchProc["<span style='color:#000'>Span Processor Production<br/>BatchSpanProcessor<br/>Span Queue max: 2048<br/>Contains: Span, Span, Span, ...<br/>Schedule every 500ms or when<br/>batch size 512 reached<br/>OTLP Exporter:<br/>• Endpoint: localhost:4317<br/>• Protocol: gRPC<br/>• Timeout: 30s</span>"]

    ProdProcessor -->|Test| SimpleProc["<span style='color:#000'>Span Processor Test Mode<br/>SimpleSpanProcessor SYNC<br/>• Immediate export no queue<br/>• 5s timeout to prevent hanging<br/>• Enabled via BatchExportConfig use_sync_export=True</span>"]

    BatchProc --> Phoenix["<span style='color:#000'>Phoenix Backend</span>"]
    SimpleProc --> Phoenix

    style AppCode fill:#90caf9,stroke:#1565c0,color:#000
    style TelManager fill:#ffcc80,stroke:#ef6c00,color:#000
    style ProdProcessor fill:#b0bec5,stroke:#546e7a,color:#000
    style BatchProc fill:#ce93d8,stroke:#7b1fa2,color:#000
    style SimpleProc fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Phoenix fill:#a5d6a7,stroke:#388e3c,color:#000

Export Modes:

  1. Production (Batch): Async export with queue, scheduled at intervals

  2. Testing (Sync): Immediate export for integration tests

Configuration:

batch_config = BatchExportConfig(
    max_queue_size=2048,           # Max spans in queue
    max_export_batch_size=512,     # Spans per export
    export_timeout_millis=30_000,  # Export timeout
    schedule_delay_millis=500,     # Export interval
)

Queue-full drop behaviour is handled natively by OTel's BatchSpanProcessor (created by phoenix.otel.register(batch=True) inside PhoenixProvider): when the queue reaches max_queue_size the processor drops new spans rather than blocking the calling thread.


3. Phoenix Project Isolation

flowchart TB
    subgraph Phoenix["<span style='color:#000'>Phoenix Backend</span>"]
        subgraph Project1["<span style='color:#000'>Project: cogniverse-acme-corp-routing</span>"]
            Acme1Spans["<span style='color:#000'>Spans from Tenant: acme-corp<br/>Service: routing<br/>Span 1: routing | Span 2: search | Span 3: enhance | ...</span>"]
        end

        subgraph Project2["<span style='color:#000'>Project: cogniverse-acme-corp-video-search</span>"]
            Acme2Spans["<span style='color:#000'>Spans from Tenant: acme-corp<br/>Service: video-search<br/>Span 1: encode | Span 2: search | Span 3: results | ...</span>"]
        end

        subgraph Project3["<span style='color:#000'>Project: cogniverse-techstart-routing</span>"]
            TechSpans["<span style='color:#000'>Spans from Tenant: techstart<br/>Service: routing<br/>Span 1: routing | Span 2: search | Span 3: enhance | ...</span>"]
        end

        Info["<span style='color:#000'>Project Naming:<br/>Template: cogniverse-tenant_id-service<br/><br/>Benefits:<br/>✓ Complete data isolation per tenant<br/>✓ Independent analytics per tenant<br/>✓ Separate retention policies per tenant<br/>✓ Easy tenant data deletion</span>"]
    end

    style Phoenix fill:#a5d6a7,stroke:#388e3c,color:#000
    style Project1 fill:#90caf9,stroke:#1565c0,color:#000
    style Project2 fill:#ffcc80,stroke:#ef6c00,color:#000
    style Project3 fill:#ce93d8,stroke:#7b1fa2,color:#000
    style Acme1Spans fill:#b0bec5,stroke:#546e7a,color:#000
    style Acme2Spans fill:#b0bec5,stroke:#546e7a,color:#000
    style TechSpans fill:#b0bec5,stroke:#546e7a,color:#000
    style Info fill:#b0bec5,stroke:#546e7a,color:#000

Project Name Resolution:

def get_project_name(tenant_id: str, service: str) -> str:
    return f"cogniverse-{tenant_id}-{service}"

# Examples:
# tenant="acme-corp", service="routing" → "cogniverse-acme-corp-routing"
# tenant="default", service="video-search" → "cogniverse-default-video-search"


Provider Abstraction

File: libs/foundation/cogniverse_foundation/telemetry/providers/base.py

The telemetry system uses a provider abstraction that defines interfaces for telemetry operations. This allows swapping backends (Phoenix, LangSmith, Datadog) without changing application code.

Store Interfaces

flowchart TB
    subgraph Abstraction["<span style='color:#000'>Foundation Layer - Provider Interfaces</span>"]
        TelemetryProvider["<span style='color:#000'>TelemetryProvider<br/><br/>• initialize(config)<br/>• configure_span_export()<br/>• session_context()</span>"]
        TraceStore["<span style='color:#000'>TraceStore<br/><br/>• get_spans()<br/>• iter_spans()<br/>• get_all_spans()<br/>• get_span_by_id()<br/>• span_projects()</span>"]
        AnnotationStore["<span style='color:#000'>AnnotationStore<br/><br/>• add_annotation()<br/>• get_annotations()<br/>• log_evaluations()</span>"]
        DatasetStore["<span style='color:#000'>DatasetStore<br/><br/>• create_dataset()<br/>• get_dataset()<br/>• append_to_dataset()<br/>• list_datasets()<br/>• delete_dataset()</span>"]
    end

    TelemetryProvider --> TraceStore
    TelemetryProvider --> AnnotationStore
    TelemetryProvider --> DatasetStore

    subgraph Implementation["<span style='color:#000'>Plugin Layer - Phoenix Implementation</span>"]
        PhoenixProvider["<span style='color:#000'>PhoenixProvider</span>"]
        PhoenixTraceStore["<span style='color:#000'>PhoenixTraceStore</span>"]
        PhoenixAnnotationStore["<span style='color:#000'>PhoenixAnnotationStore</span>"]
        PhoenixDatasetStore["<span style='color:#000'>PhoenixDatasetStore</span>"]
    end

    PhoenixProvider -.->|implements| TelemetryProvider
    PhoenixTraceStore -.->|implements| TraceStore
    PhoenixAnnotationStore -.->|implements| AnnotationStore
    PhoenixDatasetStore -.->|implements| DatasetStore

    style Abstraction fill:#a5d6a7,stroke:#388e3c,color:#000
    style Implementation fill:#ce93d8,stroke:#7b1fa2,color:#000
    style TelemetryProvider fill:#a5d6a7,stroke:#388e3c,color:#000
    style TraceStore fill:#a5d6a7,stroke:#388e3c,color:#000
    style AnnotationStore fill:#a5d6a7,stroke:#388e3c,color:#000
    style DatasetStore fill:#a5d6a7,stroke:#388e3c,color:#000
    style PhoenixProvider fill:#ce93d8,stroke:#7b1fa2,color:#000
    style PhoenixTraceStore fill:#ce93d8,stroke:#7b1fa2,color:#000
    style PhoenixAnnotationStore fill:#ce93d8,stroke:#7b1fa2,color:#000
    style PhoenixDatasetStore fill:#ce93d8,stroke:#7b1fa2,color:#000

TraceStore Interface

Query traces and spans from telemetry backend.

class TraceStore(ABC):
    @abstractmethod
    async def get_spans(
        self,
        project: str,
        start_time: Optional[datetime] = None,
        end_time: Optional[datetime] = None,
        filters: Optional[Dict[str, Any]] = None,
        limit: int = 1000,
        columns: Optional[Sequence[str]] = None,
    ) -> pd.DataFrame:
        """
        Returns DataFrame with standardized columns:
        - context.span_id: Unique span identifier
        - name: Span operation name
        - attributes.*: Span attributes (flattened)
        - start_time, end_time: Timestamps
        """
        pass

    @abstractmethod
    async def iter_spans(
        self,
        project: str,
        start_time: Optional[datetime] = None,
        end_time: Optional[datetime] = None,
        filters: Optional[Dict[str, Any]] = None,
        page_size: int = 1000,
        columns: Optional[Sequence[str]] = None,
    ) -> AsyncIterator[pd.DataFrame]:
        """Yield page-sized DataFrames for matching spans."""
        pass

    @abstractmethod
    async def get_all_spans(
        self,
        project: str,
        start_time: Optional[datetime] = None,
        end_time: Optional[datetime] = None,
        filters: Optional[Dict[str, Any]] = None,
    ) -> pd.DataFrame:
        """Return every matching span through provider-native pagination."""
        pass

    @abstractmethod
    async def get_span_by_id(
        self, span_id: str, project: str
    ) -> Optional[Dict[str, Any]]:
        """Get single span by ID."""
        pass

    @abstractmethod
    async def span_projects(self, span_ids: Sequence[str]) -> Dict[str, Optional[str]]:
        """The project holding each span, or None for one no project holds."""
        pass

Use get_spans for deliberately bounded windows. Use iter_spans when the caller needs page-sized frames and bounded memory while walking the full result set. columns projects the requested standardized fields server-side when supported; the Phoenix implementation walks the requested time range with adaptive windows, splitting any full window in half and only falling back to span_id exclusions when a window can no longer be split. Consumers whose correctness depends on complete history use get_all_spans; the Phoenix implementation builds it from iter_spans and raises if any page fails, so callers never interpret a partial history as the complete project. Its start_time and end_time columns are normalized to timezone-aware UTC datetimes; an invalid timestamp raises with project and column context before the frame is returned. iter_spans and get_all_spans take filters={"name": ..., "roots_only": True}: name keeps spans of one name or several, roots_only keeps spans without a parent (one per trace); any other key raises ValueError.

AnnotationStore Interface

Manage human/LLM annotations on spans for approval workflows and evaluations.

class AnnotationStore(ABC):
    @abstractmethod
    async def add_annotation(
        self,
        span_id: str,
        name: str,           # e.g., "human_approval", "relevance_score"
        label: str,          # e.g., "approved", "rejected"
        score: float,        # 0.0-1.0
        metadata: Dict[str, Any],
        project: str,
        identifier: Optional[str] = None,  # several annotations of one name per span
    ) -> str:
        """Add annotation to a span; replaces the span's annotation of the
        same name and identifier."""
        pass

    @abstractmethod
    async def get_annotations(
        self,
        spans_df: pd.DataFrame,
        project: str,
        annotation_names: Optional[List[str]] = None,
    ) -> pd.DataFrame:
        """Get annotations for spans."""
        pass

    @abstractmethod
    async def log_evaluations(
        self,
        eval_name: str,
        evaluations_df: pd.DataFrame,  # Must contain: span_id, score, label
        project: str,
    ) -> None:
        """Bulk upload evaluation results as span annotations."""
        pass

DatasetStore Interface

Manage training datasets for model fine-tuning.

class DatasetStore(ABC):
    @abstractmethod
    async def create_dataset(
        self, name: str, data: pd.DataFrame, metadata: Optional[Dict[str, Any]] = None
    ) -> str:
        """Create new dataset."""
        pass

    @abstractmethod
    async def get_dataset(self, name: str) -> pd.DataFrame:
        """Load dataset by name.

        Raises DatasetNotFoundError (a ValueError subclass) if no dataset by
        that name exists; a backend outage raises the underlying transport
        error instead of returning an empty/no-data result.
        """
        pass

    @abstractmethod
    async def append_to_dataset(
        self,
        name: str,
        data: pd.DataFrame,
        metadata: Optional[Dict[str, Any]] = None,
    ) -> None:
        """Append records to an existing dataset as a new version.

        ``metadata`` carries input_keys/output_keys/metadata_keys (same
        shape as ``create_dataset``). Raises DatasetNotFoundError if the
        dataset does not exist.
        """
        pass

    @abstractmethod
    async def list_datasets(self) -> List[Dict[str, Any]]:
        """Every dataset, as ``{name, example_count, created_at,
        description}``. Raises DatasetStoreUnavailableError when the store
        cannot answer."""
        pass

    async def delete_dataset(self, name: str) -> bool:
        """Delete a dataset by name (last-write-wins blob storage deletes
        before create). Returns True if one was deleted, False if none
        existed. Not ``@abstractmethod`` — backends that can't delete raise
        NotImplementedError from the default.
        """
        raise NotImplementedError

    async def describe_datasets(self) -> List[DatasetSummary]:
        """Every stored dataset, newest first, as DatasetSummary(id, name,
        example_count, created_at, description, tenant_id, metadata).
        Raises DatasetStoreUnavailableError when the store cannot answer.
        Not ``@abstractmethod``; the default raises NotImplementedError.
        """
        raise NotImplementedError

    async def replace_dataset(
        self, name: str, data: pd.DataFrame, metadata: Optional[Dict[str, Any]] = None
    ) -> str:
        """Replace a dataset's contents.

        Same-name writers are serialized per store instance. If the delete /
        create path faults after the old dataset is gone, the previous frame is
        recreated before the exception propagates.
        """

A dataset created with metadata["tenant_id"] (DATASET_TENANT_KEY) is owned by that tenant: PhoenixDatasetStore creates it with the owner in the dataset's Phoenix metadata (GraphQL createDataset) and then uploads its rows, deleting it again when the upload fails; describe_datasets reports the owner as DatasetSummary.tenant_id (None for a dataset created without one). Creating a name that exists appends to it only for its owner; any other tenant gets DatasetOwnedByAnotherTenantError (a ValueError).

DatasetNotFoundError (cogniverse_foundation.telemetry.providers.base) subclasses ValueError so existing except ValueError callers keep working, while letting callers distinguish a genuinely missing dataset from a backend outage.

PhoenixDatasetStore uploads each row through upload_dataset_rows (cogniverse_telemetry_phoenix.provider), storing every value as the string the frame's CSV rendering holds for it (an empty cell is ""). A value has no size limit: optimizer versions whose content or ledger exceeds 128 KiB (a week of consumed span ids) round-trip byte-for-byte.

The telemetry provider exposes these store interfaces. replace_dataset is the safe helper for stable-name overwrites: same-name writers are serialized only within one store instance on one event loop, and a torn delete/create that cannot restore the prior frame raises DatasetReplaceRestoreFailedError with the dataset name, the original error, and the restore error. That protection is in-process only; a hard kill between delete and create can still lose the dataset because Phoenix exposes no atomic swap primitive. Experiment tracking lives on the separate EvaluationProvider stack (PhoenixEvaluationProvider.create_experiment), and metric aggregation via PhoenixAnalytics — not on the telemetry provider.

TelemetryProvider Base Class

Providers implement all store interfaces and handle backend-specific initialization.

class TelemetryProvider(ABC):
    def __init__(self, name: str):
        self.name = name
        self._trace_store: Optional[TraceStore] = None
        self._annotation_store: Optional[AnnotationStore] = None
        self._dataset_store: Optional[DatasetStore] = None

    @abstractmethod
    def initialize(self, config: Dict[str, Any]) -> None:
        """
        Initialize provider with configuration.
        Provider extracts its own keys from generic config dict.

        Example config for Phoenix:
            {"tenant_id": "acme", "http_endpoint": "http://localhost:6006"}

        Example config for LangSmith:
            {"tenant_id": "acme", "api_key": "xxx", "project": "my-project"}
        """
        pass

    @abstractmethod
    def configure_span_export(
        self,
        endpoint: str,
        project_name: str,
        use_batch_export: bool = True,
        batch_config: Optional[Any] = None,
        resource_attributes: Optional[Dict[str, str]] = None,
        *,
        raise_on_export_failure: bool,
    ) -> Any:
        """Configure OTLP span export, returns TracerProvider.

        ``batch_config`` carries the BatchExportConfig queue/batch/timeout
        knobs; ``resource_attributes`` (service.version + extras) land on
        the provider's OTel Resource. ``raise_on_export_failure`` selects the
        checked synchronous processor used only by spans whose records cannot
        be lost; normal synchronous spans keep OpenTelemetry's non-raising
        processor.
        """
        pass

    def preload_span_export(self) -> None:
        """Import what configure_span_export loads on first use (Phoenix:
        phoenix.otel and the OTLP gRPC exporter). The runtime calls it through
        TelemetryManager.preload_span_export before serving, so a tenant's
        first span does not import them on the event loop. Default: no-op."""

    @abstractmethod
    async def list_projects(self, name_contains: str) -> List[str]:
        """Names of the projects whose name contains ``name_contains``.
        Raises when the backend does not answer, never an empty list."""

    @abstractmethod
    async def project_id(self, name: str) -> Optional[str]:
        """The backend's id for the project (the one its UI addresses it
        by), or None when it does not exist. Raises when the backend does not
        answer. GET /admin/tenant/{tenant}/telemetry/phoenix links to it."""

    @abstractmethod
    async def delete_project(self, name: str) -> bool:
        """Delete a project and its spans; False when it does not exist.
        Raises when the backend does not answer or refuses. A tenant delete
        removes the tenant's projects through these two."""

    @property
    def traces(self) -> TraceStore:
        """Get trace store (query spans). Raises RuntimeError if not initialized."""
        if self._trace_store is None:
            raise RuntimeError(f"{self.name} provider not initialized")
        return self._trace_store

    @property
    def annotations(self) -> AnnotationStore:
        """Get annotation store (manage annotations). Raises RuntimeError if not initialized."""
        if self._annotation_store is None:
            raise RuntimeError(f"{self.name} provider not initialized")
        return self._annotation_store

    @property
    def datasets(self) -> DatasetStore:
        """Get dataset store (manage training datasets). Raises RuntimeError if not initialized."""
        if self._dataset_store is None:
            raise RuntimeError(f"{self.name} provider not initialized")
        return self._dataset_store

    @abstractmethod
    @contextmanager
    def session_context(self, session_id: str) -> Generator[None, None, None]:
        """Context manager for session tracking."""
        yield

Phoenix Implementation

Package: cogniverse-telemetry-phoenix Location: libs/telemetry-phoenix/cogniverse_telemetry_phoenix/

Phoenix is the default telemetry provider, implementing all store interfaces using Phoenix AsyncClient.

PhoenixProvider

File: libs/telemetry-phoenix/cogniverse_telemetry_phoenix/provider.py

from cogniverse_telemetry_phoenix import PhoenixProvider

# Initialize Phoenix provider (provider.initialize() is called by TelemetryManager)
# Required keys: tenant_id, http_endpoint, grpc_endpoint
provider = PhoenixProvider()
provider.initialize({
    "tenant_id": "acme",
    "http_endpoint": "http://localhost:6006",     # HTTP API endpoint (required)
    "grpc_endpoint": "http://localhost:4317",     # gRPC OTLP endpoint (required)
})

# Query spans
spans_df = await provider.traces.get_spans(
    project="cogniverse-acme-routing",
    start_time=datetime.now() - timedelta(hours=1),
    limit=100
)

# Add annotation
await provider.annotations.add_annotation(
    span_id="abc123",
    name="human_approval",
    label="approved",
    score=1.0,
    metadata={"reviewer": "user@example.com"},
    project="cogniverse-acme-routing"
)

# Create dataset
dataset_id = await provider.datasets.create_dataset(
    name="golden_queries_v1",
    data=approved_queries_df,
    metadata={"version": "1.0"}
)

Plugin Registration

Phoenix provider is auto-discovered via Python entry points in pyproject.toml:

[project.entry-points."cogniverse.telemetry.providers"]
phoenix = "cogniverse_telemetry_phoenix:PhoenixProvider"

[project.entry-points."cogniverse.evaluation.providers"]
phoenix = "cogniverse_telemetry_phoenix.evaluation.evaluation_provider:PhoenixEvaluationProvider"

Implementing a Custom Provider

To implement a new telemetry provider (e.g., LangSmith, Datadog):

from cogniverse_foundation.telemetry.providers.base import (
    TelemetryProvider, TraceStore, AnnotationStore
)

class LangSmithTraceStore(TraceStore):
    async def get_spans(self, project, start_time=None, end_time=None,
                        filters=None, limit=1000) -> pd.DataFrame:
        # LangSmith-specific implementation
        runs = await self.client.list_runs(project_name=project)
        return self._convert_to_dataframe(runs)

    async def get_all_spans(self, project, start_time=None, end_time=None,
                            filters=None) -> pd.DataFrame:
        # Follow the provider's cursor until there is no next page.
        runs = await self.client.list_all_runs(project_name=project)
        return self._convert_to_dataframe(runs)

    async def get_span_by_id(self, span_id, project) -> Optional[Dict]:
        return await self.client.read_run(run_id=span_id)

class LangSmithProvider(TelemetryProvider):
    def __init__(self):
        super().__init__("langsmith")

    def initialize(self, config: Dict[str, Any]) -> None:
        self.client = LangSmithClient(api_key=config["api_key"])
        self._trace_store = LangSmithTraceStore(self.client)
        # Initialize other stores...

    def configure_span_export(
        self,
        endpoint,
        project_name,
        use_batch_export=True,
        batch_config=None,
        resource_attributes=None,
        *,
        raise_on_export_failure: bool,
    ):
        # LangSmith uses different export mechanism
        if raise_on_export_failure:
            raise RuntimeError(
                "LangSmithProvider cannot guarantee synchronous export"
            )
        return LangSmithTracerProvider(project=project_name)

Register via entry points:

[project.entry-points."cogniverse.telemetry.providers"]
langsmith = "my_package:LangSmithProvider"

Core Components

1. TelemetryManager

File: libs/foundation/cogniverse_foundation/telemetry/manager.py

Purpose: Singleton manager for multi-tenant tracer providers with lazy initialization and LRU caching.

Singleton Pattern: TelemetryManager implements a thread-safe singleton pattern using __new__ and a class-level lock. Only one instance exists per Python process, and subsequent instantiations return the same instance.

Key Attributes:

_instance: Optional[TelemetryManager]  # Singleton instance (class variable)
_lock: threading.Lock                   # Thread-safe singleton creation (class variable)
_initialized: bool                      # Tracks if __init__ has run
config: TelemetryConfig                # Telemetry configuration
_tenant_providers: Dict[str, TracerProvider]  # Cached providers per tenant
_tenant_tracers: Dict[str, Tracer]     # Cached tracers per (tenant, service)
_tracer_provider_keys: Dict[str, str]  # cache_key -> _tenant_providers key it was built from
_tracer_created_at: Dict[str, float]   # cache_key -> time.monotonic() at insert (TTL tracking)
_project_configs: Dict[str, Dict]      # Per-project configuration overrides
_cache_hits: int                        # Cache performance metrics
_cache_misses: int
_failed_initializations: int

Cache Eviction (count + TTL): _evict_old_tracers() runs after every cache miss. It first enforces the LRU count cap (max_cached_tenants), then detaches providers no cached tracer references. One background worker drains detached providers outside the manager lock after all their recording spans end. TelemetryManager.span() acquires its tracer and starts recording under the same manager lock; callers that hold a tracer across that lock (get_tracer(), TenantRoutingTracerProvider) can start a span on a provider that has since drained, and the lease processor records no lease for it. Cancellation releases the lease when the span scope exits. A retired provider whose spans have not ended within retirement_lease_timeout_seconds is shut down anyway, with a warning naming the provider and the outstanding span count. _cached_tracer() checks tenant_cache_ttl_seconds on each lookup and retires expired entries through the same worker. A TTL of 0 or less disables expiry.

Retired providers are capped at max(1, max_cached_tenants). At capacity, new optional tracer creation emits a warning and yields no recording tracer; cached tracers remain usable. This bounds exporter threads during collector outages. get_stats() exposes retired_providers and retirement_workers. shutdown() closes admission, retires all cached providers, and waits at most 30 seconds for draining. Leased providers continue draining in the background if that deadline expires.

Initialization Pattern:

The singleton uses double-checked locking for thread safety:

class TelemetryManager:
    _instance = None
    _lock = threading.Lock()

    def __new__(cls, config: Optional[TelemetryConfig] = None):
        if cls._instance is None:
            with cls._lock:
                if cls._instance is None:
                    cls._instance = super().__new__(cls)
                    cls._instance._initialized = False
        return cls._instance

    def __init__(self, config: Optional[TelemetryConfig] = None):
        if self._initialized:
            return  # Skip re-initialization
        # ... initialization code ...
        self._initialized = True

Recommended Usage:

from cogniverse_foundation.telemetry.manager import TelemetryManager, get_telemetry_manager

# Option 1: Direct instantiation (singleton pattern handles multiple calls)
telemetry = TelemetryManager()  # Creates instance on first call
telemetry2 = TelemetryManager()  # Returns same instance
assert telemetry is telemetry2  # True

# Option 2: Use the convenience function (preferred)
telemetry = get_telemetry_manager()  # Gets or creates singleton

# Configuration is only applied on FIRST instantiation
# To reconfigure, use reset() first (TESTS ONLY):
TelemetryManager.reset()  # Clears singleton, shuts down providers
telemetry = TelemetryManager(new_config)  # Fresh instance with new config

Main Methods:

__new__(cls, config: Optional[TelemetryConfig] = None) -> TelemetryManager

Thread-safe singleton constructor.

Parameters:

  • config: Required telemetry configuration (only applied on first instantiation; raises ValueError if None). Use get_telemetry_manager() to auto-load from ConfigManager.

Returns: Singleton TelemetryManager instance

Example:

# First call creates instance with config
config = TelemetryConfig(enabled=True, level=TelemetryLevel.DETAILED)
manager1 = TelemetryManager(config)

# Subsequent calls return same instance (config parameter ignored)
manager2 = TelemetryManager()
assert manager1 is manager2  # True
assert manager2.config is config  # True - uses original config

reset() -> None (class method)

Reset singleton instance - FOR TESTS ONLY.

Shuts down all tracer providers, clears caches, resets the class singleton, and clears the module-global that get_telemetry_manager() short-circuits on — so the next get_telemetry_manager() rebuilds a fresh, live instance rather than returning the shut-down one. The global is cleared under the same lock get_telemetry_manager() uses for its cold build, so a rebuild racing a reset is single-flight.

Example:

# In test setup/teardown
TelemetryManager.reset()  # Clear singleton state

# Now can create fresh instance with test config
test_config = TelemetryConfig(
    otlp_endpoint="localhost:24317",
    batch_config=BatchExportConfig(use_sync_export=True)
)
manager = TelemetryManager(test_config)

get_telemetry_manager(config_manager=None) -> TelemetryManager (module function)

Get the global telemetry manager instance. On first call, loads config from ConfigManager automatically under SYSTEM_TENANT_ID (telemetry config is cluster-wide — OTLP endpoint, batch/export settings). Per-request tenant scoping happens inside TelemetryManager.span(tenant_id=...), not here. This is the preferred way to access the singleton.

Parameters: - config_manager: Optional ConfigManager instance. If None on first call, creates one via create_default_config_manager().

On first call this also applies a TELEMETRY_OTLP_ENDPOINT env var override (set by the Helm chart in k3d deployments) if present and different from the loaded config's otlp_endpoint, clearing cached tenant providers/tracers so the new endpoint takes effect.

configure_telemetry_endpoints(*, otlp_endpoint, http_endpoint) -> None (module function)

Records the Phoenix endpoints a deployment names for the process, for entrypoints that must not build the manager eagerly (the runtime, the ingestion worker and cogniverse-eval pass TELEMETRY_OTLP_ENDPOINT / TELEMETRY_HTTP_ENDPOINT). They are applied to the singleton when it is built, or at once to one already built; None leaves the stored config's value. otlp_endpoint sets config.otlp_endpoint, http_endpoint sets provider_config["http_endpoint"].

Example:

from cogniverse_foundation.telemetry.manager import get_telemetry_manager

# Throughout application code (auto-loads config from ConfigManager)
telemetry = get_telemetry_manager()
with telemetry.span("operation", tenant_id="acme") as span:
    span.set_attribute("key", "value")


get_tracer(tenant_id: str, project_name: Optional[str] = None) -> Optional[Tracer]

Get or create tracer for a specific tenant. Note: This is LEGACY - use span() instead.

Parameters:

  • tenant_id: Tenant identifier for project isolation

  • project_name: Optional project name for management operations (e.g., "experiments", "synthetic_data", "system"). If None, uses tenant-only project.

Returns:

  • Tracer instance if telemetry enabled

  • None if telemetry disabled or initialization failed

Caching Logic:

  1. Check LRU cache for {tenant_id}:{project_name} key

  2. If cached, return tracer (cache hit)

  3. If not cached:

  4. Create TracerProvider for tenant if needed
  5. Get tracer from provider
  6. Cache tracer with LRU eviction
  7. Return tracer or None on error

Example:

manager = TelemetryManager()

# Get tracer for tenant
tracer = manager.get_tracer("acme-corp", "routing")

# Use tracer
with tracer.start_as_current_span("process_query") as span:
    span.set_attribute("query", "test")


span(name: str, tenant_id: str, project_name: Optional[str] = None, attributes: Optional[Dict[str, Any]] = None, component: str = "agents") -> ContextManager

Context manager for creating tenant-specific spans.

Parameters:

  • name: Span name (e.g., "search", "cogniverse.routing")

  • tenant_id: Tenant identifier

  • project_name: Optional project name for span isolation

  • attributes: Optional span attributes

  • component: Telemetry-level tag — one of search_service / agents / backend / pipeline / encoder. TelemetryConfig.should_instrument_component() checks this against config.level; if the level doesn't admit the component, span() yields a NoOpSpan without ever creating a tracer. Default agents emits at DETAILED and above. Pass encoder for per-inference model detail spans, pipeline for ingestion-stage spans, search_service for the top-level search HTTP entry (admitted even at BASIC).

Automatic Attributes:

  • tenant.id: Tenant identifier

  • service.name: Service name

  • environment: Environment (development, production)

The three are exported together as SPAN_ENVELOPE_ATTRIBUTES in cogniverse_foundation.telemetry.manager, so a consumer that pins a span's attribute set names the envelope by import instead of restating it.

Graceful Degradation:

  • If telemetry disabled or fails, yields no-op span

  • Application continues normally even when telemetry fails

Example:

manager = TelemetryManager()

# Basic span
with manager.span("search", tenant_id="acme-corp") as span:
    span.set_attribute("query", "test")
    results = search(query)
    span.set_attribute("num_results", len(results))

# Span with canonical input/output contract
with manager.span(
    "cogniverse.routing",
    tenant_id="acme-corp",
    attributes={
        "openinference.project.name": "cogniverse-acme-corp-routing",
    }
) as span:
    decision = route_query(query)
    # output.value = json.dumps({"chosen_agent": "video_search", ...})
    record_span_io(
        span,
        input_value=query,
        output={"chosen_agent": "video_search", "confidence": decision.confidence},
        operation="routing",
    )


session_span(name: str, tenant_id: str, session_id: str, project_name: Optional[str] = None, attributes: Optional[Dict[str, Any]] = None, component: str = "search_service") -> ContextManager

Context manager for creating spans within a session context.

Purpose: Links multiple requests in a multi-turn conversation session. All spans created within a session share the same session.id attribute and are grouped together in Phoenix Sessions view.

Parameters:

  • name: Span name

  • tenant_id: Tenant identifier

  • session_id: Session identifier (client-generated UUID)

  • project_name: Optional Phoenix project name

  • attributes: Optional span attributes

Automatic Attributes:

  • session.id: Session identifier (for grouping in Phoenix)

  • All attributes from span() method

Example:

manager = TelemetryManager()

# Multi-turn conversation with session tracking
session_id = "user-session-uuid-123"

# Turn 1: First search
with manager.session_span("search", tenant_id="acme", session_id=session_id) as span:
    span.set_attribute("query", "find basketball videos")
    results = search("find basketball videos")

# Turn 2: Follow-up search (same session)
with manager.session_span("search", tenant_id="acme", session_id=session_id) as span:
    span.set_attribute("query", "show me dunks")
    results = search("show me dunks")

# Both spans will be grouped in Phoenix Sessions view under session_id

Phoenix Sessions View:

  • Traces with the same session.id are grouped together

  • Enables viewing complete conversation trajectories

  • Supports session-level evaluation and fine-tuning data extraction


register_project(tenant_id: str, project_name: str, **kwargs) -> None

Register a project with optional config overrides.

Purpose: Allows per-project configuration (e.g., different telemetry endpoints for tests). Single source of truth for project settings.

Parameters:

  • tenant_id: Tenant identifier

  • project_name: Project name (e.g., "search", "synthetic_data", "routing")

  • **kwargs: Optional overrides:

  • otlp_endpoint: Override OTLP gRPC endpoint for span export
  • http_endpoint: Override HTTP endpoint used by the provider for span queries
  • grpc_endpoint: Override gRPC endpoint used by the provider (distinct from otlp_endpoint, which configures the OTel exporter itself)
  • use_sync_export: Override batch/sync export mode

Project Naming: cogniverse-{tenant_id}-{project_name}

Example:

manager = TelemetryManager()

# Use defaults from config
manager.register_project(tenant_id="customer-123", project_name="search")

# Override endpoints for tests
manager.register_project(
    tenant_id="test-tenant1",
    project_name="synthetic_data",
    otlp_endpoint="http://localhost:24317",
    http_endpoint="http://localhost:26006",
    use_sync_export=True
)


Note: TelemetryManager has no bare session() method (no context manager that establishes a session without also creating a span). To share a session_id across multiple spans, either call session_span() for each one (it re-enters provider.session_context(session_id) on every call, which is idempotent), or use provider.session_context() directly around plain span() calls:

manager = TelemetryManager()

# Wrap multiple spans in a session context via the provider directly
provider = manager.get_provider(tenant_id="acme")
with provider.session_context("user-session-abc"):
    # Spans created here share the session_id
    with manager.span("operation1", tenant_id="acme") as span:
        pass
    with manager.span("operation2", tenant_id="acme") as span:
        pass

shutdown() -> None

Shutdown all tracer providers gracefully.

Purpose: Flushes pending spans (force_flush(timeout_millis=5000)) and calls shutdown() on every cached TracerProvider, then clears the tenant provider/tracer caches. Used internally by reset(); call directly at application shutdown to ensure clean teardown.

Example:

manager = TelemetryManager()
# ... application runs ...
manager.shutdown()  # Flush + close all tracer providers


get_provider(tenant_id: str, project_name: Optional[str] = None) -> TelemetryProvider

Get telemetry provider for querying spans/annotations/datasets.

Purpose: This is separate from span export (which uses OpenTelemetry OTLP). Providers are used for reading data from the telemetry backend.

Parameters:

  • tenant_id: Tenant identifier

  • project_name: Optional project name to get project-specific config

Returns: TelemetryProvider instance

Raises: ValueError if no providers available or provider initialization fails

Endpoint Derivation: grpc_endpoint and http_endpoint are derived from config.otlp_endpoint when not explicitly set in config.provider_config or via a registered project override — e.g. otlp_endpoint="localhost:4317" yields grpc_endpoint="http://localhost:4317" and http_endpoint="http://localhost:6006" (the :4317 gRPC port is swapped for the :6006 HTTP port). This lets providers like Phoenix initialize without requiring manual provider_config entries in the common case. TelemetryManager.provider_endpoints() returns that {"grpc_endpoint", "http_endpoint"} pair before any project override; PhoenixEvaluationProvider.initialize() builds on it when its config names no endpoints.

Example:

manager = TelemetryManager()

# Get provider for querying
provider = manager.get_provider(tenant_id="customer-123")

# Query spans
spans_df = await provider.traces.get_spans(
    project="cogniverse-customer-123-search",
    start_time=datetime(2025, 1, 1),
    limit=1000
)

# Add annotations
await provider.annotations.add_annotation(
    span_id="abc123",
    name="human_approval",
    label="approved",
    score=1.0,
    metadata={"reviewer": "alice"},
    project="cogniverse-customer-123-synthetic_data"
)


force_flush(timeout_millis: int = 10000) -> bool

Force flush all spans for all tenants.

Parameters:

  • timeout_millis: Timeout in milliseconds for flush operation

Returns: True if all flushes succeeded, False otherwise

Usage: Call before application shutdown or after critical operations in tests

Example:

manager = TelemetryManager()

# ... application code ...

# Ensure all spans exported before shutdown
success = manager.force_flush(timeout_millis=5000)
if not success:
    logger.warning("Some spans may not have been exported")


get_stats() -> Dict[str, Any]

Get telemetry manager statistics.

Returns:

{
    "cache_hits": 150,          # Number of cache hits
    "cache_misses": 10,         # Number of cache misses
    "failed_initializations": 0,  # Failed tracer creations
    "cached_tenants": 5,        # Number of cached providers
    "cached_tracers": 8,        # Number of cached tracers
    "config": {
        "enabled": True,
        "level": "detailed",
        "environment": "production"
    }
}

Example:

manager = TelemetryManager()
stats = manager.get_stats()
print(f"Cache hit rate: {stats['cache_hits'] / (stats['cache_hits'] + stats['cache_misses']):.2%}")


2. TelemetryConfig

File: libs/foundation/cogniverse_foundation/telemetry/config.py

Purpose: Configuration for telemetry system with environment variable support.

Key Attributes:

# Core settings
enabled: bool                       # Enable/disable telemetry (default: true)
level: TelemetryLevel              # DISABLED, BASIC, DETAILED, VERBOSE
environment: str                    # development, production

# OpenTelemetry span export settings (generic OTLP - backend-agnostic)
otlp_enabled: bool                  # Enable OTLP span export (default: true)
otlp_endpoint: str                  # OTLP collector endpoint (default: localhost:4317)
otlp_use_tls: bool                 # Use TLS for OTLP connection (default: false)

# Provider selection (for querying spans/annotations/datasets)
provider: Optional[str]             # Provider name ("phoenix", "langsmith", None=auto-detect)
provider_config: Dict[str, Any]     # Provider-specific config (dict interpreted by provider)

# Multi-tenant settings
tenant_project_template: str       # "cogniverse-{tenant_id}"
tenant_service_template: str       # "cogniverse-{tenant_id}-{service}"
max_cached_tenants: int            # LRU cache size (default: 100)
tenant_cache_ttl_seconds: int      # Cache TTL (default: 3600)
retirement_lease_timeout_seconds: float  # Retired-provider lease wait (default: 30.0)

# Batch export settings
batch_config: BatchExportConfig    # Batch export configuration

# Service identification
service_name: str                  # Service name (default: "video-search")
service_version: str               # Service version (default: "1.0.0")

# Resource attributes
extra_resource_attributes: Dict[str, str]  # Additional resource attributes

TelemetryLevel Enum:

DISABLED = "disabled"   # No telemetry
BASIC = "basic"         # Only search operations
DETAILED = "detailed"   # Search + encoders + backend
VERBOSE = "verbose"     # Everything including internal operations

BatchExportConfig:

max_queue_size: int = 2048           # Max spans in queue
max_export_batch_size: int = 512     # Spans per export batch
export_timeout_millis: int = 30_000  # Export timeout
schedule_delay_millis: int = 500     # Export interval

# Test mode
use_sync_export: bool = False        # Must be set explicitly on BatchExportConfig

Note: use_sync_export has no environment-variable binding — despite a TELEMETRY_SYNC_EXPORT env var appearing in some test fixtures, nothing in cogniverse_foundation.telemetry reads it; only TELEMETRY_OTLP_ENDPOINT is read (by get_telemetry_manager()). To enable sync export, construct TelemetryConfig(batch_config=BatchExportConfig(use_sync_export=True)) explicitly.

Queue-full drop behaviour is handled natively by OTel's BatchSpanProcessor (created by phoenix.otel.register(batch=True) inside PhoenixProvider): when the queue reaches max_queue_size the processor drops new spans rather than blocking the calling thread.

Main Methods:

from_dict(data: Dict[str, Any]) -> TelemetryConfig

Deserialize from dictionary (used by ConfigManager for persistence).

to_dict() -> Dict[str, Any]

Serialize to dictionary for persistence (used by ConfigManager).

Example:

# Get default config
config = TelemetryConfig()

# Custom: override defaults
config = TelemetryConfig(
    enabled=True,
    level=TelemetryLevel.DETAILED,
    otlp_endpoint="phoenix.internal:4317",
    max_cached_tenants=200
)


get_project_name(tenant_id: str, service: Optional[str] = None) -> str

Generate project name for a tenant.

Parameters:

  • tenant_id: Tenant identifier

  • service: Optional service name. If omitted, returns the tenant-only project name (tenant_project_template) — there is no fallback to self.service_name.

Returns: Project name formatted with the matching template

Example:

config = TelemetryConfig()
project = config.get_project_name("acme-corp", "routing")
# Returns: "cogniverse-acme-corp-routing"

project_default = config.get_project_name("acme-corp")
# Returns: "cogniverse-acme-corp" (tenant_project_template, no service suffix)


is_tenant_project(name: str, tenant_id: str) -> bool

Whether project name is the tenant's own project or one of its service projects, by tenant_project_template and tenant_service_template. tenant_id is canonical. A tenant whose id begins another's (acme:prod and acme:prod2) never claims the other's projects.

config = TelemetryConfig()
config.is_tenant_project("cogniverse-acme:prod-routing", "acme:prod")   # True
config.is_tenant_project("cogniverse-acme:prod2", "acme:prod")          # False

should_instrument_component(component: str) -> bool

Check if a component should be instrumented based on the configured level.

Parameters:

  • component: Component name ("search_service", "agents", "backend", "pipeline", "encoder")

Returns: True if component should be instrumented

Level-Based Components:

  • DISABLED: No components

  • BASIC: search_service only

  • DETAILED (default): search_service, agents, backend, pipeline

  • VERBOSE: All components (adds encoder)

Example:

config = TelemetryConfig(level=TelemetryLevel.DETAILED)

config.should_instrument_component("search_service")  # True
config.should_instrument_component("backend")         # True
config.should_instrument_component("encoder")         # False (VERBOSE only)

Span Name Constants (config.py): standardized span/service names used across agents and the runtime, so callers don't hardcode string literals:

SPAN_NAME_REQUEST = "cogniverse.request"
SPAN_NAME_ROUTING = "cogniverse.routing"
SPAN_NAME_ORCHESTRATION = "cogniverse.orchestration"
SPAN_NAME_QUERY_ENHANCEMENT = "cogniverse.query_enhancement"
SPAN_NAME_GATEWAY = "cogniverse.gateway"
SPAN_NAME_PROFILE_SELECTION = "cogniverse.profile_selection"
SPAN_NAME_ENTITY_EXTRACTION = "cogniverse.entity_extraction"

SERVICE_NAME_ORCHESTRATION = "cogniverse.orchestration"

SPAN_NAME_PROFILE_SELECTION is the constant the runtime's profile-selection metrics route (routers/telemetry_metrics.py) filters on server-side via provider.traces.get_spans(..., filters={"name": SPAN_NAME_PROFILE_SELECTION}).


3. Telemetry Registry

File: libs/foundation/cogniverse_foundation/telemetry/registry.py

Purpose: Entry-point auto-discovery and tenant-scoped caching for TelemetryProvider implementations. A thin subclass of cogniverse_foundation.registry.EntryPointRegistry — the base class handles discovery, manual registration, conflict detection, tenant-scoped caching, and lifecycle-style initialization (klass() + .initialize(config)).

class TelemetryRegistry(EntryPointRegistry[TelemetryProvider]):
    _entry_point_group = "cogniverse.telemetry.providers"
    _label = "telemetry provider"
    _tenant_scoped = True

    @classmethod
    def _cache_key(cls, name, config, tenant_id):
        """Key telemetry providers per (tenant, project), not just tenant."""
        base = super()._cache_key(name, config, tenant_id)
        project = (config or {}).get("project_name")
        return f"{base}_{project}" if project else base

Why the custom cache key: a tenant can register distinct endpoints per project via manager.register_project(). The base EntryPointRegistry caches providers per tenant only, so without this override a second project for the same tenant would silently reuse the first project's cached provider (and its endpoints).

Usage:

from cogniverse_foundation.telemetry.registry import get_telemetry_registry

registry = get_telemetry_registry()
provider = registry.get(
    name=None,  # None = auto-detect first discovered provider (e.g., "phoenix")
    tenant_id="acme-corp",
    config={"tenant_id": "acme-corp", "http_endpoint": "...", "grpc_endpoint": "..."},
)

TelemetryManager.get_provider() is the primary caller — application code normally goes through the manager rather than the registry directly.


4. Telemetry Context Helpers

File: libs/foundation/cogniverse_foundation/telemetry/context.py

Purpose: Pre-built span creators matching old instrumentation format with OpenInference semantic conventions.

Main Functions:

search_span(tenant_id: str, query: str, top_k: int = 10, ranking_strategy: str = "default", profile: str = "unknown", backend: str = "vespa") -> ContextManager

Create search service span.

Attributes Set:

  • openinference.span.kind: "CHAIN"

  • operation.name: "search"

  • backend, query, strategy, top_k, profile

  • result_granularity: source or segment

  • input.value: JSON-encoded query parameters

  • latency_ms: Search latency (automatic)

Example:

from cogniverse_foundation.telemetry.context import search_span

with search_span(
    tenant_id="acme-corp",
    query="Marie Curie radioactivity",
    top_k=10,
    ranking_strategy="HYBRID_FLOAT_BM25"
) as span:
    results = vespa_client.search(query)
    span.set_attribute("num_results", len(results))


encode_span(tenant_id: str, encoder_type: str, query_length: int = 0, query: str = "") -> ContextManager

Create encoder span.

Attributes Set:

  • openinference.span.kind: "EMBEDDING"

  • operation.name: f"encode.{encoder_type.lower()}"

  • encoder_type, query_length

  • input.value: Query text

  • encoding_time_ms: Encoding time (automatic)

Example:

from cogniverse_foundation.telemetry.context import encode_span

with encode_span(
    tenant_id="acme-corp",
    encoder_type="ColPali",
    query="test query"
) as span:
    embeddings = encoder.encode(query)
    span.set_attribute("embedding_dim", embeddings.shape[-1])


backend_search_span(tenant_id: str, backend_type: str = "vespa", schema_name: str = "unknown", ranking_strategy: str = "default", top_k: int = 10, has_embeddings: bool = False, query_text: str = "") -> ContextManager

Create backend search span.

Attributes Set:

  • openinference.span.kind: "RETRIEVER"

  • operation.name: "search.execute"

  • backend, query, strategy, top_k, schema, has_embeddings

  • result_granularity: source or segment

  • input.value: JSON-encoded search parameters

  • latency_ms: Search latency (automatic)

Example:

from cogniverse_foundation.telemetry.context import backend_search_span

with backend_search_span(
    tenant_id="acme-corp",
    schema_name="video_colpali_mv",
    ranking_strategy="BINARY_BINARY",
    top_k=10,
    has_embeddings=True
) as span:
    results = vespa_client.execute_query(yql, embeddings)
    span.set_attribute("num_results", len(results))


add_search_results_to_span(span, results, output_value=None)

Add search results details to span.

Attributes Set:

  • num_results: Number of results

  • output.value: Canonical JSON list of result rows (the shape every search consumer reads): document_id, video_id, source_id, source_title (None when the hit stores no title), id, score, content

  • top_score: Score of top result

  • num_collapsed_documents: Number of documents collapsed into the returned source-level results when result_granularity="source"

  • source_search_incomplete: True when a source-granularity search returned fewer than top_k sources from a saturated nearest-neighbor candidate budget. SearchService.search sets it with num_collapsed_documents on the search and backend spans, and POST /search sets it on api.search.request

Events Added:

  • search_results: Top 3 results with rank, document_id, video_id, score, content_type

Example:

from cogniverse_foundation.telemetry.context import search_span, add_search_results_to_span

with search_span(tenant_id="acme", query="test") as span:
    results = search(query)
    add_search_results_to_span(span, results)

serialize_search_results(results) -> str

Serialize the result rows to the canonical output.value JSON once. A search records the same result set on both its RETRIEVER (backend) and CHAIN spans; serialize once with this helper and pass the string as output_value to both add_search_results_to_span calls so the O(N) row-build runs a single time per query rather than per span.

from cogniverse_foundation.telemetry.context import (
    add_search_results_to_span,
    serialize_search_results,
)

output_value = serialize_search_results(results)
add_search_results_to_span(backend_span, results, output_value=output_value)
add_search_results_to_span(search_span, results, output_value=output_value)

Usage Examples

Example 1: Basic Multi-Tenant Telemetry Setup

"""
Initialize telemetry manager and use it across application.
"""
from cogniverse_foundation.telemetry.manager import TelemetryManager
from cogniverse_foundation.telemetry.config import TelemetryConfig, TelemetryLevel

# Initialize once at application startup
config = TelemetryConfig(
    enabled=True,
    level=TelemetryLevel.DETAILED,
    otlp_endpoint="localhost:4317",
    max_cached_tenants=100
)
telemetry = TelemetryManager(config)

# Use in request handlers
def handle_search_request(tenant_id: str, query: str):
    """Process search request with telemetry."""
    with telemetry.span(
        "search_service.search",
        tenant_id=tenant_id,
        project_name="video-search",
        component="search_service",
    ) as span:
        span.set_attribute("query", query)
        span.set_attribute("user_agent", "mobile-app")

        # Business logic
        results = perform_search(query)

        # Add result metrics
        span.set_attribute("num_results", len(results))
        span.set_attribute("top_score", results[0].score if results else 0)

        return results

# Different tenants automatically isolated. project_name="video-search" combined
# with tenant_service_template ("cogniverse-{tenant_id}-{service}") produces:
results_acme = handle_search_request("acme-corp", "test query")  # → cogniverse-acme-corp-video-search
results_tech = handle_search_request("techstart", "test query")   # → cogniverse-techstart-video-search
# Omitting project_name entirely resolves to the tenant-only project
# ("cogniverse-{tenant_id}", via tenant_project_template) instead.

Example 2: Nested Spans with Context Propagation

"""
Create nested span hierarchy for complex operations.
"""
from cogniverse_foundation.telemetry.manager import TelemetryManager

telemetry = TelemetryManager()

def process_video_search(tenant_id: str, query: str):
    """Multi-step video search with nested spans."""

    # Parent span for entire search operation
    with telemetry.span(
        "video_search.process",
        tenant_id=tenant_id,
        attributes={"query": query}
    ) as parent_span:

        # Child span: Query encoding
        with telemetry.span(
            "video_search.encode",
            tenant_id=tenant_id,
            attributes={"encoder": "ColPali"}
        ) as encode_span:
            embeddings = encode_query(query)
            encode_span.set_attribute("embedding_dim", embeddings.shape[-1])

        # Child span: Vespa search
        with telemetry.span(
            "video_search.vespa_search",
            tenant_id=tenant_id,
            attributes={
                "schema": "video_colpali_mv",
                "ranking": "BINARY_BINARY"
            }
        ) as search_span:
            results = vespa_client.search(embeddings)
            search_span.set_attribute("num_results", len(results))

        # Child span: Post-processing
        with telemetry.span(
            "video_search.postprocess",
            tenant_id=tenant_id
        ) as postprocess_span:
            filtered_results = filter_results(results)
            postprocess_span.set_attribute("filtered_count", len(filtered_results))

        parent_span.set_attribute("total_results", len(filtered_results))
        return filtered_results

Resulting Phoenix trace hierarchy:

video_search.process (parent)
├── video_search.encode
├── video_search.vespa_search
└── video_search.postprocess


Example 3: Using Context Helpers for Standard Operations

"""
Use pre-built context helpers for common operations.
"""
from cogniverse_foundation.telemetry.context import search_span, encode_span, backend_search_span

def search_videos_with_telemetry(tenant_id: str, query: str):
    """Search with standardized telemetry spans."""

    # High-level search span
    with search_span(
        tenant_id=tenant_id,
        query=query,
        top_k=10,
        ranking_strategy="HYBRID_FLOAT_BM25",
        profile="video_colpali_mv"
    ) as span:

        # Encoding span
        with encode_span(
            tenant_id=tenant_id,
            encoder_type="ColPali",
            query=query
        ) as enc_span:
            embeddings = colpali_encoder.encode(query)

        # Backend search span
        with backend_search_span(
            tenant_id=tenant_id,
            schema_name="video_colpali_mv_frame",
            ranking_strategy="BINARY_BINARY",
            top_k=10,
            has_embeddings=True,
            query_text=query
        ) as backend_span:
            results = vespa_client.search(query, embeddings)

        # Add results to parent span
        from cogniverse_foundation.telemetry.context import add_search_results_to_span
        add_search_results_to_span(span, results)

        return results

Generated Spans:

  • search_service.search (openinference.span.kind=CHAIN)

  • encoder.colpali.encode (openinference.span.kind=EMBEDDING)

  • search.execute (openinference.span.kind=RETRIEVER)


Example 4: Per-Modality Observability in the Web Client

Per-modality runtime metrics (P50/P95/P99 latency, success rate, request count) are available in the web client's Profile metrics view (GET /admin/tenant/{tenant}/telemetry/profile-selection).

The runtime queries cogniverse.profile_selection spans from Phoenix for the selected tenant and aggregates them by the profile_selection.modality attribute that ProfileSelectionAgent emits on every dispatch. Each span starts when the agent began selecting (before it reads the candidate profiles and calls the LM) and ends once the answer is built, so its duration is the selection latency. No additional code is needed in the application — drive traffic through the routing agent and the view reflects real-time modality breakdown.

To view, open the web client (http://localhost:28400 under cogniverse up), choose Profile metrics under Operations, and pick the tenant and window.


Example 5: Production Configuration with Batch Export

"""
Production telemetry setup with optimized batch export.
"""
from cogniverse_foundation.telemetry.manager import TelemetryManager
from cogniverse_foundation.telemetry.config import TelemetryConfig, TelemetryLevel, BatchExportConfig

# Production batch export configuration
batch_config = BatchExportConfig(
    max_queue_size=4096,           # Large queue for high throughput
    max_export_batch_size=1024,    # Large batches for efficiency
    export_timeout_millis=60_000,  # 60s timeout
    schedule_delay_millis=1000,    # Export every 1s
)

# Production telemetry configuration
production_config = TelemetryConfig(
    enabled=True,
    level=TelemetryLevel.DETAILED,
    environment="production",
    otlp_enabled=True,
    otlp_endpoint="phoenix.internal:4317",
    otlp_use_tls=True,
    max_cached_tenants=500,        # Support many tenants
    batch_config=batch_config,
    service_name="video-search",
    service_version="2.1.0"
)

# Initialize manager
telemetry = TelemetryManager(production_config)

# Use in high-throughput application
def handle_requests():
    """Process many requests with efficient batching."""
    for request in incoming_requests:
        tenant_id = extract_tenant_id(request)

        with telemetry.span("process_request", tenant_id=tenant_id) as span:
            span.set_attribute("request_id", request.id)
            process_request(request)

    # Flush before shutdown
    telemetry.force_flush(timeout_millis=10000)

# Monitor cache performance
stats = telemetry.get_stats()
cache_hit_rate = stats['cache_hits'] / (stats['cache_hits'] + stats['cache_misses'])
print(f"Tracer cache hit rate: {cache_hit_rate:.2%}")
print(f"Cached tenants: {stats['cached_tenants']}")

Production Checklist:

  • ✅ Large queue size (4096+) for high throughput

  • ✅ Batch exports (1024 spans/batch) for efficiency

  • ✅ Large queue size prevents back-pressure (OTel BatchSpanProcessor drops on full)

  • ✅ TLS enabled for production Phoenix endpoint

  • ✅ Large tenant cache (500+) for multi-tenant apps

  • ✅ Force flush before shutdown


Example 6: Test Mode with Synchronous Export

"""
Testing configuration with immediate span export.
"""
from cogniverse_foundation.telemetry.manager import TelemetryManager
from cogniverse_foundation.telemetry.config import TelemetryConfig, TelemetryLevel, BatchExportConfig

# Test configuration — use_sync_export must be set explicitly on
# BatchExportConfig; there is no environment-variable equivalent.
test_config = TelemetryConfig(
    enabled=True,
    level=TelemetryLevel.VERBOSE,
    environment="test",
    otlp_enabled=True,
    otlp_endpoint="localhost:4317",
    otlp_use_tls=False,
    batch_config=BatchExportConfig(
        use_sync_export=True  # Synchronous export
    )
)

telemetry = TelemetryManager(test_config)

def test_search_telemetry():
    """Test that spans are immediately exported."""
    tenant_id = "test-tenant"

    # Create span - will be exported immediately
    with telemetry.span("test.search", tenant_id=tenant_id) as span:
        span.set_attribute("query", "test")

    # Force flush to ensure export (returns immediately in sync mode)
    success = telemetry.force_flush(timeout_millis=5000)
    assert success, "Flush should succeed in sync mode"

    # Verify span in Phoenix (query Phoenix API)
    # ... verification logic ...

# Sync export ensures spans available immediately for assertions
test_search_telemetry()

Test Mode Benefits:

  • Immediate export (no batching delay)

  • Predictable span timing for assertions

  • Simpler debugging (spans appear immediately in Phoenix)


Production Considerations

1. Multi-Tenant Isolation

Complete Data Separation:

# Each tenant gets isolated Phoenix project
tenant_a_project = "cogniverse-acme-corp-routing"
tenant_b_project = "cogniverse-techstart-routing"

# Benefits:
# ✓ Independent analytics per tenant
# ✓ Separate retention policies
# ✓ Easy tenant data deletion (GDPR compliance)
# ✓ No cross-tenant data leakage

Tenant Context Propagation:

  • Tenant ID automatically added to all spans

  • Service name identifies span origin

  • Environment tag for development/staging/production separation

Best Practices:

  • Use consistent tenant ID format across application

  • Include tenant ID in all telemetry spans

  • Monitor per-tenant span volume to detect anomalies


2. Performance and Scalability

LRU Cache Configuration:

# Adjust cache size based on tenant count
config = TelemetryConfig(
    max_cached_tenants=500  # Increase for many tenants
)

# Monitor cache performance
stats = telemetry.get_stats()
hit_rate = stats['cache_hits'] / (stats['cache_hits'] + stats['cache_misses'])

# Target: >95% cache hit rate
if hit_rate < 0.95:
    logger.warning(f"Low cache hit rate: {hit_rate:.2%}")

Batch Export Tuning:

# High-throughput configuration
batch_config = BatchExportConfig(
    max_queue_size=8192,          # Larger queue absorbs burst traffic
    max_export_batch_size=2048,   # Larger batches
    schedule_delay_millis=2000,   # Less frequent exports
)

# Low-latency configuration
batch_config = BatchExportConfig(
    max_queue_size=1024,
    max_export_batch_size=256,
    schedule_delay_millis=100,    # Frequent exports
)

Queue-full drop behaviour is handled natively by OTel's BatchSpanProcessor (created by phoenix.otel.register(batch=True) inside PhoenixProvider). Tune max_queue_size to control the point at which spans start being dropped.

Throughput Targets (illustrative — no benchmark suite ships in this repo; validate against your own deployment before treating these as SLAs):

  • Single tenant: 10,000+ spans/second

  • Multi-tenant (100 tenants): 5,000+ spans/second per tenant

  • LRU cache hit rate: >95% target once tenant working set is warm


3. Graceful Degradation

No-Op Spans:

# If telemetry fails, application continues with no-op spans
with telemetry.span("search", tenant_id="acme") as span:
    # Span might be NoOpSpan if telemetry disabled/failed
    span.set_attribute("query", "test")  # Safe to call
    results = search(query)  # Business logic always executes

Error Handling:

# Manager catches and logs telemetry errors
try:
    tracer = manager.get_tracer("tenant-123")
except Exception as e:
    logger.warning(f"Failed to create tracer: {e}")
    tracer = None  # Returns None instead of crashing

# Application handles None gracefully
if tracer:
    with tracer.start_as_current_span("operation"):
        do_work()
else:
    # No telemetry, but work continues
    do_work()

Configuration Validation:

config = TelemetryConfig()
try:
    config.validate()
except ValueError as e:
    logger.error(f"Invalid telemetry config: {e}")
    # Disable telemetry gracefully
    config.enabled = False


4. Monitoring and Alerting

Key Metrics to Monitor:

  1. Cache Performance:

    stats = telemetry.get_stats()
    cache_hit_rate = stats['cache_hits'] / (stats['cache_hits'] + stats['cache_misses'])
    
    # Alert if hit rate < 95%
    if cache_hit_rate < 0.95:
        alert("Low telemetry cache hit rate", severity="warning")
    

  2. Failed Initializations:

    if stats['failed_initializations'] > 0:
        alert("Telemetry provider initialization failures", severity="error")
    

  3. Per-Modality Performance: Monitor via the web client's "Profile metrics" view, which aggregates P95 latency and success rate per modality from cogniverse.profile_selection spans. Set Phoenix alerts for cogniverse.profile_selection span P95 duration or error rate thresholds directly in Phoenix.

  4. Export Queue Health:

    # Monitor queue size via Phoenix metrics
    # Alert if queue consistently full (indicates export bottleneck)
    

Recommended Alerts:

  • Cache hit rate < 95%

  • Failed initializations > 0

  • Per-modality P95 latency > 1000ms

  • Per-modality error rate > 10%

  • Queue drops > 100/minute


5. Security and Compliance

Tenant Data Isolation:

  • Each tenant's spans in separate Phoenix project

  • No cross-tenant span visibility

  • Per-tenant retention policies

PII Handling:

# Avoid logging sensitive data in spans
with telemetry.span("search", tenant_id=tenant_id) as span:
    # ✅ Safe: Metadata only
    span.set_attribute("query_length", len(query))
    span.set_attribute("has_filters", bool(filters))

    # ❌ Unsafe: PII in span
    # span.set_attribute("user_email", user.email)
    # span.set_attribute("full_query", query)  # May contain PII

Data Retention:

# extra_resource_attributes only tags spans with an informational label —
# it does NOT configure Phoenix's actual retention/GC policy, which is a
# separate Phoenix server-side setting. Use it to make the intended
# retention tier queryable/auditable alongside spans, not to enforce it.
config = TelemetryConfig(
    extra_resource_attributes={
        "retention_tier": "production" if environment == "production" else "development"
    }
)

GDPR Compliance:

  • Easy tenant data deletion (delete Phoenix project)

  • Audit log of span exports (Phoenix provides this)

  • Minimal PII in spans (query length, not content)


6. Testing Strategies

Unit Testing with Mocked Telemetry:

from unittest.mock import Mock, patch

def test_search_with_telemetry():
    """Test search with mocked telemetry."""
    mock_manager = Mock(spec=TelemetryManager)
    mock_span = Mock()
    mock_manager.span.return_value.__enter__.return_value = mock_span

    with patch('cogniverse_foundation.telemetry.get_telemetry_manager', return_value=mock_manager):
        results = search_with_telemetry("test query")

    # Verify span created
    mock_manager.span.assert_called_once_with(
        "search",
        tenant_id="test-tenant",
        attributes=ANY
    )

    # Verify attributes set
    mock_span.set_attribute.assert_any_call("query", "test query")

Integration Testing with Real Phoenix:

def test_telemetry_integration():
    """Test telemetry with real Phoenix instance."""
    # Enable sync export for tests via BatchExportConfig — TelemetryManager()
    # with no config only works if the singleton was already configured
    # elsewhere (e.g. by a fixture); otherwise pass a TelemetryConfig here.
    config = TelemetryConfig(batch_config=BatchExportConfig(use_sync_export=True))
    telemetry = TelemetryManager(config)
    tenant_id = "test-tenant"

    # Create span
    with telemetry.span("test.search", tenant_id=tenant_id) as span:
        span.set_attribute("query", "test")

    # Force flush
    success = telemetry.force_flush(timeout_millis=5000)
    assert success

    # Query Phoenix to verify span
    # ... Phoenix API query ...

Performance Testing:

import time

def test_telemetry_throughput():
    """Test telemetry can handle high throughput."""
    telemetry = TelemetryManager()

    start = time.time()
    span_count = 10000

    for i in range(span_count):
        with telemetry.span(f"test.span.{i}", tenant_id="perf-test") as span:
            span.set_attribute("iteration", i)

    duration = time.time() - start
    throughput = span_count / duration

    print(f"Throughput: {throughput:.0f} spans/second")
    assert throughput > 5000, "Should handle 5k+ spans/second"


Testing

Key Test Files

Unit Tests (tests/telemetry/unit/):

  • test_tracer_cache_eviction.py — LRU count cap, orphaned-provider cleanup, and TTL-based expiry
  • test_telemetry_level_filter.py — should_instrument_component() filter logic and span()'s NoOpSpan short-circuit when the level doesn't admit a component
  • test_session_tracking.py — session_context() on TelemetryProvider/PhoenixProvider, session_span(), and session ID propagation to nested spans
  • test_provider_project_cache.py — TelemetryRegistry caches providers per (tenant, project), not just tenant
  • test_span_export_config.py — BatchExportConfig knobs and resource attributes actually reach the live TracerProvider/exporter, not just round-trip through serialization
  • test_analytics_timestamp_tz.py — PhoenixAnalytics.get_traces produces timezone-aware timestamps (including UTC-normalizing naive start_time/end_time values)
  • test_phoenix_circuit_breaker.py — a down Phoenix trips the shared circuit breaker: PhoenixAnalytics.get_traces and PhoenixTraceStore.get_spans both RAISE (the transport error while the breaker is closed, CircuitOpenError once it trips open) rather than degrading to [] — an empty list would read as "no traces in range", indistinguishable from genuine empty data — both without repeatedly dialing a dead Phoenix

Integration Tests:

  • tests/telemetry/integration/test_multi_tenant_telemetry.py
  • Multi-tenant tracer provider isolation
  • LRU cache behavior
  • Batch vs sync export modes
  • Force flush functionality
  • Graceful degradation
  • tests/telemetry/integration/test_analytics_root_spans_real_phoenix.py — PhoenixAnalytics.get_traces passes root_spans_only=True server-side so limit bounds traces (root spans), not raw spans; a trace's root would otherwise be crowded out of a small limit slice by its own newer children
  • tests/telemetry/integration/test_dataset_store_event_loop.py — PhoenixDatasetStore.create_dataset/get_dataset/append_to_dataset run their synchronous Phoenix HTTP calls via asyncio.to_thread, off the event-loop thread
  • tests/telemetry/integration/test_session_tracking_real.py — session_span() against a real managed Phoenix container: real (non-NoOpSpan) spans carry session.id, and multiple requests sharing a session_id each propagate it into OTel context without leaking across requests

Test Scenarios:

  1. Tenant Isolation:

    def test_tenant_isolation():
        """Verify tenants get separate Phoenix projects."""
        manager = TelemetryManager()
    
        # Create tracers for different tenants
        tracer_a = manager.get_tracer("tenant-a")
        tracer_b = manager.get_tracer("tenant-b")
    
        # Verify different providers
        assert manager._tenant_providers["tenant-a"] != manager._tenant_providers["tenant-b"]
    

  2. LRU Cache Eviction:

    def test_lru_eviction():
        """Verify old tracers evicted when cache full."""
        config = TelemetryConfig(max_cached_tenants=2)
        manager = TelemetryManager(config)
    
        # Create 3 tracers (exceeds cache limit)
        manager.get_tracer("tenant-1")
        manager.get_tracer("tenant-2")
        manager.get_tracer("tenant-3")
    
        # Oldest (tenant-1) should be evicted
        assert len(manager._tenant_tracers) == 2
    

  3. Sync Export:

    def test_sync_export():
        """Verify synchronous export in test mode."""
        config = TelemetryConfig(batch_config=BatchExportConfig(use_sync_export=True))
        manager = TelemetryManager(config)
    
        with manager.span("test", tenant_id="test") as span:
            span.set_attribute("key", "value")
    
        # Span should be exported immediately
        success = manager.force_flush(timeout_millis=1000)
        assert success
    

manager.span(name, tenant_id=..., start_time=time.time_ns()) starts the span at a moment already past (nanoseconds since the epoch) and ends it when the context exits, so a span written after the work it records still lasts as long as that work. The gateway records each routing decision this way, from when classification began.

Normal spans remain non-raising when the exporter is unavailable. Synchronous worker code emits durable records with manager.span(..., require_export=True). Checked exporters use always-on sampling, reject batch mode, suppress recursive exporter instrumentation, and raise with project and endpoint context if the backend rejects the record. Each required span owns an isolated exporter lease, so a timeout or cancellation cannot shut down another concurrent durable write or poison a shared cache.

Async code must use the bounded async context manager so exporter I/O runs off the event loop. Lease construction is a local OpenTelemetry object setup and does not contact the collector; the first collector operation is the bounded span export:

async with manager.required_span(
    "adapter.publication",
    tenant_id=tenant_id,
    project_name="experiments",
    attributes={"run.id": run_id},
) as span:
    span.set_attribute("publication.state", "committed")

The bound is batch_config.export_timeout_millis. The exporter transport uses a shorter deadline so a rejected or unreachable collector raises with tenant, project, endpoint, and the original exporter failure while the caller's budget is still active. If the transport itself remains hung, the manager shuts the provider down, waits for the exporter thread to finish, and then raises TimeoutError with the same context. Calling the synchronous required form while an event loop is active raises immediately instead of blocking that loop. Task cancellation also shuts down and joins that span's exporter before the original cancellation propagates.

  1. Graceful Degradation:
    def test_graceful_degradation():
        """Verify app continues when telemetry fails."""
        # Disable OTLP span export
        config = TelemetryConfig(otlp_enabled=False)
        manager = TelemetryManager(config)
    
        # Should return no-op span
        with manager.span("test", tenant_id="test") as span:
            # Should not crash
            span.set_attribute("key", "value")
    

Test Coverage:

  • Multi-tenant isolation: ✅

  • LRU cache behavior: ✅

  • Batch export configuration: ✅

  • Sync export mode: ✅

  • Graceful degradation: ✅

  • Performance metrics tracking: ✅

  • Component-level telemetry filtering: ✅

  • Tracer cache TTL expiry: ✅

  • Provider caching per (tenant, project): ✅

  • Session ID propagation: ✅

  • Circuit breaker degradation on Phoenix outage: ✅

  • Root-span filtering for trace metrics: ✅

  • Dataset store event-loop offload: ✅


Summary

The Telemetry Module provides production-ready, multi-tenant observability with:

Core Features:

  • ✅ Multi-tenant isolation via Phoenix projects

  • ✅ Lazy initialization with LRU caching

  • ✅ Configurable batch export (async) or sync export (tests)

  • ✅ Per-modality performance tracking

  • ✅ Graceful degradation

  • ✅ OpenTelemetry standards compliance

Production Strengths:

  • Designed for 10,000+ spans/second per tenant (illustrative target, not a measured benchmark in this repo)

  • Complete tenant data isolation

  • Minimal performance overhead

  • Robust error handling

Integration Points:

  • All agents use telemetry.span() context manager

  • Search operations use pre-built context helpers

  • Per-modality observability flows through cogniverse.profile_selection spans and the web client's Profile metrics view

  • Phoenix provides analytics and visualization


For detailed examples and production configurations, see:

  • Architecture Overview: docs/architecture/overview.md

  • Agents Module: docs/modules/agents.md

  • Common Module (Config): docs/modules/common.md

Source Files:

  • Manager: libs/foundation/cogniverse_foundation/telemetry/manager.py

  • Config: libs/foundation/cogniverse_foundation/telemetry/config.py

  • Context: libs/foundation/cogniverse_foundation/telemetry/context.py