Class: Wurk::Metrics::Rollup
- Inherits:
-
Object
- Object
- Wurk::Metrics::Rollup
- Includes:
- Component
- Defined in:
- lib/wurk/metrics/rollup.rb
Overview
Leader-only background thread that rolls the per-class minute buckets
written by Wurk::Metrics::History (j|YYYYMMDD|H:M) up into compact,
cluster-total time-series buckets the dashboard "throughput" / "failures"
charts read directly:
jr|1m|<epoch> HASH {p,f,ms} TTL 24h (1-minute resolution)
jr|5m|<epoch> HASH {p,f,ms} TTL 7d (5-minute resolution)
jr|1h|<epoch> HASH {p,f,ms} TTL 30d (1-hour resolution)
<epoch> is the UTC start-of-bucket as integer seconds. Totals are summed
across every job class, so a 30-day chart reads ~720 small hashes instead
of fanning out over ~43k per-class minute keys.
Every write is an idempotent HSET. Each tick recomputes the trailing few 1m buckets from the source minute hash, then recomputes the coarse buckets from their 1m children. Re-running a tick (a missed tick, a leadership change, a late metric write) converges to the same totals — it never double-counts. Storage is bounded purely by the per-bucket TTLs; see docs/metrics-history.md for the retention math.
Constant Summary collapse
- PREFIX =
'jr'- BUCKETS =
bucket => [step_seconds, ttl_seconds]. The retention is the issue's spec: 1m kept 24h, 5m kept 7d, 1h kept 30d.
{ '1m' => [60, 24 * 60 * 60], '5m' => [300, 7 * 24 * 60 * 60], '1h' => [3600, 30 * 24 * 60 * 60] }.freeze
- COARSE =
%w[5m 1h].freeze
- DEFAULT_TICK_SECONDS =
60- LOOKBACK_MINUTES =
Re-roll the last N completed minutes from source on every tick (idempotent). This self-heals a leadership failover / restart or a late metric write up to N minutes old — the source
j|…buckets live 3 days, so re-reading them folds the gap back in. Only outages longer than this leave a hole that ages out with the bucket TTL (best-effort metrics). 15
Instance Attribute Summary collapse
-
#config ⇒ Object
included
from Component
readonly
Returns the value of attribute config.
Class Method Summary collapse
Instance Method Summary collapse
- #default_tag(dir = Dir.pwd) ⇒ Object included from Component
-
#fire_event(event, oneshot: true, reverse: false, reraise: false) ⇒ Object
included
from Component
Invokes lifecycle hooks for
event. - #handle_exception(ex, ctx = {}) ⇒ Object included from Component
- #hostname ⇒ Object included from Component
- #identity ⇒ Object included from Component
-
#initialize(config) ⇒ Rollup
constructor
A new instance of Rollup.
-
#leader? ⇒ Boolean
included
from Component
True iff this process currently holds the cluster
dear-leaderlock. -
#logger ⇒ Object
included
from Component
--- delegated to config -------------------------------------------.
- #mono_ms ⇒ Object included from Component
- #process_nonce ⇒ Object included from Component
-
#real_ms ⇒ Object
included
from Component
--- clocks ---------------------------------------------------------.
- #redis(idempotent: false) ⇒ Object included from Component
-
#roll(now = ::Time.now) ⇒ Object
One rollup pass, bypassing the leader gate and the sleep loop.
-
#safe_thread(name, priority: nil, &block) ⇒ Object
included
from Component
Spawns a named thread that runs
blockunderwatchdog(name). - #start ⇒ Object
-
#terminate ⇒ Object
Blocks until the thread is really gone: the launcher releases the cluster lock immediately after this returns, and a roll still in flight would write the same buckets as the next leader's first one.
-
#tick(now: ::Time.now) ⇒ Object
Leader-gated: only the elected leader writes the cluster-total series, so N workers don't each HSET the same buckets every minute.
- #tid ⇒ Object included from Component
-
#watchdog(last_words) ⇒ Object
included
from 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.
Constructor Details
#initialize(config) ⇒ Rollup
Returns a new instance of Rollup.
55 56 57 58 59 |
# File 'lib/wurk/metrics/rollup.rb', line 55 def initialize(config) @config = config @timer = TimerLoop.new(config[:metrics_rollup_interval] || DEFAULT_TICK_SECONDS) @thread = nil end |
Instance Attribute Details
#config ⇒ Object (readonly) Originally defined in module Component
Returns the value of attribute config.
Class Method Details
.bucket_key(bucket, epoch) ⇒ Object
51 52 53 |
# File 'lib/wurk/metrics/rollup.rb', line 51 def self.bucket_key(bucket, epoch) "#{PREFIX}|#{bucket}|#{epoch}" end |
Instance Method Details
#default_tag(dir = Dir.pwd) ⇒ Object Originally defined in module Component
#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
#hostname ⇒ Object Originally defined in module Component
#identity ⇒ Object 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.
#logger ⇒ Object Originally defined in module Component
--- delegated to config -------------------------------------------
#mono_ms ⇒ Object Originally defined in module Component
#process_nonce ⇒ Object Originally defined in module Component
#real_ms ⇒ Object Originally defined in module Component
--- clocks ---------------------------------------------------------
#redis(idempotent: false) ⇒ Object Originally defined in module Component
#roll(now = ::Time.now) ⇒ Object
One rollup pass, bypassing the leader gate and the sleep loop. Public so deterministic specs and a manual "roll now" can drive it directly.
93 94 95 96 97 98 99 |
# File 'lib/wurk/metrics/rollup.rb', line 93 def roll(now = ::Time.now) cur_min = floor_min(now) minutes = (1..LOOKBACK_MINUTES).map { |i| cur_min - (i * 60) } minutes.each { |epoch_min| write_minute_bucket(epoch_min) } recompute_coarse(minutes) nil end |
#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.
#start ⇒ Object
61 62 63 64 65 66 |
# File 'lib/wurk/metrics/rollup.rb', line 61 def start return @thread if @thread @timer.reset @thread = safe_thread('metrics-rollup') { @timer.run { tick } } end |
#terminate ⇒ Object
Blocks until the thread is really gone: the launcher releases the cluster lock immediately after this returns, and a roll still in flight would write the same buckets as the next leader's first one.
Cleared only on a confirmed join (Thread#join returns nil on timeout): a wedged thread must stay tracked so #start's guard returns it instead of calling @timer.reset, which would un-terminate the loop it is still inside and leave two threads HSETting the same buckets.
76 77 78 79 |
# File 'lib/wurk/metrics/rollup.rb', line 76 def terminate @timer.terminate @thread = nil if @thread&.join(TimerLoop::JOIN_TIMEOUT) end |
#tick(now: ::Time.now) ⇒ Object
Leader-gated: only the elected leader writes the cluster-total series, so N workers don't each HSET the same buckets every minute.
83 84 85 86 87 88 89 |
# File 'lib/wurk/metrics/rollup.rb', line 83 def tick(now: ::Time.now) return unless leader? roll(now) rescue StandardError => e handle_exception(e, { context: 'metrics-rollup' }) end |
#tid ⇒ Object 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.