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

Constants included from Component

Component::DEFAULT_THREAD_PRIORITY, Component::LEADER_CACHE_TTL_MS, Component::PROCESS_NONCE

Instance Attribute Summary collapse

Attributes included from Component

#config

Instance Method Summary collapse

Methods included from Component

#default_tag, #fire_event, #handle_exception, #hostname, #identity, #leader?, #logger, #mono_ms, #process_nonce, #real_ms, #redis, #safe_thread, #tid, #watchdog

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

#rndObject

Returns the value of attribute rnd.



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

def rnd
  @rnd
end

Instance Method Details

#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

#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