Class: FiberAudit::Runtime::Watchdog

Inherits:
Object
  • Object
show all
Defined in:
lib/fiber_audit/runtime/watchdog.rb

Overview

Monitors scheduler-owned heartbeat fibers from one dedicated process thread. rubocop:disable Metrics/ClassLength

Defined Under Namespace

Classes: Channel, Stall

Constant Summary collapse

SOURCE =
:scheduler_watchdog
STOP_TIMEOUT_SECONDS =
1

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(policy:, recorder:, redactor:, active_operations:, clock: Clock.new, thread_factory: ->(&block) { Thread.new(&block) }) ⇒ Watchdog

Returns a new instance of Watchdog.



24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
# File 'lib/fiber_audit/runtime/watchdog.rb', line 24

def initialize(
  policy:,
  recorder:,
  redactor:,
  active_operations:,
  clock: Clock.new,
  thread_factory: ->(&block) { Thread.new(&block) }
)
  validate_dependencies!(policy, recorder, redactor, active_operations, clock, thread_factory)
  @policy = policy
  @recorder = recorder
  @redactor = redactor
  @active_operations = active_operations
  @clock = clock
  @thread_factory = thread_factory
  @mutex = Mutex.new
  @wait_mutex = Mutex.new
  @condition = ConditionVariable.new
  @channels = {}.compare_by_identity
  @stall_sequence = 0
  @unsupported_seen = false
  @monitor_thread = nil
  @stopping = false
  @stopped = false
  emit_state(policy.enabled? ? :watchdog_absent : :watchdog_disabled)
end

Instance Attribute Details

#active_operationsObject (readonly)

Returns the value of attribute active_operations.



22
23
24
# File 'lib/fiber_audit/runtime/watchdog.rb', line 22

def active_operations
  @active_operations
end

#policyObject (readonly)

Returns the value of attribute policy.



22
23
24
# File 'lib/fiber_audit/runtime/watchdog.rb', line 22

def policy
  @policy
end

#recorderObject (readonly)

Returns the value of attribute recorder.



22
23
24
# File 'lib/fiber_audit/runtime/watchdog.rb', line 22

def recorder
  @recorder
end

Instance Method Details

#enabled?Boolean

Returns:

  • (Boolean)


51
52
53
# File 'lib/fiber_audit/runtime/watchdog.rb', line 51

def enabled?
  policy.enabled?
end

#fail_open?Boolean

Returns:

  • (Boolean)


55
56
57
# File 'lib/fiber_audit/runtime/watchdog.rb', line 55

def fail_open?
  recorder.session.policy.fail_open?
end

#poll(now_ns: nil) ⇒ Object



137
138
139
140
141
142
143
144
145
146
147
148
149
150
# File 'lib/fiber_audit/runtime/watchdog.rb', line 137

def poll(now_ns: nil)
  return state unless enabled?

  observed_now = now_ns.nil? ? @clock.monotonic_ns : Validation.integer(now_ns, 'watchdog poll time')
  @mutex.synchronize do
    @channels.each_value { |channel| poll_channel(channel, observed_now) if channel.active }
  end
  state
rescue StandardError => e
  handle_failure(e)
  raise e unless fail_open?

  :unsupported
end

#scheduler_closing(thread: Thread.current) ⇒ Object



113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
# File 'lib/fiber_audit/runtime/watchdog.rb', line 113

def scheduler_closing(thread: Thread.current)
  return self unless enabled?

  now_ns = safe_monotonic_ns
  @mutex.synchronize do
    channel = @channels.delete(thread)
    next unless channel

    channel.heartbeat.request_stop
    complete_stall(channel, now_ns: now_ns, resumed: false) if channel.stall
    unless channel.active || channel.unsupported
      @unsupported_seen = true
      emit_state(:watchdog_unsupported, thread: thread)
    end
  end
  wake_monitor
  self
rescue StandardError => e
  handle_failure(e)
  raise e unless fail_open?

  self
end

#scheduler_installed(thread: Thread.current, schedule: Fiber.method(:schedule), sleeper: Kernel.method(:sleep)) ⇒ Object



70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
# File 'lib/fiber_audit/runtime/watchdog.rb', line 70

def scheduler_installed(
  thread: Thread.current,
  schedule: Fiber.method(:schedule),
  sleeper: Kernel.method(:sleep)
)
  return self unless enabled?
  raise RuntimeContractError, 'scheduler thread must be a Thread' unless thread.is_a?(Thread)

  heartbeat = Heartbeat.new(
    clock: @clock,
    interval_ns: policy.heartbeat_interval_ns,
    owner_thread: thread,
    on_tick: method(:heartbeat_ticked),
    on_error: method(:heartbeat_failed)
  )
  previous = @mutex.synchronize do
    prior = @channels[thread]
    @channels[thread] = Channel.new(
      thread: thread,
      heartbeat: heartbeat,
      active: false,
      unsupported: false,
      stall: nil
    )
    prior
  end
  previous&.heartbeat&.request_stop
  heartbeat.start(schedule: schedule, sleeper: sleeper)
  self
rescue StandardError => e
  channel_failure(thread, e)
  raise e unless fail_open?

  self
end

#scheduler_unsupported(thread: Thread.current) ⇒ Object



106
107
108
109
110
111
# File 'lib/fiber_audit/runtime/watchdog.rb', line 106

def scheduler_unsupported(thread: Thread.current)
  return self unless enabled?

  channel_failure(thread, nil)
  self
end

#stateObject



59
60
61
62
63
64
65
66
67
68
# File 'lib/fiber_audit/runtime/watchdog.rb', line 59

def state
  return :disabled unless enabled?

  @mutex.synchronize do
    return :active if @channels.values.any?(&:active)
    return :unsupported if @unsupported_seen || @channels.values.any?(&:unsupported)

    :absent
  end
end

#stopObject



152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
# File 'lib/fiber_audit/runtime/watchdog.rb', line 152

def stop
  return self if @stopped

  @wait_mutex.synchronize do
    @stopping = true
    @condition.broadcast
  end
  now_ns = safe_monotonic_ns
  @mutex.synchronize do
    @channels.each_value do |channel|
      channel.heartbeat.request_stop
      complete_stall(channel, now_ns: now_ns, resumed: false) if channel.stall
      if !channel.active && !channel.unsupported
        @unsupported_seen = true
        emit_state(:watchdog_unsupported, thread: channel.thread)
      end
    end
    @channels.clear
  end
  stop_monitor_thread
  @stopped = true
  self
rescue StandardError => e
  handle_failure(e)
  @stopped = true
  raise e unless fail_open?

  self
end