Module: Datadog::Core::Utils::AtForkMonkeyPatch

Defined in:
lib/datadog/core/utils/at_fork_monkey_patch.rb

Overview

Monkey patches Kernel#fork and similar functions, adding an at_fork callback mechanism which is used to restart observability after the VM forks (e.g. in multiprocess Ruby apps).

Defined Under Namespace

Modules: KernelMonkeyPatch, ProcessMonkeyPatch

Class Method Summary collapse

Class Method Details

.apply! ⇒ Object



38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 38

def self.apply!
  return false unless supported?

  if RubyVersion.is?("< 3.1")
    [
      ::Process.singleton_class, # Process.fork
      ::Kernel.singleton_class,  # Kernel.fork
      ::Object,                  # fork without explicit receiver (it's defined as a method in ::Kernel)
      # Note: Modifying Object as we do here is irreversible. During tests, this
      # change will stick around even if we otherwise stub `Process` and `Kernel`
    ].each { |target| target.prepend(KernelMonkeyPatch) }
  end

  ::Process.singleton_class.prepend(ProcessMonkeyPatch)

  true
end

.at_fork(stage, &block) ⇒ Object

Registers a block to run at the given fork stage (+:before+, :parent, or :child).

Returns the registered block so callers can keep a handle to it and later deregister it via remove_at_fork.

Raises:

  • (ArgumentError)


94
95
96
97
98
99
100
101
102
103
104
105
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 94

def self.at_fork(stage, &block)
  raise(ArgumentError, "Missing block argument") unless block

  AT_FORK_REGISTRY_MUTEX.synchronize do
    current = @at_fork_blocks
    @at_fork_blocks = current.merge(
      stage => (blocks_for(current, stage) + [Callback.new(block)]).freeze
    ).freeze
  end

  block
end

.at_fork_blocks(before:, parent:, child:) ⇒ Object

Registers one before/parent/child callback triplet atomically relative to the snapshots taken by fork dispatch.



109
110
111
112
113
114
115
116
117
118
119
120
121
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 109

def self.at_fork_blocks(before:, parent:, child:)
  blocks = {before: before, parent: parent, child: child}
  group = Object.new.freeze
  AT_FORK_REGISTRY_MUTEX.synchronize do
    current = @at_fork_blocks
    @at_fork_blocks = {
      before: (current.fetch(:before) + [Callback.new(before, group)]).freeze,
      parent: (current.fetch(:parent) + [Callback.new(parent, group)]).freeze,
      child: (current.fetch(:child) + [Callback.new(child, group)]).freeze,
    }.freeze
  end
  blocks
end

.remove_at_fork(stage, block) ⇒ Object

Deregisters a block previously registered with at_fork for the given stage. It is a no-op (does not raise) when block was never registered (or was already removed). Raises ArgumentError for an unknown stage, matching the at_fork contract.



127
128
129
130
131
132
133
134
135
136
137
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 127

def self.remove_at_fork(stage, block)
  AT_FORK_REGISTRY_MUTEX.synchronize do
    current = @at_fork_blocks
    callbacks = blocks_for(current, stage)
    @at_fork_blocks = current.merge(
      stage => callbacks.reject { |callback| callback.block.equal?(block) }.freeze
    ).freeze
  end

  nil
end

.run_at_fork_blocks(stage, snapshot:, started: nil) ⇒ Object

Runs the callbacks copied for one fork lifecycle. Before callbacks stop at the first failure; parent and child callbacks all run before the first failure is re-raised.



59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 59

def self.run_at_fork_blocks(stage, snapshot:, started: nil)
  callbacks = blocks_for(snapshot, stage)
  if stage == :before
    return callbacks.each do |callback|
      started[callback.group] = true if started && callback.group
      callback.block.call
    end
  end

  error = nil
  callbacks.each do |callback|
    next if started && callback.group && !started.key?(callback.group)

    callback.block.call
  rescue Exception => e # rubocop:disable Lint/RescueException -- finish cleanup, then re-raise the first failure
    error ||= e
  end
  raise error if error
end

.snapshot_at_fork_blocks ⇒ Object

Returns the immutable callback registry for one fork lifecycle. Copy-on-write updates publish a new reference, so registrations made afterwards apply only to the next lifecycle without locking this read.



142
143
144
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 142

def self.snapshot_at_fork_blocks
  @at_fork_blocks
end

.supported? ⇒ Boolean

Returns:

  • (Boolean)


34
35
36
# File 'lib/datadog/core/utils/at_fork_monkey_patch.rb', line 34

def self.supported?
  Process.respond_to?(:fork)
end