Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
147 changes: 129 additions & 18 deletions tests/test_tracing.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand All @@ -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,
):
Expand All @@ -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()
Expand All @@ -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
Expand All @@ -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
6 changes: 3 additions & 3 deletions veadk/agent.py
Original file line number Diff line number Diff line change
Expand Up @@ -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(
Expand Down
22 changes: 21 additions & 1 deletion veadk/tracing/telemetry/exporters/base_exporter.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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
83 changes: 37 additions & 46 deletions veadk/tracing/telemetry/opentelemetry_tracer.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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.
Expand Down
Loading