Class: Checkend::Worker
- Inherits:
-
Object
- Object
- Checkend::Worker
- Defined in:
- lib/checkend/worker.rb
Overview
Worker handles async sending of notices via a background thread.
It maintains a queue of notices and sends them in the background, implementing throttling on errors and graceful shutdown.
Constant Summary collapse
- SHUTDOWN =
Object.new.freeze
- FLUSH =
Object.new.freeze
- BASE_THROTTLE =
Exponential backoff base for throttling
1.05- MAX_THROTTLE =
100
Instance Method Summary collapse
-
#flush(timeout: nil) ⇒ Object
Flush the queue, blocking until all current notices are sent.
-
#initialize(config) ⇒ Worker
constructor
A new instance of Worker.
-
#push(notice) ⇒ Boolean
Push a notice onto the queue for async sending.
-
#queue_size ⇒ Integer
Get the current queue size.
-
#running? ⇒ Boolean
Check if the worker is running.
-
#shutdown(timeout: nil) ⇒ Object
Shutdown the worker, waiting for pending notices.
Constructor Details
Instance Method Details
#flush(timeout: nil) ⇒ Object
Flush the queue, blocking until all current notices are sent
58 59 60 61 62 63 64 65 66 |
# File 'lib/checkend/worker.rb', line 58 def flush(timeout: nil) timeout ||= @config.timeout cv = ConditionVariable.new @mutex.synchronize do @queue.push(cv) cv.wait(@mutex, timeout) end end |
#push(notice) ⇒ Boolean
Push a notice onto the queue for async sending
31 32 33 34 35 36 37 |
# File 'lib/checkend/worker.rb', line 31 def push(notice) return false if @shutdown return false if @queue.size >= @config.max_queue_size @queue.push(notice) true end |
#queue_size ⇒ Integer
Get the current queue size
78 79 80 |
# File 'lib/checkend/worker.rb', line 78 def queue_size @queue.size end |
#running? ⇒ Boolean
Check if the worker is running
71 72 73 |
# File 'lib/checkend/worker.rb', line 71 def running? @thread&.alive? && !@shutdown end |
#shutdown(timeout: nil) ⇒ Object
Shutdown the worker, waiting for pending notices
42 43 44 45 46 47 48 49 50 51 52 53 |
# File 'lib/checkend/worker.rb', line 42 def shutdown(timeout: nil) timeout ||= @config.shutdown_timeout @mutex.synchronize do return if @shutdown @shutdown = true @queue.push(SHUTDOWN) end @thread&.join(timeout) end |