Module: Wurk::Metrics::Query
- Defined in:
- lib/wurk/metrics/query.rb
Overview
Read-side for the per-class HASH bucket schema written by
Wurk::Metrics::History. Backs the Web UI's history pane and any
external dashboarding code; both rely on the same HGETALL fan-out
over a contiguous range of minute / hour keys.
Window caps are spec-enforced (§20 "DoS caps"), one per argument:
minutes ≤ 480 (8h)
hours ≤ 72 (3d — same as MID_TERM retention)
hours: is validated against MAX_HOURS only, never re-checked against
MAX_MINUTES after the ×60 conversion: the minute buckets it reads live
for MID_TERM (3 days), so the whole 72h window is backed by real data.
A wider window has no data to read anyway (the buckets are TTL'd out), so we fail loudly rather than silently returning sparse results.
Defined Under Namespace
Classes: WindowTooWide
Constant Summary collapse
- MAX_MINUTES =
rubocop:disable Metrics/ModuleLength
480- MAX_HOURS =
72- MAX_QUEUE_SERIES =
25
Class Method Summary collapse
- .accumulate!(totals, hash) ⇒ Object
- .aggregate_minutes(now, minutes) ⇒ Object
- .bucket_spec!(bucket) ⇒ Object
-
.bucket_starts(now, step, window) ⇒ Object
The last
window/stepstep-aligned bucket starts, oldest→newest, so they match the keys the rollup writes. - .cap_hours!(hours) ⇒ Object
- .cap_minutes!(minutes) ⇒ Object
- .check_window!(value, max, label) ⇒ Object
- .clamp_history_window!(window_seconds, ttl) ⇒ Object
- .floor_to(time, unit) ⇒ Object
-
.for_job(klass, minutes: nil, hours: nil, now: ::Time.now) ⇒ Object
Per-class time-series.
-
.history(bucket, window_seconds, now: ::Time.now) ⇒ Object
Cluster-total time-series for the dashboard throughput/failures charts, read from the compact buckets written by Wurk::Metrics::Rollup.
- .hour_series(klass, now, hours) ⇒ Object
- .hour_timestamps(now, hours) ⇒ Object
- .minute_keys(now, minutes) ⇒ Object
- .minute_series(klass, now, minutes) ⇒ Object
-
.minute_timestamps(now, minutes) ⇒ Object
Truncate to the minute so the bucket boundary matches what the writer used.
- .pipeline_hgetall(keys) ⇒ Object
- .pipeline_hmget(keys, fields) ⇒ Object
- .queue_bucket_hashes(bucket, starts) ⇒ Object
-
.queue_history(bucket, window_seconds, queues: nil, now: ::Time.now) ⇒ Object
Per-queue size/latency gauge time-series written by Metrics::QueueRollup.
- .queue_names(queues) ⇒ Object
- .queue_points(name, starts, hashes) ⇒ Object
-
.top_jobs(class_filter: nil, minutes: 60, hours: nil, now: ::Time.now) ⇒ Object
Aggregate per-job-class totals over a recent window of minute buckets.
- .validate_for_job!(klass, minutes, hours) ⇒ Object
- .zip_rows(timestamps, rows) ⇒ Object
Class Method Details
.accumulate!(totals, hash) ⇒ Object
151 152 153 154 155 156 157 158 159 160 |
# File 'lib/wurk/metrics/query.rb', line 151 def accumulate!(totals, hash) return if hash.nil? || hash.empty? hash.each do |field, value| klass, kind = field.split('|', 2) next unless kind && TOTAL_FIELDS.include?(kind) totals[klass][kind.to_sym] += Integer(value) end end |
.aggregate_minutes(now, minutes) ⇒ Object
145 146 147 148 149 |
# File 'lib/wurk/metrics/query.rb', line 145 def aggregate_minutes(now, minutes) totals = ::Hash.new { |h, k| h[k] = { p: 0, f: 0, ms: 0 } } pipeline_hgetall(minute_keys(now, minutes)).each { |hash| accumulate!(totals, hash) } totals end |
.bucket_spec!(bucket) ⇒ Object
104 105 106 107 108 |
# File 'lib/wurk/metrics/query.rb', line 104 def bucket_spec!(bucket) Wurk::Metrics::Rollup::BUCKETS.fetch(bucket) do raise ArgumentError, "bucket must be one of #{Wurk::Metrics::Rollup::BUCKETS.keys.inspect}" end end |
.bucket_starts(now, step, window) ⇒ Object
The last window/step step-aligned bucket starts, oldest→newest, so
they match the keys the rollup writes.
119 120 121 122 |
# File 'lib/wurk/metrics/query.rb', line 119 def bucket_starts(now, step, window) last = (now.to_i / step) * step (0...(window / step)).map { |i| last - (i * step) }.reverse end |
.cap_hours!(hours) ⇒ Object
134 135 136 |
# File 'lib/wurk/metrics/query.rb', line 134 def cap_hours!(hours) check_window!(Integer(hours), MAX_HOURS, 'hours') end |
.cap_minutes!(minutes) ⇒ Object
130 131 132 |
# File 'lib/wurk/metrics/query.rb', line 130 def cap_minutes!(minutes) check_window!(Integer(minutes), MAX_MINUTES, 'minutes') end |
.check_window!(value, max, label) ⇒ Object
138 139 140 141 142 143 |
# File 'lib/wurk/metrics/query.rb', line 138 def check_window!(value, max, label) raise ArgumentError, "#{label} must be positive" if value <= 0 raise WindowTooWide, "#{label} must be <= #{max} (got #{value})" if value > max value end |
.clamp_history_window!(window_seconds, ttl) ⇒ Object
110 111 112 113 114 115 |
# File 'lib/wurk/metrics/query.rb', line 110 def clamp_history_window!(window_seconds, ttl) window = Integer(window_seconds) raise ArgumentError, 'window must be positive' if window <= 0 [window, ttl].min end |
.floor_to(time, unit) ⇒ Object
194 195 196 197 198 199 200 |
# File 'lib/wurk/metrics/query.rb', line 194 def floor_to(time, unit) t = time.utc case unit when :min then ::Time.utc(t.year, t.month, t.day, t.hour, t.min) when :hour then ::Time.utc(t.year, t.month, t.day, t.hour) end end |
.for_job(klass, minutes: nil, hours: nil, now: ::Time.now) ⇒ Object
Per-class time-series. minutes reads the per-minute bucket; hours
reads the per-class hourly bucket (separate keys per spec, so a long
window doesn't fan out over 4320 minute hashes).
50 51 52 53 |
# File 'lib/wurk/metrics/query.rb', line 50 def for_job(klass, minutes: nil, hours: nil, now: ::Time.now) validate_for_job!(klass, minutes, hours) minutes ? minute_series(klass, now, cap_minutes!(minutes)) : hour_series(klass, now, cap_hours!(hours)) end |
.history(bucket, window_seconds, now: ::Time.now) ⇒ Object
Cluster-total time-series for the dashboard throughput/failures charts,
read from the compact buckets written by Wurk::Metrics::Rollup. bucket
is '1m'/'5m'/'1h'; window_seconds is clamped to that bucket's
retention. Returns [{at:, p:, f:, ms:}, ...] oldest→newest, gap-filled
with zeros so a chart has a continuous x-axis.
60 61 62 63 64 65 |
# File 'lib/wurk/metrics/query.rb', line 60 def history(bucket, window_seconds, now: ::Time.now) step, ttl = bucket_spec!(bucket) starts = bucket_starts(now, step, clamp_history_window!(window_seconds, ttl)) rows = pipeline_hmget(starts.map { |s| Wurk::Metrics::Rollup.bucket_key(bucket, s) }, %w[p f ms]) starts.zip(rows).map { |at, (p, f, ms)| { at: at, p: p.to_i, f: f.to_i, ms: ms.to_i } } end |
.hour_series(klass, now, hours) ⇒ Object
167 168 169 170 171 |
# File 'lib/wurk/metrics/query.rb', line 167 def hour_series(klass, now, hours) = (now, hours) keys = .map { |t| Wurk::Metrics::History.hour_key(klass, t) } zip_rows(, pipeline_hmget(keys, %w[p f ms])) end |
.hour_timestamps(now, hours) ⇒ Object
189 190 191 192 |
# File 'lib/wurk/metrics/query.rb', line 189 def (now, hours) floor = floor_to(now, :hour) (0...hours).map { |i| floor - (i * 3600) }.reverse end |
.minute_keys(now, minutes) ⇒ Object
177 178 179 |
# File 'lib/wurk/metrics/query.rb', line 177 def minute_keys(now, minutes) (now, minutes).map { |t| Wurk::Metrics::History.minute_key(t) } end |
.minute_series(klass, now, minutes) ⇒ Object
162 163 164 165 |
# File 'lib/wurk/metrics/query.rb', line 162 def minute_series(klass, now, minutes) rows = pipeline_hmget(minute_keys(now, minutes), %W[#{klass}|p #{klass}|f #{klass}|ms]) zip_rows((now, minutes), rows) end |
.minute_timestamps(now, minutes) ⇒ Object
Truncate to the minute so the bucket boundary matches what the writer used. Fractional-second drift would otherwise pull in an unrelated minute on the edge of the window.
184 185 186 187 |
# File 'lib/wurk/metrics/query.rb', line 184 def (now, minutes) floor = floor_to(now, :min) (0...minutes).map { |i| floor - (i * 60) }.reverse end |
.pipeline_hgetall(keys) ⇒ Object
202 203 204 205 206 |
# File 'lib/wurk/metrics/query.rb', line 202 def pipeline_hgetall(keys) return [] if keys.empty? Wurk.redis { |c| c.pipelined { |p| keys.each { |k| p.call('HGETALL', k) } } } end |
.pipeline_hmget(keys, fields) ⇒ Object
208 209 210 211 212 |
# File 'lib/wurk/metrics/query.rb', line 208 def pipeline_hmget(keys, fields) return [] if keys.empty? Wurk.redis { |c| c.pipelined { |p| keys.each { |k| p.call('HMGET', k, *fields) } } } end |
.queue_bucket_hashes(bucket, starts) ⇒ Object
84 85 86 87 |
# File 'lib/wurk/metrics/query.rb', line 84 def queue_bucket_hashes(bucket, starts) pipeline_hgetall(starts.map { |s| Wurk::Metrics::QueueRollup.bucket_key(bucket, s) }) .map { |h| h.is_a?(::Array) ? h.each_slice(2).to_h : (h || {}) } end |
.queue_history(bucket, window_seconds, queues: nil, now: ::Time.now) ⇒ Object
Per-queue size/latency gauge time-series written by
Metrics::QueueRollup. bucket is '1m'/'5m'/'1h'; window_seconds is
clamped to the bucket's retention. Returns one entry per live queue
(or the explicit queues: list) — [{name:, points: [{at:, size:, latency:}, …]}, …] — oldest→newest, gap-filled with zeros so a chart
has a continuous x-axis. Capped at MAX_QUEUE_SERIES queues to bound the
payload; the cap is logged-free because queue cardinality is small.
74 75 76 77 78 79 80 81 82 |
# File 'lib/wurk/metrics/query.rb', line 74 def queue_history(bucket, window_seconds, queues: nil, now: ::Time.now) step, ttl = bucket_spec!(bucket) starts = bucket_starts(now, step, clamp_history_window!(window_seconds, ttl)) names = queue_names(queues) return [] if names.empty? hashes = queue_bucket_hashes(bucket, starts) names.map { |name| { name: name, points: queue_points(name, starts, hashes) } } end |
.queue_names(queues) ⇒ Object
91 92 93 94 |
# File 'lib/wurk/metrics/query.rb', line 91 def queue_names(queues) names = queues || Wurk.redis { |c| c.call('SMEMBERS', Wurk::Keys::QUEUES_SET) } names.sort.first(MAX_QUEUE_SERIES) end |
.queue_points(name, starts, hashes) ⇒ Object
96 97 98 99 100 101 102 |
# File 'lib/wurk/metrics/query.rb', line 96 def queue_points(name, starts, hashes) size_field = "#{name}|#{Wurk::Metrics::QueueRollup::SIZE_KIND}" lat_field = "#{name}|#{Wurk::Metrics::QueueRollup::LAT_KIND}" starts.zip(hashes).map do |at, hash| { at: at, size: hash[size_field].to_i, latency: (hash[lat_field] || 0).to_f } end end |
.top_jobs(class_filter: nil, minutes: 60, hours: nil, now: ::Time.now) ⇒ Object
Aggregate per-job-class totals over a recent window of minute
buckets. Returns array of [class_name, {p:, f:, ms:}] tuples
sorted by volume (p + f) descending so the UI's "top jobs" table
renders without a second sort pass.
40 41 42 43 44 45 |
# File 'lib/wurk/metrics/query.rb', line 40 def top_jobs(class_filter: nil, minutes: 60, hours: nil, now: ::Time.now) window = hours ? cap_hours!(hours) * 60 : cap_minutes!(minutes) rows = aggregate_minutes(now, window).to_a rows = rows.select { |(k, _)| k.start_with?(class_filter) } if class_filter && !class_filter.empty? rows.sort_by { |(_k, s)| -(s[:p] + s[:f]) } end |
.validate_for_job!(klass, minutes, hours) ⇒ Object
124 125 126 127 128 |
# File 'lib/wurk/metrics/query.rb', line 124 def validate_for_job!(klass, minutes, hours) raise ArgumentError, 'klass required' if klass.nil? || klass.empty? raise ArgumentError, 'pass exactly one of minutes: or hours:' if minutes && hours raise ArgumentError, 'pass minutes: or hours:' if minutes.nil? && hours.nil? end |
.zip_rows(timestamps, rows) ⇒ Object
173 174 175 |
# File 'lib/wurk/metrics/query.rb', line 173 def zip_rows(, rows) .zip(rows).map { |at, (p, f, ms)| { at: at, p: p.to_i, f: f.to_i, ms: ms.to_i } } end |