From f2d28445e986561cbe0570e1b00e5555d5eb2d3a Mon Sep 17 00:00:00 2001 From: "fengduzhen.666" Date: Tue, 4 Aug 2026 22:31:44 +0800 Subject: [PATCH 1/4] fix(tracing): register exporters added after initialization --- tests/test_tracing.py | 147 +++++++++++++++--- veadk/agent.py | 6 +- .../telemetry/exporters/base_exporter.py | 22 ++- .../tracing/telemetry/opentelemetry_tracer.py | 83 +++++----- 4 files changed, 190 insertions(+), 68 deletions(-) diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 5f7b2c5be..1365b280f 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -148,26 +148,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,29 +186,72 @@ 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 == [] + assert telemetry_module.meter_uploader is None 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 + 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 expected_meter_uploader = ( - constructed_exporters[-1].meter_uploader if enable_apmplus else None + constructed_exporters[0].meter_uploader if should_construct_apmplus else None ) assert telemetry_module.meter_uploader is expected_meter_uploader +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_tracing_registers_apmplus_without_global_provider( fresh_global_tracer_provider, ): @@ -214,6 +265,66 @@ 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 + assert telemetry_module.meter_uploader is apmplus_exporter.meter_uploader + + +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 + assert telemetry_module.meter_uploader is apmplus_exporter.meter_uploader + + +def test_agent_env_keeps_metrics_only_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 telemetry_module.meter_uploader is exporter.meter_uploader + 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 +343,10 @@ 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 for metrics but its span processor is not registered. + assert len(tracer.exporters) == 4 @pytest.mark.asyncio @@ -249,5 +360,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/veadk/agent.py b/veadk/agent.py index fd83dd305..d5e0e2145 100644 --- a/veadk/agent.py +++ b/veadk/agent.py @@ -671,17 +671,17 @@ 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( diff --git a/veadk/tracing/telemetry/exporters/base_exporter.py b/veadk/tracing/telemetry/exporters/base_exporter.py index 8ec2ac072..9fe0d7e61 100644 --- a/veadk/tracing/telemetry/exporters/base_exporter.py +++ b/veadk/tracing/telemetry/exporters/base_exporter.py @@ -12,7 +12,8 @@ # See the License for the specific language governing permissions and # limitations under the License. -from opentelemetry.sdk.trace import SpanProcessor +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 +33,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..282ee63b2 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 @@ -36,21 +35,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 +156,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._register_exporter(exporter) self._inmemory_exporter = InMemoryExporter() if self._inmemory_exporter.processor: @@ -232,11 +188,46 @@ def _init_global_tracer_provider(self) -> None: 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) + + registered = self._register_exporter(exporter) + + from veadk.tracing.telemetry.telemetry import ( + init_global_meter_uploader_from_exporters, + ) + + init_global_meter_uploader_from_exporters(self.exporters) + return registered + @property def trace_file_path(self) -> str: """Get the file path of the most recent trace dump. From 01806d5e684e8db1d87764212e382a98cac32367 Mon Sep 17 00:00:00 2001 From: "fengduzhen.666" Date: Tue, 4 Aug 2026 23:56:56 +0800 Subject: [PATCH 2/4] refactor(metrics): decouple APMPlus meter uploader --- tests/test_metric_uploader.py | 169 ++++++++++++++++++ tests/test_tracing.py | 61 +++++-- tests/test_tracing_content.py | 3 +- veadk/agent.py | 7 - veadk/tools/skills_tools/skills_tool.py | 45 ++--- .../telemetry/exporters/apmplus_exporter.py | 136 ++++++++++++-- .../telemetry/exporters/base_exporter.py | 11 ++ veadk/tracing/telemetry/metric_uploader.py | 151 ++++++++++++++++ .../tracing/telemetry/opentelemetry_tracer.py | 24 +-- veadk/tracing/telemetry/telemetry.py | 49 +---- 10 files changed, 535 insertions(+), 121 deletions(-) create mode 100644 tests/test_metric_uploader.py create mode 100644 veadk/tracing/telemetry/metric_uploader.py diff --git a/tests/test_metric_uploader.py b/tests/test_metric_uploader.py new file mode 100644 index 000000000..6ceee0105 --- /dev/null +++ b/tests/test_metric_uploader.py @@ -0,0 +1,169 @@ +# 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 telemetry +from veadk.tracing.telemetry.exporters import ( + apmplus_exporter as apmplus_exporter_module, +) +from veadk.tracing.telemetry.exporters.apmplus_exporter import MeterUploader +from veadk.tracing.telemetry.metric_uploader import ( + MetricUploaderRegistry, + metric_uploader_registry, +) + + +class FakeMetricUploader: + def __init__(self, registration_key): + self.registration_key = registration_key + self.llm_calls = [] + self.tool_calls = [] + self.skill_calls = [] + self.force_flush_calls = 0 + self.shutdown_calls = 0 + + 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) + + def force_flush(self): + self.force_flush_calls += 1 + return True + + def shutdown(self): + self.shutdown_calls += 1 + + +@pytest.fixture(autouse=True) +def fresh_metric_uploader_registry(): + metric_uploader_registry.clear() + yield + metric_uploader_registry.clear() + + +@pytest.fixture +def preconfigured_zero_reader_meter_provider(monkeypatch): + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) + provider = metrics_sdk.MeterProvider() + metrics_api.set_meter_provider(provider) + yield provider + provider.shutdown() + + +def test_meter_uploader_uses_private_provider_when_global_provider_exists( + monkeypatch, + preconfigured_zero_reader_meter_provider, +): + 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, + ) + + uploader = MeterUploader( + name="test-meter", + endpoint="http://localhost:4319", + headers={"x-byteapm-appkey": "test"}, + resource_attributes={"service.name": "test-service"}, + ) + uploader.llm_invoke_counter.add(1) + + assert metrics_api.get_meter_provider() is preconfigured_zero_reader_meter_provider + assert uploader.provider is not preconfigured_zero_reader_meter_provider + assert len(uploader.provider._sdk_config.metric_readers) == 1 + assert exporter_kwargs["insecure"] is True + assert uploader.force_flush() + assert reader.get_metrics_data() is not None + + uploader.shutdown() + + +def test_registry_deduplicates_by_key_without_constructing_a_duplicate(): + registry = MetricUploaderRegistry() + first = FakeMetricUploader(("apmplus", "same-destination")) + + assert registry.register(first) is first + + factory_calls = [] + resolved = registry.get_or_create( + first.registration_key, + lambda: factory_calls.append(True) + or FakeMetricUploader(first.registration_key), + ) + + assert resolved is first + assert factory_calls == [] + assert registry.uploaders == (first,) + + +def test_registry_fans_out_to_distinct_destinations_once(): + registry = MetricUploaderRegistry() + first = FakeMetricUploader(("apmplus", "destination-a")) + second = FakeMetricUploader(("apmplus", "destination-b")) + duplicate = FakeMetricUploader(first.registration_key) + + registry.register(first) + registry.register(second) + assert registry.register(duplicate) is first + + registry.record_call_llm("context", "event", "request", "response") + registry.record_tool_call("tool", {}, "event") + registry.record_skill_call("span", "skill", "tool", None, "ok") + + assert len(first.llm_calls) == len(second.llm_calls) == 1 + assert len(first.tool_calls) == len(second.tool_calls) == 1 + assert len(first.skill_calls) == len(second.skill_calls) == 1 + assert duplicate.llm_calls == [] + assert duplicate.shutdown_calls == 1 + assert registry.force_flush() + assert first.force_flush_calls == second.force_flush_calls == 1 + + +def test_telemetry_and_skill_metrics_use_the_registry(): + uploader = FakeMetricUploader(("apmplus", "destination")) + metric_uploader_registry.register(uploader) + + 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(uploader.llm_calls) == 1 + assert len(uploader.tool_calls) == 1 + assert len(uploader.skill_calls) == 1 diff --git a/tests/test_tracing.py b/tests/test_tracing.py index 1365b280f..2a7971f36 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, ) @@ -47,6 +46,7 @@ TLSExporter, TLSExporterConfig, ) +from veadk.tracing.telemetry.metric_uploader import metric_uploader_registry from veadk.tracing.telemetry.opentelemetry_tracer import OpentelemetryTracer APP_NAME = "app" @@ -98,6 +98,7 @@ def gen_span_processor(endpoint: str): @pytest.fixture def fresh_global_tracer_provider(monkeypatch): """Give each test an isolated OpenTelemetry global provider.""" + metric_uploader_registry.clear() monkeypatch.setattr(trace_api, "_TRACER_PROVIDER", None) monkeypatch.setattr(trace_api, "_TRACER_PROVIDER_SET_ONCE", Once()) @@ -106,13 +107,31 @@ def fresh_global_tracer_provider(monkeypatch): tracer_provider = trace_api.get_tracer_provider() if isinstance(tracer_provider, trace_sdk.TracerProvider): tracer_provider.shutdown() + metric_uploader_registry.clear() @pytest.fixture 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 ControlledMetricUploader: + registration_key = ("apmplus", "controlled") + + def __init__(self): + self.force_flush_calls = 0 + + def record_call_llm(self, *args): ... + + def record_tool_call(self, *args): ... + + def record_skill_call(self, *args): ... + + def force_flush(self): + self.force_flush_calls += 1 + return True + + def shutdown(self): ... class ControlledAPMPlusExporter(APMPlusExporter): def __init__(self): @@ -127,9 +146,12 @@ def __init__(self): def model_post_init(self, context): self._exporter = OTelInMemorySpanExporter() self.processor = SimpleSpanProcessor(self._exporter) - self.meter_uploader = object() + self._controlled_meter_uploader = ControlledMetricUploader() constructed_exporters.append(self) + def get_metric_uploader(self): + return self._controlled_meter_uploader + monkeypatch.setattr( apmplus_exporter_module, "APMPlusExporter", @@ -191,7 +213,7 @@ def test_apmplus_enable_provider_and_manual_exporter_matrix( if not should_create_tracer: assert final_provider is initial_provider assert constructed_exporters == [] - assert telemetry_module.meter_uploader is None + assert metric_uploader_registry.uploaders == () return tracer = agent.tracers[0] @@ -218,10 +240,12 @@ def test_apmplus_enable_provider_and_manual_exporter_matrix( 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 - expected_meter_uploader = ( - constructed_exporters[0].meter_uploader if should_construct_apmplus else None + expected_uploaders = ( + (constructed_exporters[0].get_metric_uploader(),) + if should_construct_apmplus + else () ) - assert telemetry_module.meter_uploader is expected_meter_uploader + assert metric_uploader_registry.uploaders == expected_uploaders def test_add_exporter_registers_after_tracer_initialization( @@ -252,6 +276,19 @@ def test_add_exporter_is_idempotent(fresh_global_tracer_provider): assert sum(processor is exporter.processor for processor in tracer._processors) == 1 +def test_force_export_flushes_metric_uploader_registry( + fresh_global_tracer_provider, + controlled_apmplus_exporter, +): + controlled_exporter_class, _ = controlled_apmplus_exporter + exporter = controlled_exporter_class() + tracer = OpentelemetryTracer(exporters=[exporter]) + + tracer.force_export() + + assert exporter.get_metric_uploader().force_flush_calls == 1 + + def test_tracing_registers_apmplus_without_global_provider( fresh_global_tracer_provider, ): @@ -280,7 +317,9 @@ def test_tracing_skips_apmplus_for_preconfigured_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 - assert telemetry_module.meter_uploader is apmplus_exporter.meter_uploader + assert metric_uploader_registry.uploaders == ( + apmplus_exporter.get_metric_uploader(), + ) def test_add_exporter_skips_apmplus_for_preconfigured_provider( @@ -297,7 +336,9 @@ def test_add_exporter_skips_apmplus_for_preconfigured_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 - assert telemetry_module.meter_uploader is apmplus_exporter.meter_uploader + assert metric_uploader_registry.uploaders == ( + apmplus_exporter.get_metric_uploader(), + ) def test_agent_env_keeps_metrics_only_apmplus_for_preconfigured_provider( @@ -321,7 +362,7 @@ def test_agent_env_keeps_metrics_only_apmplus_for_preconfigured_provider( assert isinstance(exporter, controlled_exporter_class) assert exporter in tracer.exporters assert exporter.processor not in span_processors - assert telemetry_module.meter_uploader is exporter.meter_uploader + assert metric_uploader_registry.uploaders == (exporter.get_metric_uploader(),) assert len(span_processors) == 1 # VeADK in-memory processor only diff --git a/tests/test_tracing_content.py b/tests/test_tracing_content.py index 333f7bd81..3964abe2c 100644 --- a/tests/test_tracing_content.py +++ b/tests/test_tracing_content.py @@ -22,6 +22,7 @@ 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.metric_uploader import metric_uploader_registry @dataclass @@ -144,7 +145,7 @@ def _event_names(span): def setup_function(): - telemetry.meter_uploader = None + metric_uploader_registry.clear() def test_trace_call_llm_records_content_by_default(monkeypatch): diff --git a/veadk/agent.py b/veadk/agent.py index d5e0e2145..287643ed7 100644 --- a/veadk/agent.py +++ b/veadk/agent.py @@ -688,13 +688,6 @@ def _prepare_tracers(self): 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..80a8b3bd3 100644 --- a/veadk/tools/skills_tools/skills_tool.py +++ b/veadk/tools/skills_tools/skills_tool.py @@ -530,39 +530,16 @@ 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.metric_uploader import ( + metric_uploader_registry, + ) + + metric_uploader_registry.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..d80d4afa9 100644 --- a/veadk/tracing/telemetry/exporters/apmplus_exporter.py +++ b/veadk/tracing/telemetry/exporters/apmplus_exporter.py @@ -12,7 +12,10 @@ # See the License for the specific language governing permissions and # limitations under the License. +import hashlib +import json import time +from collections.abc import Hashable from dataclasses import dataclass from typing import Any @@ -22,7 +25,7 @@ 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 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 @@ -36,6 +39,10 @@ from veadk.config import settings from veadk.tracing.telemetry.exporters.base_exporter import BaseExporter +from veadk.tracing.telemetry.metric_uploader import ( + MetricUploader as MetricUploaderProtocol, + metric_uploader_registry, +) from veadk.utils.logger import get_logger logger = get_logger(__name__) @@ -111,6 +118,25 @@ ] +def _metric_registration_key( + endpoint: str, + headers: dict, + resource_attributes: dict, +) -> tuple[str, str]: + """Build a stable digest without exposing authentication headers in logs.""" + payload = json.dumps( + { + "endpoint": endpoint, + "headers": headers, + "resource_attributes": resource_attributes, + }, + sort_keys=True, + default=str, + separators=(",", ":"), + ) + return "apmplus", hashlib.sha256(payload.encode()).hexdigest() + + @dataclass class Meters: """Metric names and identifiers for OpenTelemetry instrumentation. @@ -182,7 +208,7 @@ def __init__( ) -> None: """Initialize the meter uploader with APMPlus configuration. - Sets up the global metrics provider, creates metric instruments, + Sets up a private metrics provider, creates metric instruments, and configures OTLP export to APMPlus endpoints with proper resource attribution and authentication. @@ -192,9 +218,14 @@ def __init__( 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 + self._registration_key = _metric_registration_key( + endpoint, headers, resource_attributes + ) + self._shutdown = False + + # Reuse global resource attributes when available, but keep the reader + # and provider private so an existing zero-reader global provider cannot + # suppress APMPlus metrics. global_metrics_provider = metrics_api.get_meter_provider() # 1. init resource @@ -206,15 +237,19 @@ def __init__( resource = global_resource.merge(Resource.create(resource_attributes)) # 2. init exporter and reader - exporter = OTLPMetricExporter(endpoint=endpoint, headers=headers) + exporter = OTLPMetricExporter( + endpoint=endpoint, + headers=headers, + insecure=True, + ) metric_reader = PeriodicExportingMetricReader(exporter) - metrics_api.set_meter_provider( - metrics_sdk.MeterProvider(metric_readers=[metric_reader], resource=resource) + self._provider = metrics_sdk.MeterProvider( + metric_readers=[metric_reader], resource=resource ) # 3. init meter - self.meter: Meter = metrics.get_meter(name=name) + self.meter: Meter = self._provider.get_meter(name=name) # create meter attributes self.llm_invoke_counter = self.meter.create_counter( @@ -278,6 +313,24 @@ def __init__( explicit_bucket_boundaries_advisory=_GEN_AI_CLIENT_OPERATION_DURATION_BUCKETS, ) + @property + def registration_key(self) -> Hashable: + return self._registration_key + + @property + def provider(self) -> metrics_sdk.MeterProvider: + return self._provider + + def force_flush(self) -> bool: + if self._shutdown: + return False + return self._provider.force_flush() + + def shutdown(self) -> None: + if not self._shutdown: + self._provider.shutdown() + self._shutdown = True + def record_call_llm( self, invocation_context: InvocationContext, @@ -452,6 +505,36 @@ def record_tool_call( 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) + class APMPlusExporterConfig(BaseModel): """Configuration model for APMPlus exporter settings. @@ -545,12 +628,27 @@ def model_post_init(self, context: Any) -> None: ) self.processor = BatchSpanProcessor(self._exporter) - self.meter_uploader = MeterUploader( - name="apmplus_meter", - endpoint=self.config.endpoint, - headers=self.headers, - resource_attributes=self.resource_attributes, + def get_metric_uploader(self) -> MetricUploaderProtocol: + """Lazily create and process-deduplicate the APMPlus metric pipeline.""" + registration_key = _metric_registration_key( + self.config.endpoint, + self.headers, + self.resource_attributes, ) + return metric_uploader_registry.get_or_create( + registration_key, + lambda: MeterUploader( + name="apmplus_meter", + endpoint=self.config.endpoint, + headers=self.headers, + resource_attributes=self.resource_attributes, + ), + ) + + @property + def meter_uploader(self) -> MetricUploaderProtocol: + """Compatibility accessor for the lazily registered uploader.""" + return self.get_metric_uploader() @override def export(self) -> None: @@ -568,6 +666,16 @@ def export(self) -> None: if self._exporter: self._exporter.force_flush() + metric_uploader = metric_uploader_registry.get( + _metric_registration_key( + self.config.endpoint, + self.headers, + self.resource_attributes, + ) + ) + if metric_uploader: + metric_uploader.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 9fe0d7e61..f9d158dd5 100644 --- a/veadk/tracing/telemetry/exporters/base_exporter.py +++ b/veadk/tracing/telemetry/exporters/base_exporter.py @@ -12,11 +12,18 @@ # See the License for the specific language governing permissions and # limitations under the License. +from __future__ import annotations + +from typing import TYPE_CHECKING + 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 +if TYPE_CHECKING: + from veadk.tracing.telemetry.metric_uploader import MetricUploader + class BaseExporter(BaseModel): """Abstract base class for OpenTelemetry span exporters in VeADK tracing system. @@ -57,3 +64,7 @@ def register(self, provider: TracerProvider) -> bool: def export(self) -> None: """Force export of telemetry data.""" pass + + def get_metric_uploader(self) -> MetricUploader | None: + """Return this exporter's optional metric uploader.""" + return None diff --git a/veadk/tracing/telemetry/metric_uploader.py b/veadk/tracing/telemetry/metric_uploader.py new file mode 100644 index 000000000..4237ee2a4 --- /dev/null +++ b/veadk/tracing/telemetry/metric_uploader.py @@ -0,0 +1,151 @@ +# 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. + +from __future__ import annotations + +from collections.abc import Callable, Hashable +from threading import RLock +from typing import Any, Protocol + +from veadk.utils.logger import get_logger + +logger = get_logger(__name__) + + +class MetricUploader(Protocol): + """Interface implemented by exporter-specific metric uploaders.""" + + @property + def registration_key(self) -> Hashable: + """Return a stable, non-logged key for one metric destination.""" + + def record_call_llm(self, *args: Any) -> None: ... + + def record_tool_call(self, *args: Any) -> None: ... + + def record_skill_call(self, *args: Any) -> None: ... + + def force_flush(self) -> bool: ... + + def shutdown(self) -> None: ... + + +class MetricUploaderRegistry: + """Process-level registry for active metric destinations. + + Uploaders are deduplicated by their registration key. Different keys remain + active together, so one process can intentionally publish metrics to more + than one destination without repeatedly scanning tracers and exporters. + """ + + def __init__(self) -> None: + self._lock = RLock() + self._uploaders: dict[Hashable, MetricUploader] = {} + + @property + def uploaders(self) -> tuple[MetricUploader, ...]: + with self._lock: + return tuple(self._uploaders.values()) + + def get(self, registration_key: Hashable) -> MetricUploader | None: + with self._lock: + return self._uploaders.get(registration_key) + + def get_or_create( + self, + registration_key: Hashable, + factory: Callable[[], MetricUploader], + ) -> MetricUploader: + """Return the registered uploader, constructing it only when absent.""" + with self._lock: + uploader = self._uploaders.get(registration_key) + if uploader is None: + uploader = factory() + self._uploaders[registration_key] = uploader + logger.debug( + "Registered metric uploader `%s`.", + uploader.__class__.__name__, + ) + return uploader + + def register(self, uploader: MetricUploader) -> MetricUploader: + """Register an uploader, closing a duplicate instance if necessary.""" + with self._lock: + existing = self._uploaders.get(uploader.registration_key) + if existing is None: + self._uploaders[uploader.registration_key] = uploader + logger.debug( + "Registered metric uploader `%s`.", + uploader.__class__.__name__, + ) + return uploader + + if existing is not uploader: + uploader.shutdown() + return existing + + def record_call_llm(self, *args: Any) -> None: + self._fan_out("record_call_llm", *args) + + def record_tool_call(self, *args: Any) -> None: + self._fan_out("record_tool_call", *args) + + def record_skill_call(self, *args: Any) -> None: + self._fan_out("record_skill_call", *args) + + def force_flush(self) -> bool: + flushed = True + for uploader in self.uploaders: + try: + flushed = bool(uploader.force_flush()) and flushed + except Exception as e: + flushed = False + logger.warning( + "Failed to flush metric uploader `%s`: %s", + uploader.__class__.__name__, + e, + ) + return flushed + + def clear(self, *, shutdown: bool = True) -> None: + """Remove all uploaders, optionally shutting down their SDK providers.""" + with self._lock: + uploaders = tuple(self._uploaders.values()) + self._uploaders.clear() + + if shutdown: + for uploader in uploaders: + try: + uploader.shutdown() + except Exception as e: + logger.warning( + "Failed to shut down metric uploader `%s`: %s", + uploader.__class__.__name__, + e, + ) + + def _fan_out(self, method_name: str, *args: Any) -> None: + for uploader in self.uploaders: + try: + getattr(uploader, method_name)(*args) + except Exception as e: + logger.warning( + "Metric uploader `%s` failed in %s: %s", + uploader.__class__.__name__, + method_name, + e, + ) + + +metric_uploader_registry = MetricUploaderRegistry() diff --git a/veadk/tracing/telemetry/opentelemetry_tracer.py b/veadk/tracing/telemetry/opentelemetry_tracer.py index 282ee63b2..439b966a9 100644 --- a/veadk/tracing/telemetry/opentelemetry_tracer.py +++ b/veadk/tracing/telemetry/opentelemetry_tracer.py @@ -28,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.metric_uploader import metric_uploader_registry from veadk.utils.logger import get_logger from veadk.utils.misc import get_agent_dir from veadk.utils.patches import patch_google_adk_telemetry @@ -160,7 +161,7 @@ def _init_global_tracer_provider(self) -> None: self._apmplus_managed_externally = have_global_tracer_provider for exporter in self.exporters: - self._register_exporter(exporter) + self._activate_exporter(exporter) self._inmemory_exporter = InMemoryExporter() if self._inmemory_exporter.processor: @@ -181,12 +182,11 @@ 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, - ) - - init_global_meter_uploader_from_exporters(self.exporters) + def _activate_exporter(self, exporter: BaseExporter) -> bool: + metric_uploader = exporter.get_metric_uploader() + if metric_uploader: + metric_uploader_registry.register(metric_uploader) + return self._register_exporter(exporter) def _register_exporter(self, exporter: BaseExporter) -> bool: if isinstance(exporter, APMPlusExporter) and self.apmplus_managed_externally: @@ -219,14 +219,7 @@ def add_exporter(self, exporter: BaseExporter) -> bool: if not any(existing is exporter for existing in self.exporters): self.exporters.append(exporter) - registered = self._register_exporter(exporter) - - from veadk.tracing.telemetry.telemetry import ( - init_global_meter_uploader_from_exporters, - ) - - init_global_meter_uploader_from_exporters(self.exporters) - return registered + return self._activate_exporter(exporter) @property def trace_file_path(self) -> str: @@ -266,6 +259,7 @@ def force_export(self) -> None: for processor in self._processors: time.sleep(0.05) processor.force_flush() + metric_uploader_registry.force_flush() @override def dump( diff --git a/veadk/tracing/telemetry/telemetry.py b/veadk/tracing/telemetry/telemetry.py index fcf30811f..c0d1dfd5e 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.metric_uploader import metric_uploader_registry 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, @@ -63,8 +45,7 @@ def _upload_call_llm_metrics( ) -> None: """Upload LLM call metrics to configured meter uploaders. - This function extracts meter uploaders from agent tracers and records - LLM call metrics including token usage, latency, and request/response details. + This function records metrics through the process-level uploader registry. Args: invocation_context: Context containing agent, session, and user information @@ -72,18 +53,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 - ) + metric_uploader_registry.record_call_llm( + invocation_context, event_id, llm_request, llm_response + ) def _upload_tool_call_metrics( @@ -91,7 +63,7 @@ def _upload_tool_call_metrics( args: dict[str, Any], function_response_event: Event, ): - """Upload tool call metrics to the global meter uploader. + """Upload tool call metrics to all registered meter uploaders. Records tool execution metrics including function name, arguments, execution time, and response details for observability and debugging. @@ -101,15 +73,12 @@ 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) + if metric_uploader_registry.uploaders: + metric_uploader_registry.record_tool_call(tool, args, function_response_event) else: logger.debug( - "Meter uploader is not initialized yet. Skip recording tool call metrics." + "No meter uploader is registered. Skip recording tool call metrics." ) From ad67b7695683ef49a05b480bd6e9abb9be39c486 Mon Sep 17 00:00:00 2001 From: "fengduzhen.666" Date: Wed, 5 Aug 2026 15:31:22 +0800 Subject: [PATCH 3/4] fix(metrics): record APMPlus metrics through global provider --- tests/test_metric_uploader.py | 81 ++++++++++++++- .../telemetry/exporters/apmplus_exporter.py | 99 +++++++------------ 2 files changed, 112 insertions(+), 68 deletions(-) diff --git a/tests/test_metric_uploader.py b/tests/test_metric_uploader.py index 6ceee0105..228f81759 100644 --- a/tests/test_metric_uploader.py +++ b/tests/test_metric_uploader.py @@ -77,10 +77,86 @@ def preconfigured_zero_reader_meter_provider(monkeypatch): provider.shutdown() -def test_meter_uploader_uses_private_provider_when_global_provider_exists( +@pytest.fixture +def preconfigured_meter_provider(monkeypatch): + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) + reader = InMemoryMetricReader() + provider = metrics_sdk.MeterProvider(metric_readers=[reader]) + metrics_api.set_meter_provider(provider) + yield provider, reader + provider.shutdown() + + +def test_meter_uploader_records_to_preconfigured_global_provider( + monkeypatch, + preconfigured_meter_provider, +): + provider, reader = preconfigured_meter_provider + + def fail_if_apmplus_reader_is_created(**kwargs): + raise AssertionError( + "APMPlus must not add a metric pipeline to an existing global provider" + ) + + monkeypatch.setattr( + apmplus_exporter_module, + "OTLPMetricExporter", + fail_if_apmplus_reader_is_created, + ) + + uploader = MeterUploader( + name="test-meter", + endpoint="http://localhost:4319", + headers={"x-byteapm-appkey": "test"}, + resource_attributes={"service.name": "test-service"}, + ) + uploader.llm_invoke_counter.add(1) + + assert metrics_api.get_meter_provider() is provider + assert uploader.provider is provider + assert uploader.force_flush() + assert reader.get_metrics_data() is not None + + uploader.shutdown() + assert provider._shutdown is False + + +def test_meter_uploader_keeps_preconfigured_zero_reader_global_provider( monkeypatch, preconfigured_zero_reader_meter_provider, ): + def fail_if_apmplus_reader_is_created(**kwargs): + raise AssertionError( + "APMPlus must not add a metric pipeline to an existing global provider" + ) + + monkeypatch.setattr( + apmplus_exporter_module, + "OTLPMetricExporter", + fail_if_apmplus_reader_is_created, + ) + + uploader = MeterUploader( + name="test-meter", + endpoint="http://localhost:4319", + headers={"x-byteapm-appkey": "test"}, + resource_attributes={"service.name": "test-service"}, + ) + uploader.llm_invoke_counter.add(1) + + assert uploader.provider is preconfigured_zero_reader_meter_provider + assert uploader.provider._sdk_config.metric_readers == () + + uploader.shutdown() + assert preconfigured_zero_reader_meter_provider._shutdown is False + + +def test_meter_uploader_installs_apmplus_global_provider_when_none_exists( + monkeypatch, +): + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) + monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) reader = InMemoryMetricReader() exporter_kwargs = {} monkeypatch.setattr( @@ -102,8 +178,7 @@ def test_meter_uploader_uses_private_provider_when_global_provider_exists( ) uploader.llm_invoke_counter.add(1) - assert metrics_api.get_meter_provider() is preconfigured_zero_reader_meter_provider - assert uploader.provider is not preconfigured_zero_reader_meter_provider + assert metrics_api.get_meter_provider() is uploader.provider assert len(uploader.provider._sdk_config.metric_readers) == 1 assert exporter_kwargs["insecure"] is True assert uploader.force_flush() diff --git a/veadk/tracing/telemetry/exporters/apmplus_exporter.py b/veadk/tracing/telemetry/exporters/apmplus_exporter.py index d80d4afa9..093b6a007 100644 --- a/veadk/tracing/telemetry/exporters/apmplus_exporter.py +++ b/veadk/tracing/telemetry/exporters/apmplus_exporter.py @@ -12,8 +12,6 @@ # See the License for the specific language governing permissions and # limitations under the License. -import hashlib -import json import time from collections.abc import Hashable from dataclasses import dataclass @@ -29,7 +27,7 @@ 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 Meter, _ProxyMeterProvider from opentelemetry.sdk import metrics as metrics_sdk from opentelemetry.sdk.metrics.export import PeriodicExportingMetricReader from opentelemetry.sdk.resources import Resource @@ -118,23 +116,7 @@ ] -def _metric_registration_key( - endpoint: str, - headers: dict, - resource_attributes: dict, -) -> tuple[str, str]: - """Build a stable digest without exposing authentication headers in logs.""" - payload = json.dumps( - { - "endpoint": endpoint, - "headers": headers, - "resource_attributes": resource_attributes, - }, - sort_keys=True, - default=str, - separators=(",", ":"), - ) - return "apmplus", hashlib.sha256(payload.encode()).hexdigest() +_APMPLUS_PORTAL_METRIC_KEY = ("apmplus", "portal-metrics") @dataclass @@ -208,9 +190,8 @@ def __init__( ) -> None: """Initialize the meter uploader with APMPlus configuration. - Sets up a private metrics provider, creates metric instruments, - and configures OTLP export to APMPlus endpoints with proper - resource attribution and authentication. + Creates Portal metric instruments on the global MeterProvider. If no + global provider exists yet, it installs one with an APMPlus OTLP reader. Args: name: Meter name for identification and organization @@ -218,37 +199,34 @@ def __init__( headers: Authentication headers including APMPlus app key resource_attributes: Service metadata for metric attribution """ - self._registration_key = _metric_registration_key( - endpoint, headers, resource_attributes - ) + self._registration_key = _APMPLUS_PORTAL_METRIC_KEY self._shutdown = False + self._owns_provider = False - # Reuse global resource attributes when available, but keep the reader - # and provider private so an existing zero-reader global provider cannot - # suppress APMPlus metrics. 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, - insecure=True, - ) - metric_reader = PeriodicExportingMetricReader(exporter) + if isinstance(global_metrics_provider, _ProxyMeterProvider): + exporter = OTLPMetricExporter( + endpoint=endpoint, + headers=headers, + insecure=True, + ) + metric_reader = PeriodicExportingMetricReader(exporter) + provider = metrics_sdk.MeterProvider( + metric_readers=[metric_reader], + resource=Resource.create(resource_attributes), + ) + metrics_api.set_meter_provider(provider) + global_metrics_provider = metrics_api.get_meter_provider() - self._provider = metrics_sdk.MeterProvider( - metric_readers=[metric_reader], resource=resource - ) + if global_metrics_provider is provider: + self._owns_provider = True + else: + # Another component won the set-once race. Its global provider + # owns metric export, so close the unused APMPlus pipeline. + provider.shutdown() - # 3. init meter + self._provider = global_metrics_provider self.meter: Meter = self._provider.get_meter(name=name) # create meter attributes @@ -318,17 +296,19 @@ def registration_key(self) -> Hashable: return self._registration_key @property - def provider(self) -> metrics_sdk.MeterProvider: + def provider(self) -> metrics_api.MeterProvider: return self._provider def force_flush(self) -> bool: if self._shutdown: return False - return self._provider.force_flush() + force_flush = getattr(self._provider, "force_flush", None) + return bool(force_flush()) if force_flush else True def shutdown(self) -> None: if not self._shutdown: - self._provider.shutdown() + if self._owns_provider: + self._provider.shutdown() # type: ignore[attr-defined] self._shutdown = True def record_call_llm( @@ -629,14 +609,9 @@ def model_post_init(self, context: Any) -> None: self.processor = BatchSpanProcessor(self._exporter) def get_metric_uploader(self) -> MetricUploaderProtocol: - """Lazily create and process-deduplicate the APMPlus metric pipeline.""" - registration_key = _metric_registration_key( - self.config.endpoint, - self.headers, - self.resource_attributes, - ) + """Enable Portal metric recording once for the process.""" return metric_uploader_registry.get_or_create( - registration_key, + _APMPLUS_PORTAL_METRIC_KEY, lambda: MeterUploader( name="apmplus_meter", endpoint=self.config.endpoint, @@ -666,13 +641,7 @@ def export(self) -> None: if self._exporter: self._exporter.force_flush() - metric_uploader = metric_uploader_registry.get( - _metric_registration_key( - self.config.endpoint, - self.headers, - self.resource_attributes, - ) - ) + metric_uploader = metric_uploader_registry.get(_APMPLUS_PORTAL_METRIC_KEY) if metric_uploader: metric_uploader.force_flush() From e064701de0f3e40d6746cf0c51bdcdce56e5f2e5 Mon Sep 17 00:00:00 2001 From: "fengduzhen.666" Date: Wed, 5 Aug 2026 16:01:23 +0800 Subject: [PATCH 4/4] refactor(metrics): decouple Portal recording from exporters --- tests/test_metric_uploader.py | 244 -------- tests/test_portal_metrics.py | 221 +++++++ tests/test_tracing.py | 60 +- tests/test_tracing_content.py | 19 +- veadk/tools/skills_tools/skills_tool.py | 6 +- .../telemetry/exporters/apmplus_exporter.py | 539 ++---------------- .../telemetry/exporters/base_exporter.py | 9 - veadk/tracing/telemetry/metric_uploader.py | 151 ----- .../tracing/telemetry/opentelemetry_tracer.py | 7 +- veadk/tracing/telemetry/portal_metrics.py | 450 +++++++++++++++ veadk/tracing/telemetry/telemetry.py | 20 +- 11 files changed, 742 insertions(+), 984 deletions(-) delete mode 100644 tests/test_metric_uploader.py create mode 100644 tests/test_portal_metrics.py delete mode 100644 veadk/tracing/telemetry/metric_uploader.py create mode 100644 veadk/tracing/telemetry/portal_metrics.py diff --git a/tests/test_metric_uploader.py b/tests/test_metric_uploader.py deleted file mode 100644 index 228f81759..000000000 --- a/tests/test_metric_uploader.py +++ /dev/null @@ -1,244 +0,0 @@ -# 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 telemetry -from veadk.tracing.telemetry.exporters import ( - apmplus_exporter as apmplus_exporter_module, -) -from veadk.tracing.telemetry.exporters.apmplus_exporter import MeterUploader -from veadk.tracing.telemetry.metric_uploader import ( - MetricUploaderRegistry, - metric_uploader_registry, -) - - -class FakeMetricUploader: - def __init__(self, registration_key): - self.registration_key = registration_key - self.llm_calls = [] - self.tool_calls = [] - self.skill_calls = [] - self.force_flush_calls = 0 - self.shutdown_calls = 0 - - 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) - - def force_flush(self): - self.force_flush_calls += 1 - return True - - def shutdown(self): - self.shutdown_calls += 1 - - -@pytest.fixture(autouse=True) -def fresh_metric_uploader_registry(): - metric_uploader_registry.clear() - yield - metric_uploader_registry.clear() - - -@pytest.fixture -def preconfigured_zero_reader_meter_provider(monkeypatch): - monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) - monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) - provider = metrics_sdk.MeterProvider() - metrics_api.set_meter_provider(provider) - yield provider - provider.shutdown() - - -@pytest.fixture -def preconfigured_meter_provider(monkeypatch): - monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) - monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) - reader = InMemoryMetricReader() - provider = metrics_sdk.MeterProvider(metric_readers=[reader]) - metrics_api.set_meter_provider(provider) - yield provider, reader - provider.shutdown() - - -def test_meter_uploader_records_to_preconfigured_global_provider( - monkeypatch, - preconfigured_meter_provider, -): - provider, reader = preconfigured_meter_provider - - def fail_if_apmplus_reader_is_created(**kwargs): - raise AssertionError( - "APMPlus must not add a metric pipeline to an existing global provider" - ) - - monkeypatch.setattr( - apmplus_exporter_module, - "OTLPMetricExporter", - fail_if_apmplus_reader_is_created, - ) - - uploader = MeterUploader( - name="test-meter", - endpoint="http://localhost:4319", - headers={"x-byteapm-appkey": "test"}, - resource_attributes={"service.name": "test-service"}, - ) - uploader.llm_invoke_counter.add(1) - - assert metrics_api.get_meter_provider() is provider - assert uploader.provider is provider - assert uploader.force_flush() - assert reader.get_metrics_data() is not None - - uploader.shutdown() - assert provider._shutdown is False - - -def test_meter_uploader_keeps_preconfigured_zero_reader_global_provider( - monkeypatch, - preconfigured_zero_reader_meter_provider, -): - def fail_if_apmplus_reader_is_created(**kwargs): - raise AssertionError( - "APMPlus must not add a metric pipeline to an existing global provider" - ) - - monkeypatch.setattr( - apmplus_exporter_module, - "OTLPMetricExporter", - fail_if_apmplus_reader_is_created, - ) - - uploader = MeterUploader( - name="test-meter", - endpoint="http://localhost:4319", - headers={"x-byteapm-appkey": "test"}, - resource_attributes={"service.name": "test-service"}, - ) - uploader.llm_invoke_counter.add(1) - - assert uploader.provider is preconfigured_zero_reader_meter_provider - assert uploader.provider._sdk_config.metric_readers == () - - uploader.shutdown() - assert preconfigured_zero_reader_meter_provider._shutdown is False - - -def test_meter_uploader_installs_apmplus_global_provider_when_none_exists( - monkeypatch, -): - monkeypatch.setattr(metrics_internal, "_METER_PROVIDER", None) - monkeypatch.setattr(metrics_internal, "_METER_PROVIDER_SET_ONCE", Once()) - 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, - ) - - uploader = MeterUploader( - name="test-meter", - endpoint="http://localhost:4319", - headers={"x-byteapm-appkey": "test"}, - resource_attributes={"service.name": "test-service"}, - ) - uploader.llm_invoke_counter.add(1) - - assert metrics_api.get_meter_provider() is uploader.provider - assert len(uploader.provider._sdk_config.metric_readers) == 1 - assert exporter_kwargs["insecure"] is True - assert uploader.force_flush() - assert reader.get_metrics_data() is not None - - uploader.shutdown() - - -def test_registry_deduplicates_by_key_without_constructing_a_duplicate(): - registry = MetricUploaderRegistry() - first = FakeMetricUploader(("apmplus", "same-destination")) - - assert registry.register(first) is first - - factory_calls = [] - resolved = registry.get_or_create( - first.registration_key, - lambda: factory_calls.append(True) - or FakeMetricUploader(first.registration_key), - ) - - assert resolved is first - assert factory_calls == [] - assert registry.uploaders == (first,) - - -def test_registry_fans_out_to_distinct_destinations_once(): - registry = MetricUploaderRegistry() - first = FakeMetricUploader(("apmplus", "destination-a")) - second = FakeMetricUploader(("apmplus", "destination-b")) - duplicate = FakeMetricUploader(first.registration_key) - - registry.register(first) - registry.register(second) - assert registry.register(duplicate) is first - - registry.record_call_llm("context", "event", "request", "response") - registry.record_tool_call("tool", {}, "event") - registry.record_skill_call("span", "skill", "tool", None, "ok") - - assert len(first.llm_calls) == len(second.llm_calls) == 1 - assert len(first.tool_calls) == len(second.tool_calls) == 1 - assert len(first.skill_calls) == len(second.skill_calls) == 1 - assert duplicate.llm_calls == [] - assert duplicate.shutdown_calls == 1 - assert registry.force_flush() - assert first.force_flush_calls == second.force_flush_calls == 1 - - -def test_telemetry_and_skill_metrics_use_the_registry(): - uploader = FakeMetricUploader(("apmplus", "destination")) - metric_uploader_registry.register(uploader) - - 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(uploader.llm_calls) == 1 - assert len(uploader.tool_calls) == 1 - assert len(uploader.skill_calls) == 1 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 2a7971f36..0519b6aeb 100644 --- a/tests/test_tracing.py +++ b/tests/test_tracing.py @@ -46,7 +46,6 @@ TLSExporter, TLSExporterConfig, ) -from veadk.tracing.telemetry.metric_uploader import metric_uploader_registry from veadk.tracing.telemetry.opentelemetry_tracer import OpentelemetryTracer APP_NAME = "app" @@ -98,7 +97,6 @@ def gen_span_processor(endpoint: str): @pytest.fixture def fresh_global_tracer_provider(monkeypatch): """Give each test an isolated OpenTelemetry global provider.""" - metric_uploader_registry.clear() monkeypatch.setattr(trace_api, "_TRACER_PROVIDER", None) monkeypatch.setattr(trace_api, "_TRACER_PROVIDER_SET_ONCE", Once()) @@ -107,7 +105,6 @@ def fresh_global_tracer_provider(monkeypatch): tracer_provider = trace_api.get_tracer_provider() if isinstance(tracer_provider, trace_sdk.TracerProvider): tracer_provider.shutdown() - metric_uploader_registry.clear() @pytest.fixture @@ -115,24 +112,6 @@ def controlled_apmplus_exporter(monkeypatch): """Provide an APMPlus exporter without network or global meter side effects.""" constructed_exporters = [] - class ControlledMetricUploader: - registration_key = ("apmplus", "controlled") - - def __init__(self): - self.force_flush_calls = 0 - - def record_call_llm(self, *args): ... - - def record_tool_call(self, *args): ... - - def record_skill_call(self, *args): ... - - def force_flush(self): - self.force_flush_calls += 1 - return True - - def shutdown(self): ... - class ControlledAPMPlusExporter(APMPlusExporter): def __init__(self): super().__init__( @@ -146,12 +125,8 @@ def __init__(self): def model_post_init(self, context): self._exporter = OTelInMemorySpanExporter() self.processor = SimpleSpanProcessor(self._exporter) - self._controlled_meter_uploader = ControlledMetricUploader() constructed_exporters.append(self) - def get_metric_uploader(self): - return self._controlled_meter_uploader - monkeypatch.setattr( apmplus_exporter_module, "APMPlusExporter", @@ -213,7 +188,6 @@ def test_apmplus_enable_provider_and_manual_exporter_matrix( if not should_create_tracer: assert final_provider is initial_provider assert constructed_exporters == [] - assert metric_uploader_registry.uploaders == () return tracer = agent.tracers[0] @@ -240,12 +214,6 @@ def test_apmplus_enable_provider_and_manual_exporter_matrix( 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 - expected_uploaders = ( - (constructed_exporters[0].get_metric_uploader(),) - if should_construct_apmplus - else () - ) - assert metric_uploader_registry.uploaders == expected_uploaders def test_add_exporter_registers_after_tracer_initialization( @@ -276,17 +244,21 @@ def test_add_exporter_is_idempotent(fresh_global_tracer_provider): assert sum(processor is exporter.processor for processor in tracer._processors) == 1 -def test_force_export_flushes_metric_uploader_registry( +def test_force_export_flushes_portal_metrics( fresh_global_tracer_provider, - controlled_apmplus_exporter, + monkeypatch, ): - controlled_exporter_class, _ = controlled_apmplus_exporter - exporter = controlled_exporter_class() - tracer = OpentelemetryTracer(exporters=[exporter]) + flush_calls = [] + monkeypatch.setattr( + opentelemetry_tracer_module, + "portal_metric_recorder", + SimpleNamespace(force_flush=lambda: flush_calls.append(True) or True), + ) + tracer = OpentelemetryTracer() tracer.force_export() - assert exporter.get_metric_uploader().force_flush_calls == 1 + assert flush_calls == [True] def test_tracing_registers_apmplus_without_global_provider( @@ -317,9 +289,6 @@ def test_tracing_skips_apmplus_for_preconfigured_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 - assert metric_uploader_registry.uploaders == ( - apmplus_exporter.get_metric_uploader(), - ) def test_add_exporter_skips_apmplus_for_preconfigured_provider( @@ -336,12 +305,9 @@ def test_add_exporter_skips_apmplus_for_preconfigured_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 - assert metric_uploader_registry.uploaders == ( - apmplus_exporter.get_metric_uploader(), - ) -def test_agent_env_keeps_metrics_only_apmplus_for_preconfigured_provider( +def test_agent_env_keeps_apmplus_for_preconfigured_provider( fresh_global_tracer_provider, controlled_apmplus_exporter, monkeypatch, @@ -362,7 +328,6 @@ def test_agent_env_keeps_metrics_only_apmplus_for_preconfigured_provider( assert isinstance(exporter, controlled_exporter_class) assert exporter in tracer.exporters assert exporter.processor not in span_processors - assert metric_uploader_registry.uploaders == (exporter.get_metric_uploader(),) assert len(span_processors) == 1 # VeADK in-memory processor only @@ -386,7 +351,8 @@ async def test_tracing_with_global_provider(fresh_global_tracer_provider): trace_api.set_tracer_provider(tracer_provider) tracer = OpentelemetryTracer(exporters=exporters) - # APMPlus is retained for metrics but its span processor is not registered. + # APMPlus is retained so it can ensure a global metric pipeline, while its + # span processor is not registered. assert len(tracer.exporters) == 4 diff --git a/tests/test_tracing_content.py b/tests/test_tracing_content.py index 3964abe2c..4b01edf3c 100644 --- a/tests/test_tracing_content.py +++ b/tests/test_tracing_content.py @@ -21,8 +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.metric_uploader import metric_uploader_registry +from veadk.tracing.telemetry.portal_metrics import PortalMetricRecorder @dataclass @@ -144,10 +143,6 @@ def _event_names(span): return [event.name for event in span.events] -def setup_function(): - metric_uploader_registry.clear() - - def test_trace_call_llm_records_content_by_default(monkeypatch): monkeypatch.delenv("OBSERVABILITY_OPENTELEMETRY_TRACE_CONTENT", raising=False) @@ -218,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/tools/skills_tools/skills_tool.py b/veadk/tools/skills_tools/skills_tool.py index 80a8b3bd3..7b491b8a8 100644 --- a/veadk/tools/skills_tools/skills_tool.py +++ b/veadk/tools/skills_tools/skills_tool.py @@ -530,11 +530,9 @@ 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: - from veadk.tracing.telemetry.metric_uploader import ( - metric_uploader_registry, - ) + from veadk.tracing.telemetry import portal_metrics - metric_uploader_registry.record_skill_call( + portal_metrics.portal_metric_recorder.record_skill_call( span, skill_name, self.name, diff --git a/veadk/tracing/telemetry/exporters/apmplus_exporter.py b/veadk/tracing/telemetry/exporters/apmplus_exporter.py index 093b6a007..f8f4791fb 100644 --- a/veadk/tracing/telemetry/exporters/apmplus_exporter.py +++ b/veadk/tracing/telemetry/exporters/apmplus_exporter.py @@ -12,22 +12,12 @@ # See the License for the specific language governing permissions and # limitations under the License. -import time -from collections.abc import Hashable -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 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, _ProxyMeterProvider +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 @@ -37,483 +27,40 @@ from veadk.config import settings from veadk.tracing.telemetry.exporters.base_exporter import BaseExporter -from veadk.tracing.telemetry.metric_uploader import ( - MetricUploader as MetricUploaderProtocol, - metric_uploader_registry, -) 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, -] - - -_APMPLUS_PORTAL_METRIC_KEY = ("apmplus", "portal-metrics") - - -@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() - Creates Portal metric instruments on the global MeterProvider. If no - global provider exists yet, it installs one with an APMPlus OTLP reader. - - 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 - """ - self._registration_key = _APMPLUS_PORTAL_METRIC_KEY - self._shutdown = False - self._owns_provider = False - - global_metrics_provider = metrics_api.get_meter_provider() - - if isinstance(global_metrics_provider, _ProxyMeterProvider): - exporter = OTLPMetricExporter( - endpoint=endpoint, - headers=headers, - insecure=True, - ) - metric_reader = PeriodicExportingMetricReader(exporter) - provider = metrics_sdk.MeterProvider( - metric_readers=[metric_reader], - resource=Resource.create(resource_attributes), - ) - metrics_api.set_meter_provider(provider) - global_metrics_provider = metrics_api.get_meter_provider() - - if global_metrics_provider is provider: - self._owns_provider = True - else: - # Another component won the set-once race. Its global provider - # owns metric export, so close the unused APMPlus pipeline. - provider.shutdown() - - self._provider = global_metrics_provider - self.meter: Meter = self._provider.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 registration_key(self) -> Hashable: - return self._registration_key - - @property - def provider(self) -> metrics_api.MeterProvider: - return self._provider - - def force_flush(self) -> bool: - if self._shutdown: - return False - force_flush = getattr(self._provider, "force_flush", None) - return bool(force_flush()) if force_flush else True - - def shutdown(self) -> None: - if not self._shutdown: - if self._owns_provider: - self._provider.shutdown() # type: ignore[attr-defined] - self._shutdown = 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) + return resolved_provider class APMPlusExporterConfig(BaseModel): @@ -581,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}") @@ -608,23 +154,12 @@ def model_post_init(self, context: Any) -> None: ) self.processor = BatchSpanProcessor(self._exporter) - def get_metric_uploader(self) -> MetricUploaderProtocol: - """Enable Portal metric recording once for the process.""" - return metric_uploader_registry.get_or_create( - _APMPLUS_PORTAL_METRIC_KEY, - lambda: MeterUploader( - name="apmplus_meter", - endpoint=self.config.endpoint, - headers=self.headers, - resource_attributes=self.resource_attributes, - ), + ensure_apmplus_meter_provider( + endpoint=self.config.endpoint, + headers=self.headers, + resource_attributes=self.resource_attributes, ) - @property - def meter_uploader(self) -> MetricUploaderProtocol: - """Compatibility accessor for the lazily registered uploader.""" - return self.get_metric_uploader() - @override def export(self) -> None: """Force immediate export of pending telemetry data to APMPlus. @@ -641,9 +176,11 @@ def export(self) -> None: if self._exporter: self._exporter.force_flush() - metric_uploader = metric_uploader_registry.get(_APMPLUS_PORTAL_METRIC_KEY) - if metric_uploader: - metric_uploader.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 f9d158dd5..323763c71 100644 --- a/veadk/tracing/telemetry/exporters/base_exporter.py +++ b/veadk/tracing/telemetry/exporters/base_exporter.py @@ -14,16 +14,11 @@ from __future__ import annotations -from typing import TYPE_CHECKING - 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 -if TYPE_CHECKING: - from veadk.tracing.telemetry.metric_uploader import MetricUploader - class BaseExporter(BaseModel): """Abstract base class for OpenTelemetry span exporters in VeADK tracing system. @@ -64,7 +59,3 @@ def register(self, provider: TracerProvider) -> bool: def export(self) -> None: """Force export of telemetry data.""" pass - - def get_metric_uploader(self) -> MetricUploader | None: - """Return this exporter's optional metric uploader.""" - return None diff --git a/veadk/tracing/telemetry/metric_uploader.py b/veadk/tracing/telemetry/metric_uploader.py deleted file mode 100644 index 4237ee2a4..000000000 --- a/veadk/tracing/telemetry/metric_uploader.py +++ /dev/null @@ -1,151 +0,0 @@ -# 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. - -from __future__ import annotations - -from collections.abc import Callable, Hashable -from threading import RLock -from typing import Any, Protocol - -from veadk.utils.logger import get_logger - -logger = get_logger(__name__) - - -class MetricUploader(Protocol): - """Interface implemented by exporter-specific metric uploaders.""" - - @property - def registration_key(self) -> Hashable: - """Return a stable, non-logged key for one metric destination.""" - - def record_call_llm(self, *args: Any) -> None: ... - - def record_tool_call(self, *args: Any) -> None: ... - - def record_skill_call(self, *args: Any) -> None: ... - - def force_flush(self) -> bool: ... - - def shutdown(self) -> None: ... - - -class MetricUploaderRegistry: - """Process-level registry for active metric destinations. - - Uploaders are deduplicated by their registration key. Different keys remain - active together, so one process can intentionally publish metrics to more - than one destination without repeatedly scanning tracers and exporters. - """ - - def __init__(self) -> None: - self._lock = RLock() - self._uploaders: dict[Hashable, MetricUploader] = {} - - @property - def uploaders(self) -> tuple[MetricUploader, ...]: - with self._lock: - return tuple(self._uploaders.values()) - - def get(self, registration_key: Hashable) -> MetricUploader | None: - with self._lock: - return self._uploaders.get(registration_key) - - def get_or_create( - self, - registration_key: Hashable, - factory: Callable[[], MetricUploader], - ) -> MetricUploader: - """Return the registered uploader, constructing it only when absent.""" - with self._lock: - uploader = self._uploaders.get(registration_key) - if uploader is None: - uploader = factory() - self._uploaders[registration_key] = uploader - logger.debug( - "Registered metric uploader `%s`.", - uploader.__class__.__name__, - ) - return uploader - - def register(self, uploader: MetricUploader) -> MetricUploader: - """Register an uploader, closing a duplicate instance if necessary.""" - with self._lock: - existing = self._uploaders.get(uploader.registration_key) - if existing is None: - self._uploaders[uploader.registration_key] = uploader - logger.debug( - "Registered metric uploader `%s`.", - uploader.__class__.__name__, - ) - return uploader - - if existing is not uploader: - uploader.shutdown() - return existing - - def record_call_llm(self, *args: Any) -> None: - self._fan_out("record_call_llm", *args) - - def record_tool_call(self, *args: Any) -> None: - self._fan_out("record_tool_call", *args) - - def record_skill_call(self, *args: Any) -> None: - self._fan_out("record_skill_call", *args) - - def force_flush(self) -> bool: - flushed = True - for uploader in self.uploaders: - try: - flushed = bool(uploader.force_flush()) and flushed - except Exception as e: - flushed = False - logger.warning( - "Failed to flush metric uploader `%s`: %s", - uploader.__class__.__name__, - e, - ) - return flushed - - def clear(self, *, shutdown: bool = True) -> None: - """Remove all uploaders, optionally shutting down their SDK providers.""" - with self._lock: - uploaders = tuple(self._uploaders.values()) - self._uploaders.clear() - - if shutdown: - for uploader in uploaders: - try: - uploader.shutdown() - except Exception as e: - logger.warning( - "Failed to shut down metric uploader `%s`: %s", - uploader.__class__.__name__, - e, - ) - - def _fan_out(self, method_name: str, *args: Any) -> None: - for uploader in self.uploaders: - try: - getattr(uploader, method_name)(*args) - except Exception as e: - logger.warning( - "Metric uploader `%s` failed in %s: %s", - uploader.__class__.__name__, - method_name, - e, - ) - - -metric_uploader_registry = MetricUploaderRegistry() diff --git a/veadk/tracing/telemetry/opentelemetry_tracer.py b/veadk/tracing/telemetry/opentelemetry_tracer.py index 439b966a9..fdae30d7b 100644 --- a/veadk/tracing/telemetry/opentelemetry_tracer.py +++ b/veadk/tracing/telemetry/opentelemetry_tracer.py @@ -28,7 +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.metric_uploader import metric_uploader_registry +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 @@ -183,9 +183,6 @@ def _init_global_tracer_provider(self) -> None: ) def _activate_exporter(self, exporter: BaseExporter) -> bool: - metric_uploader = exporter.get_metric_uploader() - if metric_uploader: - metric_uploader_registry.register(metric_uploader) return self._register_exporter(exporter) def _register_exporter(self, exporter: BaseExporter) -> bool: @@ -259,7 +256,7 @@ def force_export(self) -> None: for processor in self._processors: time.sleep(0.05) processor.force_flush() - metric_uploader_registry.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 c0d1dfd5e..dba7d4886 100644 --- a/veadk/tracing/telemetry/telemetry.py +++ b/veadk/tracing/telemetry/telemetry.py @@ -30,7 +30,7 @@ ToolAttributesParams, ) from veadk.tracing.telemetry.content_tracing import should_trace_content -from veadk.tracing.telemetry.metric_uploader import metric_uploader_registry +from veadk.tracing.telemetry import portal_metrics from veadk.utils.logger import get_logger from veadk.utils.misc import safe_json_serialize @@ -43,9 +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 records metrics through the process-level uploader registry. + 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 @@ -53,7 +54,7 @@ def _upload_call_llm_metrics( llm_request: The request sent to the language model llm_response: The response received from the language model """ - metric_uploader_registry.record_call_llm( + portal_metrics.portal_metric_recorder.record_call_llm( invocation_context, event_id, llm_request, llm_response ) @@ -63,7 +64,7 @@ def _upload_tool_call_metrics( args: dict[str, Any], function_response_event: Event, ): - """Upload tool call metrics to all registered meter uploaders. + """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. @@ -74,12 +75,9 @@ def _upload_tool_call_metrics( function_response_event: Event containing the tool's response data """ - if metric_uploader_registry.uploaders: - metric_uploader_registry.record_tool_call(tool, args, function_response_event) - else: - logger.debug( - "No meter uploader is registered. Skip recording tool call metrics." - ) + portal_metrics.portal_metric_recorder.record_tool_call( + tool, args, function_response_event + ) def _set_agent_input_attribute(