diff --git a/.changesets/identify-each-process-that-reports.md b/.changesets/identify-each-process-that-reports.md new file mode 100644 index 00000000..28e1128b --- /dev/null +++ b/.changesets/identify-each-process-that-reports.md @@ -0,0 +1,6 @@ +--- +bump: patch +type: fix +--- + +Metrics are reported correctly when several processes of your application send the same metric, including processes forked after AppSignal starts. Each process reports under a `service.instance.id` and a `process.pid` of its own, which is what AppSignal uses to tell them apart. diff --git a/pyproject.toml b/pyproject.toml index 686c41b9..574f9e36 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -26,8 +26,8 @@ classifiers = [ "Development Status :: 5 - Production/Stable" ] dependencies = [ - "opentelemetry-api>=1.42.0", - "opentelemetry-sdk>=1.42.0", + "opentelemetry-api>=1.44.0", + "opentelemetry-sdk>=1.44.0", "opentelemetry-exporter-otlp-proto-http", "opentelemetry-instrumentation-logging>=0.63b0", "requests", diff --git a/src/appsignal/opentelemetry.py b/src/appsignal/opentelemetry.py index b0cc8528..6d9a954d 100644 --- a/src/appsignal/opentelemetry.py +++ b/src/appsignal/opentelemetry.py @@ -501,6 +501,12 @@ def _silence_internal_loggers() -> None: def _resource(config: Config) -> Resource: + # Ask for the resource detector that reports the process, so that the + # process id is reported alongside the instance id. Naming the detector in + # this variable is what makes each provider run it again when it refreshes + # the resource after a fork. + os.environ["OTEL_EXPERIMENTAL_RESOURCE_DETECTORS"] = "process" + attributes = { key: value for key, value in { @@ -513,7 +519,6 @@ def _resource(config: Config) -> Resource: "appsignal.config.language_integration": "python", "service.name": config.options.get("service_name") or "app", "host.name": config.options.get("hostname") or "unknown", - "appsignal.service.process_id": os.getpid(), "appsignal.config.filter_attributes": config.options.get( "filter_attributes" ), @@ -557,7 +562,12 @@ def _resource(config: Config) -> Resource: if value is not None } - return Resource(attributes=cast(Mapping[str, str | list[str]], attributes)) + # `Resource.create` runs the resource detectors, one of which gives the + # process a `service.instance.id`. Each provider re-runs them after a fork, + # so a forked worker reports under an identity of its own. The attributes + # passed here win over what a detector or `OTEL_RESOURCE_ATTRIBUTES` + # provides. + return Resource.create(cast(Mapping[str, str | list[str]], attributes)) def _opentelemetry_endpoint(config: Config) -> str: diff --git a/tests/test_config.py b/tests/test_config.py index 3e660b9f..89ea9e4a 100644 --- a/tests/test_config.py +++ b/tests/test_config.py @@ -1,10 +1,15 @@ from __future__ import annotations import os +import select import socket -from typing import cast +from typing import Any, cast +from uuid import UUID import pytest +from opentelemetry.sdk._logs import LoggerProvider +from opentelemetry.sdk.metrics import MeterProvider +from opentelemetry.sdk.trace import TracerProvider from appsignal.__about__ import __version__ from appsignal.config import Config, Options @@ -481,8 +486,76 @@ def test_set_private_environ(): assert os.environ["_APP_REVISION"] == "abc123" +def test_opentelemetry_resource_identifies_the_process(): + from appsignal.opentelemetry import _resource + + resource = _resource(Config()) + + instance_id = str(resource.attributes["service.instance.id"]) + + assert UUID(instance_id).version == 4 + assert resource.attributes["process.pid"] == os.getpid() + + +# Every provider re-runs the resource detectors after a fork, so a forked +# worker reports under an identity of its own rather than the one its parent +# started with. +# +# Forking is what this asserts, so the warning Python raises about forking a +# process that has threads is expected here. +@pytest.mark.filterwarnings("ignore:This process .* is multi-threaded") +def test_opentelemetry_resource_identifies_a_forked_process(): + from appsignal.opentelemetry import _resource + + resource = _resource(Config()) + providers = [ + TracerProvider(resource=resource), + MeterProvider(resource=resource), + LoggerProvider(resource=resource), + ] + + # The meter provider keeps its resource on its config rather than on + # itself, unlike the other two. + def identity(provider: Any) -> str: + sdk_config = getattr(provider, "_sdk_config", None) + attributes = ( + sdk_config.resource if sdk_config else provider.resource + ).attributes + return f"{attributes['service.instance.id']} {attributes['process.pid']}" + + def identities() -> list[str]: + return [identity(provider) for provider in providers] + + parent = identities() + read, write = os.pipe() + + if os.fork() == 0: # pragma: no cover + os.close(read) + os.write(write, "\n".join(identities()).encode()) + os._exit(0) + + os.close(write) + # Read with a timeout. A child that forked while another thread held a + # lock can deadlock, and waiting for it forever would hang the test run + # rather than fail it. + readable, _, _ = select.select([read], [], [], 30) + assert readable, "the forked process wrote nothing" + with os.fdopen(read) as pipe: + child = pipe.read().split("\n") + os.wait() + + # Every provider reports the same identity, and the child's is its own. + assert len(set(parent)) == 1 + assert len(set(child)) == 1 + assert child[0] != parent[0] + + child_instance_id, child_pid = child[0].split(" ") + + assert UUID(child_instance_id).version == 4 + assert child_pid != str(os.getpid()) + + def test_opentelemetry_resource(): - import os from appsignal.opentelemetry import _resource @@ -526,7 +599,6 @@ def test_opentelemetry_resource(): assert resource.attributes["appsignal.config.language_integration"] == "python" assert resource.attributes["service.name"] == "test-service" assert resource.attributes["host.name"] == "test-host" - assert resource.attributes["appsignal.service.process_id"] == os.getpid() # Test filter attributes assert resource.attributes["appsignal.config.filter_attributes"] == (