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 providersTelemetryConfig: Configuration for telemetry systems with BatchExportConfigTelemetryManager: Singleton manager for multi-tenant tracer providers- Context helpers for common operations (search, encode, backend)
span_metrics: aggregates over aTraceStorespan 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 theABCompareAggregateofrlm.ab_comparespans (AB_COMPARE_SPAN_NAME), with averages,per_datasetandper_row(newest first) frames.span_succeeded(status)is false only forERROR, 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_SCORESandpersist_result_relevance(provider, project, span_id, result_id, label, readable_within_s=), which checks the span is inprojectand stores the rating under the result's own identifier. A span no project holds yet is looked up again with backoff forreadable_within_s—span_readable_within_s(batch_config), the exporter'sschedule_delay_millisplusSPAN_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, raisesSpanNotInProjectError.persist_session_evaluation(provider, project, session_id, span_ids, outcome, score, readable_within_s=)stores a conversation's verdict (SESSION_OUTCOMES) as aSESSION_EVALUATIONannotation on each of its spans, after finding every one inprojectthe same way, keyed by the session so a new verdict replaces it.record_span_io(span, input_value=, output=, operation=, modality=)writes the input oninput.value, the output as JSON onoutput.value, and the type onoperation;read_span_io(row)reads{input, output, operation, modality}back andread_span_attributes(row)returns every attribute as a flat dotted-key dict from a Phoenix span row. Search (list output) and domain spans likequery_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), andrecord_search_io_on_current_span(query, results, modality)records a search on the active span (onesearch_result_rowper result); the search and document agents stamp both on their process span, whose id a client annotates. Query-enhancement spans also carryenhancement.pathwithlmfor a genuine enhancement andheuristic_fallbackfor 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) thatTripletExtractor,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-foundationdoesn'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¶
- Module Overview
- Architecture Diagrams
- Provider Abstraction
- Phoenix Implementation
- Core Components
- Usage Examples
- Production Considerations
- 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_selectionspans; surfaced in the web client's "Profile metrics" view - OpenTelemetry Integration: Standards-based distributed tracing with Phoenix backend
Key Features¶
- Singleton Manager Pattern
- Thread-safe singleton TelemetryManager
- Global instance accessible throughout application
-
Configuration from environment variables
-
Tenant-Aware Tracing
- Project-based separation:
cogniverse-{tenant_id}-{service} - LRU cache for tenant providers (configurable max tenants)
-
Automatic tenant context propagation
-
Flexible Export Modes
- Production: Batch export with queue management (async)
- Testing: Synchronous export for immediate flush
-
Configurable queue size, batch size, timeout
-
Per-Modality Observability
cogniverse.profile_selectionspans emitted byProfileSelectionAgentcarry aprofile_selection.modalityattribute- The runtime's
GET /admin/tenant/{tenant_id}/telemetry/profile-selectionroute (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 -
No separate metrics class required — all observability flows through standard OTel spans
-
Context Helpers
- Pre-built span creators for search, encoding, backend operations
- Exact format matching with old instrumentation system
- OpenInference semantic conventions
Dependencies¶
Foundation Package (cogniverse-foundation):
opentelemetry-api: OpenTelemetry APIopentelemetry-sdk: OpenTelemetry Python SDKpydantic: Data validationsqlalchemy: Database ORMpandas: Data structures
Phoenix Plugin (cogniverse-telemetry-phoenix):
arize-phoenix-client: Phoenix clientarize-phoenix-otel: Phoenix OpenTelemetry setuppandas: 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:
-
Production (Batch): Async export with queue, scheduled at intervals
-
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:
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; raisesValueErrorif None). Useget_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:
-
Tracerinstance if telemetry enabled -
Noneif telemetry disabled or initialization failed
Caching Logic:
-
Check LRU cache for
{tenant_id}:{project_name}key -
If cached, return tracer (cache hit)
-
If not cached:
- Create TracerProvider for tenant if needed
- Get tracer from provider
- Cache tracer with LRU eviction
- 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 ofsearch_service/agents/backend/pipeline/encoder.TelemetryConfig.should_instrument_component()checks this againstconfig.level; if the level doesn't admit the component,span()yields aNoOpSpanwithout ever creating a tracer. Defaultagentsemits atDETAILEDand above. Passencoderfor per-inference model detail spans,pipelinefor ingestion-stage spans,search_servicefor the top-level search HTTP entry (admitted even atBASIC).
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.idare 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 exporthttp_endpoint: Override HTTP endpoint used by the provider for span queriesgrpc_endpoint: Override gRPC endpoint used by the provider (distinct fromotlp_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 toself.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:sourceorsegment -
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:sourceorsegment -
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(Nonewhen 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 whenresult_granularity="source" -
source_search_incomplete:Truewhen a source-granularity search returned fewer thantop_ksources from a saturated nearest-neighbor candidate budget.SearchService.searchsets it withnum_collapsed_documentson the search and backend spans, andPOST /searchsets it onapi.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:
-
Cache Performance:
-
Failed Initializations:
-
Per-Modality Performance: Monitor via the web client's "Profile metrics" view, which aggregates P95 latency and success rate per modality from
cogniverse.profile_selectionspans. Set Phoenix alerts forcogniverse.profile_selectionspan P95 duration or error rate thresholds directly in Phoenix. -
Export Queue Health:
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 expirytest_telemetry_level_filter.py—should_instrument_component()filter logic andspan()'s NoOpSpan short-circuit when the level doesn't admit a componenttest_session_tracking.py—session_context()onTelemetryProvider/PhoenixProvider,session_span(), and session ID propagation to nested spanstest_provider_project_cache.py—TelemetryRegistrycaches providers per (tenant, project), not just tenanttest_span_export_config.py—BatchExportConfigknobs and resource attributes actually reach the liveTracerProvider/exporter, not just round-trip through serializationtest_analytics_timestamp_tz.py—PhoenixAnalytics.get_tracesproduces timezone-aware timestamps (including UTC-normalizing naivestart_time/end_timevalues)test_phoenix_circuit_breaker.py— a down Phoenix trips the shared circuit breaker:PhoenixAnalytics.get_tracesandPhoenixTraceStore.get_spansboth RAISE (the transport error while the breaker is closed,CircuitOpenErroronce 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_tracespassesroot_spans_only=Trueserver-side solimitbounds traces (root spans), not raw spans; a trace's root would otherwise be crowded out of a smalllimitslice by its own newer childrentests/telemetry/integration/test_dataset_store_event_loop.py—PhoenixDatasetStore.create_dataset/get_dataset/append_to_datasetrun their synchronous Phoenix HTTP calls viaasyncio.to_thread, off the event-loop threadtests/telemetry/integration/test_session_tracking_real.py—session_span()against a real managed Phoenix container: real (non-NoOpSpan) spans carrysession.id, and multiple requests sharing asession_ideach propagate it into OTel context without leaking across requests
Test Scenarios:
-
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"] -
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 -
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.
- 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_selectionspans 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