switch OTEL metrics to prometheus (#344)

* feat: replace OTEL with Prometheus

* fix: second pass of docs and cleanup
This commit is contained in:
Rajat Ahuja 2026-01-25 17:26:42 -05:00 committed by GitHub
parent d5e66dd565
commit 2270e5666f
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
28 changed files with 487 additions and 969 deletions

View File

@ -229,14 +229,10 @@ LLM_ANTHROPIC_API_KEY=your-anthropic-api-key-here
# SENTRY_PROFILES_SAMPLE_RATE=0.1
# =============================================================================
# OpenTelemetry Settings (Push-based metrics via OTLP)
# Prometheus Metrics Settings (Pull-based metrics)
# =============================================================================
# OTEL_ENABLED=false
# OTEL_ENDPOINT=https://mimir.example.com/otlp/v1/metrics
# OTEL_HEADERS={"X-Scope-OrgID": "honcho"} # JSON string for auth headers
# OTEL_EXPORT_INTERVAL_MILLIS=60000
# OTEL_SERVICE_NAME=honcho
# OTEL_SERVICE_NAMESPACE=honcho # Inherits from NAMESPACE if not set
# METRICS_ENABLED=false
# METRICS_NAMESPACE=honcho # Inherits from NAMESPACE if not set
# =============================================================================
# CloudEvents Telemetry Settings (Analytics events)

2
.gitignore vendored
View File

@ -190,6 +190,4 @@ CRUSH.md
metrics.jsonl
AGENTS.md
lancedb_data/
mimir-data/
grafana-data/

View File

@ -411,7 +411,7 @@ Then modify the values as needed. The TOML file is organized into sections:
- `[summary]` - Session summarization settings
- `[dream]` - Dream processing configuration (including specialist models and surprisal settings)
- `[webhook]` - Webhook configuration
- `[otel]` - OpenTelemetry push-based metrics via OTLP
- `[metrics]` - Prometheus pull-based metrics
- `[telemetry]` - CloudEvents telemetry for analytics
- `[vector_store]` - Vector store configuration (pgvector, turbopuffer, or lancedb)
- `[sentry]` - Error tracking and monitoring settings
@ -431,7 +431,7 @@ Examples:
- `DERIVER_PROVIDER` - Provider for background deriver
- `SUMMARY_PROVIDER` - Summary generation provider
- `LOG_LEVEL` - Application log level
- `OTEL_ENABLED` - Enable OpenTelemetry metrics
- `METRICS_ENABLED` - Enable Prometheus metrics
- `TELEMETRY_ENABLED` - Enable CloudEvents telemetry
### Configuration Priority

View File

@ -187,14 +187,10 @@ INCLUDE_LEVELS = ["explicit", "deductive"]
SECRET = ""
MAX_WORKSPACE_LIMIT = 10
# OpenTelemetry settings (push-based metrics via OTLP)
[otel]
# Prometheus metrics settings (pull-based metrics)
[metrics]
ENABLED = false
# ENDPOINT = "https://mimir.example.com/otlp/v1/metrics"
# HEADERS = '{"X-Scope-OrgID": "honcho"}' # JSON string for auth headers
EXPORT_INTERVAL_MILLIS = 60000
SERVICE_NAME = "honcho"
# SERVICE_NAMESPACE = "honcho" # Inherits from app.NAMESPACE if not set
# NAMESPACE = "honcho" # Inherits from app.NAMESPACE if not set
# CloudEvents telemetry settings (analytics events)
[telemetry]

View File

@ -59,17 +59,6 @@ services:
interval: 5s
timeout: 5s
retries: 5
mimir:
image: grafana/mimir:2.14.0
command:
- -config.file=/etc/mimir/mimir.yaml
ports:
- 9009:9009
volumes:
- ./mimir-data:/data
configs:
- source: mimir_config
target: /etc/mimir/mimir.yaml
grafana:
image: grafana/grafana:11.4.0
ports:
@ -81,66 +70,6 @@ services:
- GF_AUTH_ANONYMOUS_ORG_ROLE=Viewer
volumes:
- ./grafana-data:/var/lib/grafana
depends_on:
- mimir
volumes:
pgdata:
venv:
configs:
mimir_config:
content: |
multitenancy_enabled: false
server:
http_listen_port: 9009
log_level: warn
common:
storage:
backend: filesystem
filesystem:
dir: /data
blocks_storage:
backend: filesystem
filesystem:
dir: /data/blocks
bucket_store:
sync_dir: /data/tsdb-sync
tsdb:
dir: /data/tsdb
compactor:
data_dir: /data/compactor
sharding_ring:
kvstore:
store: memberlist
distributor:
ring:
instance_addr: 127.0.0.1
kvstore:
store: memberlist
ingester:
ring:
instance_addr: 127.0.0.1
kvstore:
store: memberlist
replication_factor: 1
ruler_storage:
backend: filesystem
filesystem:
dir: /data/rules
alertmanager_storage:
backend: filesystem
filesystem:
dir: /data/alertmanager
store_gateway:
sharding_ring:
replication_factor: 1

View File

@ -57,7 +57,7 @@ Then modify the values as needed. The TOML file is organized into sections:
- `[summary]` - Session summarization settings (frequency thresholds, provider, model, token limits for short and long summaries)
- `[dream]` - Dream processing configuration (enable/disable, thresholds, idle timeouts, dream types, LLM settings, surprisal sampling)
- `[webhook]` - Webhook configuration (webhook secret, workspace limits)
- `[otel]` - OpenTelemetry settings for push-based metrics via OTLP
- `[metrics]` - Prometheus pull-based metrics settings
- `[telemetry]` - CloudEvents telemetry settings for analytics
- `[vector_store]` - Vector store configuration (pgvector, Turbopuffer, LanceDB)
- `[sentry]` - Error tracking and monitoring settings (enable/disable, DSN, environment, sample rates)
@ -540,29 +540,19 @@ VECTOR_STORE_LANCEDB_PATH=./lancedb_data
## Monitoring Configuration
### OpenTelemetry (Push-based Metrics)
### Prometheus Metrics (Pull-based)
Honcho supports push-based metrics via OpenTelemetry Protocol (OTLP) to any compatible backend (Mimir, Grafana Cloud, etc.).
Honcho exposes Prometheus metrics via `/metrics` endpoints for scraping:
- **API process**: Port 8000 at `/metrics`
- **Deriver process**: Port 9090 at `/metrics`
**OpenTelemetry Settings:**
**Metrics Settings:**
```bash
# Enable/disable OTel metrics
OTEL_ENABLED=false
# Enable/disable Prometheus metrics
METRICS_ENABLED=false
# OTLP HTTP endpoint for metrics
# For Mimir: <mimir-url>/otlp/v1/metrics
# For Grafana Cloud: https://otlp-gateway-<region>.grafana.net/otlp/v1/metrics
OTEL_ENDPOINT=https://mimir.example.com/otlp/v1/metrics
# Optional auth headers (JSON format in env var)
OTEL_HEADERS='{"X-Scope-OrgID": "honcho"}'
# Export interval in milliseconds (default: 60 seconds)
OTEL_EXPORT_INTERVAL_MILLIS=60000
# Service identification
OTEL_SERVICE_NAME=honcho
OTEL_SERVICE_NAMESPACE=honcho # Inherits from app.NAMESPACE if not set
# Namespace label for all metrics (inherits from app.NAMESPACE if not set)
METRICS_NAMESPACE=honcho
```
### CloudEvents Telemetry (Analytics)
@ -691,7 +681,7 @@ MODEL = "claude-sonnet-4-20250514"
[webhook]
MAX_WORKSPACE_LIMIT = 10
[otel]
[metrics]
ENABLED = false
[telemetry]
@ -798,7 +788,7 @@ MODEL = "claude-sonnet-4-20250514"
[webhook]
MAX_WORKSPACE_LIMIT = 10
[otel]
[metrics]
ENABLED = true
[telemetry]
@ -838,7 +828,7 @@ LLM_GROQ_API_KEY=your-prod-groq-key
WEBHOOK_SECRET=your-webhook-signing-secret
# Monitoring
OTEL_ENDPOINT=https://mimir.example.com/otlp/v1/metrics
METRICS_ENABLED=true
TELEMETRY_ENDPOINT=https://telemetry.honcho.dev/v1/events
SENTRY_DSN=https://your-sentry-dsn@sentry.io/project-id
SENTRY_ENVIRONMENT=production

View File

@ -28,3 +28,13 @@ kill_timeout = '5s'
cpu_kind = 'shared'
cpus = 1
processes = ['api', 'deriver']
[[metrics]]
port = 8000
path = "/metrics"
processes = ["api"]
[[metrics]]
port = 9090
path = "/metrics"
processes = ["deriver"]

View File

@ -37,8 +37,7 @@ dependencies = [
"redis>=7.0.0,<8.0.0",
"cashews[redis]==7.4.4",
"scikit-learn>=1.6.0",
"opentelemetry-sdk>=1.36.0",
"opentelemetry-exporter-otlp-proto-http>=1.36.0",
"prometheus_client>=0.21.0",
"cloudevents>=1.12.0",
]
[dependency-groups]
@ -116,7 +115,7 @@ reportUnusedCallResult = false
reportCallInDefaultInitializer = false
reportAny = false
reportExplicitAny = false
allowedUntypedLibraries = ["langfuse", "lancedb", "pyarrow", "opentelemetry"]
allowedUntypedLibraries = ["langfuse", "lancedb", "pyarrow"]
reportImplicitOverride = false
reportImportCycles = false

View File

@ -68,7 +68,7 @@ class TomlConfigSettingsSource(PydanticBaseSettingsSource):
"WEBHOOK": "webhook",
"DREAM": "dream",
"VECTOR_STORE": "vector_store",
"OTEL": "otel",
"METRICS": "metrics",
"TELEMETRY": "telemetry",
"": "app", # For AppSettings with no prefix
}
@ -445,33 +445,10 @@ class WebhookSettings(HonchoSettings):
MAX_WORKSPACE_LIMIT: int = 10
class OpenTelemetrySettings(HonchoSettings):
"""OpenTelemetry settings for push-based metrics via OTLP.
These settings configure the OTel SDK to push metrics via OTLP HTTP
to any compatible backend (Mimir, Grafana Cloud, etc.).
"""
model_config = SettingsConfigDict(env_prefix="OTEL_", extra="ignore") # pyright: ignore
# Master toggle for OTel metrics
class MetricsSettings(HonchoSettings):
model_config = SettingsConfigDict(env_prefix="METRICS_", extra="ignore") # pyright: ignore
ENABLED: bool = False
# OTLP HTTP endpoint for metrics (e.g., "https://mimir.example.com/otlp/v1/metrics")
# For Mimir, the endpoint is typically: <mimir-url>/otlp/v1/metrics
# For Grafana Cloud: https://otlp-gateway-<region>.grafana.net/otlp/v1/metrics
ENDPOINT: str | None = None
HEADERS: dict[str, str] | None = None
# Export interval in milliseconds (default: 60 seconds)
EXPORT_INTERVAL_MILLIS: int = 60000
# Service name for resource attributes (identifies what this service is)
SERVICE_NAME: str = "honcho"
# Service namespace for resource attributes (defaults to top-level NAMESPACE if not set)
SERVICE_NAMESPACE: str | None = None
NAMESPACE: str | None = None
class TelemetrySettings(HonchoSettings):
@ -676,7 +653,7 @@ class AppSettings(HonchoSettings):
PEER_CARD: PeerCardSettings = Field(default_factory=PeerCardSettings)
SUMMARY: SummarySettings = Field(default_factory=SummarySettings)
WEBHOOK: WebhookSettings = Field(default_factory=WebhookSettings)
OTEL: OpenTelemetrySettings = Field(default_factory=OpenTelemetrySettings)
METRICS: MetricsSettings = Field(default_factory=MetricsSettings)
TELEMETRY: TelemetrySettings = Field(default_factory=TelemetrySettings)
CACHE: CacheSettings = Field(default_factory=CacheSettings)
DREAM: DreamSettings = Field(default_factory=DreamSettings)
@ -691,20 +668,15 @@ class AppSettings(HonchoSettings):
@model_validator(mode="after")
def propagate_namespace(self) -> "AppSettings":
"""Propagate top-level NAMESPACE to nested settings if not explicitly set.
After this validator runs, CACHE.NAMESPACE,
VECTOR_STORE.NAMESPACE, TELEMETRY.NAMESPACE, and OTEL.SERVICE_NAMESPACE
are guaranteed to exist. Explicitly provided nested namespaces are preserved.
"""
"""Propagate top-level NAMESPACE to nested settings if not explicitly set."""
if "NAMESPACE" not in self.CACHE.model_fields_set:
self.CACHE.NAMESPACE = self.NAMESPACE
if "NAMESPACE" not in self.VECTOR_STORE.model_fields_set:
self.VECTOR_STORE.NAMESPACE = self.NAMESPACE
if "NAMESPACE" not in self.TELEMETRY.model_fields_set:
self.TELEMETRY.NAMESPACE = self.NAMESPACE
if "SERVICE_NAMESPACE" not in self.OTEL.model_fields_set:
self.OTEL.SERVICE_NAMESPACE = self.NAMESPACE
if "NAMESPACE" not in self.METRICS.model_fields_set:
self.METRICS.NAMESPACE = self.NAMESPACE
return self

View File

@ -3,17 +3,20 @@ import logging
import os
import uvloop
from prometheus_client import start_http_server
from src.config import settings
from src.telemetry import (
initialize_telemetry,
initialize_telemetry_async,
shutdown_telemetry,
)
from src.telemetry import initialize_telemetry_async, shutdown_telemetry
from .queue_manager import main
def start_metrics_server() -> None:
"""Start the Prometheus metrics HTTP server on port 9090."""
start_http_server(9090)
print("[DERIVER] Prometheus metrics server started on port 9090")
def setup_logging():
"""
Configure logging for the deriver process.
@ -55,7 +58,7 @@ async def run_deriver():
try:
await main()
finally:
# Shutdown telemetry (flush CloudEvents buffer, shutdown OTel metrics)
# Shutdown telemetry (flush CloudEvents buffer)
await shutdown_telemetry()
@ -67,8 +70,10 @@ if __name__ == "__main__":
asyncio.set_event_loop_policy(uvloop.EventLoopPolicy())
try:
# Initialize sync telemetry (OTel metrics)
initialize_telemetry()
# Start Prometheus metrics server if enabled
if settings.METRICS.ENABLED:
start_metrics_server()
print("[DERIVER] Running main loop")
asyncio.run(run_deriver())
except KeyboardInterrupt:

View File

@ -7,10 +7,14 @@ from src.crud.representation import RepresentationManager
from src.dependencies import tracked_db
from src.models import Message
from src.schemas import ResolvedConfiguration
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.telemetry.events import RepresentationCompletedEvent, emit
from src.telemetry.logging import accumulate_metric, log_performance_metrics
from src.telemetry.otel.metrics import DeriverComponents, DeriverTaskTypes, TokenTypes
from src.telemetry.prometheus.metrics import (
DeriverComponents,
DeriverTaskTypes,
TokenTypes,
)
from src.telemetry.sentry import with_sentry_transaction
from src.utils.clients import honcho_llm_call
from src.utils.config_helpers import get_configuration
@ -141,9 +145,9 @@ async def process_representation_tasks_batch(
"ms",
)
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_deriver_tokens(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_deriver_tokens(
count=response.output_tokens,
task_type=DeriverTaskTypes.INGESTION.value,
token_type=TokenTypes.OUTPUT.value,

View File

@ -35,7 +35,7 @@ from src.reconciler import (
set_reconciler_scheduler,
)
from src.schemas import ResolvedConfiguration
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.telemetry.sentry import initialize_sentry
from src.utils.work_unit import parse_work_unit_key
from src.webhooks.events import (
@ -791,9 +791,9 @@ class QueueManager:
if (
work_unit.task_type in ["representation", "summary"]
and work_unit.workspace_name is not None
and settings.OTEL.ENABLED
and settings.METRICS.ENABLED
):
otel_metrics.record_deriver_queue_item(
prometheus_metrics.record_deriver_queue_item(
count=len(items),
workspace_name=work_unit.workspace_name,
task_type=work_unit.task_type,

View File

@ -16,14 +16,14 @@ from sqlalchemy.ext.asyncio import AsyncSession
from src import crud
from src.config import ReasoningLevel, settings
from src.dialectic import prompts
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.telemetry.events import DialecticCompletedEvent, emit
from src.telemetry.logging import (
accumulate_metric,
log_performance_metrics,
log_token_usage_metrics,
)
from src.telemetry.otel.metrics import DialecticComponents, TokenTypes
from src.telemetry.prometheus.metrics import DialecticComponents, TokenTypes
from src.utils.agent_tools import (
DIALECTIC_TOOLS,
DIALECTIC_TOOLS_MINIMAL,
@ -337,15 +337,15 @@ class DialecticAgent:
if not self.metric_key and run_id is not None:
log_performance_metrics("dialectic_chat", run_id)
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_dialectic_tokens(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_dialectic_tokens(
count=input_tokens,
token_type=TokenTypes.INPUT.value,
component=DialecticComponents.TOTAL.value,
reasoning_level=self.reasoning_level,
)
otel_metrics.record_dialectic_tokens(
prometheus_metrics.record_dialectic_tokens(
count=output_tokens,
token_type=TokenTypes.OUTPUT.value,
component=DialecticComponents.TOTAL.value,

View File

@ -22,10 +22,10 @@ from sqlalchemy.ext.asyncio import AsyncSession
from src.config import settings
from src.schemas import ResolvedConfiguration
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.telemetry.events import DreamSpecialistEvent, emit
from src.telemetry.logging import accumulate_metric, log_performance_metrics
from src.telemetry.otel.metrics import TokenTypes
from src.telemetry.prometheus.metrics import TokenTypes
from src.utils.agent_tools import (
DEDUCTION_SPECIALIST_TOOLS,
INDUCTION_SPECIALIST_TOOLS,
@ -186,14 +186,14 @@ class BaseSpecialist(ABC):
accumulate_metric(task_name, "input_tokens", response.input_tokens, "count")
accumulate_metric(task_name, "output_tokens", response.output_tokens, "count")
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_dreamer_tokens(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_dreamer_tokens(
count=response.input_tokens,
specialist_name=self.name,
token_type=TokenTypes.INPUT.value,
)
otel_metrics.record_dreamer_tokens(
prometheus_metrics.record_dreamer_tokens(
count=response.output_tokens,
specialist_name=self.name,
token_type=TokenTypes.OUTPUT.value,

View File

@ -30,9 +30,9 @@ from src.routers import (
)
from src.security import create_admin_jwt
from src.telemetry import (
initialize_telemetry,
initialize_telemetry_async,
otel_metrics,
metrics_endpoint,
prometheus_metrics,
shutdown_telemetry,
)
from src.telemetry.logging import get_route_template
@ -120,8 +120,7 @@ if SENTRY_ENABLED:
@asynccontextmanager
async def lifespan(_: FastAPI):
# Initialize telemetry (OTel metrics + CloudEvents emitter)
initialize_telemetry()
# Initialize CloudEvents telemetry
await initialize_telemetry_async()
try:
@ -140,7 +139,7 @@ async def lifespan(_: FastAPI):
await close_external_vector_store()
await close_cache()
await engine.dispose()
# Shutdown telemetry (flush CloudEvents buffer, shutdown OTel metrics)
# Shutdown telemetry (flush CloudEvents buffer)
await shutdown_telemetry()
@ -191,6 +190,9 @@ app.include_router(conclusions.router, prefix="/v3")
app.include_router(keys.router, prefix="/v3")
app.include_router(webhooks.router, prefix="/v3")
# Prometheus metrics endpoint
app.add_route("/metrics", metrics_endpoint, methods=["GET"])
# Global exception handlers
@app.exception_handler(HonchoException)
@ -233,9 +235,9 @@ async def track_request(
response = await call_next(request)
# Track metrics if enabled
if settings.OTEL.ENABLED:
if settings.METRICS.ENABLED:
template = get_route_template(request)
otel_metrics.record_api_request(
prometheus_metrics.record_api_request(
method=request.method,
endpoint=template,
status_code=str(response.status_code),

View File

@ -22,7 +22,7 @@ from src.dependencies import db
from src.deriver import enqueue
from src.exceptions import FileTooLargeError, ResourceNotFoundException
from src.security import require_auth
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.utils.files import process_file_uploads_for_messages
logger = logging.getLogger(__name__)
@ -100,9 +100,9 @@ async def create_messages_for_session(
session_name=session_id,
)
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_messages_created(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_messages_created(
count=len(created_messages),
workspace_name=workspace_id,
)
@ -199,9 +199,9 @@ async def create_messages_with_file(
len(created_messages),
)
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_messages_created(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_messages_created(
count=len(created_messages),
workspace_name=workspace_id,
)

View File

@ -14,7 +14,7 @@ from src.dependencies import db, tracked_db
from src.dialectic.chat import agentic_chat, agentic_chat_stream
from src.exceptions import AuthenticationException, ResourceNotFoundException
from src.security import JWTParams, require_auth
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.utils.search import search
logger = logging.getLogger(__name__)
@ -184,9 +184,9 @@ async def chat(
yield f"data: {json.dumps({'delta': {'content': chunk}, 'done': False})}\n\n"
yield f"data: {json.dumps({'done': True})}\n\n"
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_dialectic_call(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_dialectic_call(
workspace_name=workspace_id,
reasoning_level=options.reasoning_level,
)
@ -216,9 +216,9 @@ async def chat(
reasoning_level=options.reasoning_level,
)
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_dialectic_call(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_dialectic_call(
workspace_name=workspace_id,
reasoning_level=options.reasoning_level,
)

View File

@ -3,7 +3,7 @@ Telemetry module for Honcho.
This module consolidates all telemetry, metrics, and observability functionality:
- Sentry: Error tracking and performance tracing
- OTel: Push-based metrics via OTLP to any compatible backend (e.g., Mimir)
- Prometheus: Pull-based metrics scraped by Fly.io
- CloudEvents: Structured events for analytics (push-based)
- Logging: Langfuse integration, Rich console output, metric accumulation
- Tracing: Sentry transaction decorators
@ -12,54 +12,27 @@ This module consolidates all telemetry, metrics, and observability functionality
"""
from src.telemetry.events import emit
from src.telemetry.otel import get_meter, initialize_otel_metrics, shutdown_otel_metrics
from src.telemetry.otel.metrics import otel_metrics
from src.telemetry.prometheus import metrics_endpoint, prometheus_metrics
__all__ = [
"emit",
"get_meter",
"initialize_otel_metrics",
"initialize_telemetry",
"initialize_telemetry_async",
"otel_metrics",
"shutdown_otel_metrics",
"metrics_endpoint",
"prometheus_metrics",
"shutdown_telemetry",
]
def initialize_telemetry() -> None:
"""
Initialize all telemetry systems based on configuration.
This should be called once at application startup (in main.py lifespan).
It reads configuration from settings and initializes:
- OTel metrics (if OTEL_ENABLED=true)
Note: CloudEvents telemetry requires async initialization and should be
initialized separately using initialize_telemetry_async().
Sentry is initialized separately in sentry.py as it has its own lifecycle.
"""
from src.config import settings
if settings.OTEL.ENABLED:
initialize_otel_metrics(
endpoint=settings.OTEL.ENDPOINT,
headers=settings.OTEL.HEADERS,
export_interval_millis=settings.OTEL.EXPORT_INTERVAL_MILLIS,
service_name=settings.OTEL.SERVICE_NAME,
service_namespace=settings.OTEL.SERVICE_NAMESPACE,
enabled=True,
)
async def initialize_telemetry_async() -> None:
"""
Initialize async telemetry systems based on configuration.
This should be called once at application startup (in main.py lifespan),
after initialize_telemetry(). It initializes:
This should be called once at application startup (in main.py lifespan).
It initializes:
- CloudEvents emitter (if TELEMETRY_ENABLED=true)
Note: Prometheus metrics are pull-based and require no initialization.
Sentry is initialized separately in sentry.py as it has its own lifecycle.
"""
from src.config import settings
from src.telemetry.events import initialize_telemetry_events
@ -73,13 +46,9 @@ async def shutdown_telemetry() -> None:
Shutdown all telemetry systems gracefully.
This should be called during application shutdown to ensure:
- OTel metrics are flushed
- CloudEvents buffer is flushed
"""
from src.telemetry.events import shutdown_telemetry_events
# Shutdown CloudEvents emitter (flushes buffer)
await shutdown_telemetry_events()
# Shutdown OTel metrics
shutdown_otel_metrics()

View File

@ -1,26 +0,0 @@
"""
OpenTelemetry metrics module for Honcho.
This module provides push-based metrics using OpenTelemetry SDK
with OTLP HTTP export to any compatible backend (Mimir, Grafana Cloud, etc.).
"""
from src.telemetry.otel.metrics import (
DeriverComponents,
DeriverTaskTypes,
DialecticComponents,
TokenTypes,
get_meter,
initialize_otel_metrics,
shutdown_otel_metrics,
)
__all__ = [
"DeriverComponents",
"DeriverTaskTypes",
"DialecticComponents",
"TokenTypes",
"get_meter",
"initialize_otel_metrics",
"shutdown_otel_metrics",
]

View File

@ -1,481 +0,0 @@
"""
OpenTelemetry metrics implementation with OTLP export.
This module provides push-based metrics that replace the pull-based Prometheus
/metrics endpoint. Metrics are pushed via OTLP to any compatible backend
(Mimir, Grafana Cloud, etc.) at configurable intervals.
Usage:
from src.telemetry.otel import get_meter, initialize_otel_metrics
# Initialize once at startup
initialize_otel_metrics()
# Get a meter for your component
meter = get_meter("honcho.deriver")
# Create instruments
counter = meter.create_counter("tokens_processed", unit="tokens")
counter.add(100, {"task_type": "ingestion"})
"""
from __future__ import annotations
import atexit
import logging
from enum import Enum
from typing import TYPE_CHECKING, final
from opentelemetry import metrics
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader
from opentelemetry.sdk.resources import Resource
if TYPE_CHECKING:
from opentelemetry.metrics import Counter, Meter
logger = logging.getLogger(__name__)
# =============================================================================
# Metric label enums
# =============================================================================
class TokenTypes(Enum):
INPUT = "input"
OUTPUT = "output"
class DeriverTaskTypes(Enum):
INGESTION = "ingestion"
SUMMARY = "summary"
class DeriverComponents(Enum):
PROMPT = "prompt" # used in ingestion and summary
MESSAGES = "messages" # used in ingestion and summary
PREVIOUS_SUMMARY = "previous_summary" # only used for summary
OUTPUT_TOTAL = "output_total"
class DialecticComponents(Enum):
TOTAL = "total"
# =============================================================================
# OTel metrics infrastructure
# =============================================================================
# Global state
_meter_provider: MeterProvider | None = None
_initialized: bool = False
def initialize_otel_metrics(
*,
endpoint: str | None = None,
headers: dict[str, str] | None = None,
export_interval_millis: int = 60000,
service_name: str = "honcho",
service_namespace: str | None = None,
enabled: bool = True,
) -> None:
"""
Initialize OpenTelemetry metrics with OTLP export.
This should be called once at application startup. If already initialized,
subsequent calls are no-ops.
Args:
endpoint: OTLP HTTP endpoint URL (e.g., "https://mimir.example.com/otlp/v1/metrics").
If None, metrics are collected but not exported (useful for testing).
headers: Optional headers to include in requests (e.g., {"X-Scope-OrgID": "tenant"}).
export_interval_millis: How often to export metrics (default: 60 seconds).
service_name: Service name for resource attributes (default: "honcho").
service_namespace: Optional namespace for the service.
enabled: If False, metrics are no-ops (default: True).
"""
global _meter_provider, _initialized
if _initialized:
logger.debug("OTel metrics already initialized, skipping")
return
if not enabled:
logger.info("OTel metrics disabled")
_initialized = True
return
# Build resource attributes
resource_attributes = {
"service.name": service_name,
}
if service_namespace:
resource_attributes["service.namespace"] = service_namespace
resource = Resource.create(resource_attributes)
# Create metric reader
readers: list[PeriodicExportingMetricReader] = []
if endpoint:
try:
from opentelemetry.exporter.otlp.proto.http.metric_exporter import (
OTLPMetricExporter,
)
exporter = OTLPMetricExporter(
endpoint=endpoint,
headers=headers or {},
)
reader = PeriodicExportingMetricReader(
exporter,
export_interval_millis=export_interval_millis,
)
readers.append(reader)
logger.info(f"OTel metrics configured to push via OTLP to {endpoint}")
except Exception as e:
logger.error(f"Failed to configure OTLP metrics exporter: {e}")
# Continue without exporter - metrics still work locally
else:
logger.info(
"OTel metrics initialized without remote export (no endpoint configured)"
)
# Create and set the meter provider
# Note: empty list is valid for metric_readers (metrics still work, just not exported)
_meter_provider = MeterProvider(
resource=resource,
metric_readers=readers,
)
metrics.set_meter_provider(_meter_provider)
# Register shutdown handler
atexit.register(shutdown_otel_metrics)
_initialized = True
logger.info("OTel metrics initialized successfully")
def shutdown_otel_metrics() -> None:
"""
Shutdown the OTel metrics provider, flushing any pending metrics.
This is automatically called at process exit via atexit, but can be
called manually for graceful shutdown.
"""
global _meter_provider, _initialized
if _meter_provider is not None:
try:
_meter_provider.shutdown()
logger.info("OTel metrics shutdown complete")
except Exception as e:
logger.error(f"Error during OTel metrics shutdown: {e}")
finally:
_meter_provider = None
_initialized = False
def get_meter(name: str, version: str = "") -> Meter:
"""
Get an OTel Meter for creating instruments.
Args:
name: The name of the instrumentation scope (e.g., "honcho.deriver").
version: Optional version of the instrumentation scope.
Returns:
An OTel Meter instance for creating counters, histograms, etc.
Example:
meter = get_meter("honcho.deriver")
counter = meter.create_counter("tokens_processed", unit="tokens")
counter.add(100, {"task_type": "ingestion"})
"""
return metrics.get_meter(name, version)
# =============================================================================
# Pre-defined metrics that mirror existing Prometheus counters
# =============================================================================
# These are created lazily on first use to avoid issues with initialization order
@final
class OTelMetrics:
"""
Container for OTel metrics that mirror existing Prometheus counters.
This class provides a bridge during migration - the same metrics are
available via both Prometheus (pull) and OTel (push).
Namespace is managed at the instance level (from settings), not per-call.
"""
_instance: OTelMetrics | None = None
_is_initialized: bool = False
_namespace: str = "honcho"
# Meters (lazily initialized)
_api_meter: Meter | None = None
_deriver_meter: Meter | None = None
_dialectic_meter: Meter | None = None
_dreamer_meter: Meter | None = None
# Counters (lazily initialized)
_api_requests: Counter | None = None
_messages_created: Counter | None = None
_dialectic_calls: Counter | None = None
_deriver_queue_items: Counter | None = None
_deriver_tokens: Counter | None = None
_dialectic_tokens: Counter | None = None
_dreamer_tokens: Counter | None = None
def __new__(cls) -> OTelMetrics:
if cls._instance is None:
cls._instance = super().__new__(cls)
return cls._instance
def _ensure_initialized(self) -> None:
"""Lazily initialize meters and instruments."""
if self._is_initialized:
return
# Get namespace from settings (same as Prometheus)
from src.config import settings
self._namespace = settings.OTEL.SERVICE_NAMESPACE or "honcho"
# Get meters for different components
self._api_meter = get_meter("honcho.api")
self._deriver_meter = get_meter("honcho.deriver")
self._dialectic_meter = get_meter("honcho.dialectic")
self._dreamer_meter = get_meter("honcho.dreamer")
# Create counters that mirror Prometheus metrics (_total suffix follows OpenMetrics convention)
# API requests
self._api_requests = self._api_meter.create_counter(
name="api_requests_total",
unit="requests",
description="Total API requests",
)
# Messages created
self._messages_created = self._api_meter.create_counter(
name="messages_created_total",
unit="messages",
description="Total messages created",
)
# Dialectic calls
self._dialectic_calls = self._dialectic_meter.create_counter(
name="dialectic_calls_total",
unit="calls",
description="Total dialectic calls",
)
# Deriver queue items processed
self._deriver_queue_items = self._deriver_meter.create_counter(
name="deriver_queue_items_processed_total",
unit="items",
description="Total deriver queue items processed",
)
# Token counters
self._deriver_tokens = self._deriver_meter.create_counter(
name="deriver_tokens_processed_total",
unit="tokens",
description="Total tokens processed by the deriver",
)
self._dialectic_tokens = self._dialectic_meter.create_counter(
name="dialectic_tokens_processed_total",
unit="tokens",
description="Total tokens processed by the dialectic",
)
self._dreamer_tokens = self._dreamer_meter.create_counter(
name="dreamer_tokens_processed_total",
unit="tokens",
description="Total tokens processed by the dreamer",
)
self._is_initialized = True
def _handle_metric_error(self, method_name: str, error: Exception) -> None:
"""Handle errors from metric recording by logging to Sentry."""
import sentry_sdk
sentry_sdk.capture_exception(error)
logger.warning(
"Failed to record OTel metric in %s: %s", method_name, str(error)
)
def record_api_request(
self,
*,
method: str,
endpoint: str,
status_code: str,
) -> None:
"""Record an API request metric."""
try:
self._ensure_initialized()
if self._api_requests is None:
return # Not initialized, skip silently
self._api_requests.add(
1,
{
"method": method,
"endpoint": endpoint,
"status_code": status_code,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_api_request", e)
def record_messages_created(
self,
*,
count: int,
workspace_name: str,
) -> None:
"""Record messages created metric."""
try:
self._ensure_initialized()
if self._messages_created is None:
return # Not initialized, skip silently
self._messages_created.add(
count,
{
"workspace_name": workspace_name,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_messages_created", e)
def record_dialectic_call(
self,
*,
workspace_name: str,
reasoning_level: str,
) -> None:
"""Record a dialectic call metric."""
try:
self._ensure_initialized()
if self._dialectic_calls is None:
return # Not initialized, skip silently
self._dialectic_calls.add(
1,
{
"workspace_name": workspace_name,
"reasoning_level": reasoning_level,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_dialectic_call", e)
def record_deriver_queue_item(
self,
*,
count: int,
workspace_name: str,
task_type: str,
) -> None:
"""Record deriver queue items processed metric."""
try:
self._ensure_initialized()
if self._deriver_queue_items is None:
return # Not initialized, skip silently
self._deriver_queue_items.add(
count,
{
"workspace_name": workspace_name,
"task_type": task_type,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_deriver_queue_item", e)
def record_deriver_tokens(
self,
*,
count: int,
task_type: str,
token_type: str,
component: str,
) -> None:
"""Record deriver token usage metric."""
try:
self._ensure_initialized()
if self._deriver_tokens is None:
return # Not initialized, skip silently
self._deriver_tokens.add(
count,
{
"task_type": task_type,
"token_type": token_type,
"component": component,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_deriver_tokens", e)
def record_dialectic_tokens(
self,
*,
count: int,
token_type: str,
component: str,
reasoning_level: str,
) -> None:
"""Record dialectic token usage metric."""
try:
self._ensure_initialized()
if self._dialectic_tokens is None:
return # Not initialized, skip silently
self._dialectic_tokens.add(
count,
{
"token_type": token_type,
"component": component,
"reasoning_level": reasoning_level,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_dialectic_tokens", e)
def record_dreamer_tokens(
self,
*,
count: int,
specialist_name: str,
token_type: str,
) -> None:
"""Record dreamer token usage metric."""
try:
self._ensure_initialized()
if self._dreamer_tokens is None:
return # Not initialized, skip silently
self._dreamer_tokens.add(
count,
{
"specialist_name": specialist_name,
"token_type": token_type,
"namespace": self._namespace,
},
)
except Exception as e:
self._handle_metric_error("record_dreamer_tokens", e)
# Singleton instance
otel_metrics = OTelMetrics()

View File

@ -0,0 +1,26 @@
"""
Prometheus telemetry module.
Exports:
- prometheus_metrics: Singleton for recording metrics
- metrics_endpoint: Async endpoint for /metrics route
- Label enums for metric values
"""
from src.telemetry.prometheus.metrics import (
DeriverComponents,
DeriverTaskTypes,
DialecticComponents,
TokenTypes,
metrics_endpoint,
prometheus_metrics,
)
__all__ = [
"DeriverComponents",
"DeriverTaskTypes",
"DialecticComponents",
"TokenTypes",
"metrics_endpoint",
"prometheus_metrics",
]

View File

@ -0,0 +1,234 @@
"""Prometheus metrics for Honcho."""
from __future__ import annotations
import logging
from enum import Enum
from typing import cast, final
from prometheus_client import (
CONTENT_TYPE_LATEST,
REGISTRY,
Counter,
disable_created_metrics,
generate_latest,
)
from starlette.requests import Request
from starlette.responses import Response
from src.config import settings
disable_created_metrics()
logger = logging.getLogger(__name__)
class NamespacedCounter(Counter):
def labels(self, **kwargs: str) -> NamespacedCounter:
kwargs["namespace"] = cast(str, settings.METRICS.NAMESPACE)
return super().labels(**kwargs) # type: ignore[return-value]
class TokenTypes(Enum):
INPUT = "input"
OUTPUT = "output"
class DeriverTaskTypes(Enum):
INGESTION = "ingestion"
SUMMARY = "summary"
class DeriverComponents(Enum):
PROMPT = "prompt"
MESSAGES = "messages"
PREVIOUS_SUMMARY = "previous_summary"
OUTPUT_TOTAL = "output_total"
class DialecticComponents(Enum):
TOTAL = "total"
api_requests_counter = NamespacedCounter(
"api_requests",
"Total API requests",
["namespace", "method", "endpoint", "status_code"],
)
messages_created_counter = NamespacedCounter(
"messages_created",
"Total messages created",
["namespace", "workspace_name"],
)
dialectic_calls_counter = NamespacedCounter(
"dialectic_calls",
"Total dialectic calls",
["namespace", "workspace_name", "reasoning_level"],
)
deriver_queue_items_processed_counter = NamespacedCounter(
"deriver_queue_items_processed",
"Total deriver queue items processed",
["namespace", "workspace_name", "task_type"],
)
deriver_tokens_processed_counter = NamespacedCounter(
"deriver_tokens_processed",
"Total tokens processed by the deriver",
["namespace", "task_type", "token_type", "component"],
)
dialectic_tokens_processed_counter = NamespacedCounter(
"dialectic_tokens_processed",
"Total tokens processed by the dialectic",
["namespace", "token_type", "component", "reasoning_level"],
)
dreamer_tokens_processed_counter = NamespacedCounter(
"dreamer_tokens_processed",
"Total tokens processed by the dreamer",
["namespace", "specialist_name", "token_type"],
)
@final
class PrometheusMetrics:
_instance: PrometheusMetrics | None = None
def __new__(cls) -> PrometheusMetrics:
if cls._instance is None:
cls._instance = super().__new__(cls)
return cls._instance
def _handle_metric_error(self, method_name: str, error: Exception) -> None:
import sentry_sdk
sentry_sdk.capture_exception(error)
logger.warning(
"Failed to record Prometheus metric in %s: %s", method_name, str(error)
)
def record_api_request(
self,
*,
method: str,
endpoint: str,
status_code: str,
) -> None:
try:
api_requests_counter.labels(
method=method,
endpoint=endpoint,
status_code=status_code,
).inc()
except Exception as e:
self._handle_metric_error("record_api_request", e)
def record_messages_created(
self,
*,
count: int,
workspace_name: str,
) -> None:
try:
messages_created_counter.labels(
workspace_name=workspace_name,
).inc(count)
except Exception as e:
self._handle_metric_error("record_messages_created", e)
def record_dialectic_call(
self,
*,
workspace_name: str,
reasoning_level: str,
) -> None:
try:
dialectic_calls_counter.labels(
workspace_name=workspace_name,
reasoning_level=reasoning_level,
).inc()
except Exception as e:
self._handle_metric_error("record_dialectic_call", e)
def record_deriver_queue_item(
self,
*,
count: int,
workspace_name: str,
task_type: str,
) -> None:
try:
deriver_queue_items_processed_counter.labels(
workspace_name=workspace_name,
task_type=task_type,
).inc(count)
except Exception as e:
self._handle_metric_error("record_deriver_queue_item", e)
def record_deriver_tokens(
self,
*,
count: int,
task_type: str,
token_type: str,
component: str,
) -> None:
try:
deriver_tokens_processed_counter.labels(
task_type=task_type,
token_type=token_type,
component=component,
).inc(count)
except Exception as e:
self._handle_metric_error("record_deriver_tokens", e)
def record_dialectic_tokens(
self,
*,
count: int,
token_type: str,
component: str,
reasoning_level: str,
) -> None:
try:
dialectic_tokens_processed_counter.labels(
token_type=token_type,
component=component,
reasoning_level=reasoning_level,
).inc(count)
except Exception as e:
self._handle_metric_error("record_dialectic_tokens", e)
def record_dreamer_tokens(
self,
*,
count: int,
specialist_name: str,
token_type: str,
) -> None:
try:
dreamer_tokens_processed_counter.labels(
specialist_name=specialist_name,
token_type=token_type,
).inc(count)
except Exception as e:
self._handle_metric_error("record_dreamer_tokens", e)
prometheus_metrics = PrometheusMetrics()
async def metrics_endpoint(_request: Request) -> Response:
if not settings.METRICS.ENABLED:
return Response("Metrics are disabled", status_code=404)
try:
return Response(
content=generate_latest(REGISTRY),
media_type=CONTENT_TYPE_LATEST,
)
except Exception as e:
logger.error(f"Failed to generate metrics: {e}", exc_info=True)
return Response("Failed to generate metrics", status_code=500)

View File

@ -16,10 +16,14 @@ from src.crud.session import session_cache_key
from src.dependencies import tracked_db
from src.exceptions import ResourceNotFoundException
from src.models import Message
from src.telemetry import otel_metrics
from src.telemetry import prometheus_metrics
from src.telemetry.events import AgentToolSummaryCreatedEvent, emit
from src.telemetry.logging import accumulate_metric, conditional_observe
from src.telemetry.otel.metrics import DeriverComponents, DeriverTaskTypes, TokenTypes
from src.telemetry.prometheus.metrics import (
DeriverComponents,
DeriverTaskTypes,
TokenTypes,
)
from src.utils.clients import HonchoLLMCallResponse, honcho_llm_call
from src.utils.formatting import utc_now_iso
from src.utils.tokens import estimate_tokens, track_deriver_input_tokens
@ -449,8 +453,8 @@ async def _create_and_save_summary(
)
# Track output tokens
if settings.OTEL.ENABLED:
otel_metrics.record_deriver_tokens(
if settings.METRICS.ENABLED:
prometheus_metrics.record_deriver_tokens(
count=new_summary["token_count"],
task_type=DeriverTaskTypes.SUMMARY.value,
token_type=TokenTypes.OUTPUT.value,

View File

@ -1,8 +1,12 @@
import tiktoken
from src.config import settings
from src.telemetry import otel_metrics
from src.telemetry.otel.metrics import DeriverComponents, DeriverTaskTypes, TokenTypes
from src.telemetry import prometheus_metrics
from src.telemetry.prometheus.metrics import (
DeriverComponents,
DeriverTaskTypes,
TokenTypes,
)
tokenizer = tiktoken.get_encoding("o200k_base")
@ -31,9 +35,9 @@ def track_deriver_input_tokens(
components: Dict mapping component names to token counts
"""
for component, token_count in components.items():
# OTel metrics (push-based)
if settings.OTEL.ENABLED:
otel_metrics.record_deriver_tokens(
# Prometheus metrics
if settings.METRICS.ENABLED:
prometheus_metrics.record_deriver_tokens(
count=token_count,
task_type=task_type.value,
token_type=TokenTypes.INPUT.value,

View File

@ -1,23 +1,21 @@
"""Integration tests for OpenTelemetry token metrics tracking.
# pyright: reportPrivateUsage=false, reportUnknownVariableType=false
"""Integration tests for Prometheus token metrics tracking.
These tests verify that deriver and dialectic token metrics are correctly
emitted with accurate token counts when processing messages and dialectic queries.
The approach uses delta-based verification with OTel's InMemoryMetricReader:
The approach uses delta-based verification by directly accessing Prometheus counters:
1. Capture counter values before test execution
2. Run the code under test (with mocked LLM)
3. Verify deltas match expected values
"""
from collections.abc import Iterator
from typing import Any, cast
from unittest.mock import AsyncMock, patch
import pytest
from nanoid import generate as generate_nanoid
from opentelemetry.metrics import Meter
from opentelemetry.sdk.metrics import MeterProvider
from opentelemetry.sdk.metrics.export import InMemoryMetricReader
from prometheus_client import Counter
from sqlalchemy.ext.asyncio import AsyncSession
from src import crud, models, schemas
@ -29,12 +27,15 @@ from src.schemas import (
ResolvedReasoningConfiguration,
ResolvedSummaryConfiguration,
)
from src.telemetry.otel.metrics import otel_metrics
from src.telemetry.prometheus.metrics import (
deriver_tokens_processed_counter,
dialectic_tokens_processed_counter,
)
from src.utils.clients import HonchoLLMCallResponse
from src.utils.representation import ExplicitObservationBase, PromptRepresentation
from src.utils.summarizer import (
SummaryType,
_create_and_save_summary, # pyright: ignore[reportPrivateUsage]
_create_and_save_summary,
estimate_short_summary_prompt_tokens,
)
@ -43,137 +44,70 @@ from src.utils.summarizer import (
# =============================================================================
class OTelMetricChecker:
"""Utility class to capture and verify OTel counter deltas."""
class PrometheusMetricChecker:
"""Utility class to capture and verify Prometheus counter deltas."""
_reader: InMemoryMetricReader
def __init__(self, reader: InMemoryMetricReader):
self._reader = reader
def _get_metric_value(self, metric_name: str, labels: dict[str, str]) -> float:
def _get_counter_value(self, counter: Counter, labels: dict[str, str]) -> float:
"""Get current value of a counter with specific labels.
Note: The OTel SDK's MetricsData types are not fully typed, so we cast
to Any to avoid type warnings when traversing the metrics data structure.
Note: For prometheus_client counters, we access the internal _value
of the labeled metric. This is implementation-specific but works for testing.
"""
raw_data = self._reader.get_metrics_data() # pyright: ignore[reportUnknownVariableType]
if raw_data is None:
try:
labeled_counter = counter.labels(**labels)
return labeled_counter._value.get()
except Exception:
return 0.0
data = cast(Any, raw_data)
for resource_metrics in data.resource_metrics:
for scope_metrics in resource_metrics.scope_metrics:
for metric in scope_metrics.metrics:
if metric.name == metric_name and hasattr(
metric.data, "data_points"
):
for point in metric.data.data_points:
# Check if labels match
point_attrs: dict[str, str] = (
dict(point.attributes) if point.attributes else {}
)
if all(point_attrs.get(k) == v for k, v in labels.items()):
return float(point.value)
return 0.0
def capture(self, metric_name: str, labels: dict[str, str]) -> float:
def capture(self, counter: Counter, labels: dict[str, str]) -> float:
"""Capture current value of a counter with specific labels."""
return self._get_metric_value(metric_name, labels)
return self._get_counter_value(counter, labels)
def get_delta(
self, metric_name: str, labels: dict[str, str], before: float
self, counter: Counter, labels: dict[str, str], before: float
) -> float:
"""Get the delta between a before value and current."""
return self.capture(metric_name, labels) - before
return self.capture(counter, labels) - before
def assert_delta(
self,
metric_name: str,
counter: Counter,
labels: dict[str, str],
before: float,
expected: int | float,
message: str = "",
) -> None:
"""Assert that the delta matches expected value."""
delta = self.get_delta(metric_name, labels, before)
delta = self.get_delta(counter, labels, before)
assert (
delta == expected
), f"{message}: expected delta {expected}, got {delta}. Labels: {labels}"
def _reset_otel_metrics_singleton() -> None:
"""Reset the otel_metrics singleton instance so it reinitializes with a new provider.
This clears all instance attributes so _ensure_initialized() will create
new meters tied to the current global MeterProvider.
"""
# Reset initialization flag (instance attribute shadows class attribute)
otel_metrics._is_initialized = False # pyright: ignore[reportPrivateUsage]
# Reset all meters
otel_metrics._api_meter = None # pyright: ignore[reportPrivateUsage]
otel_metrics._deriver_meter = None # pyright: ignore[reportPrivateUsage]
otel_metrics._dialectic_meter = None # pyright: ignore[reportPrivateUsage]
otel_metrics._dreamer_meter = None # pyright: ignore[reportPrivateUsage]
# Reset all counters
otel_metrics._api_requests = None # pyright: ignore[reportPrivateUsage]
otel_metrics._messages_created = None # pyright: ignore[reportPrivateUsage]
otel_metrics._dialectic_calls = None # pyright: ignore[reportPrivateUsage]
otel_metrics._deriver_queue_items = None # pyright: ignore[reportPrivateUsage]
otel_metrics._deriver_tokens = None # pyright: ignore[reportPrivateUsage]
otel_metrics._dialectic_tokens = None # pyright: ignore[reportPrivateUsage]
otel_metrics._dreamer_tokens = None # pyright: ignore[reportPrivateUsage]
@pytest.fixture
def otel_test_setup(
def prometheus_test_setup(
monkeypatch: pytest.MonkeyPatch,
) -> Iterator[tuple[InMemoryMetricReader, OTelMetricChecker]]:
"""Set up OTel metrics with in-memory reader for testing.
The OTel SDK only allows set_meter_provider() to be called once per process.
To work around this for testing, we patch get_meter() to return meters
from our test provider directly.
) -> Iterator[PrometheusMetricChecker]:
"""Set up Prometheus metrics for testing.
Yields:
Tuple of (reader, checker) for verifying metrics
A PrometheusMetricChecker for verifying metrics
"""
# Create in-memory reader
reader = InMemoryMetricReader()
# Enable METRICS in settings and set namespace for test assertions
monkeypatch.setattr("src.config.settings.METRICS.ENABLED", True)
monkeypatch.setattr("src.config.settings.METRICS.NAMESPACE", "test")
# Create a test meter provider
provider = MeterProvider(metric_readers=[reader])
checker = PrometheusMetricChecker()
# Patch get_meter to return meters from our test provider
# This is necessary because set_meter_provider() can only be called once per process
def test_get_meter(name: str, version: str = "") -> Meter:
return provider.get_meter(name, version)
monkeypatch.setattr("src.telemetry.otel.metrics.get_meter", test_get_meter)
# Reset the otel_metrics singleton instance so it reinitializes with our test provider
_reset_otel_metrics_singleton()
# Enable OTEL in settings and set namespace for test assertions
monkeypatch.setattr("src.config.settings.OTEL.ENABLED", True)
monkeypatch.setattr("src.config.settings.OTEL.SERVICE_NAMESPACE", "test")
checker = OTelMetricChecker(reader)
yield reader, checker
# Cleanup: reset singleton so other tests aren't affected
_reset_otel_metrics_singleton()
yield checker
@pytest.fixture
def metric_checker(
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
) -> OTelMetricChecker:
prometheus_test_setup: PrometheusMetricChecker,
) -> PrometheusMetricChecker:
"""Fixture providing a metric checker instance."""
return otel_test_setup[1]
return prometheus_test_setup
# =============================================================================
@ -285,12 +219,12 @@ class TestDeriverIngestionMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify OUTPUT_TOTAL tokens match response.output_tokens from LLM."""
from src.deriver.deriver import process_representation_tasks_batch
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
messages = await create_test_messages(
@ -309,7 +243,7 @@ class TestDeriverIngestionMetrics:
"token_type": "output",
"component": "output_total",
}
before = metric_checker.capture("deriver_tokens_processed_total", labels)
before = metric_checker.capture(deriver_tokens_processed_counter, labels)
# Mock the LLM call and save_representation (we're testing metrics, not DB writes)
with (
@ -332,7 +266,7 @@ class TestDeriverIngestionMetrics:
# Verify output tokens metric
metric_checker.assert_delta(
"deriver_tokens_processed_total",
deriver_tokens_processed_counter,
labels,
before,
expected_output_tokens,
@ -343,13 +277,13 @@ class TestDeriverIngestionMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify PROMPT component is tracked for ingestion input."""
from src.deriver.deriver import process_representation_tasks_batch
from src.deriver.prompts import estimate_minimal_deriver_prompt_tokens
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
messages = await create_test_messages(
@ -367,7 +301,7 @@ class TestDeriverIngestionMetrics:
"token_type": "input",
"component": "prompt",
}
before = metric_checker.capture("deriver_tokens_processed_total", labels)
before = metric_checker.capture(deriver_tokens_processed_counter, labels)
with (
patch(
@ -388,7 +322,7 @@ class TestDeriverIngestionMetrics:
)
metric_checker.assert_delta(
"deriver_tokens_processed_total",
deriver_tokens_processed_counter,
labels,
before,
expected_prompt_tokens,
@ -399,12 +333,12 @@ class TestDeriverIngestionMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify MESSAGES component is tracked for ingestion input."""
from src.deriver.deriver import process_representation_tasks_batch
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
messages = await create_test_messages(
@ -424,7 +358,7 @@ class TestDeriverIngestionMetrics:
"token_type": "input",
"component": "messages",
}
before = metric_checker.capture("deriver_tokens_processed_total", labels)
before = metric_checker.capture(deriver_tokens_processed_counter, labels)
with (
patch(
@ -446,7 +380,7 @@ class TestDeriverIngestionMetrics:
# Verify messages tokens were tracked (should be > 0)
delta = metric_checker.get_delta(
"deriver_tokens_processed_total", labels, before
deriver_tokens_processed_counter, labels, before
)
assert delta > 0, f"Expected messages input tokens > 0, got {delta}"
@ -464,11 +398,11 @@ class TestDeriverSummaryMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify OUTPUT_TOTAL tokens are tracked for summary."""
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
@ -493,7 +427,7 @@ class TestDeriverSummaryMetrics:
"token_type": "output",
"component": "output_total",
}
before = metric_checker.capture("deriver_tokens_processed_total", labels)
before = metric_checker.capture(deriver_tokens_processed_counter, labels)
with (
patch(
@ -519,7 +453,7 @@ class TestDeriverSummaryMetrics:
# Verify output tokens match the summary token_count
metric_checker.assert_delta(
"deriver_tokens_processed_total",
deriver_tokens_processed_counter,
labels,
before,
expected_output_tokens,
@ -530,11 +464,11 @@ class TestDeriverSummaryMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify PROMPT component is tracked for summary input."""
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
messages = await create_test_messages(
@ -557,7 +491,7 @@ class TestDeriverSummaryMetrics:
"token_type": "input",
"component": "prompt",
}
before = metric_checker.capture("deriver_tokens_processed_total", labels)
before = metric_checker.capture(deriver_tokens_processed_counter, labels)
with (
patch(
@ -582,7 +516,7 @@ class TestDeriverSummaryMetrics:
)
metric_checker.assert_delta(
"deriver_tokens_processed_total",
deriver_tokens_processed_counter,
labels,
before,
expected_prompt_tokens,
@ -593,11 +527,11 @@ class TestDeriverSummaryMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify MESSAGES component is tracked for summary input."""
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
messages = await create_test_messages(
@ -629,7 +563,7 @@ class TestDeriverSummaryMetrics:
"token_type": "input",
"component": "messages",
}
before = metric_checker.capture("deriver_tokens_processed_total", labels)
before = metric_checker.capture(deriver_tokens_processed_counter, labels)
with (
patch(
@ -655,7 +589,7 @@ class TestDeriverSummaryMetrics:
# Verify messages tokens match what summarizer actually computed
delta = metric_checker.get_delta(
"deriver_tokens_processed_total", labels, before
deriver_tokens_processed_counter, labels, before
)
assert (
delta == expected_messages_tokens
@ -666,11 +600,11 @@ class TestDeriverSummaryMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify metrics are NOT emitted when _create_summary returns is_fallback=True."""
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
messages = await create_test_messages(
@ -699,10 +633,10 @@ class TestDeriverSummaryMetrics:
"component": "prompt",
}
before_output = metric_checker.capture(
"deriver_tokens_processed_total", output_labels
deriver_tokens_processed_counter, output_labels
)
before_prompt = metric_checker.capture(
"deriver_tokens_processed_total", prompt_labels
deriver_tokens_processed_counter, prompt_labels
)
with patch(
@ -723,10 +657,10 @@ class TestDeriverSummaryMetrics:
# Verify NO change in metrics when fallback
output_delta = metric_checker.get_delta(
"deriver_tokens_processed_total", output_labels, before_output
deriver_tokens_processed_counter, output_labels, before_output
)
prompt_delta = metric_checker.get_delta(
"deriver_tokens_processed_total", prompt_labels, before_prompt
deriver_tokens_processed_counter, prompt_labels, before_prompt
)
assert (
@ -750,12 +684,12 @@ class TestDialecticTokenMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify INPUT tokens are tracked from LLM response."""
from src.dialectic.core import DialecticAgent
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
@ -770,7 +704,7 @@ class TestDialecticTokenMetrics:
"component": "total",
"reasoning_level": "low",
}
before = metric_checker.capture("dialectic_tokens_processed_total", labels)
before = metric_checker.capture(dialectic_tokens_processed_counter, labels)
agent = DialecticAgent(
db=db_session,
@ -787,7 +721,7 @@ class TestDialecticTokenMetrics:
await agent.answer("What do you know about this user?")
metric_checker.assert_delta(
"dialectic_tokens_processed_total",
dialectic_tokens_processed_counter,
labels,
before,
expected_input_tokens,
@ -798,12 +732,12 @@ class TestDialecticTokenMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
):
"""Verify OUTPUT tokens are tracked from LLM response."""
from src.dialectic.core import DialecticAgent
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
@ -818,7 +752,7 @@ class TestDialecticTokenMetrics:
"component": "total",
"reasoning_level": "low",
}
before = metric_checker.capture("dialectic_tokens_processed_total", labels)
before = metric_checker.capture(dialectic_tokens_processed_counter, labels)
agent = DialecticAgent(
db=db_session,
@ -835,7 +769,7 @@ class TestDialecticTokenMetrics:
await agent.answer("What do you know about this user?")
metric_checker.assert_delta(
"dialectic_tokens_processed_total",
dialectic_tokens_processed_counter,
labels,
before,
expected_output_tokens,
@ -846,16 +780,16 @@ class TestDialecticTokenMetrics:
self,
db_session: AsyncSession,
sample_data: tuple[Workspace, Peer],
otel_test_setup: tuple[InMemoryMetricReader, OTelMetricChecker],
prometheus_test_setup: PrometheusMetricChecker,
monkeypatch: pytest.MonkeyPatch,
):
"""Verify metrics are NOT emitted when OTEL.ENABLED=False."""
"""Verify metrics are NOT emitted when PROMETHEUS.ENABLED=False."""
from src.dialectic.core import DialecticAgent
_, metric_checker = otel_test_setup
metric_checker = prometheus_test_setup
# Explicitly disable OTEL metrics
monkeypatch.setattr("src.config.settings.OTEL.ENABLED", False)
# Explicitly disable Prometheus metrics
monkeypatch.setattr("src.config.settings.METRICS.ENABLED", False)
workspace, peer = sample_data
session = await create_test_session_with_peer(db_session, workspace, peer)
@ -878,10 +812,10 @@ class TestDialecticTokenMetrics:
"reasoning_level": "low",
}
before_input = metric_checker.capture(
"dialectic_tokens_processed_total", input_labels
dialectic_tokens_processed_counter, input_labels
)
before_output = metric_checker.capture(
"dialectic_tokens_processed_total", output_labels
dialectic_tokens_processed_counter, output_labels
)
agent = DialecticAgent(
@ -900,10 +834,10 @@ class TestDialecticTokenMetrics:
# Verify NO change in metrics
input_delta = metric_checker.get_delta(
"dialectic_tokens_processed_total", input_labels, before_input
dialectic_tokens_processed_counter, input_labels, before_input
)
output_delta = metric_checker.get_delta(
"dialectic_tokens_processed_total", output_labels, before_output
dialectic_tokens_processed_counter, output_labels, before_output
)
assert (

View File

@ -337,7 +337,7 @@ def mock_telemetry_settings():
mock_settings.TELEMETRY.MAX_RETRIES = max_retries
mock_settings.TELEMETRY.MAX_BUFFER_SIZE = max_buffer_size
mock_settings.TELEMETRY.HEADERS = headers
mock_settings.OTEL.ENABLED = False
mock_settings.METRICS.ENABLED = False
return patch("src.telemetry.emitter.settings", mock_settings)
return _configure

View File

@ -350,56 +350,6 @@ class TestShutdownTelemetryEvents:
mock_shutdown.assert_called_once()
# =============================================================================
# Tests for initialize_telemetry() function
# =============================================================================
class TestInitializeTelemetry:
"""Tests for initialize_telemetry() in src.telemetry."""
def test_initialize_otel_when_enabled(self):
"""initialize_telemetry() initializes OTel metrics when enabled."""
import src.telemetry as telemetry_module
with (
patch("src.config.settings") as mock_settings,
# Patch where it's used (in src.telemetry), not where it's defined
patch.object(telemetry_module, "initialize_otel_metrics") as mock_otel_init,
):
mock_settings.OTEL.ENABLED = True
mock_settings.OTEL.ENDPOINT = "http://otel:9009/metrics"
mock_settings.OTEL.HEADERS = None
mock_settings.OTEL.EXPORT_INTERVAL_MILLIS = 60000
mock_settings.OTEL.SERVICE_NAME = "honcho"
mock_settings.OTEL.SERVICE_NAMESPACE = "test"
telemetry_module.initialize_telemetry()
mock_otel_init.assert_called_once_with(
endpoint="http://otel:9009/metrics",
headers=None,
export_interval_millis=60000,
service_name="honcho",
service_namespace="test",
enabled=True,
)
def test_skip_otel_when_disabled(self):
"""initialize_telemetry() skips OTel when disabled."""
import src.telemetry as telemetry_module
with (
patch("src.config.settings") as mock_settings,
patch.object(telemetry_module, "initialize_otel_metrics") as mock_otel_init,
):
mock_settings.OTEL.ENABLED = False
telemetry_module.initialize_telemetry()
mock_otel_init.assert_not_called()
# =============================================================================
# Tests for initialize_telemetry_async() function
# =============================================================================
@ -454,17 +404,13 @@ class TestShutdownTelemetry:
"""Tests for shutdown_telemetry() in src.telemetry."""
@pytest.mark.asyncio
async def test_shutdown_calls_all_subsystems(self):
"""shutdown_telemetry() shuts down all telemetry subsystems."""
async def test_shutdown_calls_cloudevents(self):
"""shutdown_telemetry() shuts down CloudEvents emitter."""
from src.telemetry import shutdown_telemetry
with (
patch(
"src.telemetry.events.shutdown_telemetry_events", new_callable=AsyncMock
) as mock_ce_shutdown,
patch("src.telemetry.shutdown_otel_metrics") as mock_otel_shutdown,
):
with patch(
"src.telemetry.events.shutdown_telemetry_events", new_callable=AsyncMock
) as mock_ce_shutdown:
await shutdown_telemetry()
mock_ce_shutdown.assert_called_once()
mock_otel_shutdown.assert_called_once()

15
uv.lock
View File

@ -1042,10 +1042,9 @@ dependencies = [
{ name = "langfuse" },
{ name = "nanoid" },
{ name = "openai" },
{ name = "opentelemetry-exporter-otlp-proto-http" },
{ name = "opentelemetry-sdk" },
{ name = "pdfplumber" },
{ name = "pgvector" },
{ name = "prometheus-client" },
{ name = "psycopg", extra = ["binary"] },
{ name = "pydantic" },
{ name = "pydantic-settings" },
@ -1098,10 +1097,9 @@ requires-dist = [
{ name = "langfuse", specifier = ">=3.3.2" },
{ name = "nanoid", specifier = ">=2.0.0" },
{ name = "openai", specifier = ">=1.99.7" },
{ name = "opentelemetry-exporter-otlp-proto-http", specifier = ">=1.36.0" },
{ name = "opentelemetry-sdk", specifier = ">=1.36.0" },
{ name = "pdfplumber", specifier = ">=0.11.7" },
{ name = "pgvector", specifier = ">=0.2.5" },
{ name = "prometheus-client", specifier = ">=0.21.0" },
{ name = "psycopg", extras = ["binary"], specifier = ">=3.1.19" },
{ name = "pydantic", specifier = ">=2.11.7" },
{ name = "pydantic-settings", specifier = ">=2.10.1" },
@ -2127,6 +2125,15 @@ wheels = [
{ url = "https://files.pythonhosted.org/packages/88/74/a88bf1b1efeae488a0c0b7bdf71429c313722d1fc0f377537fbe554e6180/pre_commit-4.2.0-py2.py3-none-any.whl", hash = "sha256:a009ca7205f1eb497d10b845e52c838a98b6cdd2102a6c8e4540e94ee75c58bd", size = 220707, upload-time = "2025-03-18T21:35:19.343Z" },
]
[[package]]
name = "prometheus-client"
version = "0.24.1"
source = { registry = "https://pypi.org/simple" }
sdist = { url = "https://files.pythonhosted.org/packages/f0/58/a794d23feb6b00fc0c72787d7e87d872a6730dd9ed7c7b3e954637d8f280/prometheus_client-0.24.1.tar.gz", hash = "sha256:7e0ced7fbbd40f7b84962d5d2ab6f17ef88a72504dcf7c0b40737b43b2a461f9", size = 85616, upload-time = "2026-01-14T15:26:26.965Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/74/c3/24a2f845e3917201628ecaba4f18bab4d18a337834c1df2a159ee9d22a42/prometheus_client-0.24.1-py3-none-any.whl", hash = "sha256:150db128af71a5c2482b36e588fc8a6b95e498750da4b17065947c16070f4055", size = 64057, upload-time = "2026-01-14T15:26:24.42Z" },
]
[[package]]
name = "propcache"
version = "0.4.1"