diff --git a/CHANGELOG.md b/CHANGELOG.md index 32ca2729..98a797c1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,10 @@ # Release History # Unreleased ### Features Added +- Add network-isolated upstream compatibility tests for Azure resource detectors, + OpenTelemetry SDK telemetry contracts, dependency plugins, exporters, and + `service.instance.id` precedence and fork-refresh regressions. + ([#279](https://github.com/microsoft/opentelemetry-distro-python/pull/279)) - Add typed Agent365 execute-tool argument and result schema models with `schema_version: "1.0"` serialization, `ToolCallAction`/`ToolCallOutcomeStatus`/`ToolPolicyDecision` enums, provider extension data wrapped under the JSON `metadata` property, public exports, and `ExecuteToolScope` support while preserving raw diff --git a/tests/upstream/__init__.py b/tests/upstream/__init__.py new file mode 100644 index 00000000..e69de29b diff --git a/tests/upstream/conftest.py b/tests/upstream/conftest.py new file mode 100644 index 00000000..d5b07ddd --- /dev/null +++ b/tests/upstream/conftest.py @@ -0,0 +1,13 @@ +import socket + +import pytest + + +@pytest.fixture(autouse=True) +def block_external_network(monkeypatch): + def unexpected_network_call(*args, **kwargs): + pytest.fail("upstream compatibility tests must not perform external network calls") + + monkeypatch.delenv("OTEL_EXPERIMENTAL_RESOURCE_DETECTORS", raising=False) + monkeypatch.setattr(socket, "create_connection", unexpected_network_call) + monkeypatch.setattr("requests.sessions.Session.request", unexpected_network_call) diff --git a/tests/upstream/test_azure_resource_detectors.py b/tests/upstream/test_azure_resource_detectors.py new file mode 100644 index 00000000..7b86a453 --- /dev/null +++ b/tests/upstream/test_azure_resource_detectors.py @@ -0,0 +1,247 @@ +import json +import os +from importlib.metadata import entry_points +from urllib.error import URLError + +import pytest +from opentelemetry.resource.detector.azure.app_service import AzureAppServiceResourceDetector +from opentelemetry.resource.detector.azure.functions import AzureFunctionsResourceDetector +from opentelemetry.resource.detector.azure.vm import AzureVMResourceDetector +from opentelemetry.sdk.environment_variables import OTEL_EXPERIMENTAL_RESOURCE_DETECTORS +from opentelemetry.sdk.resources import Resource, ResourceDetector, get_aggregated_resources + +from microsoft.opentelemetry._azure_monitor._constants import RESOURCE_ARG +from microsoft.opentelemetry._azure_monitor._utils.configurations import _default_resource + +_AZURE_ENVIRONMENT_VARIABLES = ( + "AKS_ARM_NAMESPACE_ID", + "FUNCTIONS_WORKER_RUNTIME", + "REGION_NAME", + "WEBSITE_HOME_STAMPNAME", + "WEBSITE_HOSTNAME", + "WEBSITE_INSTANCE_ID", + "WEBSITE_MEMORY_LIMIT_MB", + "WEBSITE_OWNER_NAME", + "WEBSITE_RESOURCE_GROUP", + "WEBSITE_SITE_NAME", + "WEBSITE_SLOT_NAME", +) + + +@pytest.fixture(autouse=True) +def clear_azure_environment(monkeypatch): + for name in _AZURE_ENVIRONMENT_VARIABLES: + monkeypatch.delenv(name, raising=False) + monkeypatch.delenv(OTEL_EXPERIMENTAL_RESOURCE_DETECTORS, raising=False) + + +def test_azure_resource_detector_entry_points_remain_loadable(): + detector_entry_points = { + entry_point.name: entry_point + for entry_point in entry_points(group="opentelemetry_resource_detector") + if entry_point.name.startswith("azure_") + } + + required_detectors = {"azure_app_service", "azure_functions", "azure_vm"} + assert required_detectors <= detector_entry_points.keys() + for name in required_detectors: + entry_point = detector_entry_points[name] + detector_class = entry_point.load() + assert issubclass(detector_class, ResourceDetector) + assert isinstance(detector_class(), ResourceDetector) + + +def test_app_service_detector_maps_azure_environment_to_resource(monkeypatch): + monkeypatch.setenv("WEBSITE_SITE_NAME", "orders-api") + monkeypatch.setenv("REGION_NAME", "westus2") + monkeypatch.setenv("WEBSITE_SLOT_NAME", "staging") + monkeypatch.setenv("WEBSITE_HOSTNAME", "orders-api.azurewebsites.net") + monkeypatch.setenv("WEBSITE_INSTANCE_ID", "instance-123") + monkeypatch.setenv("WEBSITE_HOME_STAMPNAME", "waws-prod-bay-001") + monkeypatch.setenv("WEBSITE_OWNER_NAME", "subscription-id+resource-owner") + monkeypatch.setenv("WEBSITE_RESOURCE_GROUP", "observability-rg") + + attributes = AzureAppServiceResourceDetector().detect().attributes + + assert attributes == { + "azure.app.service.stamp": "waws-prod-bay-001", + "cloud.platform": "azure_app_service", + "cloud.provider": "azure", + "cloud.region": "westus2", + "cloud.resource_id": ( + "/subscriptions/subscription-id/resourceGroups/observability-rg/" "providers/Microsoft.Web/sites/orders-api" + ), + "deployment.environment": "staging", + "host.id": "orders-api.azurewebsites.net", + "service.instance.id": "instance-123", + "service.name": "orders-api", + } + + +def test_app_service_detector_returns_empty_resource_outside_app_service(): + assert not AzureAppServiceResourceDetector().detect().attributes + + +def test_app_service_detector_defers_service_identity_to_functions(monkeypatch): + monkeypatch.setenv("FUNCTIONS_WORKER_RUNTIME", "python") + monkeypatch.setenv("WEBSITE_SITE_NAME", "function-app") + monkeypatch.setenv("WEBSITE_INSTANCE_ID", "function-instance") + + attributes = AzureAppServiceResourceDetector().detect().attributes + + assert attributes["cloud.provider"] == "azure" + assert attributes["service.instance.id"] == "function-instance" + assert "cloud.platform" not in attributes + assert "service.name" not in attributes + + +def test_functions_detector_maps_azure_environment_to_resource(monkeypatch): + monkeypatch.setenv("FUNCTIONS_WORKER_RUNTIME", "python") + monkeypatch.setenv("WEBSITE_SITE_NAME", "queue-processor") + monkeypatch.setenv("REGION_NAME", "eastus") + monkeypatch.setenv("WEBSITE_INSTANCE_ID", "function-instance") + monkeypatch.setenv("WEBSITE_MEMORY_LIMIT_MB", "1536") + monkeypatch.setenv("WEBSITE_OWNER_NAME", "subscription-id") + monkeypatch.setenv("WEBSITE_RESOURCE_GROUP", "functions-rg") + + attributes = AzureFunctionsResourceDetector().detect().attributes + + assert attributes["cloud.platform"] == "azure_functions" + assert attributes["cloud.provider"] == "azure" + assert attributes["cloud.region"] == "eastus" + assert attributes["cloud.resource_id"] == ( + "/subscriptions/subscription-id/resourceGroups/functions-rg/" "providers/Microsoft.Web/sites/queue-processor" + ) + assert attributes["faas.instance"] == "function-instance" + assert attributes["faas.max_memory"] == 1536 + process_pid = attributes["process.pid"] + assert isinstance(process_pid, int) + assert process_pid > 0 + assert attributes["service.name"] == "queue-processor" + + +def test_functions_detector_returns_empty_resource_outside_functions(): + assert not AzureFunctionsResourceDetector().detect().attributes + + +def test_functions_detector_ignores_invalid_memory_limit(monkeypatch): + monkeypatch.setenv("FUNCTIONS_WORKER_RUNTIME", "python") + monkeypatch.setenv("WEBSITE_MEMORY_LIMIT_MB", "not-an-integer") + + attributes = AzureFunctionsResourceDetector().detect().attributes + + assert attributes["cloud.platform"] == "azure_functions" + assert "faas.max_memory" not in attributes + + +class _MetadataResponse: + def __init__(self, metadata): + self._payload = json.dumps(metadata).encode() + + def __enter__(self): + return self + + def __exit__(self, exc_type, exc_value, traceback): + return False + + def read(self): + return self._payload + + +def test_vm_detector_maps_metadata_service_response_to_resource(monkeypatch): + metadata = { + "location": "centralus", + "name": "worker-01", + "osType": "Linux", + "resourceId": "/subscriptions/sub/resourceGroups/rg/providers/Microsoft.Compute/virtualMachines/worker-01", + "sku": "2022-datacenter", + "version": "22.04", + "vmId": "vm-id-123", + "vmScaleSetName": "worker-pool", + "vmSize": "Standard_D4s_v5", + } + + def fake_urlopen(request, timeout): + assert request.full_url.endswith("api-version=2021-12-13&format=json") + assert request.get_header("Metadata") == "True" + assert timeout == 0.2 + return _MetadataResponse(metadata) + + monkeypatch.setattr("opentelemetry.resource.detector.azure.vm.urlopen", fake_urlopen) + + attributes = AzureVMResourceDetector().detect().attributes + + assert attributes == { + "azure.vm.scaleset.name": "worker-pool", + "azure.vm.sku": "2022-datacenter", + "cloud.platform": "azure_vm", + "cloud.provider": "azure", + "cloud.region": "centralus", + "cloud.resource_id": metadata["resourceId"], + "host.id": "vm-id-123", + "host.name": "worker-01", + "host.type": "Standard_D4s_v5", + "os.type": "Linux", + "os.version": "22.04", + "service.instance.id": "vm-id-123", + } + + +def test_vm_detector_handles_non_azure_hosts_without_failing(monkeypatch): + def unavailable_metadata_service(request, timeout): + raise URLError("metadata service unavailable") + + monkeypatch.setattr("opentelemetry.resource.detector.azure.vm.urlopen", unavailable_metadata_service) + + assert AzureVMResourceDetector().detect().attributes == Resource.get_empty().attributes + + +@pytest.mark.parametrize( + ("environment_variable", "value"), + [ + ("AKS_ARM_NAMESPACE_ID", "cluster-resource-id"), + ("FUNCTIONS_WORKER_RUNTIME", "python"), + ("WEBSITE_SITE_NAME", "app-service"), + ], +) +def test_vm_detector_skips_metadata_request_on_other_azure_platforms(monkeypatch, environment_variable, value): + monkeypatch.setenv(environment_variable, value) + + def unexpected_request(request, timeout): + pytest.fail("VM metadata must not be queried on another Azure platform") + + monkeypatch.setattr("opentelemetry.resource.detector.azure.vm.urlopen", unexpected_request) + + assert not AzureVMResourceDetector().detect().attributes + + +def test_resource_aggregation_ignores_failing_detector(): + class FailingDetector(ResourceDetector): + def detect(self): + raise RuntimeError("detector failure") + + class WorkingDetector(ResourceDetector): + def detect(self): + return Resource({"service.name": "working-detector"}) + + resource = get_aggregated_resources( + [FailingDetector(), WorkingDetector()], + initial_resource=Resource({"deployment.environment": "test"}), + ) + + assert resource.attributes["deployment.environment"] == "test" + assert resource.attributes["service.name"] == "working-detector" + + +def test_distro_default_resource_runs_registered_azure_detectors(monkeypatch): + monkeypatch.setenv("WEBSITE_SITE_NAME", "inventory-api") + monkeypatch.setenv("WEBSITE_INSTANCE_ID", "inventory-instance") + configurations = {} + + _default_resource(configurations) + + assert configurations[RESOURCE_ARG].attributes["service.name"] == "inventory-api" + assert configurations[RESOURCE_ARG].attributes["service.instance.id"] + assert configurations[RESOURCE_ARG].attributes["cloud.platform"] == "azure_app_service" + assert configurations[RESOURCE_ARG].attributes["cloud.provider"] == "azure" + assert os.environ[OTEL_EXPERIMENTAL_RESOURCE_DETECTORS] == "azure_app_service,azure_vm" diff --git a/tests/upstream/test_dependency_plugins.py b/tests/upstream/test_dependency_plugins.py new file mode 100644 index 00000000..e1e484ff --- /dev/null +++ b/tests/upstream/test_dependency_plugins.py @@ -0,0 +1,129 @@ +import importlib +import importlib.util +from importlib.metadata import entry_points, version +from unittest.mock import patch + +import pytest +from azure.monitor.opentelemetry.exporter import ( + AzureMonitorLogExporter, + AzureMonitorMetricExporter, + AzureMonitorTraceExporter, +) +from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter +from opentelemetry.exporter.otlp.proto.http.metric_exporter import OTLPMetricExporter +from opentelemetry.exporter.otlp.proto.http.trace_exporter import OTLPSpanExporter + +_LOADABLE_INSTRUMENTORS = { + "logging", + "requests", + "urllib", + "urllib3", +} + +_OPTIONAL_RUNTIME_INSTRUMENTORS = { + "django": ("opentelemetry.instrumentation.django", "DjangoInstrumentor"), + "fastapi": ("opentelemetry.instrumentation.fastapi", "FastAPIInstrumentor"), + "flask": ("opentelemetry.instrumentation.flask", "FlaskInstrumentor"), + "httpx": ("opentelemetry.instrumentation.httpx", "HTTPXClientInstrumentor"), + "httpx2": ("opentelemetry.instrumentation.httpx", "HTTPX2ClientInstrumentor"), + "psycopg2": ("opentelemetry.instrumentation.psycopg2", "Psycopg2Instrumentor"), +} + + +@pytest.mark.parametrize( + ("distribution", "module"), + [ + ("azure-core", "azure.core"), + ("azure-core-tracing-opentelemetry", "azure.core.tracing.ext.opentelemetry_span"), + ("azure-monitor-opentelemetry-exporter", "azure.monitor.opentelemetry.exporter"), + ("opentelemetry-api", "opentelemetry.trace"), + ("opentelemetry-exporter-otlp-proto-http", "opentelemetry.exporter.otlp.proto.http"), + ("opentelemetry-instrumentation", "opentelemetry.instrumentation"), + ("opentelemetry-resource-detector-azure", "opentelemetry.resource.detector.azure"), + ("opentelemetry-sdk", "opentelemetry.sdk"), + ("opentelemetry-util-genai", "opentelemetry.util.genai"), + ], +) +def test_direct_dependency_distribution_and_module_are_available(distribution, module): + assert version(distribution) + assert importlib.import_module(module) + + +@pytest.mark.parametrize( + ("group", "required_names"), + [ + ( + "opentelemetry_instrumentor", + _LOADABLE_INSTRUMENTORS | _OPTIONAL_RUNTIME_INSTRUMENTORS.keys(), + ), + ("opentelemetry_logs_exporter", {"azure_monitor_opentelemetry_exporter", "otlp_proto_http"}), + ("opentelemetry_metrics_exporter", {"azure_monitor_opentelemetry_exporter", "otlp_proto_http"}), + ("opentelemetry_propagator", {"baggage", "tracecontext"}), + ("opentelemetry_traces_exporter", {"azure_monitor_opentelemetry_exporter", "otlp_proto_http"}), + ("opentelemetry_traces_sampler", {"always_off", "always_on", "parentbased_traceidratio", "traceidratio"}), + ], +) +def test_required_plugin_entry_points_load(group, required_names): + plugins = {plugin.name: plugin for plugin in entry_points(group=group)} + + assert required_names <= plugins.keys() + names_to_load = _LOADABLE_INSTRUMENTORS if group == "opentelemetry_instrumentor" else required_names + for name in names_to_load: + assert plugins[name].load() + + +@pytest.mark.parametrize( + ("name", "expected_target"), + _OPTIONAL_RUNTIME_INSTRUMENTORS.items(), +) +def test_optional_runtime_instrumentor_entry_points_target_installed_modules(name, expected_target): + plugins = {plugin.name: plugin for plugin in entry_points(group="opentelemetry_instrumentor")} + module, attribute = expected_target + plugin = plugins[name] + + # These instrumentors import optional user frameworks or drivers when loaded. + # Validate their targets without requiring those applications in the distro test environment. + assert plugin.module == module + assert plugin.attr == attribute + assert importlib.util.find_spec(module) is not None + + +@pytest.mark.parametrize( + ("exporter_class", "endpoint"), + [ + (OTLPSpanExporter, "http://localhost:4318/v1/traces"), + (OTLPMetricExporter, "http://localhost:4318/v1/metrics"), + (OTLPLogExporter, "http://localhost:4318/v1/logs"), + ], +) +def test_otlp_http_exporters_construct_and_shutdown(exporter_class, endpoint): + exporter = exporter_class(endpoint=endpoint) + + exporter.shutdown() + + +@pytest.mark.parametrize( + "exporter_class", + [ + AzureMonitorTraceExporter, + AzureMonitorMetricExporter, + AzureMonitorLogExporter, + ], +) +def test_azure_monitor_exporters_construct_and_shutdown(exporter_class): + with ( + patch( + "azure.monitor.opentelemetry.exporter.export._base.get_configuration_manager", + return_value=None, + ), + patch.object(exporter_class, "_should_collect_stats", return_value=False), + patch.object(exporter_class, "_should_collect_customer_sdkstats", return_value=False), + ): + exporter = exporter_class( + connection_string=( + "InstrumentationKey=00000000-0000-0000-0000-000000000000;" "IngestionEndpoint=https://example.test/" + ), + disable_offline_storage=True, + ) + + exporter.shutdown() diff --git a/tests/upstream/test_opentelemetry_sdk_contracts.py b/tests/upstream/test_opentelemetry_sdk_contracts.py new file mode 100644 index 00000000..c5978c5f --- /dev/null +++ b/tests/upstream/test_opentelemetry_sdk_contracts.py @@ -0,0 +1,267 @@ +import logging +from unittest.mock import patch + +import pytest +from opentelemetry import baggage, trace +from opentelemetry.baggage.propagation import W3CBaggagePropagator +from opentelemetry.instrumentation.logging.handler import LoggingHandler +from opentelemetry.propagators.composite import CompositePropagator +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk._logs.export import InMemoryLogRecordExporter, SimpleLogRecordProcessor +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.metrics.export import InMemoryMetricReader +from opentelemetry.sdk.resources import Resource +from opentelemetry.sdk.trace import TracerProvider +from opentelemetry.sdk.trace.export import SimpleSpanProcessor +from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter +from opentelemetry.sdk.trace.sampling import ALWAYS_OFF +from opentelemetry.trace import StatusCode +from opentelemetry.trace.propagation.tracecontext import TraceContextTextMapPropagator + + +def test_resource_create_adds_sdk_attributes_and_preserves_user_attributes(): + resource = Resource.create( + { + "deployment.environment": "test", + "service.name": "compatibility-suite", + } + ) + + assert resource.attributes["deployment.environment"] == "test" + assert resource.attributes["service.name"] == "compatibility-suite" + assert resource.attributes["service.instance.id"] + assert resource.attributes["telemetry.sdk.language"] == "python" + assert resource.attributes["telemetry.sdk.name"] == "opentelemetry" + assert resource.attributes["telemetry.sdk.version"] + + +def test_resource_create_honors_standard_environment_attributes(monkeypatch): + monkeypatch.setenv("OTEL_SERVICE_NAME", "environment-service") + monkeypatch.setenv("OTEL_RESOURCE_ATTRIBUTES", "deployment.environment=production,service.version=2.1") + + resource = Resource.create({"service.name": "explicit-service"}) + + assert resource.attributes["service.name"] == "explicit-service" + assert resource.attributes["deployment.environment"] == "production" + assert resource.attributes["service.version"] == "2.1" + + +def test_tracing_exports_parent_child_spans_with_events_and_attributes(): + exporter = InMemorySpanExporter() + provider = TracerProvider( + resource=Resource.create({"service.name": "compatibility-suite"}), + shutdown_on_exit=False, + ) + provider.add_span_processor(SimpleSpanProcessor(exporter)) + tracer = provider.get_tracer("compatibility-tests", "1.0") + + with tracer.start_as_current_span("parent") as parent: + parent.set_attribute("request.id", "request-123") + parent.add_event("request.received", {"payload.size": 42}) + with tracer.start_as_current_span("child") as child: + child.set_attribute("component", "dependency") + + spans = {span.name: span for span in exporter.get_finished_spans()} + parent_span = spans["parent"] + child_span = spans["child"] + + assert child_span.context.trace_id == parent_span.context.trace_id + assert child_span.parent.span_id == parent_span.context.span_id + assert parent_span.attributes is not None + assert parent_span.attributes["request.id"] == "request-123" + assert parent_span.events[0].name == "request.received" + assert parent_span.events[0].attributes is not None + assert parent_span.events[0].attributes["payload.size"] == 42 + assert parent_span.resource.attributes["service.name"] == "compatibility-suite" + assert parent_span.instrumentation_scope.name == "compatibility-tests" + provider.shutdown() + + +def test_always_off_sampler_creates_non_recording_span_and_exports_nothing(): + exporter = InMemorySpanExporter() + provider = TracerProvider(sampler=ALWAYS_OFF, shutdown_on_exit=False) + provider.add_span_processor(SimpleSpanProcessor(exporter)) + + with provider.get_tracer("compatibility-tests").start_as_current_span("dropped") as span: + assert not span.is_recording() + assert not span.get_span_context().trace_flags.sampled + + assert not exporter.get_finished_spans() + provider.shutdown() + + +def test_record_exception_sets_error_status_and_exception_event(): + exporter = InMemorySpanExporter() + provider = TracerProvider(shutdown_on_exit=False) + provider.add_span_processor(SimpleSpanProcessor(exporter)) + tracer = provider.get_tracer("compatibility-tests") + + with pytest.raises(ValueError, match="bad input"): + with tracer.start_as_current_span("failing-operation"): + raise ValueError("bad input") + + span = exporter.get_finished_spans()[0] + assert span.status.status_code is StatusCode.ERROR + assert span.events[0].name == "exception" + assert span.events[0].attributes is not None + assert span.events[0].attributes["exception.type"] == "ValueError" + assert span.events[0].attributes["exception.message"] == "bad input" + provider.shutdown() + + +def test_tracer_provider_force_flush_reaches_span_processor(): + exporter = InMemorySpanExporter() + provider = TracerProvider(shutdown_on_exit=False) + processor = SimpleSpanProcessor(exporter) + provider.add_span_processor(processor) + + with provider.get_tracer("compatibility-tests").start_as_current_span("flush-me"): + pass + + with patch.object(processor, "force_flush", wraps=processor.force_flush) as force_flush: + assert provider.force_flush(timeout_millis=1_000) + + force_flush.assert_called_once() + assert [span.name for span in exporter.get_finished_spans()] == ["flush-me"] + provider.shutdown() + + +def test_metrics_collect_counter_and_histogram_measurements(): + reader = InMemoryMetricReader() + provider = MeterProvider( + metric_readers=[reader], + resource=Resource.create({"service.name": "compatibility-suite"}), + shutdown_on_exit=False, + ) + meter = provider.get_meter("compatibility-tests", "1.0") + + meter.create_counter("requests").add(3, {"route": "/orders"}) + meter.create_histogram("request.duration", unit="ms").record(12.5, {"route": "/orders"}) + + metrics_data = reader.get_metrics_data() + resource_metrics = metrics_data.resource_metrics[0] + metrics = { + metric.name: metric for scope_metrics in resource_metrics.scope_metrics for metric in scope_metrics.metrics + } + + counter_point = metrics["requests"].data.data_points[0] + histogram_point = metrics["request.duration"].data.data_points[0] + assert counter_point.value == 3 + assert counter_point.attributes == {"route": "/orders"} + assert histogram_point.count == 1 + assert histogram_point.sum == 12.5 + assert histogram_point.min == 12.5 + assert histogram_point.max == 12.5 + assert histogram_point.attributes == {"route": "/orders"} + assert resource_metrics.resource.attributes["service.name"] == "compatibility-suite" + provider.shutdown() + + +def test_metrics_collect_up_down_counter_measurements(): + reader = InMemoryMetricReader() + provider = MeterProvider(metric_readers=[reader], shutdown_on_exit=False) + counter = provider.get_meter("compatibility-tests").create_up_down_counter("active.requests") + + counter.add(3, {"queue": "primary"}) + counter.add(-1, {"queue": "primary"}) + + metric = reader.get_metrics_data().resource_metrics[0].scope_metrics[0].metrics[0] + point = metric.data.data_points[0] + assert point.value == 2 + assert point.attributes == {"queue": "primary"} + provider.shutdown() + + +def test_logging_handler_exports_formatted_body_attributes_and_scope(): + exporter = InMemoryLogRecordExporter() + provider = LoggerProvider( + resource=Resource.create({"service.name": "compatibility-suite"}), + shutdown_on_exit=False, + ) + provider.add_log_record_processor(SimpleLogRecordProcessor(exporter)) + handler = LoggingHandler(logger_provider=provider) + logger = logging.getLogger("compatibility-tests") + original_handlers = logger.handlers[:] + original_level = logger.level + original_propagate = logger.propagate + + try: + logger.handlers = [handler] + logger.setLevel(logging.INFO) + logger.propagate = False + logger.info("processed %s items", 5, extra={"request_id": "request-123"}) + + exported = exporter.get_finished_logs() + assert len(exported) == 1 + record = exported[0] + assert record.log_record.body == "processed 5 items" + assert record.log_record.attributes is not None + assert record.log_record.attributes["request_id"] == "request-123" + assert record.instrumentation_scope.name == "compatibility-tests" + assert record.resource.attributes["service.name"] == "compatibility-suite" + finally: + logger.handlers = original_handlers + logger.setLevel(original_level) + logger.propagate = original_propagate + provider.shutdown() + + +def test_log_record_carries_active_trace_context(): + span_exporter = InMemorySpanExporter() + tracer_provider = TracerProvider(shutdown_on_exit=False) + tracer_provider.add_span_processor(SimpleSpanProcessor(span_exporter)) + log_exporter = InMemoryLogRecordExporter() + logger_provider = LoggerProvider(shutdown_on_exit=False) + logger_provider.add_log_record_processor(SimpleLogRecordProcessor(log_exporter)) + logger = logging.getLogger("trace-correlated-logger") + handler = LoggingHandler(logger_provider=logger_provider) + original_handlers = logger.handlers[:] + original_propagate = logger.propagate + logger.handlers = [handler] + logger.propagate = False + + try: + with tracer_provider.get_tracer("compatibility-tests").start_as_current_span("operation") as span: + logger.warning("correlated") + + record = log_exporter.get_finished_logs()[0].log_record + assert record.trace_id == span.get_span_context().trace_id + assert record.span_id == span.get_span_context().span_id + assert record.trace_flags == span.get_span_context().trace_flags + finally: + logger.handlers = original_handlers + logger.propagate = original_propagate + logger_provider.shutdown() + tracer_provider.shutdown() + + +def test_w3c_trace_context_and_baggage_round_trip(): + provider = TracerProvider(shutdown_on_exit=False) + tracer = provider.get_tracer("compatibility-tests") + propagator = CompositePropagator( + [ + TraceContextTextMapPropagator(), + W3CBaggagePropagator(), + ] + ) + + with tracer.start_as_current_span("producer") as span: + context = baggage.set_baggage("tenant.id", "tenant-123") + carrier = {} + propagator.inject(carrier, context=context) + + extracted_context = propagator.extract(carrier) + extracted_span_context = trace.get_current_span(extracted_context).get_span_context() + + assert carrier["traceparent"].startswith("00-") + assert carrier["baggage"] == "tenant.id=tenant-123" + assert extracted_span_context.trace_id == span.get_span_context().trace_id + assert extracted_span_context.span_id == span.get_span_context().span_id + assert baggage.get_baggage("tenant.id", context=extracted_context) == "tenant-123" + provider.shutdown() + + +def test_invalid_traceparent_is_ignored(): + extracted_context = TraceContextTextMapPropagator().extract({"traceparent": "invalid"}) + + assert not trace.get_current_span(extracted_context).get_span_context().is_valid diff --git a/tests/upstream/test_service_instance_id_regressions.py b/tests/upstream/test_service_instance_id_regressions.py new file mode 100644 index 00000000..c2e21922 --- /dev/null +++ b/tests/upstream/test_service_instance_id_regressions.py @@ -0,0 +1,168 @@ +import os + +import pytest +import opentelemetry.sdk.resources as resources_module +from opentelemetry.resource.detector.azure.vm import AzureVMResourceDetector +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.resources import ( + Resource, + ServiceInstanceIdResourceDetector, +) +from opentelemetry.sdk.trace import TracerProvider + + +@pytest.fixture(autouse=True) +def reset_generated_service_instance_id(monkeypatch): + monkeypatch.setattr(resources_module, "_service_instance_id", None) + monkeypatch.setattr(resources_module, "_service_instance_id_pid", None) + monkeypatch.delenv("OTEL_RESOURCE_ATTRIBUTES", raising=False) + monkeypatch.delenv("OTEL_SERVICE_NAME", raising=False) + + +def test_generated_service_instance_id_is_stable_within_one_process(): + detector = ServiceInstanceIdResourceDetector() + + first = detector.detect().attributes["service.instance.id"] + second = detector.detect().attributes["service.instance.id"] + + assert first == second + + +def test_generated_service_instance_id_changes_when_process_identity_changes(monkeypatch): + detector = ServiceInstanceIdResourceDetector() + parent_id = detector.detect().attributes["service.instance.id"] + child_pid = os.getpid() + 1 + + monkeypatch.setattr(resources_module.os, "getpid", lambda: child_pid) + child_id = detector.detect().attributes["service.instance.id"] + + assert child_id != parent_id + + +def test_resource_create_does_not_generate_unbounded_ids_in_one_process(): + instance_ids = {Resource.create().attributes["service.instance.id"] for _ in range(100)} + + assert len(instance_ids) == 1 + + +def test_explicit_resource_service_instance_id_overrides_generated_id(): + resource = Resource.create({"service.instance.id": "explicit-instance"}) + + assert resource.attributes["service.instance.id"] == "explicit-instance" + + +def test_environment_service_instance_id_overrides_generated_id(monkeypatch): + monkeypatch.setenv("OTEL_RESOURCE_ATTRIBUTES", "service.instance.id=pod-123") + + resource = Resource.create() + + assert resource.attributes["service.instance.id"] == "pod-123" + + +def test_azure_vm_detector_does_not_override_environment_service_instance_id(monkeypatch): + monkeypatch.setenv("OTEL_EXPERIMENTAL_RESOURCE_DETECTORS", "azure_vm") + monkeypatch.setenv("OTEL_RESOURCE_ATTRIBUTES", "service.instance.id=pod-123") + monkeypatch.setattr( + AzureVMResourceDetector, + "detect", + lambda self: Resource( + { + "cloud.platform": "azure_vm", + "host.id": "vm-id", + "service.instance.id": "vm-id", + } + ), + ) + + resource = Resource.create() + + assert resource.attributes["cloud.platform"] == "azure_vm" + assert resource.attributes["host.id"] == "vm-id" + assert resource.attributes["service.instance.id"] == "pod-123" + + +def test_aks_identity_is_not_replaced_by_vm_detector(monkeypatch): + monkeypatch.setenv("AKS_ARM_NAMESPACE_ID", "aks-resource-id") + monkeypatch.setenv("OTEL_EXPERIMENTAL_RESOURCE_DETECTORS", "azure_vm") + monkeypatch.setenv( + "OTEL_RESOURCE_ATTRIBUTES", + "k8s.pod.name=orders-7f8c9,service.instance.id=orders-7f8c9", + ) + + resource = Resource.create() + + assert resource.attributes["k8s.pod.name"] == "orders-7f8c9" + assert resource.attributes["service.instance.id"] == "orders-7f8c9" + assert "cloud.platform" not in resource.attributes + assert "host.id" not in resource.attributes + + +@pytest.mark.xfail( + strict=True, + reason="opentelemetry-sdk 1.44 runs the generated service_instance detector after Azure detectors", +) +def test_azure_app_service_instance_id_overrides_generated_id(monkeypatch): + monkeypatch.setenv("OTEL_EXPERIMENTAL_RESOURCE_DETECTORS", "azure_app_service") + monkeypatch.setenv("WEBSITE_SITE_NAME", "orders-api") + monkeypatch.setenv("WEBSITE_INSTANCE_ID", "app-service-worker") + + resource = Resource.create() + + assert resource.attributes["service.instance.id"] == "app-service-worker" + + +@pytest.mark.xfail( + strict=True, + reason="opentelemetry-sdk 1.44 merges process-dependent detector refreshes over explicit resources", +) +def test_initial_resource_service_instance_id_has_highest_aggregation_priority(): + from opentelemetry.sdk.resources import get_aggregated_resources + + resource = get_aggregated_resources( + [ServiceInstanceIdResourceDetector()], + initial_resource=Resource({"service.instance.id": "explicit-instance"}), + ) + + assert resource.attributes["service.instance.id"] == "explicit-instance" + + +def _provider_resource(provider): + if isinstance(provider, MeterProvider): + return provider._sdk_config.resource + return provider.resource + + +@pytest.mark.parametrize("provider_class", [TracerProvider, MeterProvider, LoggerProvider]) +@pytest.mark.xfail( + strict=True, + reason="opentelemetry-sdk 1.44 fork refresh overwrites explicit service.instance.id", +) +def test_provider_fork_refresh_preserves_explicit_instance_id(monkeypatch, provider_class): + provider = provider_class( + resource=Resource({"service.instance.id": "explicit-instance"}), + shutdown_on_exit=False, + ) + child_pid = os.getpid() + 1 + monkeypatch.setattr(resources_module.os, "getpid", lambda: child_pid) + + try: + provider._handle_fork() + assert _provider_resource(provider).attributes["service.instance.id"] == "explicit-instance" + finally: + provider.shutdown() + + +@pytest.mark.parametrize("provider_class", [TracerProvider, MeterProvider, LoggerProvider]) +def test_provider_fork_refresh_regenerates_generated_instance_id(monkeypatch, provider_class): + provider = provider_class(shutdown_on_exit=False) + parent_id = _provider_resource(provider).attributes["service.instance.id"] + child_pid = os.getpid() + 1 + monkeypatch.setattr(resources_module.os, "getpid", lambda: child_pid) + + try: + provider._handle_fork() + child_id = _provider_resource(provider).attributes["service.instance.id"] + assert child_id != parent_id + finally: + provider.shutdown()