Skip to content
Open
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
Original file line number Diff line number Diff line change
@@ -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.
6 changes: 5 additions & 1 deletion lib/appsignal/opentelemetry.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -103,6 +106,7 @@ def configure(config)
)

@started = true
ProcessResource.attach_fork_hook
rescue LoadError => e
@started = false
Appsignal::Utils::StdoutAndLoggerMessage.error(
Expand Down
137 changes: 137 additions & 0 deletions lib/appsignal/opentelemetry/process_resource.rb
Original file line number Diff line number Diff line change
@@ -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
90 changes: 90 additions & 0 deletions spec/integration/collector_mode_fork_resource_spec.rb
Original file line number Diff line number Diff line change
@@ -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
31 changes: 31 additions & 0 deletions spec/integration/runners/collector_mode_fork_resource.rb
Original file line number Diff line number Diff line change
@@ -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)
Loading
Loading