Class: Wurk::Swarm
- Inherits:
-
Object
- Object
- Wurk::Swarm
- Includes:
- Component
- Defined in:
- lib/wurk/swarm.rb,
lib/wurk/swarm/backoff.rb,
lib/wurk/swarm/restart.rb,
lib/wurk/swarm/child_boot.rb,
lib/wurk/swarm/orphan_guard.rb
Overview
Parent supervisor. Forks N children per the worker topology, monitors PIDs, relays signals, respawns crashed children with per-slot exponential backoff, handles rolling restart on SIGUSR1, recycles RSS-bloated children.
The supervise loop never sleeps on behalf of a respawn or a restart: crash backoff is tracked as per-slot due-times (Swarm::Backoff) and rolling restart / recycle run as a non-blocking state machine (Swarm::Restart) advanced one phase per tick. TERM/INT is therefore honored within a tick regardless of restart or backoff state.
Boot ordering (must be exact — see docs/idea/03-process-model.md):
1. Host app boots fully; eager loads done.
2. Railtie `after_initialize` fires.
3. `boot` closes parent-side connections (Redis, ActiveRecord).
4. `boot` forks N children.
5. Each child reconnects DB + opens a fresh Redis pool, then
installs its own signal handlers and starts the Launcher.
6. Parent calls `supervise` to enter the wait/relay loop.
Signals (see docs/idea/04-signals.md):
TERM/INT → `shutdown` (graceful drain; aborts any restart)
TSTP → relay TSTP (quiet — stop fetching; one-way, no resume)
USR1 → `rolling_restart` (zero-downtime cycle)
Defined Under Namespace
Classes: Backoff, ChildBoot, OrphanGuard, Restart
Constant Summary collapse
- SUPERVISE_TICK =
0.2- RESPAWN_BACKOFF =
1.0- HEARTBEAT_WAIT =
30- MEMORY_CHECK_INTERVAL =
10- DEFAULT_SHUTDOWN_TIMEOUT =
25- SUPERVISOR_POOL_SIZE =
The supervisor's own Redis pool: one connection, disjoint from every capsule pool. Post-boot the parent's only Redis use is the heartbeat probe (
heartbeat_seen?) — one SISMEMBER per restart tick — and routing it through the capsule main pool reopened, at capsule size (>= 10), the very sockets step 3 had just closed; the next respawn/restart fork then inherited them. Dedicated and torn down before every fork, so no child ever sees a parent socket. 1- SUPERVISOR_POOL_NAME =
'swarm-supervisor'- HARD_KILL_REAP_ATTEMPTS =
Poll budget for the post-SIGKILL reap sweep (250ms). Deliberately far shorter than any drain deadline — it runs after the fleet is already dead and races a kernel teardown measured in microseconds.
25- HARD_KILL_REAP_INTERVAL =
0.01- SHUTDOWN_GRACE =
Children each hard_shutdown after their own drain deadline (bulk_requeue
- a 3s ensure window + heartbeat cleanup); the parent must not SIGKILL them mid-tail, so its own wait always extends past theirs by this much.
5- RETERM_INTERVAL =
How often the drain re-relays TERM to whatever is still in the child table. A child forked in the instant before the drain can miss the first one outright: from
Process.forkreturning in the child until ChildBoot resets the inherited traps, the PARENT's handler is what a TERM finds there, and it writes the signal into the parent's self-pipe (or, oncefork_childhas dropped that pipe, into nothing) instead of stopping the child. The child then boots clean, drains nothing, andwait_for_childrenburns its whole budget beforehard_kill_stragglersSIGKILLs it — taking its in-flight jobs with it. The window is the_forkhook chain (ActiveSupport's ForkTracker and friends all run child-side beforeProcess.forkreturns), so it widens with load: measured 10.5s drains against a 2.6s baseline once that chain is slowed. No child-side fix can close it — the window precedes any code the child controls — so the supervisor repeats itself instead. Re-sending is safe: by the second pass the child owns its handler, TERM to a child already draining is a no-op, and the pids come from the child table, where an exited child is a zombie until reaped and so can never be a recycled pid. 1.0- SWARM_SIGNALS =
USR2 is relayed (log reopen) — without a trap, a logrotate config that signals the master pid would hit USR2's default disposition and kill the whole swarm.
{ 'TERM' => :term, 'INT' => :term, 'TSTP' => :tstp, 'USR1' => :usr1, 'USR2' => :usr2 }.freeze
Constants included from Component
Component::DEFAULT_THREAD_PRIORITY, Component::LEADER_CACHE_TTL_MS, Component::PROCESS_NONCE
Instance Attribute Summary collapse
-
#topology ⇒ Object
readonly
Returns the value of attribute topology.
Attributes included from Component
Instance Method Summary collapse
-
#boot(install_signals: true) ⇒ Object
install_signals:is false in tests so the integration suite can driveshutdown/rolling_restartdirectly without poisoning the test process's signal handlers. -
#children ⇒ Object
A snapshot, not the live table.
-
#initialize(topology:, config: Wurk.configuration, memory_limit: config.memory_limit_kb, shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT) ⇒ Swarm
constructor
A new instance of Swarm.
-
#owner? ⇒ Boolean
Only the process that forked the children may supervise or signal them.
-
#quiet_swarm ⇒ Object
TSTP quiet is one-way and GLOBAL (spec §21.3): it must survive respawns and memory recycles, or a quieted-but-crashed child's replacement would resume fetching mid-maintenance.
-
#request_shutdown ⇒ Object
Cross-thread drain request: raise a flag the supervise loop observes on its next tick instead of draining on the caller's thread.
-
#rolling_restart ⇒ Object
SIGUSR1: queue every live child for the rolling-restart state machine, which replaces one slot at a time (spawn replacement → await its heartbeat → TERM the old child → await its drain) without blocking the supervise thread, so TERM stays responsive throughout the cycle.
- #shutdown(timeout: @shutdown_timeout) ⇒ Object
- #supervise ⇒ Object
Methods included from Component
#default_tag, #fire_event, #handle_exception, hostname, #hostname, identity, #identity, #leader?, #logger, #mono_ms, #process_nonce, #real_ms, #redis, #safe_thread, tid, #tid, #watchdog
Constructor Details
#initialize(topology:, config: Wurk.configuration, memory_limit: config.memory_limit_kb, shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT) ⇒ Swarm
Returns a new instance of Swarm.
94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 |
# File 'lib/wurk/swarm.rb', line 94 def initialize(topology:, config: Wurk.configuration, memory_limit: config.memory_limit_kb, shutdown_timeout: DEFAULT_SHUTDOWN_TIMEOUT) @topology = topology @config = config @memory_limit = memory_limit @shutdown_timeout = shutdown_timeout # One reentrant lock covers the child table AND the restart machine: they # are mutually recursive (Restart#advance calls back into `describe` and # `spawn`, which read and write `@children`), so separate locks would be # taken in opposite orders — enqueue holds children→restart, advance holds # restart→children — and deadlock. Monitor's reentrancy is what lets those # callbacks re-enter. Never held across a sleep: `wait_for_children` and # the supervise tick sleep outside it. @lock = ::Monitor.new @children = {} @assignments = [] @owner_pid = nil @stopping = false @shutdown_requested = false @quieted = false @last_memory_check = 0 @respawn_backoff = Backoff.new(base: RESPAWN_BACKOFF) @restart = build_restart init_fork_unsafe_handles end |
Instance Attribute Details
#topology ⇒ Object (readonly)
Returns the value of attribute topology.
92 93 94 |
# File 'lib/wurk/swarm.rb', line 92 def topology @topology end |
Instance Method Details
#boot(install_signals: true) ⇒ Object
install_signals: is false in tests so the integration suite can
drive shutdown / rolling_restart directly without poisoning the
test process's signal handlers.
Traps go in BEFORE fork_children: a TERM landing in the (previously post-fork) window between fork and trap installation left the parent on its default disposition — it died instantly and orphaned live, fetching children. Installed first, the trap queues the TERM and the supervise loop drains it (relaying to children) even if it arrives mid-boot.
prepare_for_fork! is the last thing before the fork and stays after
close_parent_sockets: it materializes the slot-independent config so
every child inherits it copy-on-write, and none of what it touches opens
a socket — anything that did would land back on the wrong side of step 3.
Nothing else belongs between steps 3 and 4. #101 boot-audit: a pre-fork
SCRIPT LOAD was tried here (the cache is server-global, so one upload
could serve the whole fleet and children would only PING) and MEASURED
SLOWER — bench:swarm_boot 152 -> 95 i/s, a 36.7% regression. Step 3 has
just closed every parent socket, so the upload has to open its own
connection, and it lands serially ahead of the first fork: ~2ms of a
~10ms boot, on every child's path. The per-child upload it replaced cost
far less — the children reconnect in parallel, and it rides in the
pipeline of a PING each one already sends (ChildBoot#validate_redis!).
144 145 146 147 148 149 150 151 152 153 154 155 |
# File 'lib/wurk/swarm.rb', line 144 def boot(install_signals: true) raise 'Wurk::Swarm already booted' unless @assignments.empty? raise ArgumentError, 'Topology has no slots' if @topology.empty? @assignments = @topology.assignments.freeze @owner_pid = ::Process.pid install_signal_handlers if install_signals close_parent_sockets @config.prepare_for_fork! fork_children child_pids end |
#children ⇒ Object
A snapshot, not the live table. Callers read this off the supervise thread, which inserts (respawn, restart) and deletes (reap) on every tick; handing out the live Hash lets them iterate it mid-mutation.
160 161 162 |
# File 'lib/wurk/swarm.rb', line 160 def children @lock.synchronize { @children.dup } end |
#owner? ⇒ Boolean
Only the process that forked the children may supervise or signal them.
A forked child inherits this object along with the host's at_exit hooks
— and rails_boot registers one that drains the swarm — so ChildBoot's
exit reaches a full drain inside the child. There @children holds the
child's SIBLINGS, not its own children: unguarded, that TERMs them, stalls
the whole shutdown timeout waiting on PIDs it can never reap, then SIGKILLs
whatever survived.
207 208 209 |
# File 'lib/wurk/swarm.rb', line 207 def owner? ::Process.pid == @owner_pid end |
#quiet_swarm ⇒ Object
TSTP quiet is one-way and GLOBAL (spec §21.3): it must survive respawns and memory recycles, or a quieted-but-crashed child's replacement would resume fetching mid-maintenance. The flag makes every future fork boot already-quieted (see ChildBoot start_quiet).
215 216 217 218 |
# File 'lib/wurk/swarm.rb', line 215 def quiet_swarm @quieted = true relay_signal('TSTP') end |
#request_shutdown ⇒ Object
Cross-thread drain request: raise a flag the supervise loop observes on
its next tick instead of draining on the caller's thread. shutdown
walks and clears the child table that the supervise thread is
concurrently mutating (reap, respawn, restart), so two threads inside it
race. Every caller that isn't the supervise thread — rails_boot's
at_exit, which fires on the host's main thread — comes through here.
196 197 198 |
# File 'lib/wurk/swarm.rb', line 196 def request_shutdown @shutdown_requested = true end |
#rolling_restart ⇒ Object
SIGUSR1: queue every live child for the rolling-restart state machine, which replaces one slot at a time (spawn replacement → await its heartbeat → TERM the old child → await its drain) without blocking the supervise thread, so TERM stays responsive throughout the cycle.
224 225 226 |
# File 'lib/wurk/swarm.rb', line 224 def rolling_restart @lock.synchronize { @restart.enqueue(@children.keys) } end |
#shutdown(timeout: @shutdown_timeout) ⇒ Object
178 179 180 181 182 183 184 185 186 187 188 |
# File 'lib/wurk/swarm.rb', line 178 def shutdown(timeout: @shutdown_timeout) return unless owner? @stopping = true @lock.synchronize { @restart.abort } relay_signal('TERM') wait_for_children(timeout + SHUTDOWN_GRACE) hard_kill_stragglers close_supervisor_pool close_signal_pipe end |
#supervise ⇒ Object
164 165 166 167 168 169 170 171 172 173 174 175 176 |
# File 'lib/wurk/swarm.rb', line 164 def supervise return unless owner? until done? drain_signals shutdown if @shutdown_requested && !@stopping reap_children spawn_due_respawns advance_restart unless @stopping check_memory_pressure sleep SUPERVISE_TICK end end |