Class: Wurk::Scheduled::Poller

Inherits:
Object
  • Object
show all
Includes:
Component
Defined in:
lib/wurk/scheduled.rb

Overview

Single thread that wakes on a randomized interval, drains both ZSETs, then sleeps again. Random spread prevents the cluster from dogpiling Redis at the top of each cadence.

Constant Summary collapse

INITIAL_WAIT =
10

Instance Attribute Summary collapse

Instance Method Summary collapse

Constructor Details

#initialize(config) ⇒ Poller

Returns a new instance of Poller.



163
164
165
166
167
168
169
170
171
172
# File 'lib/wurk/scheduled.rb', line 163

def initialize(config)
  @config = config
  @enq = (config[:scheduled_enq] || Enq).new(config)
  @done = false
  @mutex = ::Mutex.new
  @sleeper = ::ConditionVariable.new
  @thread = nil
  @rnd = ::Random.new
  @last_cleanup_ms = 0
end

Instance Attribute Details

#configObject (readonly) Originally defined in module Component

Returns the value of attribute config.

#rndObject

Returns the value of attribute rnd.



161
162
163
# File 'lib/wurk/scheduled.rb', line 161

def rnd
  @rnd
end

Instance Method Details

#default_tag(dir = Dir.pwd) ⇒ Object Originally defined in module Component

#enqueueObject

Called on every wake. Any raise inside the Enq is reported and the loop continues — a transient Redis blip must not kill the scheduler.



211
212
213
214
215
# File 'lib/wurk/scheduled.rb', line 211

def enqueue
  @enq.enqueue_jobs
rescue StandardError => e
  handle_exception(e, { context: 'scheduler' })
end

#fire_event(event, oneshot: true, reverse: false, reraise: false) ⇒ Object Originally defined in module Component

Invokes lifecycle hooks for event. Hooks run in registration order (or LIFO when reverse: true, used for teardown). A raise in one hook is reported via handle_exception and does NOT stop the next hook unless reraise: true (used in tests / fail-fast boot). oneshot: true clears the bucket after dispatch so the event can't fire twice.

#handle_exception(ex, ctx = {}) ⇒ Object Originally defined in module Component

#hostnameObject Originally defined in module Component

#identityObject Originally defined in module Component

#leader?Boolean Originally defined in module Component

True iff this process currently holds the cluster dear-leader lock. Cached per Component instance for LEADER_CACHE_TTL_MS (~5s): cron and the metrics rollups call this every tick, and an uncached GET would double their Redis traffic at short intervals for no benefit — the lock's own renewal cadence (60s+, spec §6.1) easily tolerates a few-second-stale read. Returns false unconditionally when WURK_LEADER=false (or SIDEKIQ_LEADER=false) is set on the process (opt-out hot-standby). Any Redis error is swallowed → false, so a transient partition can't propagate as an exception into user code.

Spec: docs/target/sidekiq-ent.md §6.1.

Returns:

  • (Boolean)

#loggerObject Originally defined in module Component

--- delegated to config -------------------------------------------

#mono_msObject Originally defined in module Component

#process_nonceObject Originally defined in module Component

#real_msObject Originally defined in module Component

--- clocks ---------------------------------------------------------

#redis(idempotent: false) ⇒ Object Originally defined in module Component

#safe_thread(name, priority: nil, &block) ⇒ Object Originally defined in module Component

Spawns a named thread that runs block under watchdog(name). The parent must retain the returned Thread; otherwise GC may not, but report_on_exception is disabled so we don't double-log on death.

Priority resolution matches Sidekiq (component.rb:44-48): explicit argument, then config.thread_priority, then -1. Ruby's default of 0 buys a 100ms timeslice; each negative step halves it, so -1 keeps a CPU-heavy capsule from starving its siblings for a whole tick.

#startObject

Spawns the scheduler thread. INITIAL_WAIT delays the first sweep so a fleet-wide deploy doesn't have every freshly-booted process hit Redis simultaneously.



177
178
179
180
181
182
183
184
185
186
# File 'lib/wurk/scheduled.rb', line 177

def start
  @thread ||= safe_thread('scheduler') do # rubocop:disable Naming/MemoizedInstanceVariableName
    initial_wait
    until @done
      enqueue
      wait
    end
    logger.info('Scheduler exiting...')
  end
end

#terminateObject

Idempotent. Wakes the sleeping thread so it observes @done and exits. Also propagates the stop signal to @enq so any in-flight drain loop short-circuits instead of running to completion. Terminal, not a pause: @enq's stop flag is one-way, so this poller never polls again.

Joins before returning — the caller (Launcher#quiet, then #stop) clears the heartbeat right after, and a sweep still in flight would promote jobs on behalf of a process that no longer exists.

Cleared only on a confirmed join (Thread#join returns nil on timeout): a wedged sweep must stay tracked so #start's ||= guard returns it rather than spawning a second scheduler thread alongside it.



200
201
202
203
204
205
206
207
# File 'lib/wurk/scheduled.rb', line 200

def terminate
  @mutex.synchronize do
    @done = true
    @enq.terminate
    @sleeper.signal
  end
  @thread = nil if @thread&.join(TimerLoop::JOIN_TIMEOUT)
end

#tidObject Originally defined in module Component

#watchdog(last_words) ⇒ Object Originally defined in module Component

Wraps a block at a thread boundary: any unhandled exception is reported via handle_exception (so it lands in error_handlers / the log) and then re-raised. last_words is the component label included in the context.