diff --git a/tests/test_portal_metrics.py b/tests/test_portal_metrics.py new file mode 100644 index 000000000..0b60a2353 --- /dev/null +++ b/tests/test_portal_metrics.py @@ -0,0 +1,221 @@ +# Copyright (c) 2025 Beijing Volcano Engine Technology Co., Ltd. and/or its affiliates. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +import time +from types import SimpleNamespace + +import pytest +from opentelemetry import metrics as metrics_api +from opentelemetry.metrics import _internal as metrics_internal +from opentelemetry.sdk import metrics as metrics_sdk +from opentelemetry.sdk.metrics.export import InMemoryMetricReader +from opentelemetry.util._once import Once + +from veadk.tools.skills_tools.skills_tool import SkillsTool +from veadk.tracing.telemetry import portal_metrics, telemetry +from veadk.tracing.telemetry.exporters import ( + apmplus_exporter as apmplus_exporter_module, +) +from veadk.tracing.telemetry.exporters.apmplus_exporter import ( + APMPlusExporter, + APMPlusExporterConfig, + ensure_apmplus_meter_provider, +) +from veadk.tracing.telemetry.portal_metrics import PortalMetricRecorder + + +class FakePortalMetricRecorder: + def __init__(self): + self.llm_calls = [] + self.tool_calls = [] + self.skill_calls = [] + + def record_call_llm(self, *args): + self.llm_calls.append(args) + + def record_tool_call(self, *args): + self.tool_calls.append(args) + + def record_skill_call(self, *args): + self.skill_calls.append(args) + + +@pytest.fixture +def fresh_global_meter_provider(monkeypatch): + proxy_provider = metrics_internal._PROXY_METER_PROVIDER + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) + monkeypatch.setattr(proxy_provider, "_real_meter_provider", None) + monkeypatch.setattr(proxy_provider, "_meters", []) + + yield + + provider = metrics_internal._METER_PROVIDER + if isinstance(provider, metrics_sdk.MeterProvider): + provider.shutdown() + + +def test_portal_metrics_are_recorded_without_apmplus_exporter( + monkeypatch, +): + recorder = FakePortalMetricRecorder() + monkeypatch.setattr(portal_metrics, "portal_metric_recorder", recorder) + + telemetry._upload_call_llm_metrics(None, "event", "request", "response") + telemetry._upload_tool_call_metrics("tool", {}, "response") + + skills_tool = SkillsTool({}) + span = SimpleNamespace(start_time=time.time_ns()) + skills_tool._upload_skill_metrics(span, "missing-skill", "ok") + + assert len(recorder.llm_calls) == 1 + assert len(recorder.tool_calls) == 1 + assert len(recorder.skill_calls) == 1 + + +def test_portal_recorder_uses_default_proxy_without_installing_provider( + fresh_global_meter_provider, +): + default_provider = metrics_api.get_meter_provider() + + recorder = PortalMetricRecorder(name="test-default-proxy") + recorder.llm_invoke_counter.add(1) + + assert isinstance(default_provider, metrics_internal._ProxyMeterProvider) + assert metrics_api.get_meter_provider() is default_provider + + +def test_portal_recorder_uses_preconfigured_global_provider( + fresh_global_meter_provider, +): + reader = InMemoryMetricReader() + provider = metrics_sdk.MeterProvider(metric_readers=[reader]) + metrics_api.set_meter_provider(provider) + + recorder = PortalMetricRecorder(name="test-preconfigured-provider") + recorder.llm_invoke_counter.add(1) + provider.force_flush() + + assert recorder.provider is provider + assert reader.get_metrics_data() is not None + + +def test_proxy_instruments_follow_provider_installed_later( + fresh_global_meter_provider, +): + recorder = PortalMetricRecorder(name="test-late-provider") + recorder.llm_invoke_counter.add(1, {"phase": "before"}) + + reader = InMemoryMetricReader() + provider = metrics_sdk.MeterProvider(metric_readers=[reader]) + metrics_api.set_meter_provider(provider) + recorder.llm_invoke_counter.add(2, {"phase": "after"}) + provider.force_flush() + + metrics_data = reader.get_metrics_data() + points = [ + point + for resource_metrics in metrics_data.resource_metrics + for scope_metrics in resource_metrics.scope_metrics + for metric in scope_metrics.metrics + if metric.name == "gen_ai.chat.count" + for point in metric.data.data_points + ] + + assert [(dict(point.attributes), point.value) for point in points] == [ + ({"phase": "after"}, 2) + ] + + +def test_apmplus_reuses_preconfigured_global_provider( + fresh_global_meter_provider, + monkeypatch, +): + provider = metrics_sdk.MeterProvider() + metrics_api.set_meter_provider(provider) + + def fail_if_exporter_is_created(**kwargs): + raise AssertionError("APMPlus must not modify an existing global provider") + + monkeypatch.setattr( + apmplus_exporter_module, + "OTLPMetricExporter", + fail_if_exporter_is_created, + ) + + resolved_provider = ensure_apmplus_meter_provider( + endpoint="http://localhost:4319", + headers={"x-byteapm-appkey": "test"}, + resource_attributes={"service.name": "test-service"}, + ) + + assert resolved_provider is provider + assert provider._sdk_config.metric_readers == () + + +def test_apmplus_installs_global_provider_when_none_exists( + fresh_global_meter_provider, + monkeypatch, +): + reader = InMemoryMetricReader() + exporter_kwargs = {} + monkeypatch.setattr( + apmplus_exporter_module, + "OTLPMetricExporter", + lambda **kwargs: exporter_kwargs.update(kwargs) or object(), + ) + monkeypatch.setattr( + apmplus_exporter_module, + "PeriodicExportingMetricReader", + lambda exporter: reader, + ) + + provider = ensure_apmplus_meter_provider( + endpoint="http://localhost:4319", + headers={"x-byteapm-appkey": "test"}, + resource_attributes={"service.name": "test-service"}, + ) + + assert metrics_api.get_meter_provider() is provider + assert provider._sdk_config.metric_readers == [reader] + assert exporter_kwargs == { + "endpoint": "http://localhost:4319", + "headers": {"x-byteapm-appkey": "test"}, + "insecure": True, + } + + +def test_apmplus_exporter_only_bootstraps_meter_provider(monkeypatch): + bootstrap_calls = [] + monkeypatch.setattr( + apmplus_exporter_module, + "ensure_apmplus_meter_provider", + lambda **kwargs: bootstrap_calls.append(kwargs), + ) + + APMPlusExporter( + config=APMPlusExporterConfig( + endpoint="http://localhost:4319", + app_key="test-app-key", + service_name="test-service", + ) + ) + + assert bootstrap_calls == [ + { + "endpoint": "http://localhost:4319", + "headers": {"x-byteapm-appkey": "test-app-key"}, + "resource_attributes": {"service.name": "test-service"}, + } + ] diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 5f7b2c5be..0519b6aeb 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -31,7 +31,6 @@ from veadk.tracing.telemetry import ( opentelemetry_tracer as opentelemetry_tracer_module, ) -from veadk.tracing.telemetry import telemetry as telemetry_module from veadk.tracing.telemetry.exporters import ( apmplus_exporter as apmplus_exporter_module, ) @@ -112,7 +111,6 @@ def fresh_global_tracer_provider(monkeypatch): def controlled_apmplus_exporter(monkeypatch): """Provide an APMPlus exporter without network or global meter side effects.""" constructed_exporters = [] - monkeypatch.setattr(telemetry_module, "meter_uploader", None) class ControlledAPMPlusExporter(APMPlusExporter): def __init__(self): @@ -127,7 +125,6 @@ def __init__(self): def model_post_init(self, context): self._exporter = OTelInMemorySpanExporter() self.processor = SimpleSpanProcessor(self._exporter) - self.meter_uploader = object() constructed_exporters.append(self) monkeypatch.setattr( @@ -148,26 +145,34 @@ def model_post_init(self, context): [False, True], ids=["env-disabled", "env-enabled"], ) +@pytest.mark.parametrize( + "provider_preconfigured", + [False, True], + ids=["provider-proxy", "provider-preconfigured"], +) @pytest.mark.parametrize( "manual_exporter", [False, True], ids=["no-manual-exporter", "manual-exporter"], ) -def test_apmplus_preconfigured_provider_matrix( +def test_apmplus_enable_provider_and_manual_exporter_matrix( fresh_global_tracer_provider, controlled_apmplus_exporter, monkeypatch, enable_apmplus, + provider_preconfigured, manual_exporter, ): - """A preconfigured provider owns traces; env exporter retains metrics.""" + """Validate env, global provider, and explicit exporter independently.""" controlled_exporter_class, constructed_exporters = controlled_apmplus_exporter monkeypatch.setenv("ENABLE_APMPLUS", str(enable_apmplus).lower()) monkeypatch.setenv("ENABLE_COZELOOP", "false") monkeypatch.setenv("ENABLE_TLS", "false") - tracer_provider = trace_sdk.TracerProvider() - trace_api.set_tracer_provider(tracer_provider) + initial_provider = trace_api.get_tracer_provider() + if provider_preconfigured: + initial_provider = trace_sdk.TracerProvider() + trace_api.set_tracer_provider(initial_provider) tracers = [] if manual_exporter: @@ -178,27 +183,82 @@ def test_apmplus_preconfigured_provider_matrix( should_create_tracer = manual_exporter or enable_apmplus assert len(agent.tracers) == int(should_create_tracer) - assert trace_api.get_tracer_provider() is tracer_provider - assert len(constructed_exporters) == int(manual_exporter) + int(enable_apmplus) - span_processors = tracer_provider._active_span_processor._span_processors + final_provider = trace_api.get_tracer_provider() if not should_create_tracer: - assert span_processors == () + assert final_provider is initial_provider + assert constructed_exporters == [] return tracer = agent.tracers[0] + assert isinstance(final_provider, trace_sdk.TracerProvider) + if provider_preconfigured: + assert final_provider is initial_provider + else: + assert final_provider is not initial_provider + + should_register_apmplus = not provider_preconfigured and ( + manual_exporter or enable_apmplus + ) + should_construct_apmplus = manual_exporter or enable_apmplus + assert len(constructed_exporters) == int(should_construct_apmplus) assert sum( isinstance(exporter, controlled_exporter_class) for exporter in tracer.exporters - ) == int(enable_apmplus) - assert all( - exporter.processor not in span_processors for exporter in constructed_exporters + ) == int(should_construct_apmplus) + + span_processors = final_provider._active_span_processor._span_processors + registered_apmplus_processors = sum( + any(processor is exporter.processor for processor in span_processors) + for exporter in constructed_exporters ) - assert len(span_processors) == 1 # VeADK in-memory processor only - assert tracer.apmplus_managed_externally is True - expected_meter_uploader = ( - constructed_exporters[-1].meter_uploader if enable_apmplus else None + assert registered_apmplus_processors == int(should_register_apmplus) + assert len(span_processors) == 1 + int(should_register_apmplus) + assert tracer.apmplus_managed_externally is provider_preconfigured + + +def test_add_exporter_registers_after_tracer_initialization( + fresh_global_tracer_provider, +): + tracer = OpentelemetryTracer() + exporter = init_exporters()[1] + + tracer.add_exporter(exporter) + + tracer_provider = trace_api.get_tracer_provider() + span_processors = tracer_provider._active_span_processor._span_processors + assert exporter in tracer.exporters + assert exporter.processor in span_processors + assert exporter.processor in tracer._processors + + +def test_add_exporter_is_idempotent(fresh_global_tracer_provider): + tracer = OpentelemetryTracer() + exporter = init_exporters()[1] + + tracer.add_exporter(exporter) + tracer.add_exporter(exporter) + + tracer_provider = trace_api.get_tracer_provider() + span_processors = tracer_provider._active_span_processor._span_processors + assert sum(processor is exporter.processor for processor in span_processors) == 1 + assert sum(processor is exporter.processor for processor in tracer._processors) == 1 + + +def test_force_export_flushes_portal_metrics( + fresh_global_tracer_provider, + monkeypatch, +): + flush_calls = [] + monkeypatch.setattr( + opentelemetry_tracer_module, + "portal_metric_recorder", + SimpleNamespace(force_flush=lambda: flush_calls.append(True) or True), ) - assert telemetry_module.meter_uploader is expected_meter_uploader + tracer = OpentelemetryTracer() + + tracer.force_export() + + assert flush_calls == [True] def test_tracing_registers_apmplus_without_global_provider( @@ -214,6 +274,63 @@ def test_tracing_registers_apmplus_without_global_provider( assert apmplus_exporter.processor in span_processors +def test_tracing_skips_apmplus_for_preconfigured_provider( + fresh_global_tracer_provider, +): + tracer_provider = trace_sdk.TracerProvider() + trace_api.set_tracer_provider(tracer_provider) + apmplus_exporter = init_apmplus_exporter() + + tracer = OpentelemetryTracer(exporters=[apmplus_exporter]) + global_tracer_provider = trace_api.get_tracer_provider() + span_processors = global_tracer_provider._active_span_processor._span_processors + + assert global_tracer_provider is tracer_provider + assert apmplus_exporter in tracer.exporters + assert apmplus_exporter.processor not in span_processors + assert len(span_processors) == 1 # VeADK in-memory processor only + + +def test_add_exporter_skips_apmplus_for_preconfigured_provider( + fresh_global_tracer_provider, +): + tracer_provider = trace_sdk.TracerProvider() + trace_api.set_tracer_provider(tracer_provider) + tracer = OpentelemetryTracer() + apmplus_exporter = init_apmplus_exporter() + + tracer.add_exporter(apmplus_exporter) + + span_processors = tracer_provider._active_span_processor._span_processors + assert apmplus_exporter in tracer.exporters + assert apmplus_exporter.processor not in span_processors + assert len(span_processors) == 1 # VeADK in-memory processor only + + +def test_agent_env_keeps_apmplus_for_preconfigured_provider( + fresh_global_tracer_provider, + controlled_apmplus_exporter, + monkeypatch, +): + controlled_exporter_class, constructed_exporters = controlled_apmplus_exporter + tracer_provider = trace_sdk.TracerProvider() + trace_api.set_tracer_provider(tracer_provider) + tracer = OpentelemetryTracer() + monkeypatch.setenv("ENABLE_APMPLUS", "true") + monkeypatch.setenv("ENABLE_COZELOOP", "false") + monkeypatch.setenv("ENABLE_TLS", "false") + + Agent._prepare_tracers(SimpleNamespace(tracers=[tracer])) + + span_processors = tracer_provider._active_span_processor._span_processors + assert len(constructed_exporters) == 1 + exporter = constructed_exporters[0] + assert isinstance(exporter, controlled_exporter_class) + assert exporter in tracer.exporters + assert exporter.processor not in span_processors + assert len(span_processors) == 1 # VeADK in-memory processor only + + @pytest.mark.asyncio async def test_tracing(fresh_global_tracer_provider): exporters = init_exporters() @@ -232,10 +349,11 @@ async def test_tracing_with_global_provider(fresh_global_tracer_provider): tracer_provider = trace_api.get_tracer_provider() tracer_provider.add_span_processor(gen_span_processor("http://localhost:8000")) trace_api.set_tracer_provider(tracer_provider) - # tracer = OpentelemetryTracer(exporters=exporters) - assert len(tracer.exporters) == 3 # APMPlus is managed by the existing provider + # APMPlus is retained so it can ensure a global metric pipeline, while its + # span processor is not registered. + assert len(tracer.exporters) == 4 @pytest.mark.asyncio @@ -249,5 +367,5 @@ async def test_tracing_with_apmplus_global_provider(fresh_global_tracer_provider # init OpentelemetryTracer tracer = OpentelemetryTracer(exporters=exporters) - # apmplus exporter won't init again, so there are cozeloop, tls, in_memory exporter - assert len(tracer.exporters) == 3 # with extra 1 built-in exporters + # APMPlus is retained for metrics but its span processor is not registered. + assert len(tracer.exporters) == 4 # with extra 1 built-in exporters diff --git a/tests/test_tracing_content.py b/tests/test_tracing_content.py index 333f7bd81..4b01edf3c 100644 --- a/tests/test_tracing_content.py +++ b/tests/test_tracing_content.py @@ -21,7 +21,7 @@ from veadk.config import settings from veadk.tracing.telemetry import telemetry from veadk.tracing.telemetry.content_tracing import should_trace_content -from veadk.tracing.telemetry.exporters.apmplus_exporter import MeterUploader +from veadk.tracing.telemetry.portal_metrics import PortalMetricRecorder @dataclass @@ -143,10 +143,6 @@ def _event_names(span): return [event.name for event in span.events] -def setup_function(): - telemetry.meter_uploader = None - - def test_trace_call_llm_records_content_by_default(monkeypatch): monkeypatch.delenv("OBSERVABILITY_OPENTELEMETRY_TRACE_CONTENT", raising=False) @@ -217,19 +213,19 @@ def test_trace_tool_call_skips_content_when_env_false(monkeypatch): def test_apmplus_tool_metrics_skip_token_usage_when_tool_content_missing(): - meter_uploader = object.__new__(MeterUploader) - meter_uploader.apmplus_span_latency = _FakeMetricRecorder() - meter_uploader.apmplus_tool_token_usage = _FakeMetricRecorder() + metric_recorder = object.__new__(PortalMetricRecorder) + metric_recorder.apmplus_span_latency = _FakeMetricRecorder() + metric_recorder.apmplus_tool_token_usage = _FakeMetricRecorder() with _start_test_span("execute_tool lookup"): - meter_uploader.record_tool_call( + metric_recorder.record_tool_call( _FakeTool(), {"query": "tool input secret"}, _ExplodingFunctionResponseEvent(), ) - assert len(meter_uploader.apmplus_span_latency.records) == 1 - assert meter_uploader.apmplus_tool_token_usage.records == [] + assert len(metric_recorder.apmplus_span_latency.records) == 1 + assert metric_recorder.apmplus_tool_token_usage.records == [] def test_agent_root_span_skips_content_when_env_false(monkeypatch): diff --git a/veadk/agent.py b/veadk/agent.py index fd83dd305..287643ed7 100644 --- a/veadk/agent.py +++ b/veadk/agent.py @@ -671,30 +671,23 @@ def _prepare_tracers(self): if enable_apmplus_tracer and not any( isinstance(e, APMPlusExporter) for e in exporters ): - self.tracers[0].exporters.append(APMPlusExporter()) # type: ignore + self.tracers[0].add_exporter(APMPlusExporter()) # type: ignore logger.info("Enable APMPlus exporter by env.") if enable_cozeloop_tracer and not any( isinstance(e, CozeloopExporter) for e in exporters ): - self.tracers[0].exporters.append(CozeloopExporter()) # type: ignore + self.tracers[0].add_exporter(CozeloopExporter()) # type: ignore logger.info("Enable CozeLoop exporter by env.") if enable_tls_tracer and not any(isinstance(e, TLSExporter) for e in exporters): - self.tracers[0].exporters.append(TLSExporter()) # type: ignore + self.tracers[0].add_exporter(TLSExporter()) # type: ignore logger.info("Enable TLS exporter by env.") logger.debug( f"Opentelemetry Tracer init {len(self.tracers[0].exporters)} exporters" # type: ignore ) - # Initialize global meter_uploader from exporters - from veadk.tracing.telemetry.telemetry import ( - init_global_meter_uploader_from_exporters, - ) - - init_global_meter_uploader_from_exporters(self.tracers[0].exporters) # type: ignore - @property def _llm_flow(self) -> BaseLlmFlow: from google.adk.flows.llm_flows.auto_flow import AutoFlow diff --git a/veadk/tools/skills_tools/skills_tool.py b/veadk/tools/skills_tools/skills_tool.py index 6b4a2e40f..7b491b8a8 100644 --- a/veadk/tools/skills_tools/skills_tool.py +++ b/veadk/tools/skills_tools/skills_tool.py @@ -530,39 +530,14 @@ def _add_skill_span_attributes( def _upload_skill_metrics(self, span: _Span, skill_name: str, result: str) -> None: """Upload skill metrics to the telemetry system.""" try: - import time - from veadk.tracing.telemetry.telemetry import meter_uploader - - if meter_uploader: - # 初始化属性,包含技能相关信息 - skill = self.skills.get(skill_name) - attributes = { - "skill_name": skill_name, - "tool_name": self.name, - "skill_space_id": ( - skill.skill_space_id if skill and skill.skill_space_id else "" - ), - "skill_id": skill.id if skill and skill.id else "", - "gen_ai.operation.name": "execute_skill", - "error_type": ( - "skill_execution_error" if result.startswith("Error:") else "" - ), - } - - # 计算 span 执行耗时(秒) - latency_seconds = 0 - if hasattr(span, "start_time"): - # 计算耗时(秒) - latency_seconds = (time.time_ns() - span.start_time) / 1e9 # type: ignore - - # 记录技能执行延迟 - if hasattr(meter_uploader, "skill_invoke_latency"): - # 使用 skill_invoke_latency 记录技能执行延迟(秒) - meter_uploader.skill_invoke_latency.record( - latency_seconds, attributes - ) - logger.debug( - f"Uploaded skill metrics for {skill_name} with latency {latency_seconds:.4f}s and attributes {attributes}" - ) + from veadk.tracing.telemetry import portal_metrics + + portal_metrics.portal_metric_recorder.record_skill_call( + span, + skill_name, + self.name, + self.skills.get(skill_name), + result, + ) except Exception as e: logger.warning(f"Failed to upload skill metrics: {e}") diff --git a/veadk/tracing/telemetry/exporters/apmplus_exporter.py b/veadk/tracing/telemetry/exporters/apmplus_exporter.py index 2b096aa4a..f8f4791fb 100644 --- a/veadk/tracing/telemetry/exporters/apmplus_exporter.py +++ b/veadk/tracing/telemetry/exporters/apmplus_exporter.py @@ -12,21 +12,12 @@ # See the License for the specific language governing permissions and # limitations under the License. -import time -from dataclasses import dataclass from typing import Any -from google.adk.agents.invocation_context import InvocationContext -from google.adk.agents.run_config import StreamingMode -from google.adk.events import Event -from google.adk.models.llm_request import LlmRequest -from google.adk.models.llm_response import LlmResponse -from google.adk.tools import BaseTool -from opentelemetry import metrics, trace from opentelemetry import metrics as metrics_api from opentelemetry.exporter.otlp.proto.grpc.metric_exporter import OTLPMetricExporter from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter -from opentelemetry.metrics._internal import Meter +from opentelemetry.metrics._internal import _ProxyMeterProvider from opentelemetry.sdk import metrics as metrics_sdk from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader from opentelemetry.sdk.resources import Resource @@ -41,416 +32,35 @@ logger = get_logger(__name__) -_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS = [ - 0.01, - 0.02, - 0.04, - 0.08, - 0.16, - 0.32, - 0.64, - 1.28, - 2.56, - 5.12, - 10.24, - 20.48, - 40.96, - 81.92, -] - -_GEN_AI_SERVER_TIME_PER_OUTPUT_TOKEN_BUCKETS = [ - 0.01, - 0.025, - 0.05, - 0.075, - 0.1, - 0.15, - 0.2, - 0.3, - 0.4, - 0.5, - 0.75, - 1.0, - 2.5, -] - -_GEN_AI_SERVER_TIME_TO_FIRST_TOKEN_BUCKETS = [ - 0.001, - 0.005, - 0.01, - 0.02, - 0.04, - 0.06, - 0.08, - 0.1, - 0.25, - 0.5, - 0.75, - 1.0, - 2.5, - 5.0, - 7.5, - 10.0, -] - -_GEN_AI_CLIENT_TOKEN_USAGE_BUCKETS = [ - 1, - 4, - 16, - 64, - 256, - 1024, - 4096, - 16384, - 65536, - 262144, - 1048576, - 4194304, - 16777216, - 67108864, -] - - -@dataclass -class Meters: - """Metric names and identifiers for OpenTelemetry instrumentation. - - This class defines standardized metric names used for LLM and agent - observability. The metrics follow OpenTelemetry semantic conventions - for generative AI operations and include custom APMPlus metrics for - enhanced monitoring capabilities. - - Standard Gen AI Metrics: - - LLM_CHAT_COUNT: Counter for LLM invocation frequency - - LLM_TOKEN_USAGE: Histogram for token consumption analysis - - LLM_OPERATION_DURATION: Histogram for operation latency tracking - - LLM_COMPLETIONS_EXCEPTIONS: Counter for error rate monitoring - - Streaming metrics: Performance analysis for streaming responses - - APMPlus Custom Metrics: - - APMPLUS_SPAN_LATENCY: Span execution time for performance analysis - - APMPLUS_TOOL_TOKEN_USAGE: Tool-specific token consumption tracking - """ - - LLM_CHAT_COUNT = "gen_ai.chat.count" - LLM_TOKEN_USAGE = "gen_ai.client.token.usage" - LLM_OPERATION_DURATION = "gen_ai.client.operation.duration" - LLM_COMPLETIONS_EXCEPTIONS = "gen_ai.chat_completions.exceptions" - LLM_STREAMING_TIME_TO_FIRST_TOKEN = ( - "gen_ai.chat_completions.streaming_time_to_first_token" - ) - LLM_STREAMING_TIME_TO_GENERATE = ( - "gen_ai.chat_completions.streaming_time_to_generate" +def ensure_apmplus_meter_provider( + endpoint: str, + headers: dict, + resource_attributes: dict, +) -> metrics_api.MeterProvider: + """Install an APMPlus-backed global MeterProvider only when none exists.""" + global_provider = metrics_api.get_meter_provider() + if not isinstance(global_provider, _ProxyMeterProvider): + return global_provider + + exporter = OTLPMetricExporter( + endpoint=endpoint, + headers=headers, + insecure=True, ) - LLM_STREAMING_TIME_PER_OUTPUT_TOKEN = ( - "gen_ai.chat_completions.streaming_time_per_output_token" + metric_reader = PeriodicExportingMetricReader(exporter) + provider = metrics_sdk.MeterProvider( + metric_readers=[metric_reader], + resource=Resource.create(resource_attributes), ) + metrics_api.set_meter_provider(provider) + resolved_provider = metrics_api.get_meter_provider() - # apmplus metrics - # span duration - APMPLUS_SPAN_LATENCY = "apmplus_span_latency" - # tool token usage - APMPLUS_TOOL_TOKEN_USAGE = "apmplus_tool_token_usage" - # skill invoke latency - GEN_AI_SKILL_INVOKE_LATENCY = "gen_ai_skill_invoke_latency" - - -class MeterUploader: - """Metrics uploader for APMPlus observability platform integration. - - MeterUploader manages the collection and transmission of telemetry metrics - to Volcengine's APMPlus platform. It creates and maintains OpenTelemetry - metric instruments for comprehensive agent performance monitoring. - - Key Features: - - Automatic metric instrument creation with appropriate buckets - - LLM call metrics including token usage and latency - - Tool execution metrics for performance analysis - - Error tracking and exception monitoring - - Integration with OpenTelemetry metrics SDK - - Metrics Collected: - - LLM invocation counts and frequencies - - Token consumption (input/output) with histogram distribution - - Operation latency with performance bucket analysis - - Error rates and exception details - - Span-level performance metrics for APMPlus dashboards - """ - - def __init__( - self, name: str, endpoint: str, headers: dict, resource_attributes: dict - ) -> None: - """Initialize the meter uploader with APMPlus configuration. + if resolved_provider is not provider: + # Another component won the set-once race. Keep its provider and close + # the unused APMPlus pipeline created by this call. + provider.shutdown() - Sets up the global metrics provider, creates metric instruments, - and configures OTLP export to APMPlus endpoints with proper - resource attribution and authentication. - - Args: - name: Meter name for identification and organization - endpoint: APMPlus OTLP endpoint URL for metric transmission - headers: Authentication headers including APMPlus app key - resource_attributes: Service metadata for metric attribution - """ - # global_metrics_provider -> global_tracer_provider - # exporter -> exporter - # metric_reader -> processor - global_metrics_provider = metrics_api.get_meter_provider() - - # 1. init resource - if hasattr(global_metrics_provider, "_sdk_config"): - global_resource = global_metrics_provider._sdk_config.resource # type: ignore - else: - global_resource = Resource.create() - - resource = global_resource.merge(Resource.create(resource_attributes)) - - # 2. init exporter and reader - exporter = OTLPMetricExporter(endpoint=endpoint, headers=headers) - metric_reader = PeriodicExportingMetricReader(exporter) - - metrics_api.set_meter_provider( - metrics_sdk.MeterProvider(metric_readers=[metric_reader], resource=resource) - ) - - # 3. init meter - self.meter: Meter = metrics.get_meter(name=name) - - # create meter attributes - self.llm_invoke_counter = self.meter.create_counter( - name=Meters.LLM_CHAT_COUNT, - description="Number of LLM invocations", - unit="count", - ) - self.token_usage = self.meter.create_histogram( - name=Meters.LLM_TOKEN_USAGE, - description="Token consumption of LLM invocations", - unit="count", - explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_TOKEN_USAGE_BUCKETS, - ) - self.duration_histogram = self.meter.create_histogram( - name=Meters.LLM_OPERATION_DURATION, - unit="s", - description="GenAI operation duration", - explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, - ) - self.chat_exception_counter = self.meter.create_counter( - name=Meters.LLM_COMPLETIONS_EXCEPTIONS, - unit="time", - description="Number of exceptions occurred during chat completions", - ) - self.streaming_time_to_first_token = self.meter.create_histogram( - name=Meters.LLM_STREAMING_TIME_TO_FIRST_TOKEN, - unit="s", - description="Time to first token in streaming chat completions", - explicit_bucket_boundaries_advisory=_GEN_AI_SERVER_TIME_TO_FIRST_TOKEN_BUCKETS, - ) - self.streaming_time_to_generate = self.meter.create_histogram( - name=Meters.LLM_STREAMING_TIME_TO_GENERATE, - unit="s", - description="Time between first token and completion in streaming chat completions", - explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, - ) - self.streaming_time_per_output_token = self.meter.create_histogram( - name=Meters.LLM_STREAMING_TIME_PER_OUTPUT_TOKEN, - unit="s", - description="Time per output token in streaming chat completions", - explicit_bucket_boundaries_advisory=_GEN_AI_SERVER_TIME_PER_OUTPUT_TOKEN_BUCKETS, - ) - - # apmplus metrics for veadk dashboard - self.apmplus_span_latency = self.meter.create_histogram( - name=Meters.APMPLUS_SPAN_LATENCY, - description="Latency of span", - unit="s", - explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, - ) - self.apmplus_tool_token_usage = self.meter.create_histogram( - name=Meters.APMPLUS_TOOL_TOKEN_USAGE, - description="Token consumption of APMPlus tool token", - unit="count", - explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_TOKEN_USAGE_BUCKETS, - ) - self.skill_invoke_latency = self.meter.create_histogram( - name=Meters.GEN_AI_SKILL_INVOKE_LATENCY, - description="Latency of skill invocations", - unit="s", - explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, - ) - - def record_call_llm( - self, - invocation_context: InvocationContext, - event_id: str, - llm_request: LlmRequest, - llm_response: LlmResponse, - ) -> None: - """Record comprehensive metrics for LLM call operations. - - Captures detailed telemetry data for language model invocations - including token consumption, latency, and error information. - This data enables cost optimization, performance analysis, and - reliability monitoring in APMPlus dashboards. - - Metrics Recorded: - - Invocation count with model and operation attributes - - Input/output token usage with separate tracking - - Operation duration from span timing data - - Error counts and exception details - - Span latency for performance analysis - - Args: - invocation_context: Context with agent, session, and user information - event_id: Unique identifier for this LLM call event - llm_request: Request object with model and parameter details - llm_response: Response object with content and usage metadata - """ - is_streaming = bool( - invocation_context.run_config - and invocation_context.run_config.streaming_mode != StreamingMode.NONE - ) - server_address = getattr(invocation_context.agent, "model_api_base", "unknown") - attributes = { - "gen_ai_system": "volcengine", - "gen_ai_response_model": llm_request.model, - "gen_ai_operation_name": "chat", - "gen_ai_operation_type": "llm", - "stream": is_streaming, - "server_address": server_address, - } # required by Volcengine APMPlus - - if llm_response.usage_metadata: - # llm invocation number += 1 - self.llm_invoke_counter.add(1, attributes) - - # upload token usage - input_token = llm_response.usage_metadata.prompt_token_count - output_token = llm_response.usage_metadata.candidates_token_count - - if input_token: - token_attributes = {**attributes, "gen_ai_token_type": "input"} - self.token_usage.record(input_token, attributes=token_attributes) - if output_token: - token_attributes = {**attributes, "gen_ai_token_type": "output"} - self.token_usage.record(output_token, attributes=token_attributes) - - # Get llm duration - span = trace.get_current_span() - if span and hasattr(span, "start_time") and self.duration_histogram: - # We use span start time as the llm request start time - tik = span.start_time # type: ignore - # We use current time as the llm request end time - tok = time.time_ns() - # Calculate duration in seconds - duration = (tok - tik) / 1e9 - self.duration_histogram.record( - duration, attributes=attributes - ) # unit in seconds - - # Get model request error - if llm_response.error_code and self.chat_exception_counter: - exception_attributes = { - **attributes, - "error_type": llm_response.error_code, - } - self.chat_exception_counter.add(1, exception_attributes) - - # TODO: Get streaming time to first token - # time_to_frist_token = 0.1 - # if self.streaming_time_to_first_token: - # self.streaming_time_to_first_token.record( - # time_to_frist_token, attributes=attributes - # ) - - # TODO: Get streaming time to generate - # time_to_generate = 1.0 - # if self.streaming_time_to_generate: - # self.streaming_time_to_generate.record( - # time_to_generate, attributes=attributes - # ) - - # TODO: Get streaming time per output token - # time_per_output_token = 0.01 - # if self.streaming_time_per_output_token: - # self.streaming_time_per_output_token.record( - # time_per_output_token, attributes=attributes - # ) - - # add span name attribute - span = trace.get_current_span() - if not span: - return - - # record span latency - if hasattr(span, "start_time") and self.apmplus_span_latency: - # span 耗时 - duration = (time.time_ns() - span.start_time) / 1e9 # type: ignore - self.apmplus_span_latency.record(duration, attributes=attributes) - - def record_tool_call( - self, - tool: BaseTool, - args: dict[str, Any], - function_response_event: Event, - ): - """Record metrics for tool execution operations. - - Captures performance and usage metrics for tool invocations - including execution latency and estimated token consumption. - Enables monitoring of tool performance and resource usage patterns. - - Metrics Recorded: - - Tool execution latency from span timing - - Input/output token estimation based on text length - - Tool-specific attributes for categorization - - Args: - tool: Tool instance that was executed - args: Arguments passed to the tool function - function_response_event: Event containing execution results - """ - logger.debug(f"Record tool call work in progress. Tool: {tool.name}") - span = trace.get_current_span() - if not span: - return - operation_type = "tool" - operation_name = tool.name - operation_backend = "" - if hasattr(tool, "custom_metadata") and tool.custom_metadata: - operation_backend = tool.custom_metadata.get("backend", "") - - attributes = { - "gen_ai_operation_name": operation_name, - "gen_ai_operation_type": operation_type, - "gen_ai_operation_backend": operation_backend, - } - - if hasattr(span, "start_time") and self.apmplus_span_latency: - # span 耗时 - duration = (time.time_ns() - span.start_time) / 1e9 # type: ignore - self.apmplus_span_latency.record(duration, attributes=attributes) - - if self.apmplus_tool_token_usage and hasattr(span, "attributes"): - span_attributes = span.attributes or {} - tool_input = span_attributes.get("gen_ai.tool.input") - if tool_input: - tool_token_usage_input = ( - len(tool_input) / 4 - ) # tool token 数量,使用文本长度/4 - input_tool_token_attributes = {**attributes, "token_type": "input"} - self.apmplus_tool_token_usage.record( - tool_token_usage_input, attributes=input_tool_token_attributes - ) - - tool_output = span_attributes.get("gen_ai.tool.output") - if tool_output: - tool_token_usage_output = ( - len(tool_output) / 4 - ) # tool token 数量,使用文本长度/4 - output_tool_token_attributes = {**attributes, "token_type": "output"} - self.apmplus_tool_token_usage.record( - tool_token_usage_output, attributes=output_tool_token_attributes - ) + return resolved_provider class APMPlusExporterConfig(BaseModel): @@ -518,14 +128,13 @@ class APMPlusExporter(BaseExporter): def model_post_init(self, context: Any) -> None: """Initialize APMPlus exporter components after model construction. - Sets up OTLP span exporter, batch processor, and meter uploader - with proper authentication and resource attribution for APMPlus - integration. + Sets up the OTLP span exporter and ensures a usable global + MeterProvider without replacing one configured by the application. Components Initialized: - OTLP span exporter with APMPlus endpoint and authentication - Batch span processor for efficient data transmission - - Meter uploader for comprehensive metrics collection + - APMPlus metric pipeline only when no global MeterProvider exists - Resource attributes for service identification """ logger.info(f"APMPlusExporter sevice name: {self.config.service_name}") @@ -545,8 +154,7 @@ def model_post_init(self, context: Any) -> None: ) self.processor = BatchSpanProcessor(self._exporter) - self.meter_uploader = MeterUploader( - name="apmplus_meter", + ensure_apmplus_meter_provider( endpoint=self.config.endpoint, headers=self.headers, resource_attributes=self.resource_attributes, @@ -568,6 +176,12 @@ def export(self) -> None: if self._exporter: self._exporter.force_flush() + from veadk.tracing.telemetry.portal_metrics import ( + portal_metric_recorder, + ) + + portal_metric_recorder.force_flush() + logger.info( f"APMPlusExporter exports data to {self.config.endpoint}, service name: {self.config.service_name}" ) diff --git a/veadk/tracing/telemetry/exporters/base_exporter.py b/veadk/tracing/telemetry/exporters/base_exporter.py index 8ec2ac072..323763c71 100644 --- a/veadk/tracing/telemetry/exporters/base_exporter.py +++ b/veadk/tracing/telemetry/exporters/base_exporter.py @@ -12,7 +12,10 @@ # See the License for the specific language governing permissions and # limitations under the License. -from opentelemetry.sdk.trace import SpanProcessor +from __future__ import annotations + +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import SpanProcessor, TracerProvider from opentelemetry.sdk.trace.export import SpanExporter from pydantic import BaseModel, ConfigDict, Field @@ -32,8 +35,27 @@ class BaseExporter(BaseModel): headers: dict = Field(default_factory=dict) _exporter: SpanExporter | None = None + _registered_provider: TracerProvider | None = None processor: SpanProcessor | None = None + def register(self, provider: TracerProvider) -> bool: + """Register this exporter's processor with a tracer provider once. + + Returns: + Whether the processor was newly registered with ``provider``. + """ + if self.processor is None or self._registered_provider is provider: + return False + + if self.resource_attributes: + provider._resource = provider._resource.merge( + Resource.create(self.resource_attributes) + ) + + provider.add_span_processor(self.processor) + self._registered_provider = provider + return True + def export(self) -> None: """Force export of telemetry data.""" pass diff --git a/veadk/tracing/telemetry/opentelemetry_tracer.py b/veadk/tracing/telemetry/opentelemetry_tracer.py index 276111cb9..fdae30d7b 100644 --- a/veadk/tracing/telemetry/opentelemetry_tracer.py +++ b/veadk/tracing/telemetry/opentelemetry_tracer.py @@ -20,7 +20,6 @@ from opentelemetry import trace as trace_api from opentelemetry.sdk import trace as trace_sdk -from opentelemetry.sdk.resources import Resource from opentelemetry.sdk.trace import SpanLimits, TracerProvider from pydantic import BaseModel, ConfigDict, Field, field_validator from typing_extensions import override @@ -29,6 +28,7 @@ from veadk.tracing.telemetry.exporters.apmplus_exporter import APMPlusExporter from veadk.tracing.telemetry.exporters.base_exporter import BaseExporter from veadk.tracing.telemetry.exporters.inmemory_exporter import InMemoryExporter +from veadk.tracing.telemetry.portal_metrics import portal_metric_recorder from veadk.utils.logger import get_logger from veadk.utils.misc import get_agent_dir from veadk.utils.patches import patch_google_adk_telemetry @@ -36,21 +36,6 @@ logger = get_logger(__name__) -def _update_resource_attributions( - provider: TracerProvider, resource_attributes: dict -) -> None: - """Update the resource attributes of a TracerProvider instance. - - This function merges new resource attributes with the existing ones in the - provider, allowing dynamic configuration of telemetry metadata. - - Args: - provider: The TracerProvider instance to update - resource_attributes: Dictionary of attributes to merge with existing resources - """ - provider._resource = provider._resource.merge(Resource.create(resource_attributes)) - - class OpentelemetryTracer(BaseModel, BaseTracer): """OpenTelemetry-based tracer implementation for comprehensive agent observability. @@ -172,39 +157,11 @@ def _init_global_tracer_provider(self) -> None: global_tracer_provider = trace_api.get_tracer_provider() global_tracer_provider: TracerProvider + self._global_tracer_provider = global_tracer_provider self._apmplus_managed_externally = have_global_tracer_provider - if self._apmplus_managed_externally: - exporter_count = len(self.exporters) - self.exporters = [ - e for e in self.exporters if not isinstance(e, APMPlusExporter) - ] - if len(self.exporters) != exporter_count: - logger.info( - "Reuse existing global TracerProvider and skip registering " - "APMPlusExporter." - ) - for exporter in self.exporters: - processor = exporter.processor - resource_attributes = exporter.resource_attributes - - if resource_attributes: - _update_resource_attributions( - global_tracer_provider, resource_attributes - ) - - if processor: - global_tracer_provider.add_span_processor(processor) - self._processors.append(processor) - - logger.debug( - f"Add span processor for exporter `{exporter.__class__.__name__}` to OpentelemetryTracer." - ) - else: - logger.error( - f"Add span processor for exporter `{exporter.__class__.__name__}` to OpentelemetryTracer failed." - ) + self._activate_exporter(exporter) self._inmemory_exporter = InMemoryExporter() if self._inmemory_exporter.processor: @@ -225,18 +182,42 @@ def _init_global_tracer_provider(self) -> None: f"Init OpentelemetryTracer with {len(self._processors)} exporter(s)." ) - # Initialize global meter_uploader from exporters - from veadk.tracing.telemetry.telemetry import ( - init_global_meter_uploader_from_exporters, - ) + def _activate_exporter(self, exporter: BaseExporter) -> bool: + return self._register_exporter(exporter) - init_global_meter_uploader_from_exporters(self.exporters) + def _register_exporter(self, exporter: BaseExporter) -> bool: + if isinstance(exporter, APMPlusExporter) and self.apmplus_managed_externally: + logger.info( + "Reuse existing global TracerProvider and skip registering " + "APMPlusExporter span processor." + ) + return False + + if exporter.register(self._global_tracer_provider): + self._processors.append(exporter.processor) + logger.debug( + f"Add span processor for exporter `{exporter.__class__.__name__}` to OpentelemetryTracer." + ) + return True + + if exporter.processor is None: + logger.error( + f"Add span processor for exporter `{exporter.__class__.__name__}` to OpentelemetryTracer failed." + ) + return False @property def apmplus_managed_externally(self) -> bool: """Whether a global provider existed before this tracer was initialized.""" return self._apmplus_managed_externally + def add_exporter(self, exporter: BaseExporter) -> bool: + """Add an exporter and immediately register it with the global provider.""" + if not any(existing is exporter for existing in self.exporters): + self.exporters.append(exporter) + + return self._activate_exporter(exporter) + @property def trace_file_path(self) -> str: """Get the file path of the most recent trace dump. @@ -275,6 +256,7 @@ def force_export(self) -> None: for processor in self._processors: time.sleep(0.05) processor.force_flush() + portal_metric_recorder.force_flush() @override def dump( diff --git a/veadk/tracing/telemetry/portal_metrics.py b/veadk/tracing/telemetry/portal_metrics.py new file mode 100644 index 000000000..4842dce2a --- /dev/null +++ b/veadk/tracing/telemetry/portal_metrics.py @@ -0,0 +1,450 @@ +# Copyright (c) 2025 Beijing Volcano Engine Technology Co., Ltd. and/or its affiliates. +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +"""Always-on VeADK Portal metric instrumentation.""" + +import time +from dataclasses import dataclass +from typing import Any + +from google.adk.agents.invocation_context import InvocationContext +from google.adk.agents.run_config import StreamingMode +from google.adk.events import Event +from google.adk.models.llm_request import LlmRequest +from google.adk.models.llm_response import LlmResponse +from google.adk.tools import BaseTool +from opentelemetry import metrics as metrics_api +from opentelemetry import trace +from opentelemetry.metrics._internal import Meter + +from veadk.utils.logger import get_logger + +logger = get_logger(__name__) + +_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS = [ + 0.01, + 0.02, + 0.04, + 0.08, + 0.16, + 0.32, + 0.64, + 1.28, + 2.56, + 5.12, + 10.24, + 20.48, + 40.96, + 81.92, +] + +_GEN_AI_SERVER_TIME_PER_OUTPUT_TOKEN_BUCKETS = [ + 0.01, + 0.025, + 0.05, + 0.075, + 0.1, + 0.15, + 0.2, + 0.3, + 0.4, + 0.5, + 0.75, + 1.0, + 2.5, +] + +_GEN_AI_SERVER_TIME_TO_FIRST_TOKEN_BUCKETS = [ + 0.001, + 0.005, + 0.01, + 0.02, + 0.04, + 0.06, + 0.08, + 0.1, + 0.25, + 0.5, + 0.75, + 1.0, + 2.5, + 5.0, + 7.5, + 10.0, +] + +_GEN_AI_CLIENT_TOKEN_USAGE_BUCKETS = [ + 1, + 4, + 16, + 64, + 256, + 1024, + 4096, + 16384, + 65536, + 262144, + 1048576, + 4194304, + 16777216, + 67108864, +] + + +@dataclass +class Meters: + """Metric names and identifiers for OpenTelemetry instrumentation. + + This class defines standardized metric names used for LLM and agent + observability. The metrics follow OpenTelemetry semantic conventions + for generative AI operations and include custom APMPlus metrics for + enhanced monitoring capabilities. + + Standard Gen AI Metrics: + - LLM_CHAT_COUNT: Counter for LLM invocation frequency + - LLM_TOKEN_USAGE: Histogram for token consumption analysis + - LLM_OPERATION_DURATION: Histogram for operation latency tracking + - LLM_COMPLETIONS_EXCEPTIONS: Counter for error rate monitoring + - Streaming metrics: Performance analysis for streaming responses + + APMPlus Custom Metrics: + - APMPLUS_SPAN_LATENCY: Span execution time for performance analysis + - APMPLUS_TOOL_TOKEN_USAGE: Tool-specific token consumption tracking + """ + + LLM_CHAT_COUNT = "gen_ai.chat.count" + LLM_TOKEN_USAGE = "gen_ai.client.token.usage" + LLM_OPERATION_DURATION = "gen_ai.client.operation.duration" + LLM_COMPLETIONS_EXCEPTIONS = "gen_ai.chat_completions.exceptions" + LLM_STREAMING_TIME_TO_FIRST_TOKEN = ( + "gen_ai.chat_completions.streaming_time_to_first_token" + ) + LLM_STREAMING_TIME_TO_GENERATE = ( + "gen_ai.chat_completions.streaming_time_to_generate" + ) + LLM_STREAMING_TIME_PER_OUTPUT_TOKEN = ( + "gen_ai.chat_completions.streaming_time_per_output_token" + ) + + # apmplus metrics + # span duration + APMPLUS_SPAN_LATENCY = "apmplus_span_latency" + # tool token usage + APMPLUS_TOOL_TOKEN_USAGE = "apmplus_tool_token_usage" + # skill invoke latency + GEN_AI_SKILL_INVOKE_LATENCY = "gen_ai_skill_invoke_latency" + + +class PortalMetricRecorder: + """Record VeADK Portal metrics through the global MeterProvider. + + The recorder creates OpenTelemetry instruments and records measurements. + It does not configure a MeterProvider, reader, exporter, endpoint, or + credentials. Export is entirely owned by the global MeterProvider. + + Key Features: + - Automatic metric instrument creation with appropriate buckets + - LLM call metrics including token usage and latency + - Tool execution metrics for performance analysis + - Error tracking and exception monitoring + - Integration with OpenTelemetry metrics SDK + + Metrics Collected: + - LLM invocation counts and frequencies + - Token consumption (input/output) with histogram distribution + - Operation latency with performance bucket analysis + - Error rates and exception details + - Span-level performance metrics for APMPlus dashboards + """ + + def __init__(self, name: str = "veadk_portal") -> None: + """Create metric instruments without configuring an export pipeline.""" + self.meter: Meter = metrics_api.get_meter(name=name) + + # create meter attributes + self.llm_invoke_counter = self.meter.create_counter( + name=Meters.LLM_CHAT_COUNT, + description="Number of LLM invocations", + unit="count", + ) + self.token_usage = self.meter.create_histogram( + name=Meters.LLM_TOKEN_USAGE, + description="Token consumption of LLM invocations", + unit="count", + explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_TOKEN_USAGE_BUCKETS, + ) + self.duration_histogram = self.meter.create_histogram( + name=Meters.LLM_OPERATION_DURATION, + unit="s", + description="GenAI operation duration", + explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, + ) + self.chat_exception_counter = self.meter.create_counter( + name=Meters.LLM_COMPLETIONS_EXCEPTIONS, + unit="time", + description="Number of exceptions occurred during chat completions", + ) + self.streaming_time_to_first_token = self.meter.create_histogram( + name=Meters.LLM_STREAMING_TIME_TO_FIRST_TOKEN, + unit="s", + description="Time to first token in streaming chat completions", + explicit_bucket_boundaries_advisory=_GEN_AI_SERVER_TIME_TO_FIRST_TOKEN_BUCKETS, + ) + self.streaming_time_to_generate = self.meter.create_histogram( + name=Meters.LLM_STREAMING_TIME_TO_GENERATE, + unit="s", + description="Time between first token and completion in streaming chat completions", + explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, + ) + self.streaming_time_per_output_token = self.meter.create_histogram( + name=Meters.LLM_STREAMING_TIME_PER_OUTPUT_TOKEN, + unit="s", + description="Time per output token in streaming chat completions", + explicit_bucket_boundaries_advisory=_GEN_AI_SERVER_TIME_PER_OUTPUT_TOKEN_BUCKETS, + ) + + # apmplus metrics for veadk dashboard + self.apmplus_span_latency = self.meter.create_histogram( + name=Meters.APMPLUS_SPAN_LATENCY, + description="Latency of span", + unit="s", + explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, + ) + self.apmplus_tool_token_usage = self.meter.create_histogram( + name=Meters.APMPLUS_TOOL_TOKEN_USAGE, + description="Token consumption of APMPlus tool token", + unit="count", + explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_TOKEN_USAGE_BUCKETS, + ) + self.skill_invoke_latency = self.meter.create_histogram( + name=Meters.GEN_AI_SKILL_INVOKE_LATENCY, + description="Latency of skill invocations", + unit="s", + explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, + ) + + @property + def provider(self) -> metrics_api.MeterProvider: + return metrics_api.get_meter_provider() + + def force_flush(self) -> bool: + force_flush = getattr(self.provider, "force_flush", None) + return bool(force_flush()) if force_flush else True + + def record_call_llm( + self, + invocation_context: InvocationContext, + event_id: str, + llm_request: LlmRequest, + llm_response: LlmResponse, + ) -> None: + """Record comprehensive metrics for LLM call operations. + + Captures detailed telemetry data for language model invocations + including token consumption, latency, and error information. + This data enables cost optimization, performance analysis, and + reliability monitoring in APMPlus dashboards. + + Metrics Recorded: + - Invocation count with model and operation attributes + - Input/output token usage with separate tracking + - Operation duration from span timing data + - Error counts and exception details + - Span latency for performance analysis + + Args: + invocation_context: Context with agent, session, and user information + event_id: Unique identifier for this LLM call event + llm_request: Request object with model and parameter details + llm_response: Response object with content and usage metadata + """ + is_streaming = bool( + invocation_context.run_config + and invocation_context.run_config.streaming_mode != StreamingMode.NONE + ) + server_address = getattr(invocation_context.agent, "model_api_base", "unknown") + attributes = { + "gen_ai_system": "volcengine", + "gen_ai_response_model": llm_request.model, + "gen_ai_operation_name": "chat", + "gen_ai_operation_type": "llm", + "stream": is_streaming, + "server_address": server_address, + } # required by Volcengine APMPlus + + if llm_response.usage_metadata: + # llm invocation number += 1 + self.llm_invoke_counter.add(1, attributes) + + # upload token usage + input_token = llm_response.usage_metadata.prompt_token_count + output_token = llm_response.usage_metadata.candidates_token_count + + if input_token: + token_attributes = {**attributes, "gen_ai_token_type": "input"} + self.token_usage.record(input_token, attributes=token_attributes) + if output_token: + token_attributes = {**attributes, "gen_ai_token_type": "output"} + self.token_usage.record(output_token, attributes=token_attributes) + + # Get llm duration + span = trace.get_current_span() + if span and hasattr(span, "start_time") and self.duration_histogram: + # We use span start time as the llm request start time + tik = span.start_time # type: ignore + # We use current time as the llm request end time + tok = time.time_ns() + # Calculate duration in seconds + duration = (tok - tik) / 1e9 + self.duration_histogram.record( + duration, attributes=attributes + ) # unit in seconds + + # Get model request error + if llm_response.error_code and self.chat_exception_counter: + exception_attributes = { + **attributes, + "error_type": llm_response.error_code, + } + self.chat_exception_counter.add(1, exception_attributes) + + # TODO: Get streaming time to first token + # time_to_frist_token = 0.1 + # if self.streaming_time_to_first_token: + # self.streaming_time_to_first_token.record( + # time_to_frist_token, attributes=attributes + # ) + + # TODO: Get streaming time to generate + # time_to_generate = 1.0 + # if self.streaming_time_to_generate: + # self.streaming_time_to_generate.record( + # time_to_generate, attributes=attributes + # ) + + # TODO: Get streaming time per output token + # time_per_output_token = 0.01 + # if self.streaming_time_per_output_token: + # self.streaming_time_per_output_token.record( + # time_per_output_token, attributes=attributes + # ) + + # add span name attribute + span = trace.get_current_span() + if not span: + return + + # record span latency + if hasattr(span, "start_time") and self.apmplus_span_latency: + # span 耗时 + duration = (time.time_ns() - span.start_time) / 1e9 # type: ignore + self.apmplus_span_latency.record(duration, attributes=attributes) + + def record_tool_call( + self, + tool: BaseTool, + args: dict[str, Any], + function_response_event: Event, + ): + """Record metrics for tool execution operations. + + Captures performance and usage metrics for tool invocations + including execution latency and estimated token consumption. + Enables monitoring of tool performance and resource usage patterns. + + Metrics Recorded: + - Tool execution latency from span timing + - Input/output token estimation based on text length + - Tool-specific attributes for categorization + + Args: + tool: Tool instance that was executed + args: Arguments passed to the tool function + function_response_event: Event containing execution results + """ + logger.debug(f"Record tool call work in progress. Tool: {tool.name}") + span = trace.get_current_span() + if not span: + return + operation_type = "tool" + operation_name = tool.name + operation_backend = "" + if hasattr(tool, "custom_metadata") and tool.custom_metadata: + operation_backend = tool.custom_metadata.get("backend", "") + + attributes = { + "gen_ai_operation_name": operation_name, + "gen_ai_operation_type": operation_type, + "gen_ai_operation_backend": operation_backend, + } + + if hasattr(span, "start_time") and self.apmplus_span_latency: + # span 耗时 + duration = (time.time_ns() - span.start_time) / 1e9 # type: ignore + self.apmplus_span_latency.record(duration, attributes=attributes) + + if self.apmplus_tool_token_usage and hasattr(span, "attributes"): + span_attributes = span.attributes or {} + tool_input = span_attributes.get("gen_ai.tool.input") + if tool_input: + tool_token_usage_input = ( + len(tool_input) / 4 + ) # tool token 数量,使用文本长度/4 + input_tool_token_attributes = {**attributes, "token_type": "input"} + self.apmplus_tool_token_usage.record( + tool_token_usage_input, attributes=input_tool_token_attributes + ) + + tool_output = span_attributes.get("gen_ai.tool.output") + if tool_output: + tool_token_usage_output = ( + len(tool_output) / 4 + ) # tool token 数量,使用文本长度/4 + output_tool_token_attributes = {**attributes, "token_type": "output"} + self.apmplus_tool_token_usage.record( + tool_token_usage_output, attributes=output_tool_token_attributes + ) + + def record_skill_call( + self, + span: Any, + skill_name: str, + tool_name: str, + skill: Any, + result: str, + ) -> None: + """Record latency and result attributes for a skill invocation.""" + attributes = { + "skill_name": skill_name, + "tool_name": tool_name, + "skill_space_id": ( + skill.skill_space_id + if skill and getattr(skill, "skill_space_id", None) + else "" + ), + "skill_id": skill.id if skill and getattr(skill, "id", None) else "", + "gen_ai.operation.name": "execute_skill", + "error_type": ( + "skill_execution_error" if result.startswith("Error:") else "" + ), + } + + latency_seconds = 0.0 + if hasattr(span, "start_time"): + latency_seconds = (time.time_ns() - span.start_time) / 1e9 + + self.skill_invoke_latency.record(latency_seconds, attributes) + + +portal_metric_recorder = PortalMetricRecorder() diff --git a/veadk/tracing/telemetry/telemetry.py b/veadk/tracing/telemetry/telemetry.py index fcf30811f..dba7d4886 100644 --- a/veadk/tracing/telemetry/telemetry.py +++ b/veadk/tracing/telemetry/telemetry.py @@ -30,30 +30,12 @@ ToolAttributesParams, ) from veadk.tracing.telemetry.content_tracing import should_trace_content +from veadk.tracing.telemetry import portal_metrics from veadk.utils.logger import get_logger from veadk.utils.misc import safe_json_serialize logger = get_logger(__name__) -meter_uploader = None - - -def init_global_meter_uploader_from_exporters(exporters): - """Initialize global meter_uploader from a list of exporters. - - Args: - exporters: List of exporter instances to search for meter_uploader - """ - global meter_uploader - for exporter in exporters: - if hasattr(exporter, "meter_uploader") and exporter.meter_uploader: - meter_uploader = exporter.meter_uploader - logger.debug( - "Global meter_uploader initialized from exporter: %s", - exporter.__class__.__name__, - ) - break - def _upload_call_llm_metrics( invocation_context: InvocationContext, @@ -61,10 +43,10 @@ def _upload_call_llm_metrics( llm_request: LlmRequest, llm_response: LlmResponse, ) -> None: - """Upload LLM call metrics to configured meter uploaders. + """Record LLM call metrics through the global MeterProvider. - This function extracts meter uploaders from agent tracers and records - LLM call metrics including token usage, latency, and request/response details. + Recording is independent from exporter configuration. OpenTelemetry's + default proxy instruments remain no-op until a real provider is installed. Args: invocation_context: Context containing agent, session, and user information @@ -72,18 +54,9 @@ def _upload_call_llm_metrics( llm_request: The request sent to the language model llm_response: The response received from the language model """ - from veadk.agent import Agent - - if isinstance(invocation_context.agent, Agent): - tracers = invocation_context.agent.tracers - for tracer in tracers: - for exporter in getattr(tracer, "exporters", []): - if getattr(exporter, "meter_uploader", None): - global meter_uploader - meter_uploader = exporter.meter_uploader - exporter.meter_uploader.record_call_llm( - invocation_context, event_id, llm_request, llm_response - ) + portal_metrics.portal_metric_recorder.record_call_llm( + invocation_context, event_id, llm_request, llm_response + ) def _upload_tool_call_metrics( @@ -91,7 +64,7 @@ def _upload_tool_call_metrics( args: dict[str, Any], function_response_event: Event, ): - """Upload tool call metrics to the global meter uploader. + """Record tool call metrics through the global MeterProvider. Records tool execution metrics including function name, arguments, execution time, and response details for observability and debugging. @@ -101,16 +74,10 @@ def _upload_tool_call_metrics( args: Arguments passed to the tool function function_response_event: Event containing the tool's response data - Note: - - Requires global meter_uploader to be initialized """ - global meter_uploader - if meter_uploader: - meter_uploader.record_tool_call(tool, args, function_response_event) - else: - logger.debug( - "Meter uploader is not initialized yet. Skip recording tool call metrics." - ) + portal_metrics.portal_metric_recorder.record_tool_call( + tool, args, function_response_event + ) def _set_agent_input_attribute(