From 623d67b2289fca91d92ece970d39de39433fd0c7 Mon Sep 17 00:00:00 2001 From: Noemi Lapresta Date: Wed, 30 Sep 2026 13:03:55 +0200 Subject: [PATCH] Refresh the process resource after fork The OpenTelemetry providers are built once, and their resource carries the `process.pid` of the process that built them. A process forked after that, such as a Puma worker under `preload_app!`, exports its data with its parent's resource, so the collector computes the deltas of every worker's up-down counters as one series. The resource now also carries a `service.instance.id`, a UUIDv4 generated once per process, as the Python SDK does. A `Process._fork` hook replaces it and `process.pid` on each provider's resource in the child. The providers have no resource setter, so the hook writes `@resource`, reads the resource back, and logs a warning when the provider did not keep it. --- ...ounters-correctly-from-forked-processes.md | 6 + lib/appsignal/opentelemetry.rb | 6 +- .../opentelemetry/process_resource.rb | 137 ++++++++++++ .../collector_mode_fork_resource_spec.rb | 90 ++++++++ .../runners/collector_mode_fork_resource.rb | 31 +++ .../opentelemetry/process_resource_spec.rb | 208 ++++++++++++++++++ 6 files changed, 477 insertions(+), 1 deletion(-) create mode 100644 .changesets/report-counters-correctly-from-forked-processes.md create mode 100644 lib/appsignal/opentelemetry/process_resource.rb create mode 100644 spec/integration/collector_mode_fork_resource_spec.rb create mode 100644 spec/integration/runners/collector_mode_fork_resource.rb create mode 100644 spec/lib/appsignal/opentelemetry/process_resource_spec.rb diff --git a/.changesets/report-counters-correctly-from-forked-processes.md b/.changesets/report-counters-correctly-from-forked-processes.md new file mode 100644 index 000000000..f775a3a55 --- /dev/null +++ b/.changesets/report-counters-correctly-from-forked-processes.md @@ -0,0 +1,6 @@ +--- +bump: patch +type: fix +--- + +In collector mode, report counter metrics correctly from processes forked after AppSignal starts, such as Puma workers when using `preload_app!`. Before this change, applications running several forked processes would lose data points or report incorrect values for them. diff --git a/lib/appsignal/opentelemetry.rb b/lib/appsignal/opentelemetry.rb index 4f0285aee..04da2cb48 100644 --- a/lib/appsignal/opentelemetry.rb +++ b/lib/appsignal/opentelemetry.rb @@ -8,6 +8,7 @@ require "appsignal/opentelemetry/http_response" require "appsignal/opentelemetry/http_server_request" require "appsignal/opentelemetry/messaging" +require "appsignal/opentelemetry/process_resource" require "appsignal/opentelemetry/proxied_exporter" require "appsignal/opentelemetry/rendering" require "appsignal/opentelemetry/sql_db_system" @@ -59,7 +60,9 @@ def configure(config) # `LoggerProvider` take a `resource:` kwarg that replaces (not # merges), so we do the merge ourselves and use the same merged # resource for the tracer provider to keep all three in sync. - resource = ::OpenTelemetry::SDK::Resources::Resource.default.merge(build_resource(config)) + resource = ProcessResource.resource + .merge(::OpenTelemetry::SDK::Resources::Resource.default) + .merge(build_resource(config)) span_exporter = build_exporter( ::OpenTelemetry::Exporter::OTLP::Exporter, @@ -103,6 +106,7 @@ def configure(config) ) @started = true + ProcessResource.attach_fork_hook rescue LoadError => e @started = false Appsignal::Utils::StdoutAndLoggerMessage.error( diff --git a/lib/appsignal/opentelemetry/process_resource.rb b/lib/appsignal/opentelemetry/process_resource.rb new file mode 100644 index 000000000..149be4358 --- /dev/null +++ b/lib/appsignal/opentelemetry/process_resource.rb @@ -0,0 +1,137 @@ +# frozen_string_literal: true + +require "securerandom" + +module Appsignal + module OpenTelemetry + # @!visibility private + # + # The resource attributes that identify the current process: + # `service.instance.id` and `process.pid`. The OpenTelemetry providers are + # built once, so a process forked after that inherits its parent's values + # until {.attach_fork_hook} refreshes them in the child. + module ProcessResource + SERVICE_INSTANCE_ID = "service.instance.id" + PROCESS_PID = "process.pid" + REFRESH_FAILED = "Could not refresh the OpenTelemetry resource after fork" + + class << self + # A random UUIDv4, generated again when the process ID changes. + def service_instance_id + pid = Process.pid + unless @service_instance_id_pid == pid + @service_instance_id = SecureRandom.uuid + @service_instance_id_pid = pid + end + @service_instance_id + end + + def resource + ::OpenTelemetry::SDK::Resources::Resource.create( + SERVICE_INSTANCE_ID => service_instance_id, + PROCESS_PID => Process.pid + ) + end + + def attach_fork_hook + return if @fork_hook_attached + + Process.singleton_class.prepend(ForkHook) + @fork_hook_attached = true + end + + def before_fork + return unless Appsignal::OpenTelemetry.started? + + providers.transform_values do |provider| + current = current_resource(provider) + current && attribute(current, SERVICE_INSTANCE_ID) + end + rescue => e + log_failure("#{e.class}: #{e.message}") + nil + end + + def after_fork(instance_ids_before_fork) + return unless instance_ids_before_fork + + update = resource + providers.each do |name, provider| + refresh(name, provider, update, instance_ids_before_fork[name]) + end + rescue => e + log_failure("#{e.class}: #{e.message}") + end + + private + + def providers + { + "tracer" => ::OpenTelemetry.tracer_provider, + "meter" => ::OpenTelemetry.meter_provider, + "logger" => ::OpenTelemetry.logger_provider + } + end + + def refresh(name, provider, update, instance_id_before_fork) + current = current_resource(provider) + unless current + log_failure("the #{name} provider's resource could not be read") + return + end + return unless attribute(current, SERVICE_INSTANCE_ID) == instance_id_before_fork + + refreshed = current.merge(update) + write_resource(provider, refreshed) + return if current_resource(provider).equal?(refreshed) + + log_failure("the #{name} provider did not keep the new resource") + end + + # The SDK providers take their resource only in their constructors, and + # read it from `@resource` when they export. + def current_resource(provider) + if provider.respond_to?(:resource) + resource = provider.resource + unless provider.respond_to?(:resource=) + # Without a setter, only `@resource` can be written, and it isn't there. + return unless provider.instance_variable_defined?(:@resource) + # The reader returns something else, so writing `@resource` would change nothing. + return unless provider.instance_variable_get(:@resource).equal?(resource) + end + elsif provider.instance_variable_defined?(:@resource) + resource = provider.instance_variable_get(:@resource) + end + # If it's not a resource, it isn't what the provider exports. + resource if resource.is_a?(::OpenTelemetry::SDK::Resources::Resource) + end + + def write_resource(provider, resource) + if provider.respond_to?(:resource=) + provider.resource = resource + else + provider.instance_variable_set(:@resource, resource) + end + end + + def attribute(resource, key) + resource.attribute_enumerator.each { |k, v| return v if k == key } + nil + end + + def log_failure(reason) + Appsignal.internal_logger.warn("#{REFRESH_FAILED}: #{reason}") + end + end + + module ForkHook + def _fork + instance_ids = ProcessResource.before_fork + pid = super + ProcessResource.after_fork(instance_ids) if pid.zero? + pid + end + end + end + end +end diff --git a/spec/integration/collector_mode_fork_resource_spec.rb b/spec/integration/collector_mode_fork_resource_spec.rb new file mode 100644 index 000000000..6da3d614d --- /dev/null +++ b/spec/integration/collector_mode_fork_resource_spec.rb @@ -0,0 +1,90 @@ +# Skipped on JRuby because `Process.fork` raises NotImplementedError there. +if DependencyHelper.opentelemetry_present? && !DependencyHelper.running_jruby? + require "opentelemetry/exporter/otlp" + require "opentelemetry/proto/collector/trace/v1/trace_service_pb" + require "opentelemetry/proto/collector/metrics/v1/metrics_service_pb" + require "opentelemetry/proto/collector/logs/v1/logs_service_pb" + + describe "Collector mode resource under fork" do + before { OTLPCollectorServer.clear } + + it "exports the child's data with the child's process ID and service instance ID" do + runner = Runner.new( + "collector_mode_fork_resource", + :env => OTLPCollectorServer.env.merge("APPSIGNAL_LOG" => "stdout") + ) + runner.run + + parent_pid = runner.output[/PARENT_PID=(\d+)/, 1].to_i + child_pid = runner.output[/CHILD_PID=(\d+)/, 1].to_i + + resources = { + "span" => span_resources, + "metric" => metric_resources, + "log" => log_resources + } + + resources.each do |signal, by_process| + parent = by_process.fetch("parent") { raise "No parent #{signal} exported" } + child = by_process.fetch("child") { raise "No child #{signal} exported" } + + expect(parent["process.pid"]).to eq(parent_pid) + expect(child["process.pid"]).to eq(child_pid) + expect(parent["service.instance.id"]).to match(/\A\h{8}-\h{4}-4\h{3}-\h{4}-\h{12}\z/) + expect(child["service.instance.id"]).to match(/\A\h{8}-\h{4}-4\h{3}-\h{4}-\h{12}\z/) + expect(child["service.instance.id"]).not_to eq(parent["service.instance.id"]) + end + child_instance_ids = resources.values.map { |by_process| by_process["child"] } + .map { |attributes| attributes["service.instance.id"] } + expect(child_instance_ids.uniq.size).to eq(1) + + expect(runner.output).not_to include(Appsignal::OpenTelemetry::ProcessResource::REFRESH_FAILED) + end + + def drain(path, message_class) + queue = OTLPCollectorServer.received[path] + Array.new(queue.size) { message_class.decode(queue.pop[:body]) } + end + + def resource_attributes(resource) + resource.attributes.to_h do |kv| + value = kv.value + [kv.key, value.value == :int_value ? value.int_value : value.string_value] + end + end + + def span_resources + drain("/v1/traces", Opentelemetry::Proto::Collector::Trace::V1::ExportTraceServiceRequest) + .flat_map(&:resource_spans).each_with_object({}) do |rs, by_process| + rs.scope_spans.flat_map(&:spans).each do |span| + process = span.name[/\A(parent|child)#run\z/, 1] + by_process[process] = resource_attributes(rs.resource) if process + end + end + end + + def metric_resources + drain("/v1/metrics", Opentelemetry::Proto::Collector::Metrics::V1::ExportMetricsServiceRequest) + .flat_map(&:resource_metrics).each_with_object({}) do |rm, by_process| + rm.scope_metrics.flat_map(&:metrics).each do |metric| + next unless metric.name == "fork_resource_counter" + + metric.sum.data_points.each do |dp| + process = dp.attributes.find { |kv| kv.key == "process" }&.value&.string_value + by_process[process] = resource_attributes(rm.resource) if process + end + end + end + end + + def log_resources + drain("/v1/logs", Opentelemetry::Proto::Collector::Logs::V1::ExportLogsServiceRequest) + .flat_map(&:resource_logs).each_with_object({}) do |rl, by_process| + rl.scope_logs.flat_map(&:log_records).each do |record| + process = record.body.string_value[/\A(parent|child) log\z/, 1] + by_process[process] = resource_attributes(rl.resource) if process + end + end + end + end +end diff --git a/spec/integration/runners/collector_mode_fork_resource.rb b/spec/integration/runners/collector_mode_fork_resource.rb new file mode 100644 index 000000000..30f051c08 --- /dev/null +++ b/spec/integration/runners/collector_mode_fork_resource.rb @@ -0,0 +1,31 @@ +PROJECT_ROOT = "../../../".freeze +$LOAD_PATH.unshift(File.expand_path("ext", PROJECT_ROOT)) +$LOAD_PATH.unshift(File.expand_path("lib", PROJECT_ROOT)) + +require "appsignal" + +Appsignal.start + +def emit(process) + Appsignal.monitor(:action => "#{process}#run") do + Appsignal.increment_counter("fork_resource_counter", 1, :process => process) + Appsignal::Logger.new("fork-resource").info("#{process} log") + end +end + +emit("parent") +puts "PARENT_PID=#{Process.pid}" + +child_pid = Process.fork do + puts "CHILD_PID=#{Process.pid}" + emit("child") + Appsignal.stop("integration test") +rescue => e + warn "child failed: #{e.class}: #{e.message}" + warn e.backtrace + exit!(1) +end + +_, status = Process.waitpid2(child_pid) +Appsignal.stop("integration test") +exit(status.exitstatus || 1) diff --git a/spec/lib/appsignal/opentelemetry/process_resource_spec.rb b/spec/lib/appsignal/opentelemetry/process_resource_spec.rb new file mode 100644 index 000000000..d27f984b1 --- /dev/null +++ b/spec/lib/appsignal/opentelemetry/process_resource_spec.rb @@ -0,0 +1,208 @@ +# frozen_string_literal: true + +if DependencyHelper.opentelemetry_present? + require "opentelemetry/sdk" + require "opentelemetry-metrics-sdk" + require "opentelemetry-logs-sdk" + + describe Appsignal::OpenTelemetry::ProcessResource do + let(:config) do + build_config( + :options => { + :name => "process-resource-spec", + :push_api_key => "abc", + :collector_endpoint => "http://127.0.0.1:9090" + } + ) + end + let(:uuid_v4) { /\A\h{8}-\h{4}-4\h{3}-\h{4}-\h{12}\z/ } + + around do |example| + providers = [ + ::OpenTelemetry.tracer_provider, + ::OpenTelemetry.meter_provider, + ::OpenTelemetry.logger_provider + ] + example.run + ensure + ::OpenTelemetry.tracer_provider, ::OpenTelemetry.meter_provider, + ::OpenTelemetry.logger_provider = providers + end + + before { Appsignal::OpenTelemetry.reset! } + after { Appsignal::OpenTelemetry.reset! } + + def attributes(resource) + resource.attribute_enumerator.to_h + end + + def provider_attributes + { + "tracer" => attributes(::OpenTelemetry.tracer_provider.resource), + "meter" => attributes(::OpenTelemetry.meter_provider.resource), + "logger" => attributes(::OpenTelemetry.logger_provider.instance_variable_get(:@resource)) + } + end + + def fork_in_place + instance_ids = described_class.before_fork + allow(Process).to receive(:pid).and_return(Process.pid + 1) + described_class.after_fork(instance_ids) + end + + describe ".service_instance_id" do + it "is a UUIDv4 that stays the same within a process" do + id = described_class.service_instance_id + + expect(id).to match(uuid_v4) + expect(described_class.service_instance_id).to eq(id) + end + + it "changes when the process ID changes" do + id = described_class.service_instance_id + allow(Process).to receive(:pid).and_return(Process.pid + 1) + + expect(described_class.service_instance_id).to match(uuid_v4) + expect(described_class.service_instance_id).not_to eq(id) + end + end + + it "puts the service instance ID on every provider's resource" do + Appsignal::OpenTelemetry.configure(config) + + provider_attributes.each_value do |attrs| + expect(attrs["service.instance.id"]).to eq(described_class.service_instance_id) + end + end + + describe "after a fork" do + it "gives every provider the new process ID and service instance ID" do + Appsignal::OpenTelemetry.configure(config) + parent_id = described_class.service_instance_id + + logs = capture_logs { fork_in_place } + + expect(logs).to be_empty + provider_attributes.each_value do |attrs| + expect(attrs["process.pid"]).to eq(Process.pid) + expect(attrs["service.instance.id"]).to eq(described_class.service_instance_id) + expect(attrs["service.instance.id"]).not_to eq(parent_id) + expect(attrs["appsignal.config.name"]).to eq("process-resource-spec") + end + end + + it "does nothing when collector mode has not started" do + expect(described_class.before_fork).to be_nil + expect { described_class.after_fork(nil) }.not_to raise_error + end + + it "leaves a provider alone when its service instance ID already changed" do + Appsignal::OpenTelemetry.configure(config) + instance_ids = described_class.before_fork + refreshed = ::OpenTelemetry.meter_provider.resource.merge( + ::OpenTelemetry::SDK::Resources::Resource.create("service.instance.id" => "other") + ) + ::OpenTelemetry.meter_provider.instance_variable_set(:@resource, refreshed) + allow(Process).to receive(:pid).and_return(Process.pid + 1) + + logs = capture_logs { described_class.after_fork(instance_ids) } + + expect(logs).to be_empty + expect(::OpenTelemetry.meter_provider.resource).to equal(refreshed) + end + + it "uses a public resource setter when the provider has one" do + Appsignal::OpenTelemetry.configure(config) + provider = ::OpenTelemetry.meter_provider + provider.singleton_class.attr_writer(:resource) + allow(provider).to receive(:resource=).and_call_original + + logs = capture_logs { fork_in_place } + + expect(logs).to be_empty + expect(provider).to have_received(:resource=) + expect(attributes(provider.resource)["process.pid"]).to eq(Process.pid) + end + + it "logs a warning when a provider's resource can't be read" do + Appsignal::OpenTelemetry.configure(config) + instance_ids = described_class.before_fork + ::OpenTelemetry.meter_provider = Object.new + + logs = capture_logs { described_class.after_fork(instance_ids) } + + expect(logs).to contains_log( + :warn, + "Could not refresh the OpenTelemetry resource after fork: " \ + "the meter provider's resource could not be read" + ) + end + + it "logs a warning when the resource reader doesn't return @resource" do + Appsignal::OpenTelemetry.configure(config) + provider = ::OpenTelemetry.meter_provider + original = provider.resource + instance_ids = described_class.before_fork + copy = ::OpenTelemetry::SDK::Resources::Resource.create(attributes(original)) + provider.define_singleton_method(:resource) { copy } + + logs = capture_logs { described_class.after_fork(instance_ids) } + + expect(logs).to contains_log( + :warn, + "Could not refresh the OpenTelemetry resource after fork: " \ + "the meter provider's resource could not be read" + ) + expect(provider.instance_variable_get(:@resource)).to equal(original) + end + + it "logs a warning when the resource reader returns something other than a resource" do + Appsignal::OpenTelemetry.configure(config) + provider = ::OpenTelemetry.meter_provider + instance_ids = described_class.before_fork + provider.instance_variable_set(:@resource, "not a resource") + provider.define_singleton_method(:resource) { @resource } + + logs = capture_logs { described_class.after_fork(instance_ids) } + + expect(logs).to contains_log( + :warn, + "Could not refresh the OpenTelemetry resource after fork: " \ + "the meter provider's resource could not be read" + ) + expect(provider.instance_variable_get(:@resource)).to eq("not a resource") + end + + it "logs a warning when the provider doesn't keep the new resource" do + Appsignal::OpenTelemetry.configure(config) + provider = ::OpenTelemetry.meter_provider + original = provider.resource + provider.define_singleton_method(:resource=) { |_resource| nil } + + logs = capture_logs { fork_in_place } + + expect(logs).to contains_log( + :warn, + "Could not refresh the OpenTelemetry resource after fork: " \ + "the meter provider did not keep the new resource" + ) + expect(provider.resource).to equal(original) + end + + it "logs a warning and doesn't raise when refreshing raises" do + Appsignal::OpenTelemetry.configure(config) + instance_ids = described_class.before_fork + allow(::OpenTelemetry).to receive(:tracer_provider).and_raise(RuntimeError, "boom") + + logs = capture_logs do + expect { described_class.after_fork(instance_ids) }.not_to raise_error + end + + expect(logs).to contains_log( + :warn, + "Could not refresh the OpenTelemetry resource after fork: RuntimeError: boom" + ) + end + end + end +end