API reference¶
Generated from the apogee-ai-observability source with mkdocstrings. Every symbol below is exported from apogee_ai_observability, so it is part of the supported public surface.
Application · Use cases¶
EmitSpanUseCase
¶
EmitSpanUseCase(emitter: ITelemetryEmitter)
Sends a span to the configured emitter (and optionally a CostTracker).
Kept thin on purpose: orchestration logic lives in callers (e.g., the agent runtime); this just centralises the contract.
Source code in apogee_ai_observability/application/use_cases/emit_span_use_case.py
EndTraceUseCase
¶
EndTraceUseCase(store: ITraceStore)
Mark a trace as finished and persist it to the store.
Source code in apogee_ai_observability/application/use_cases/end_trace_use_case.py
GetTraceUseCase
¶
GetTraceUseCase(store: ITraceStore)
ListTracesUseCase
¶
ListTracesUseCase(store: ITraceStore)
RecordCostUseCase
¶
RecordCostUseCase(calculator: ICostCalculator, tracker: CostTracker)
Estimate cost from usage, persist into the tracker.
Source code in apogee_ai_observability/application/use_cases/record_cost_use_case.py
execute
async
¶
execute(usage: LLMUsage, *, trace_id: str = '', span_id: str | None = None, user_id: str | None = None, tenant_id: str | None = None) -> CostLedgerEntry
Source code in apogee_ai_observability/application/use_cases/record_cost_use_case.py
async def execute(
self,
usage: LLMUsage,
*,
trace_id: str = "",
span_id: str | None = None,
user_id: str | None = None,
tenant_id: str | None = None,
) -> CostLedgerEntry:
entry = self._calculator.estimate(
usage,
trace_id=trace_id,
span_id=span_id,
user_id=user_id,
tenant_id=tenant_id,
)
self._tracker.record(entry)
return entry
ReplayTraceUseCase
¶
ReplayTraceUseCase(store: ITraceStore, emitter: ITelemetryEmitter)
Load a stored trace and re-emit through emitter.
Source code in apogee_ai_observability/application/use_cases/replay_trace_use_case.py
execute
async
¶
StartTraceUseCase
¶
Domain¶
CostLedgerEntry
dataclass
¶
CostLedgerEntry(trace_id: str, span_id: str | None, provider: str, model: str, input_tokens: int, output_tokens: int, input_cost_usd: float, output_cost_usd: float, user_id: str | None = None, tenant_id: str | None = None, timestamp: datetime = (lambda: now(utc))())
Single cost charge attributable to a trace/span/user/tenant.
timestamp
class-attribute
instance-attribute
¶
LLMOpName
¶
Bases: str, Enum
Standardised span names so dashboards can group across emitters.
LLMUsage
dataclass
¶
LLMUsage(provider: str, model: str, prompt_tokens: int = 0, completion_tokens: int = 0, total_tokens: int = 0)
LatencyBucket
dataclass
¶
LatencyBucket(p50_ms: float, p95_ms: float, p99_ms: float, max_ms: float, samples: int, window_seconds: int = 60)
Aggregated latency stats over a window.
MetricSample
dataclass
¶
MetricSample(name: str, value: float, unit: str = '', attributes: dict[str, str] = dict(), timestamp: datetime = (lambda: now(utc))())
Span
dataclass
¶
Span(name: str, kind: SpanKind = INTERNAL, span_id: str = (lambda: token_hex(8))(), trace_id: str = (lambda: token_hex(16))(), parent_span_id: str | None = None, start_time: datetime = (lambda: now(utc))(), end_time: datetime | None = None, status: SpanStatus = UNSET, error: str | None = None, attributes: dict[str, Any] = dict(), events: list[dict[str, Any]] = list(), usage: LLMUsage | None = None)
A single operation in a trace.
Mutable because emitters often build the span before/after work runs
(set end_time, status, usage post-hoc).
span_id
class-attribute
instance-attribute
¶
trace_id
class-attribute
instance-attribute
¶
start_time
class-attribute
instance-attribute
¶
attributes
class-attribute
instance-attribute
¶
events
class-attribute
instance-attribute
¶
end
¶
end(*, status: SpanStatus | None = None, error: str | None = None) -> None
Source code in apogee_ai_observability/domain/entities/span.py
def end(self, *, status: SpanStatus | None = None, error: str | None = None) -> None:
self.end_time = datetime.now(timezone.utc)
if error is not None:
self.error = error
self.status = status if status is not None else SpanStatus.ERROR
else:
self.status = status if status is not None else SpanStatus.OK
add_event
¶
set_attribute
¶
SpanStatus
¶
Trace
dataclass
¶
Trace(trace_id: str, name: str = '', started_at: datetime = (lambda: now(utc))(), finished_at: datetime | None = None, spans: list[Span] = list(), metadata: dict[str, str] = dict())
TraceContext
dataclass
¶
TraceContext(trace_id: str = (lambda: token_hex(16))(), span_id: str = (lambda: token_hex(8))(), parent_span_id: str | None = None, sampled: bool = True, baggage: dict[str, str] = dict())
Carries the active trace + span ids across async boundaries.
Modeled after the W3C trace-context: 16-byte trace id and 8-byte span id encoded as lowercase hex.
semantic
¶
LatencyBucket
dataclass
¶
LatencyBucket(p50_ms: float, p95_ms: float, p99_ms: float, max_ms: float, samples: int, window_seconds: int = 60)
Aggregated latency stats over a window.
TraceContext
dataclass
¶
TraceContext(trace_id: str = (lambda: token_hex(16))(), span_id: str = (lambda: token_hex(8))(), parent_span_id: str | None = None, sampled: bool = True, baggage: dict[str, str] = dict())
Carries the active trace + span ids across async boundaries.
Modeled after the W3C trace-context: 16-byte trace id and 8-byte span id encoded as lowercase hex.
Domain · Enums¶
EmitterKind
¶
Bases: str, Enum
SpanKind
¶
Bases: str, Enum
Domain · Exceptions¶
EmitterFailureException
¶
Bases: ObservabilityError
Source code in apogee_ai_observability/domain/exceptions/observability_exceptions.py
ObservabilityError
¶
Bases: Exception
Base exception for apogee-ai-observability.
PricingNotAvailableException
¶
Bases: ObservabilityError
Source code in apogee_ai_observability/domain/exceptions/observability_exceptions.py
TraceNotFoundException
¶
Bases: ObservabilityError
Source code in apogee_ai_observability/domain/exceptions/observability_exceptions.py
Domain · Protocols (ports)¶
ICostCalculator
¶
Bases: Protocol
estimate
¶
estimate(usage: LLMUsage, *, trace_id: str = '', span_id: str | None = None, user_id: str | None = None, tenant_id: str | None = None) -> CostLedgerEntry
IPricingProvider
¶
Bases: Protocol
Returns (input_per_million, output_per_million) USD or None.
ITelemetryEmitter
¶
Bases: Protocol
Sink for spans and metrics produced by the agent runtime.
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
ITraceStore
¶
Infrastructure¶
CompositeEmitter
¶
CompositeEmitter(children: Iterable[ITelemetryEmitter])
Fan-out wrapper: every event goes to every child emitter.
Failures of one child do not block the others.
Source code in apogee_ai_observability/infrastructure/emitters/builtin_emitters.py
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
ConsoleEmitter
¶
Pretty single-line output, useful in dev.
Source code in apogee_ai_observability/infrastructure/emitters/builtin_emitters.py
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/builtin_emitters.py
async def emit_span(self, span: Span) -> None:
duration = f"{span.duration_ms:.1f}ms" if span.duration_ms is not None else "-"
usage = ""
if span.usage is not None:
usage = f" tokens={span.usage.total_tokens}"
print(
f"[span] {span.kind.value:8s} {span.name:24s} {duration:>10s} "
f"{span.status.value}{usage}",
file=self._stream,
)
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
CostTracker
¶
Append-only ledger with rollup helpers (by trace, user, tenant).
Source code in apogee_ai_observability/infrastructure/cost/cost_tracker.py
record
¶
record(entry: CostLedgerEntry) -> None
extend
¶
extend(entries: Iterable[CostLedgerEntry]) -> None
by_trace
¶
by_user
¶
by_tenant
¶
by_model
¶
DatadogLLMEmitter
¶
Adapter for Datadog LLM Observability via ddtrace.
Lazy import: install via pip install 'apogee-ai-observability[datadog]'.
Source code in apogee_ai_observability/infrastructure/emitters/datadog_emitter.py
def __init__(
self,
*,
service: str = "apogee-ai",
env: str | None = None,
) -> None:
try:
import ddtrace # type: ignore # noqa: F401
except ImportError as exc:
raise ImportError(
"DatadogLLMEmitter requires `ddtrace`. "
"Install with: pip install 'apogee-ai-observability[datadog]'"
) from exc
from ddtrace import tracer # type: ignore
self._tracer = tracer
self._service = service
self._env = env
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/datadog_emitter.py
async def emit_span(self, span: Span) -> None:
try:
with self._tracer.trace(
span.name,
service=self._service,
resource=span.kind.value,
span_type="llm",
) as dd_span:
for key, value in span.attributes.items():
dd_span.set_tag(key, value)
if span.usage is not None:
dd_span.set_metric("llm.prompt_tokens", span.usage.prompt_tokens)
dd_span.set_metric("llm.completion_tokens", span.usage.completion_tokens)
dd_span.set_metric("llm.total_tokens", span.usage.total_tokens)
dd_span.set_tag("llm.model", span.usage.model)
dd_span.set_tag("llm.provider", span.usage.provider)
if self._env:
dd_span.set_tag("env", self._env)
if span.error:
dd_span.set_traceback()
dd_span.set_tag("error", True)
dd_span.set_tag("error.message", span.error)
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
HeliconeEmitter
¶
Posts manual log events to Helicone's custom-events endpoint.
Helicone's primary mode is header-based proxying — this adapter targets
the supplemental /v1/log style endpoint for spans we observed
locally. Lazy import via [helicone] extra (httpx).
Source code in apogee_ai_observability/infrastructure/emitters/helicone_emitter.py
def __init__(
self,
*,
api_key: str,
base_url: str = "https://api.helicone.ai",
) -> None:
try:
import httpx # type: ignore # noqa: F401
except ImportError as exc:
raise ImportError(
"HeliconeEmitter requires `httpx`. "
"Install with: pip install 'apogee-ai-observability[helicone]'"
) from exc
self._api_key = api_key
self._base_url = base_url.rstrip("/")
self._client = None
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/helicone_emitter.py
async def emit_span(self, span: Span) -> None:
try:
payload = {
"trace_id": span.trace_id,
"span_id": span.span_id,
"name": span.name,
"kind": span.kind.value,
"status": span.status.value,
"start": span.start_time.isoformat(),
"end": span.end_time.isoformat() if span.end_time else None,
"duration_ms": span.duration_ms,
"attributes": dict(span.attributes),
"usage": (
{
"model": span.usage.model,
"prompt_tokens": span.usage.prompt_tokens,
"completion_tokens": span.usage.completion_tokens,
}
if span.usage
else None
),
}
response = await self._get_client().post("/v1/log", json=payload)
response.raise_for_status()
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
InMemoryEmitter
¶
InMemoryTraceStore
¶
Source code in apogee_ai_observability/infrastructure/replay/in_memory_trace_store.py
save
async
¶
list
async
¶
Source code in apogee_ai_observability/infrastructure/replay/in_memory_trace_store.py
async def list(
self,
*,
limit: int | None = None,
tenant_id: str | None = None,
) -> list[Trace]:
items = sorted(self._store.values(), key=lambda t: t.started_at, reverse=True)
if tenant_id is not None:
items = [
t for t in items if t.metadata.get("tenant.id") == tenant_id
]
if limit is not None:
items = items[:limit]
return items
JsonTraceStore
¶
LangSmithEmitter
¶
Adapter for LangChain LangSmith.
Lazy import: install via pip install 'apogee-ai-observability[langsmith]'.
Source code in apogee_ai_observability/infrastructure/emitters/langsmith_emitter.py
def __init__(
self,
*,
api_key: str | None = None,
project: str | None = None,
) -> None:
try:
import langsmith # type: ignore # noqa: F401
except ImportError as exc:
raise ImportError(
"LangSmithEmitter requires `langsmith`. "
"Install with: pip install 'apogee-ai-observability[langsmith]'"
) from exc
from langsmith import Client # type: ignore
self._client = Client(api_key=api_key)
self._project = project
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/langsmith_emitter.py
async def emit_span(self, span: Span) -> None:
try:
self._client.create_run(
name=span.name,
run_type=_map_kind(span.kind.value),
inputs=span.attributes,
outputs={"events": span.events} if span.events else None,
start_time=span.start_time,
end_time=span.end_time,
error=span.error,
project_name=self._project,
trace_id=span.trace_id,
parent_run_id=span.parent_span_id,
)
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
LangfuseEmitter
¶
LangfuseEmitter(*, public_key: str | None = None, secret_key: str | None = None, host: str | None = None)
Adapter for Langfuse cloud or self-hosted.
Lazy import: install via pip install 'apogee-ai-observability[langfuse]'.
Source code in apogee_ai_observability/infrastructure/emitters/langfuse_emitter.py
def __init__(
self,
*,
public_key: str | None = None,
secret_key: str | None = None,
host: str | None = None,
) -> None:
try:
import langfuse # type: ignore # noqa: F401
except ImportError as exc:
raise ImportError(
"LangfuseEmitter requires `langfuse`. "
"Install with: pip install 'apogee-ai-observability[langfuse]'"
) from exc
from langfuse import Langfuse # type: ignore
self._client = Langfuse(
public_key=public_key,
secret_key=secret_key,
host=host,
)
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/langfuse_emitter.py
async def emit_span(self, span: Span) -> None:
try:
attributes = dict(span.attributes)
usage = None
if span.usage is not None:
usage = {
"input": span.usage.prompt_tokens,
"output": span.usage.completion_tokens,
"total": span.usage.total_tokens,
"unit": "TOKENS",
}
kwargs = {
"trace_id": span.trace_id,
"id": span.span_id,
"name": span.name,
"start_time": span.start_time,
"end_time": span.end_time,
"metadata": attributes,
"level": "ERROR" if span.error else "DEFAULT",
"status_message": span.error,
}
if usage:
self._client.generation(
**kwargs,
model=span.usage.model if span.usage else None,
usage=usage,
)
else:
self._client.span(**kwargs)
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
Source code in apogee_ai_observability/infrastructure/emitters/langfuse_emitter.py
shutdown
async
¶
NoopEmitter
¶
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
OTelEmitter
¶
OTelEmitter(*, endpoint: str | None = None, service_name: str = 'apogee-ai', insecure: bool = True)
OpenTelemetry adapter — works with any OTLP backend (Jaeger, Tempo, Honeycomb, Datadog APM, Lightstep, etc).
Lazy import: install via pip install 'apogee-ai-observability[otel]'.
Source code in apogee_ai_observability/infrastructure/emitters/otel_emitter.py
def __init__(
self,
*,
endpoint: str | None = None,
service_name: str = "apogee-ai",
insecure: bool = True,
) -> None:
try:
from opentelemetry import metrics, trace # type: ignore
from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import ( # type: ignore
OTLPMetricExporter,
)
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import ( # type: ignore
OTLPSpanExporter,
)
from opentelemetry.sdk.metrics import MeterProvider # type: ignore
from opentelemetry.sdk.metrics.export import ( # type: ignore
PeriodicExportingMetricReader,
)
from opentelemetry.sdk.resources import Resource # type: ignore
from opentelemetry.sdk.trace import TracerProvider # type: ignore
from opentelemetry.sdk.trace.export import ( # type: ignore
BatchSpanProcessor,
)
except ImportError as exc:
raise ImportError(
"OTelEmitter requires opentelemetry-* packages. "
"Install with: pip install 'apogee-ai-observability[otel]'"
) from exc
resource = Resource.create({"service.name": service_name})
tracer_provider = TracerProvider(resource=resource)
span_exporter = OTLPSpanExporter(endpoint=endpoint, insecure=insecure)
tracer_provider.add_span_processor(BatchSpanProcessor(span_exporter))
trace.set_tracer_provider(tracer_provider)
self._tracer = trace.get_tracer("apogee-ai-observability")
metric_exporter = OTLPMetricExporter(endpoint=endpoint, insecure=insecure)
meter_provider = MeterProvider(
resource=resource,
metric_readers=[PeriodicExportingMetricReader(metric_exporter)],
)
metrics.set_meter_provider(meter_provider)
self._meter = metrics.get_meter("apogee-ai-observability")
self._counters: dict[str, object] = {}
self._tracer_provider = tracer_provider
self._meter_provider = meter_provider
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/otel_emitter.py
async def emit_span(self, span: Span) -> None:
try:
from opentelemetry.trace import StatusCode # type: ignore
with self._tracer.start_as_current_span(span.name) as ot_span:
for key, value in span.attributes.items():
ot_span.set_attribute(key, _coerce(value))
if span.usage is not None:
ot_span.set_attribute("llm.prompt_tokens", span.usage.prompt_tokens)
ot_span.set_attribute("llm.completion_tokens", span.usage.completion_tokens)
ot_span.set_attribute("llm.total_tokens", span.usage.total_tokens)
ot_span.set_attribute("llm.model", span.usage.model)
ot_span.set_attribute("llm.provider", span.usage.provider)
if span.error:
ot_span.record_exception(Exception(span.error))
ot_span.set_status(StatusCode.ERROR, span.error)
elif span.status == SpanStatus.OK:
ot_span.set_status(StatusCode.OK)
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
Source code in apogee_ai_observability/infrastructure/emitters/otel_emitter.py
async def emit_metric(self, sample: MetricSample) -> None:
try:
counter = self._counters.get(sample.name)
if counter is None:
counter = self._meter.create_counter(sample.name, unit=sample.unit or "1")
self._counters[sample.name] = counter
counter.add(sample.value, attributes=dict(sample.attributes)) # type: ignore[attr-defined]
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
shutdown
async
¶
PRICING_TABLE
module-attribute
¶
PRICING_TABLE: dict[str, dict[str, tuple[float, float]]] = {'openai': {'gpt-4o': (5.0, 15.0), 'gpt-4o-mini': (0.15, 0.6), 'gpt-4-turbo': (10.0, 30.0), 'gpt-4': (30.0, 60.0), 'gpt-3.5-turbo': (0.5, 1.5), 'o1': (15.0, 60.0), 'o1-mini': (3.0, 12.0), 'o3-mini': (1.1, 4.4), 'text-embedding-3-large': (0.13, 0.0), 'text-embedding-3-small': (0.02, 0.0)}, 'anthropic': {'claude-3-5-sonnet': (3.0, 15.0), 'claude-3-5-haiku': (0.8, 4.0), 'claude-3-opus': (15.0, 75.0), 'claude-3-sonnet': (3.0, 15.0), 'claude-3-haiku': (0.25, 1.25), 'claude-opus-4': (15.0, 75.0), 'claude-sonnet-4': (3.0, 15.0), 'claude-haiku-4': (0.8, 4.0)}, 'google': {'gemini-1.5-pro': (1.25, 5.0), 'gemini-1.5-flash': (0.075, 0.3), 'gemini-2.0-flash': (0.1, 0.4), 'gemini-2.5-pro': (1.25, 5.0), 'gemini-2.5-flash': (0.1, 0.4)}, 'openrouter': {'default': (1.0, 3.0)}, 'bedrock': {'anthropic.claude-3-5-sonnet': (3.0, 15.0), 'anthropic.claude-3-haiku': (0.25, 1.25), 'amazon.titan-text-express': (0.2, 0.6), 'meta.llama3-70b-instruct': (2.65, 3.5)}}
PhoenixEmitter
¶
Adapter for Arize Phoenix (open source LLM observability).
Phoenix consumes OpenTelemetry spans, so this emitter wraps an OTel tracer and tags spans with OpenInference semantic conventions.
Lazy import: install via pip install 'apogee-ai-observability[phoenix]'.
Source code in apogee_ai_observability/infrastructure/emitters/phoenix_emitter.py
def __init__(self, *, endpoint: str | None = None, project: str | None = None) -> None:
try:
import opentelemetry # type: ignore # noqa: F401
except ImportError as exc:
raise ImportError(
"PhoenixEmitter requires opentelemetry-api/sdk and arize-phoenix. "
"Install with: pip install 'apogee-ai-observability[phoenix]'"
) from exc
try:
from phoenix.otel import register # type: ignore
except ImportError as exc: # pragma: no cover
raise ImportError(
"phoenix-otel is required: pip install 'apogee-ai-observability[phoenix]'"
) from exc
self._tracer_provider = register(
project_name=project or "apogee-ai",
endpoint=endpoint,
)
self._tracer = self._tracer_provider.get_tracer("apogee-ai-observability")
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/phoenix_emitter.py
async def emit_span(self, span: Span) -> None:
try:
with self._tracer.start_as_current_span(span.name) as ot_span:
for key, value in span.attributes.items():
ot_span.set_attribute(key, _coerce(value))
if span.usage is not None:
ot_span.set_attribute("llm.token_count.prompt", span.usage.prompt_tokens)
ot_span.set_attribute("llm.token_count.completion", span.usage.completion_tokens)
ot_span.set_attribute("llm.token_count.total", span.usage.total_tokens)
if span.error:
ot_span.record_exception(Exception(span.error))
except Exception as exc: # noqa: BLE001
raise EmitterFailureException(self.name, str(exc)) from exc
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
shutdown
async
¶
SpanNode
dataclass
¶
StructuredJsonEmitter
¶
NDJSON emitter for log aggregators (Loki, OpenSearch, Datadog logs).
Source code in apogee_ai_observability/infrastructure/emitters/builtin_emitters.py
emit_span
async
¶
emit_span(span: Span) -> None
Source code in apogee_ai_observability/infrastructure/emitters/builtin_emitters.py
async def emit_span(self, span: Span) -> None:
record = {
"kind": "span",
"trace_id": span.trace_id,
"span_id": span.span_id,
"parent_span_id": span.parent_span_id,
"name": span.name,
"type": span.kind.value,
"status": span.status.value,
"start_time": span.start_time.isoformat(),
"end_time": span.end_time.isoformat() if span.end_time else None,
"duration_ms": span.duration_ms,
"attributes": dict(span.attributes),
"events": list(span.events),
"error": span.error,
}
if span.usage is not None:
record["usage"] = {
"provider": span.usage.provider,
"model": span.usage.model,
"prompt_tokens": span.usage.prompt_tokens,
"completion_tokens": span.usage.completion_tokens,
"total_tokens": span.usage.total_tokens,
}
self._stream.write(json.dumps(record, default=str) + "\n")
self._stream.flush()
emit_metric
async
¶
emit_metric(sample: MetricSample) -> None
Source code in apogee_ai_observability/infrastructure/emitters/builtin_emitters.py
async def emit_metric(self, sample: MetricSample) -> None:
record = {
"kind": "metric",
"name": sample.name,
"value": sample.value,
"unit": sample.unit,
"timestamp": sample.timestamp.isoformat(),
"attributes": dict(sample.attributes),
}
self._stream.write(json.dumps(record, default=str) + "\n")
self._stream.flush()
shutdown
async
¶
TableCostCalculator
¶
TableCostCalculator(pricing: IPricingProvider | None = None)
Source code in apogee_ai_observability/infrastructure/cost/table_cost_calculator.py
estimate
¶
estimate(usage: LLMUsage, *, trace_id: str = '', span_id: str | None = None, user_id: str | None = None, tenant_id: str | None = None) -> CostLedgerEntry
Source code in apogee_ai_observability/infrastructure/cost/table_cost_calculator.py
def estimate(
self,
usage: LLMUsage,
*,
trace_id: str = "",
span_id: str | None = None,
user_id: str | None = None,
tenant_id: str | None = None,
) -> CostLedgerEntry:
price = self._pricing.lookup(usage.provider, usage.model)
input_per_m, output_per_m = price if price is not None else (0.0, 0.0)
input_cost = (usage.prompt_tokens / 1_000_000.0) * input_per_m
output_cost = (usage.completion_tokens / 1_000_000.0) * output_per_m
return CostLedgerEntry(
trace_id=trace_id,
span_id=span_id,
provider=usage.provider,
model=usage.model,
input_tokens=usage.prompt_tokens,
output_tokens=usage.completion_tokens,
input_cost_usd=round(input_cost, 6),
output_cost_usd=round(output_cost, 6),
user_id=user_id,
tenant_id=tenant_id,
)
TablePricingProvider
¶
Source code in apogee_ai_observability/infrastructure/cost/pricing_table.py
lookup
¶
Source code in apogee_ai_observability/infrastructure/cost/pricing_table.py
def lookup(self, provider: str, model: str) -> tuple[float, float] | None:
p = provider.lower()
if p not in self._pricing:
return None
catalog = self._pricing[p]
m = model.lower()
if m in catalog:
return catalog[m]
for key in catalog:
if key in m:
return catalog[key]
if "default" in catalog:
return catalog["default"]
return None
TraceContextScope
¶
TraceContextScope(ctx: TraceContext)
TraceReplayer
¶
TraceReplayer(emitter: ITelemetryEmitter)
Re-emits a previously stored trace through any emitter.
Useful for debugging: load a production trace into an in-memory or Phoenix instance for visual inspection without re-running the agent.
Source code in apogee_ai_observability/infrastructure/replay/trace_replayer.py
replay
async
¶
replay(trace: Trace) -> int
Re-emit every span. Returns the number of spans emitted.
Source code in apogee_ai_observability/infrastructure/replay/trace_replayer.py
build_tree
¶
Reconstruct the parent/child tree from flat span list.
Source code in apogee_ai_observability/infrastructure/replay/trace_replayer.py
def build_tree(self, trace: Trace) -> SpanNode | None:
"""Reconstruct the parent/child tree from flat span list."""
nodes: dict[str, SpanNode] = {s.span_id: SpanNode(span=s) for s in trace.spans}
root: SpanNode | None = None
for s in trace.spans:
node = nodes[s.span_id]
if s.parent_span_id is None:
root = node
continue
parent = nodes.get(s.parent_span_id)
if parent is None:
# Orphan; treat as root
root = root or node
continue
parent.children.append(node)
if root is None and nodes:
root = next(iter(nodes.values()))
return root
get_active
¶
get_active() -> TraceContext | None
reset_active
¶
set_active
¶
set_active(ctx: TraceContext | None)