From 91ecca792e72595e99527e25a43119b54561b7f9 Mon Sep 17 00:00:00 2001 From: James Boyer Date: Wed, 29 Jul 2026 15:40:09 -0400 Subject: [PATCH 1/4] Adding ability to query by max_age_by_name --- README.md | 22 +++- lib/delayed/monitor.rb | 44 +++++-- .../__snapshots__/monitor_spec.rb.snap | 109 ++++++++++++++++++ spec/delayed/monitor_spec.rb | 79 ++++++++++++- 4 files changed, 244 insertions(+), 10 deletions(-) diff --git a/README.md b/README.md index 03715a1..426c6b8 100644 --- a/README.md +++ b/README.md @@ -409,8 +409,26 @@ An additional _experimental_ metric is available, intended for use with applicat - **delayed.job.alert_age_percent** - the _percent_ to which the oldest job has reached the "age alert" threshold. (See the [Alerting Threshholds](#priority-based-alerting-threshholds) section above.) -All of these events may be subscribed to via a single regular expression (again, in your application -config or in an initializer): +If your jobs table has a `name` column (see [Database Setup](#database-setup)), one additional +event will be emitted, grouped by priority name, queue name, **and job name**: + +- **delayed.job.max_age_by_name** - the age of the oldest run_at value, per job name (excludes failed jobs) + +This answers "_which_ job is stuck?" when `delayed.job.max_age` alerts. A few behavioral notes: + +- Unlike the other metrics, it is not zero-filled: job names cannot be enumerated in advance, so a + (priority, queue, name) series is only emitted while matching jobs are present in the queue. +- Jobs enqueued before the `name` column existed (i.e. mid-upgrade) are reported under the name + `"unknown"`. +- Because this metric's cardinality scales with the number of distinct job names, it can be + disabled if that's a concern for your metrics provider: + +```ruby +Delayed::Monitor.emit_max_age_by_name = false +``` + +All of these events (including `max_age_by_name`) may be subscribed to via a single regular +expression (again, in your application config or in an initializer): ```ruby ActiveSupport::Notifications.subscribe(/delayed\.job\..*_(count|age|percent)/) do |*args| diff --git a/lib/delayed/monitor.rb b/lib/delayed/monitor.rb index c2871cb..c1d74e9 100644 --- a/lib/delayed/monitor.rb +++ b/lib/delayed/monitor.rb @@ -16,6 +16,7 @@ class Monitor ).freeze cattr_accessor :sleep_delay, instance_writer: false, default: 60 + cattr_accessor :emit_max_age_by_name, instance_writer: false, default: true def initialize @jobs = Job.group(:priority, :queue) @@ -27,6 +28,7 @@ def run! @memo = {} ActiveSupport::Notifications.instrument('delayed.monitor.run', default_tags) do METRICS.each { |metric| emit_metric!(metric) } + emit_metric_by_name!('max_age') if emit_max_age_by_name && Job.name_assignable? end interruptable_sleep(sleep_delay) end @@ -75,6 +77,15 @@ def emit_metric!(metric) end end + def emit_metric_by_name!(metric) + query_for("#{metric}_by_name").each do |(priority, queue, name), value| + ActiveSupport::Notifications.instrument( + "delayed.job.#{metric}_by_name", + default_tags.merge(priority: Priority.new(priority).to_s, queue: queue, name: name || 'unknown', value: value), + ) + end + end + def default_results @default_results ||= Priority.names.values.flat_map { |priority| (Worker.queues.presence || [Worker.default_queue_name]).map do |queue| @@ -96,24 +107,30 @@ def default_tags end # This method generates a query that scans the specified scope, groups by - # priority and queue, and calculates the specified aggregates. An outer - # query is executed for priority bucketing and appending db_now_utc (to - # avoid running these computations for each tuple in the inner query). - def grouped_query(scope, include_db_time: false, **kwargs) + # priority and queue (plus any extra_group_columns), and calculates the + # specified aggregates. An outer query is executed for priority bucketing + # and appending db_now_utc (to avoid running these computations for each + # tuple in the inner query). + def grouped_query(scope, include_db_time: false, extra_group_columns: [], **kwargs) inner_selects = kwargs.map { |key, (agg, expr)| as_expression(agg, expr, key) } outer_selects = kwargs.map { |key, (agg, _)| as_expression(agg == :count ? :sum : agg, key, key) } outer_selects << "#{self.class.sql_now_in_utc} AS db_now_utc" if include_db_time Delayed::Job - .from(scope.select(:priority, :queue, *inner_selects).group(:priority, :queue)) - .group(priority_case_statement, :queue).select( + .from(scope.select(:priority, :queue, *extra_group_columns, *inner_selects).group(:priority, :queue, *extra_group_columns)) + .group(priority_case_statement, :queue, *extra_group_columns).select( *outer_selects, "#{priority_case_statement} AS priority", 'queue AS queue', - ).group_by { |j| [j.priority.to_i, j.queue] } + *extra_group_columns.map { |column| "#{column} AS #{column}" }, + ).group_by { |j| result_key(j, extra_group_columns) } .transform_values(&:first) end + def result_key(record, extra_group_columns) + [record.priority.to_i, record.queue, *extra_group_columns.map { |column| record[column] }] + end + def as_expression(aggregate_function, aggregate_expression, column_name) "#{aggregate_function.to_s.upcase}(#{aggregate_expression}) AS #{column_name}" end @@ -150,6 +167,10 @@ def max_age_grouped pending_counts.transform_values { |j| time_ago(db_now(j), j.run_at) } end + def max_age_by_name_grouped + pending_run_ats_by_name.transform_values { |j| time_ago(db_now(j), j.run_at) } + end + def alert_age_percent_grouped pending_counts.each_with_object({}) do |((priority, queue), j), metrics| max_age = time_ago(db_now(j), j.run_at) @@ -192,6 +213,15 @@ def pending_counts ) end + def pending_run_ats_by_name + @memo[:pending_run_ats_by_name] ||= grouped_query( + jobs.pending, + include_db_time: true, + extra_group_columns: %i(name), + run_at: [:min, case_when(Job.claimable_clause.to_sql, 'run_at')], + ) + end + def failed_counts @memo[:failed_counts] ||= grouped_query(jobs.failed, count: [:count, '*']) end diff --git a/spec/delayed/__snapshots__/monitor_spec.rb.snap b/spec/delayed/__snapshots__/monitor_spec.rb.snap index 8da11e2..2722581 100644 --- a/spec/delayed/__snapshots__/monitor_spec.rb.snap +++ b/spec/delayed/__snapshots__/monitor_spec.rb.snap @@ -117,6 +117,39 @@ GroupAggregate (cost=...) -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) +-- QUERIES FOR `max_age_by_name`: +--------------------------------- +SELECT MIN(run_at) AS run_at, + TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, + CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at + FROM \"delayed_jobs\" + WHERE \"delayed_jobs\".\"failed_at\" IS NULL + AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" + +GroupAggregate (cost=...) + Output: min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + -> Sort (cost=...) + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.run_at + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + -> Subquery Scan on subquery (cost=...) + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.run_at + -> GroupAggregate (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Sort (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Index Scan using idx_delayed_jobs_live on public.delayed_jobs (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Index Cond: (delayed_jobs.run_at <= '2025-11-10 17:20:13'::timestamp without time zone) +--- SNAP snapshots["[legacy index] runs the expected postgresql queries with the expected plans 1"] = <<-SNAP @@ -244,6 +277,40 @@ GroupAggregate (cost=...) -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) +-- QUERIES FOR `max_age_by_name`: +--------------------------------- +SELECT MIN(run_at) AS run_at, + TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, + CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at + FROM \"delayed_jobs\" + WHERE \"delayed_jobs\".\"failed_at\" IS NULL + AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" + +GroupAggregate (cost=...) + Output: min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + -> Sort (cost=...) + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.run_at + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + -> Subquery Scan on subquery (cost=...) + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.run_at + -> GroupAggregate (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Sort (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Index Scan using delayed_jobs_priority on public.delayed_jobs (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Index Cond: (delayed_jobs.run_at <= '2025-11-10 17:20:13'::timestamp without time zone) + Filter: (delayed_jobs.failed_at IS NULL) +--- SNAP snapshots["runs the expected sqlite3 queries with the expected plans 1"] = <<-SNAP @@ -333,6 +400,27 @@ USE TEMP B-TREE FOR GROUP BY -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) +-- QUERIES FOR `max_age_by_name`: +--------------------------------- +SELECT MIN(run_at) AS run_at, + CURRENT_TIMESTAMP AS db_now_utc, + CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at + FROM \"delayed_jobs\" + WHERE \"delayed_jobs\".\"failed_at\" IS NULL + AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" + +CO-ROUTINE subquery +SCAN delayed_jobs USING INDEX idx_delayed_jobs_live +USE TEMP B-TREE FOR GROUP BY +SCAN subquery +USE TEMP B-TREE FOR GROUP BY +--- SNAP snapshots["[legacy index] runs the expected sqlite3 queries with the expected plans 1"] = <<-SNAP @@ -423,6 +511,27 @@ USE TEMP B-TREE FOR GROUP BY -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) +-- QUERIES FOR `max_age_by_name`: +--------------------------------- +SELECT MIN(run_at) AS run_at, + CURRENT_TIMESTAMP AS db_now_utc, + CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at + FROM \"delayed_jobs\" + WHERE \"delayed_jobs\".\"failed_at\" IS NULL + AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" + +CO-ROUTINE subquery +SCAN delayed_jobs USING INDEX delayed_jobs_priority +USE TEMP B-TREE FOR GROUP BY +SCAN subquery +USE TEMP B-TREE FOR GROUP BY +--- SNAP snapshots["runs the expected mysql2 queries with the expected plans 1"] = <<-SNAP diff --git a/spec/delayed/monitor_spec.rb b/spec/delayed/monitor_spec.rb index 4e13fa9..4c41b75 100644 --- a/spec/delayed/monitor_spec.rb +++ b/spec/delayed/monitor_spec.rb @@ -5,6 +5,13 @@ described_class.sleep_delay = 0 end + def emitted_event_names(pattern, &block) + events = [] + callback = ->(name, *) { events << name } + ActiveSupport::Notifications.subscribed(callback, pattern, &block) + events + end + let(:default_payload) do { table: 'delayed_jobs', @@ -83,6 +90,10 @@ .and emit_notification("delayed.job.alert_age_percent").with_payload(default_payload.merge(priority: 'reporting')).approximately.with_value(0) end + it 'does not emit max_age_by_name when no jobs are present' do + expect(emitted_event_names("delayed.job.max_age_by_name") { subject.run! }).to be_empty + end + context 'when named priorities are customized' do around do |example| Delayed::Priority.names = { high: 0, low: 7 } @@ -201,6 +212,67 @@ .and emit_notification("delayed.job.max_age").with_payload(p30_payload.merge(queue: 'banana')).approximately.with_value(4.hours) end + it 'emits max_age_by_name grouped by job name' do + expect { subject.run! } + .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'SimpleJob')).approximately.with_value(30.seconds) + .and emit_notification("delayed.job.max_age_by_name").with_payload(p10_payload.merge(name: 'SimpleJob')).approximately.with_value(2.minutes) + .and emit_notification("delayed.job.max_age_by_name").with_payload(p20_payload.merge(name: 'SimpleJob')).approximately.with_value(1.hour) + .and emit_notification("delayed.job.max_age_by_name").with_payload(p30_payload.merge(name: 'SimpleJob')).approximately.with_value(6.hours) + .and emit_notification("delayed.job.max_age_by_name").with_payload(p30_payload.merge(queue: 'banana', name: 'SimpleJob')).approximately.with_value(4.hours) + end + + context 'when multiple job names share a priority and queue' do + let!(:other_named_job) { Delayed::Job.create! p0_attributes.merge(name: 'OtherJob', run_at: now - 10.minutes) } + + it 'emits a separate max_age_by_name series per name' do + expect { subject.run! } + .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'SimpleJob')).approximately.with_value(30.seconds) + .and emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'OtherJob')).approximately.with_value(10.minutes) + end + end + + context 'when emit_max_age_by_name is disabled' do + around do |example| + described_class.emit_max_age_by_name = false + example.run + ensure + described_class.emit_max_age_by_name = true + end + + it 'does not emit max_age_by_name' do + expect(emitted_event_names("delayed.job.max_age_by_name") { subject.run! }).to be_empty + end + end + + context 'when the delayed_jobs table has no name column' do + before do + allow(Delayed::Job).to receive(:name_assignable?).and_return(false) + end + + it 'does not emit max_age_by_name' do + expect(emitted_event_names("delayed.job.max_age_by_name") { subject.run! }).to be_empty + end + end + + context 'when a job predates the name column' do + around do |example| + ValidateRunAtAndNameNotNull.migrate(:down) + AddRunAtAndNameNotNullCheck.migrate(:down) + example.run + ensure + Delayed::Job.delete_all + AddRunAtAndNameNotNullCheck.migrate(:up) + ValidateRunAtAndNameNotNull.migrate(:up) + end + + let!(:unnamed_job) { Delayed::Job.create! p0_attributes.merge(name: nil, run_at: now - 10.minutes) } + + it "emits max_age_by_name under the name 'unknown'" do + expect { subject.run! } + .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'unknown')).approximately.with_value(10.minutes) + end + end + context 'when named priorities are customized' do around do |example| Delayed::Priority.names = { high: 0, low: 20 } @@ -369,6 +441,11 @@ .and emit_notification("delayed.job.alert_age_percent").with_payload(payload).approximately.with_value(0) end + it 'excludes the locked job from max_age_by_name' do + expect { subject.run! } + .to emit_notification("delayed.job.max_age_by_name").with_payload(payload.merge(name: 'SimpleJob')).approximately.with_value(0) + end + context 'and a workable job is also present in the same group' do # The workable job's run_at is newer than the locked job's, so max_age must # track the workable job (30s), not the locked job (1 hour). @@ -405,7 +482,7 @@ end def query_descriptions - described_class::METRICS.each do |metric| + (described_class::METRICS + %w(max_age_by_name)).each do |metric| queries << "-- QUERIES FOR `#{metric}`:" queries << "---------------------------------" monitor.query_for(metric) From 1e8829e65f92855ddbe0cb112cb621830682fec2 Mon Sep 17 00:00:00 2001 From: James Boyer Date: Thu, 30 Jul 2026 12:08:01 -0400 Subject: [PATCH 2/4] Enabling generic metric by attribute configuration --- README.md | 41 ++++++- lib/delayed/monitor.rb | 111 ++++++++++-------- .../__snapshots__/monitor_spec.rb.snap | 60 +++++++--- spec/delayed/monitor_spec.rb | 66 ++++++++++- 4 files changed, 202 insertions(+), 76 deletions(-) diff --git a/README.md b/README.md index 426c6b8..6c8a36b 100644 --- a/README.md +++ b/README.md @@ -420,14 +420,47 @@ This answers "_which_ job is stuck?" when `delayed.job.max_age` alerts. A few be (priority, queue, name) series is only emitted while matching jobs are present in the queue. - Jobs enqueued before the `name` column existed (i.e. mid-upgrade) are reported under the name `"unknown"`. -- Because this metric's cardinality scales with the number of distinct job names, it can be - disabled if that's a concern for your metrics provider: + +These by-attribute pairings are driven by `Delayed::Monitor.metrics_by_attribute`, which defaults +to `{ max_age: %i(name) }`. Any metric can be paired with any groupable column — including columns +your application adds to the jobs table (populated at enqueue time) — and each pairing emits its +own **delayed.job.<metric>_by_<column>** event, grouped by priority name, queue name, +and that column's value (`'unknown'` when NULL). For example, with a per-record `owner` column: ```ruby -Delayed::Monitor.emit_max_age_by_name = false +Delayed::Monitor.metrics_by_attribute = { + max_age: %i(name owner), + failed_count: %i(owner), +} ``` -All of these events (including `max_age_by_name`) may be subscribed to via a single regular +The above would result in the following monitors: + +- 'delayed.job.max_age_by_name' +- 'delayed.job.max_age_by_owner' +- 'delayed.job.failed_count_by_owner' + +Configured columns that do not (yet) exist on the table are skipped. For example: + +```ruby +Delayed::Monitor.metrics_by_attribute = { + max_age: %i(owner non_existent_attribute), + failed_count: %i(owner), +} +``` + +The above would result in the following monitors based on attribute column availability: + +- 'delayed.job.max_age_by_owner' +- 'delayed.job.failed_count_by_owner' + +Metric names, on the other hand, are validated when the config is assigned — pairing a metric that +does not exist (e.g. `min_age`) raises an `ArgumentError` at boot, rather than failing later in the +monitor process. + +Each pairing's cardinality scales with the number of distinct values in the column, so if series cardinality is a concern for your metrics provider, by-attribute metrics can be disabled entirely by setting the config to `{}`. + +All of these events (including the `*_by_*` pairings) may be subscribed to via a single regular expression (again, in your application config or in an initializer): ```ruby diff --git a/lib/delayed/monitor.rb b/lib/delayed/monitor.rb index c1d74e9..6827f47 100644 --- a/lib/delayed/monitor.rb +++ b/lib/delayed/monitor.rb @@ -16,7 +16,20 @@ class Monitor ).freeze cattr_accessor :sleep_delay, instance_writer: false, default: 60 - cattr_accessor :emit_max_age_by_name, instance_writer: false, default: true + + def self.metrics_by_attribute + @metrics_by_attribute ||= { max_age: %i(name) }.freeze + end + + def self.metrics_by_attribute=(value) + unknown_metrics = value.keys.map(&:to_s) - METRICS + if unknown_metrics.any? + raise ArgumentError, "Unknown metrics in metrics_by_attribute: #{unknown_metrics.join(', ')}. " \ + "Valid metrics: #{METRICS.join(', ')}" + end + + @metrics_by_attribute = value.freeze + end def initialize @jobs = Job.group(:priority, :queue) @@ -28,13 +41,13 @@ def run! @memo = {} ActiveSupport::Notifications.instrument('delayed.monitor.run', default_tags) do METRICS.each { |metric| emit_metric!(metric) } - emit_metric_by_name!('max_age') if emit_max_age_by_name && Job.name_assignable? + assignable_metrics_by_attribute.each { |(metric, attribute)| emit_metric_by_attribute!(metric, attribute) } end interruptable_sleep(sleep_delay) end - def query_for(metric) - send(:"#{metric}_grouped") + def query_for(metric, *args) + send(:"#{metric}_grouped", *args) end def self.sql_now_in_utc @@ -77,12 +90,18 @@ def emit_metric!(metric) end end - def emit_metric_by_name!(metric) - query_for("#{metric}_by_name").each do |(priority, queue, name), value| - ActiveSupport::Notifications.instrument( - "delayed.job.#{metric}_by_name", - default_tags.merge(priority: Priority.new(priority).to_s, queue: queue, name: name || 'unknown', value: value), - ) + def emit_metric_by_attribute!(metric, attribute) + query_for(metric, attribute).each do |(priority, queue, label), value| + tags = { priority: Priority.new(priority).to_s, queue: queue, attribute => label || 'unknown', value: value } + ActiveSupport::Notifications.instrument("delayed.job.#{metric}_by_#{attribute}", default_tags.merge(tags)) + end + end + + def assignable_metrics_by_attribute + self.class.metrics_by_attribute.flat_map do |metric, attributes| + Array(attributes) + .select { |attribute| Job.column_names.include?(attribute.to_s) } + .map { |attribute| [metric.to_s, attribute.to_sym] } end end @@ -135,52 +154,48 @@ def as_expression(aggregate_function, aggregate_expression, column_name) "#{aggregate_function.to_s.upcase}(#{aggregate_expression}) AS #{column_name}" end - def count_grouped - failed_count_grouped.merge(live_count_grouped) { |_, l, f| l + f } - end - - def live_count_grouped - live_counts.transform_values(&:count) + def count_grouped(attribute = nil) + failed_count_grouped(attribute).merge(live_count_grouped(attribute)) { |_, l, f| l + f } end - def future_count_grouped - live_counts.transform_values(&:future_count) + def live_count_grouped(attribute = nil) + live_counts(attribute).transform_values(&:count) end - def erroring_count_grouped - live_counts.transform_values(&:erroring_count) + def future_count_grouped(attribute = nil) + live_counts(attribute).transform_values(&:future_count) end - def locked_count_grouped - pending_counts.transform_values(&:claimed_count) + def erroring_count_grouped(attribute = nil) + live_counts(attribute).transform_values(&:erroring_count) end - def failed_count_grouped - failed_counts.transform_values(&:count) + def locked_count_grouped(attribute = nil) + pending_counts(attribute).transform_values(&:claimed_count) end - def max_lock_age_grouped - pending_counts.transform_values { |j| time_ago(db_now(j), j.locked_at) } + def failed_count_grouped(attribute = nil) + failed_counts(attribute).transform_values(&:count) end - def max_age_grouped - pending_counts.transform_values { |j| time_ago(db_now(j), j.run_at) } + def max_lock_age_grouped(attribute = nil) + pending_counts(attribute).transform_values { |j| time_ago(db_now(j), j.locked_at) } end - def max_age_by_name_grouped - pending_run_ats_by_name.transform_values { |j| time_ago(db_now(j), j.run_at) } + def max_age_grouped(attribute = nil) + pending_counts(attribute).transform_values { |j| time_ago(db_now(j), j.run_at) } end - def alert_age_percent_grouped - pending_counts.each_with_object({}) do |((priority, queue), j), metrics| + def alert_age_percent_grouped(attribute = nil) + pending_counts(attribute).each_with_object({}) do |(key, j), metrics| max_age = time_ago(db_now(j), j.run_at) - alert_age = Priority.new(priority).alert_age - metrics[[priority, queue]] = [max_age / alert_age * 100, 100].min if alert_age + alert_age = Priority.new(key.first).alert_age + metrics[key] = [max_age / alert_age * 100, 100].min if alert_age end end - def workable_count_grouped - pending_counts.transform_values(&:claimable_count) + def workable_count_grouped(attribute = nil) + pending_counts(attribute).transform_values(&:claimable_count) end alias working_count_grouped locked_count_grouped @@ -193,19 +208,21 @@ def oldest_workable_job_grouped pending_counts.transform_values(&:run_at).compact end - def live_counts - @memo[:live_counts] ||= grouped_query( + def live_counts(attribute = nil) + @memo[[:live_counts, attribute]] ||= grouped_query( jobs.live, + extra_group_columns: [attribute].compact, count: [:count, '*'], future_count: [:sum, case_when(Job.future_clause.to_sql)], erroring_count: [:sum, case_when(Job.erroring_clause.to_sql)], ) end - def pending_counts - @memo[:pending_counts] ||= grouped_query( + def pending_counts(attribute = nil) + @memo[[:pending_counts, attribute]] ||= grouped_query( jobs.pending, include_db_time: true, + extra_group_columns: [attribute].compact, claimed_count: [:sum, case_when(Job.claimed_clause.to_sql)], claimable_count: [:sum, case_when(Job.claimable_clause.to_sql)], locked_at: [:min, case_when(Job.claimed_clause.to_sql, 'locked_at')], @@ -213,17 +230,9 @@ def pending_counts ) end - def pending_run_ats_by_name - @memo[:pending_run_ats_by_name] ||= grouped_query( - jobs.pending, - include_db_time: true, - extra_group_columns: %i(name), - run_at: [:min, case_when(Job.claimable_clause.to_sql, 'run_at')], - ) - end - - def failed_counts - @memo[:failed_counts] ||= grouped_query(jobs.failed, count: [:count, '*']) + def failed_counts(attribute = nil) + @memo[[:failed_counts, attribute]] ||= + grouped_query(jobs.failed, extra_group_columns: [attribute].compact, count: [:count, '*']) end def db_now(record) diff --git a/spec/delayed/__snapshots__/monitor_spec.rb.snap b/spec/delayed/__snapshots__/monitor_spec.rb.snap index 2722581..f6abb86 100644 --- a/spec/delayed/__snapshots__/monitor_spec.rb.snap +++ b/spec/delayed/__snapshots__/monitor_spec.rb.snap @@ -119,12 +119,19 @@ GroupAggregate (cost=...) -- (no new queries) -- QUERIES FOR `max_age_by_name`: --------------------------------- -SELECT MIN(run_at) AS run_at, +SELECT SUM(claimed_count) AS claimed_count, + SUM(claimable_count) AS claimable_count, + MIN(locked_at) AS locked_at, + MIN(run_at) AS run_at, TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, queue AS queue, name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, + MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, + MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL @@ -133,15 +140,15 @@ SELECT MIN(run_at) AS run_at, GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.run_at + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.run_at + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at @@ -279,12 +286,19 @@ GroupAggregate (cost=...) -- (no new queries) -- QUERIES FOR `max_age_by_name`: --------------------------------- -SELECT MIN(run_at) AS run_at, +SELECT SUM(claimed_count) AS claimed_count, + SUM(claimable_count) AS claimable_count, + MIN(locked_at) AS locked_at, + MIN(run_at) AS run_at, TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, queue AS queue, name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, + MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, + MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL @@ -293,15 +307,15 @@ SELECT MIN(run_at) AS run_at, GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.run_at + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.run_at + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at @@ -402,12 +416,19 @@ USE TEMP B-TREE FOR GROUP BY -- (no new queries) -- QUERIES FOR `max_age_by_name`: --------------------------------- -SELECT MIN(run_at) AS run_at, +SELECT SUM(claimed_count) AS claimed_count, + SUM(claimable_count) AS claimable_count, + MIN(locked_at) AS locked_at, + MIN(run_at) AS run_at, CURRENT_TIMESTAMP AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, queue AS queue, name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, + MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, + MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL @@ -513,12 +534,19 @@ USE TEMP B-TREE FOR GROUP BY -- (no new queries) -- QUERIES FOR `max_age_by_name`: --------------------------------- -SELECT MIN(run_at) AS run_at, +SELECT SUM(claimed_count) AS claimed_count, + SUM(claimable_count) AS claimable_count, + MIN(locked_at) AS locked_at, + MIN(run_at) AS run_at, CURRENT_TIMESTAMP AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, queue AS queue, name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL + OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, + MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, + MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL diff --git a/spec/delayed/monitor_spec.rb b/spec/delayed/monitor_spec.rb index 4c41b75..7804338 100644 --- a/spec/delayed/monitor_spec.rb +++ b/spec/delayed/monitor_spec.rb @@ -231,12 +231,12 @@ def emitted_event_names(pattern, &block) end end - context 'when emit_max_age_by_name is disabled' do + context 'when metrics_by_attribute is empty' do around do |example| - described_class.emit_max_age_by_name = false + described_class.metrics_by_attribute = {} example.run ensure - described_class.emit_max_age_by_name = true + described_class.metrics_by_attribute = { max_age: %i(name) } end it 'does not emit max_age_by_name' do @@ -246,7 +246,7 @@ def emitted_event_names(pattern, &block) context 'when the delayed_jobs table has no name column' do before do - allow(Delayed::Job).to receive(:name_assignable?).and_return(false) + allow(Delayed::Job).to receive(:column_names).and_return(Delayed::Job.column_names - ['name']) end it 'does not emit max_age_by_name' do @@ -273,6 +273,54 @@ def emitted_event_names(pattern, &block) end end + context 'when metrics_by_attribute pairs metrics with a custom column' do + around do |example| + Delayed::Job.connection.add_column :delayed_jobs, :owner, :string + Delayed::Job.reset_column_information + described_class.metrics_by_attribute = { max_age: %i(name owner), failed_count: %i(owner) } + example.run + ensure + described_class.metrics_by_attribute = { max_age: %i(name) } + Delayed::Job.connection.remove_column :delayed_jobs, :owner + Delayed::Job.reset_column_information + end + + let!(:team_a_job) { Delayed::Job.create! p0_attributes.merge(owner: 'team_a', run_at: now - 10.minutes) } + let!(:team_b_job) { Delayed::Job.create! p0_attributes.merge(owner: 'team_b', run_at: now - 20.minutes) } + + it "emits a max_age_by_owner series per owner, reporting ownerless rows as 'unknown'" do + expect { subject.run! } + .to emit_notification("delayed.job.max_age_by_owner").with_payload(p0_payload.merge(owner: 'team_a')).approximately.with_value(10.minutes) + .and emit_notification("delayed.job.max_age_by_owner").with_payload(p0_payload.merge(owner: 'team_b')).approximately.with_value(20.minutes) + .and emit_notification("delayed.job.max_age_by_owner").with_payload(p0_payload.merge(owner: 'unknown')).approximately.with_value(30.seconds) + end + + it 'emits other paired metrics by the same column' do + expect { subject.run! } + .to emit_notification("delayed.job.failed_count_by_owner").with_payload(p0_payload.merge(owner: 'unknown')).with_value(1) + end + + it 'emits max_age_by_name independently, without owner tags' do + expect { subject.run! } + .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'SimpleJob')).approximately.with_value(20.minutes) + end + end + + context 'when metrics_by_attribute names a column that does not exist' do + around do |example| + described_class.metrics_by_attribute = { max_age: %i(name owner) } + example.run + ensure + described_class.metrics_by_attribute = { max_age: %i(name) } + end + + it 'skips the missing column but still emits max_age_by_name' do + events = emitted_event_names(/delayed\.job\.max_age_by_/) { subject.run! } + expect(events).to include("delayed.job.max_age_by_name") + expect(events).not_to include("delayed.job.max_age_by_owner") + end + end + context 'when named priorities are customized' do around do |example| Delayed::Priority.names = { high: 0, low: 20 } @@ -461,6 +509,13 @@ def emitted_event_names(pattern, &block) end end + describe '.metrics_by_attribute=' do + it 'rejects unknown metrics at assignment time' do + expect { described_class.metrics_by_attribute = { min_age: %i(name) } } + .to raise_error(ArgumentError, /Unknown metrics in metrics_by_attribute: min_age/) + end + end + describe 'SQL' do let(:monitor) { described_class.new } let(:queries) { [] } @@ -485,7 +540,8 @@ def query_descriptions (described_class::METRICS + %w(max_age_by_name)).each do |metric| queries << "-- QUERIES FOR `#{metric}`:" queries << "---------------------------------" - monitor.query_for(metric) + base, attribute = metric.split('_by_') + monitor.query_for(base, *[attribute&.to_sym].compact) queries << "-- (no new queries)" unless queries.last == '---' end queries.dup.map { |query| query.try(:full_description) || query } From f29aa8a47190ad53a2c1bb509569f324bcd686a3 Mon Sep 17 00:00:00 2001 From: James Boyer Date: Fri, 31 Jul 2026 09:57:24 -0400 Subject: [PATCH 3/4] ensuring untagged zero values emit and just doing tag columns --- README.md | 99 ++--- lib/delayed/monitor.rb | 100 +++-- .../__snapshots__/monitor_spec.rb.snap | 365 ++++++------------ spec/delayed/monitor_spec.rb | 131 +++---- 4 files changed, 263 insertions(+), 432 deletions(-) diff --git a/README.md b/README.md index 6c8a36b..fb00f3f 100644 --- a/README.md +++ b/README.md @@ -391,14 +391,16 @@ QUEUE=tracking rake delayed:monitor QUEUES=mailers,tasks rake delayed:monitor ``` -The following events will be emitted, grouped by priority name (e.g. "interactive") and queue name, -and the metric's "`:value`" will be available in the event's payload. **This means that there will -be one value _per_ unique combination of queue & priority**, and totals must be computed via -downstream aggregation (e.g. as a StatsD "gauge" metric). +The following events will be emitted, grouped by priority name (e.g. "interactive"), queue name, +and the values of any configured `tag_columns`. By default, job `name` is included. The +metric's "`:value`" will be available in the event's payload. **This means that there will be one +value _per_ unique combination of queue, priority, and tag column values**, and totals must be +computed via downstream aggregation (e.g. as a StatsD "gauge" metric, summed or maxed by tag). - **delayed.job.count** - the total number of jobs - **delayed.job.future_count** - jobs where run_at is in the future - **delayed.job.working_count** - jobs that are currently being worked off (excludes failed jobs) +- **delayed.job.locked_count** - jobs that are currently locked by a worker (equivalent to working_count) - **delayed.job.workable_count** - jobs that are waiting to be worked off - **delayed.job.erroring_count** - jobs where attempts > 0 - **delayed.job.failed_count** - jobs where failed_at is not nil @@ -409,59 +411,66 @@ An additional _experimental_ metric is available, intended for use with applicat - **delayed.job.alert_age_percent** - the _percent_ to which the oldest job has reached the "age alert" threshold. (See the [Alerting Threshholds](#priority-based-alerting-threshholds) section above.) -If your jobs table has a `name` column (see [Database Setup](#database-setup)), one additional -event will be emitted, grouped by priority name, queue name, **and job name**: +By default, these events are also tagged with the job's `name` (when the jobs table has a `name` +column — see [Database Setup](#database-setup)) so that downstream aggregation can answer +"_which_ job is stuck?" when `delayed.job.max_age` alerts (e.g. `max by {queue, priority, name}` +in Datadog). -- **delayed.job.max_age_by_name** - the age of the oldest run_at value, per job name (excludes failed jobs) - -This answers "_which_ job is stuck?" when `delayed.job.max_age` alerts. A few behavioral notes: - -- Unlike the other metrics, it is not zero-filled: job names cannot be enumerated in advance, so a - (priority, queue, name) series is only emitted while matching jobs are present in the queue. -- Jobs enqueued before the `name` column existed (i.e. mid-upgrade) are reported under the name - `"unknown"`. - -These by-attribute pairings are driven by `Delayed::Monitor.metrics_by_attribute`, which defaults -to `{ max_age: %i(name) }`. Any metric can be paired with any groupable column — including columns -your application adds to the jobs table (populated at enqueue time) — and each pairing emits its -own **delayed.job.<metric>_by_<column>** event, grouped by priority name, queue name, -and that column's value (`'unknown'` when NULL). For example, with a per-record `owner` column: +The set of tagged columns is driven by `Delayed::Monitor.tag_columns`, which defaults to +`%i(name)` when the jobs table has a `name` column (and to `[]` otherwise). You can include +columns your application adds to the jobs table (populated at enqueue time). For example, if your +jobs table has an `owner` column you wish to also monitor: ```ruby -Delayed::Monitor.metrics_by_attribute = { - max_age: %i(name owner), - failed_count: %i(owner), -} +Delayed::Monitor.tag_columns = %i(name owner) ``` -The above would result in the following monitors: +A few behavioral notes: -- 'delayed.job.max_age_by_name' -- 'delayed.job.max_age_by_owner' -- 'delayed.job.failed_count_by_owner' +- Rows whose value was never populated for a tagged column are reported under the value `'unset'` + (e.g. jobs enqueued before the `name` column existed, mid-upgrade). +- Configured columns must exist on the jobs table: assigning a missing column to `tag_columns` + raises an `ArgumentError` immediately, rather than the column being silently skipped. Because + the assignment validates against the schema, setting `tag_columns` in an initializer requires a + database connection at boot, in every process that loads it. See the rollout steps below. +- Tag values cannot be enumerated in advance, so a tagged series is only emitted while matching + jobs are present. Separately, an untagged zero value is always emitted for every + (priority, queue) combination, so that each metric maintains a baseline series even when no + matching jobs are enqueued. For example, `delayed.job.count` with a single enqueued job would + emit the following series: -Configured columns that do not (yet) exist on the table are skipped. For example: + ```ruby + { priority: 'interactive', queue: 'default', name: 'SimpleJob', value: 1 } + { priority: 'interactive', queue: 'default', value: 0 } + { priority: 'user_visible', queue: 'default', value: 0 } + { priority: 'eventual', queue: 'default', value: 0 } + { priority: 'reporting', queue: 'default', value: 0 } + ``` +- Each column multiplies a metric's series cardinality by its number of distinct values (though in + practice a job's `name` tends to determine its `priority` and any ownership tags, making the + number of distinct job names the effective upper bound). If cardinality is a concern for your + metrics provider, tagging can be disabled entirely with `Delayed::Monitor.tag_columns = []`. -```ruby -Delayed::Monitor.metrics_by_attribute = { - max_age: %i(owner non_existent_attribute), - failed_count: %i(owner), -} -``` +#### Rolling out a new tag column -The above would result in the following monitors based on attribute column availability: +Because assignment fails loudly on a missing column, a new tag column should be rolled out in +three separate deploys, each fully released before the next begins: -- 'delayed.job.max_age_by_owner' -- 'delayed.job.failed_count_by_owner' +1. Migrate the column onto the jobs table (nullable — no backfill required). +2. Deploy the code that populates the column at enqueue time. +3. Add the column to `Delayed::Monitor.tag_columns` in an initializer, and deploy. -Metric names, on the other hand, are validated when the config is assigned — pairing a metric that -does not exist (e.g. `min_age`) raises an `ArgumentError` at boot, rather than failing later in the -monitor process. +Adding the column to `tag_columns` before the migration has run everywhere would raise at boot in +every process that loads the initializer. Jobs enqueued before step 2 will report under the +`'unset'` tag value until they are worked off (or backfilled). -Each pairing's cardinality scales with the number of distinct values in the column, so if series cardinality is a concern for your metrics provider, by-attribute metrics can be disabled entirely by setting the config to `{}`. +The default `name` tag needs no such rollout: it applies only when the jobs table already has a +`name` column, so a monitor running against an older schema simply emits untagged metrics until +the generated migrations (see [Database Setup](#database-setup)) have run and the monitor process +has restarted (the default is resolved once per process). -All of these events (including the `*_by_*` pairings) may be subscribed to via a single regular -expression (again, in your application config or in an initializer): +All of these events may be subscribed to via a single regular expression (again, in your +application config or in an initializer): ```ruby ActiveSupport::Notifications.subscribe(/delayed\.job\..*_(count|age|percent)/) do |*args| @@ -474,7 +483,7 @@ ActiveSupport::Notifications.subscribe(/delayed\.job\..*_(count|age|percent)/) d end ``` -Additionally, the monitor process with emit a **delayed.monitor.run** event with a duration +Additionally, the monitor process will emit a **delayed.monitor.run** event with a duration attached, so that you can monitor the time it takes to emit these aggregate metrics. ```ruby diff --git a/lib/delayed/monitor.rb b/lib/delayed/monitor.rb index 6827f47..d849063 100644 --- a/lib/delayed/monitor.rb +++ b/lib/delayed/monitor.rb @@ -17,18 +17,17 @@ class Monitor cattr_accessor :sleep_delay, instance_writer: false, default: 60 - def self.metrics_by_attribute - @metrics_by_attribute ||= { max_age: %i(name) }.freeze + def self.tag_columns + @tag_columns ||= (Job.column_names.include?('name') ? %i(name) : []).freeze end - def self.metrics_by_attribute=(value) - unknown_metrics = value.keys.map(&:to_s) - METRICS - if unknown_metrics.any? - raise ArgumentError, "Unknown metrics in metrics_by_attribute: #{unknown_metrics.join(', ')}. " \ - "Valid metrics: #{METRICS.join(', ')}" + def self.tag_columns=(columns) + if columns.any? { |column| Job.column_names.exclude?(column.to_s) } + raise ArgumentError, "Delayed::Monitor.tag_columns includes columns missing from #{Job.table_name}. " \ + "Available columns: #{Job.column_names.join(', ')}" end - @metrics_by_attribute = value.freeze + @tag_columns = columns.map(&:to_sym).freeze end def initialize @@ -41,13 +40,12 @@ def run! @memo = {} ActiveSupport::Notifications.instrument('delayed.monitor.run', default_tags) do METRICS.each { |metric| emit_metric!(metric) } - assignable_metrics_by_attribute.each { |(metric, attribute)| emit_metric_by_attribute!(metric, attribute) } end interruptable_sleep(sleep_delay) end - def query_for(metric, *args) - send(:"#{metric}_grouped", *args) + def query_for(metric) + send(:"#{metric}_grouped") end def self.sql_now_in_utc @@ -82,29 +80,17 @@ def self.parse_utc_time(string) attr_reader :jobs def emit_metric!(metric) - query_for(metric).reverse_merge(default_results).each do |(priority, queue), value| + query_for(metric) + .merge!(default_results) { |_key, existing, _default| existing } + .each do |(priority, queue, *column_values), value| + tags = column_values.zip(self.class.tag_columns).to_h { |val, column| [column, val.nil? ? 'unset' : val] } ActiveSupport::Notifications.instrument( "delayed.job.#{metric}", - default_tags.merge(priority: Priority.new(priority).to_s, queue: queue, value: value), + default_tags.merge(priority: Priority.new(priority).to_s, queue: queue, **tags, value: value), ) end end - def emit_metric_by_attribute!(metric, attribute) - query_for(metric, attribute).each do |(priority, queue, label), value| - tags = { priority: Priority.new(priority).to_s, queue: queue, attribute => label || 'unknown', value: value } - ActiveSupport::Notifications.instrument("delayed.job.#{metric}_by_#{attribute}", default_tags.merge(tags)) - end - end - - def assignable_metrics_by_attribute - self.class.metrics_by_attribute.flat_map do |metric, attributes| - Array(attributes) - .select { |attribute| Job.column_names.include?(attribute.to_s) } - .map { |attribute| [metric.to_s, attribute.to_sym] } - end - end - def default_results @default_results ||= Priority.names.values.flat_map { |priority| (Worker.queues.presence || [Worker.default_queue_name]).map do |queue| @@ -154,48 +140,48 @@ def as_expression(aggregate_function, aggregate_expression, column_name) "#{aggregate_function.to_s.upcase}(#{aggregate_expression}) AS #{column_name}" end - def count_grouped(attribute = nil) - failed_count_grouped(attribute).merge(live_count_grouped(attribute)) { |_, l, f| l + f } + def count_grouped + failed_count_grouped.merge(live_count_grouped) { |_, l, f| l + f } end - def live_count_grouped(attribute = nil) - live_counts(attribute).transform_values(&:count) + def live_count_grouped + live_counts.transform_values(&:count) end - def future_count_grouped(attribute = nil) - live_counts(attribute).transform_values(&:future_count) + def future_count_grouped + live_counts.transform_values(&:future_count) end - def erroring_count_grouped(attribute = nil) - live_counts(attribute).transform_values(&:erroring_count) + def erroring_count_grouped + live_counts.transform_values(&:erroring_count) end - def locked_count_grouped(attribute = nil) - pending_counts(attribute).transform_values(&:claimed_count) + def locked_count_grouped + pending_counts.transform_values(&:claimed_count) end - def failed_count_grouped(attribute = nil) - failed_counts(attribute).transform_values(&:count) + def failed_count_grouped + failed_counts.transform_values(&:count) end - def max_lock_age_grouped(attribute = nil) - pending_counts(attribute).transform_values { |j| time_ago(db_now(j), j.locked_at) } + def max_lock_age_grouped + pending_counts.transform_values { |j| time_ago(db_now(j), j.locked_at) } end - def max_age_grouped(attribute = nil) - pending_counts(attribute).transform_values { |j| time_ago(db_now(j), j.run_at) } + def max_age_grouped + pending_counts.transform_values { |j| time_ago(db_now(j), j.run_at) } end - def alert_age_percent_grouped(attribute = nil) - pending_counts(attribute).each_with_object({}) do |(key, j), metrics| + def alert_age_percent_grouped + pending_counts.each_with_object({}) do |(key, j), metrics| max_age = time_ago(db_now(j), j.run_at) alert_age = Priority.new(key.first).alert_age metrics[key] = [max_age / alert_age * 100, 100].min if alert_age end end - def workable_count_grouped(attribute = nil) - pending_counts(attribute).transform_values(&:claimable_count) + def workable_count_grouped + pending_counts.transform_values(&:claimable_count) end alias working_count_grouped locked_count_grouped @@ -208,21 +194,21 @@ def oldest_workable_job_grouped pending_counts.transform_values(&:run_at).compact end - def live_counts(attribute = nil) - @memo[[:live_counts, attribute]] ||= grouped_query( + def live_counts + @memo[:live_counts] ||= grouped_query( jobs.live, - extra_group_columns: [attribute].compact, + extra_group_columns: self.class.tag_columns, count: [:count, '*'], future_count: [:sum, case_when(Job.future_clause.to_sql)], erroring_count: [:sum, case_when(Job.erroring_clause.to_sql)], ) end - def pending_counts(attribute = nil) - @memo[[:pending_counts, attribute]] ||= grouped_query( + def pending_counts + @memo[:pending_counts] ||= grouped_query( jobs.pending, include_db_time: true, - extra_group_columns: [attribute].compact, + extra_group_columns: self.class.tag_columns, claimed_count: [:sum, case_when(Job.claimed_clause.to_sql)], claimable_count: [:sum, case_when(Job.claimable_clause.to_sql)], locked_at: [:min, case_when(Job.claimed_clause.to_sql, 'locked_at')], @@ -230,9 +216,9 @@ def pending_counts(attribute = nil) ) end - def failed_counts(attribute = nil) - @memo[[:failed_counts, attribute]] ||= - grouped_query(jobs.failed, extra_group_columns: [attribute].compact, count: [:count, '*']) + def failed_counts + @memo[:failed_counts] ||= + grouped_query(jobs.failed, extra_group_columns: self.class.tag_columns, count: [:count, '*']) end def db_now(record) diff --git a/spec/delayed/__snapshots__/monitor_spec.rb.snap b/spec/delayed/__snapshots__/monitor_spec.rb.snap index f6abb86..b20a2cf 100644 --- a/spec/delayed/__snapshots__/monitor_spec.rb.snap +++ b/spec/delayed/__snapshots__/monitor_spec.rb.snap @@ -3,56 +3,61 @@ snapshots["runs the expected postgresql queries with the expected plans 1"] = << --------------------------------- SELECT SUM(count) AS count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NOT NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: sum(subquery.count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: sum(subquery.count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.count - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.count + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.count + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.count -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, count(*) - Group Key: delayed_jobs.priority, delayed_jobs.queue - -> Index Only Scan using idx_delayed_jobs_failed on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, count(*) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Sort (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Index Scan using idx_delayed_jobs_failed on public.delayed_jobs (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name --- SELECT SUM(count) AS count, SUM(future_count) AS future_count, SUM(erroring_count) AS erroring_count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count, SUM(CASE WHEN \"delayed_jobs\".\"run_at\" > '2025-11-10 17:20:13' THEN 1 ELSE 0 END) AS future_count, SUM(CASE WHEN \"delayed_jobs\".\"attempts\" > 0 THEN 1 ELSE 0 END) AS erroring_count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: sum(subquery.count), sum(subquery.future_count), sum(subquery.erroring_count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: sum(subquery.count), sum(subquery.future_count), sum(subquery.erroring_count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.count, subquery.future_count, subquery.erroring_count - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.count, subquery.future_count, subquery.erroring_count + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.count, subquery.future_count, subquery.erroring_count + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.count, subquery.future_count, subquery.erroring_count -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, count(*), sum(CASE WHEN (delayed_jobs.run_at > '2025-11-10 17:20:13'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN (delayed_jobs.attempts > 0) THEN 1 ELSE 0 END) - Group Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, count(*), sum(CASE WHEN (delayed_jobs.run_at > '2025-11-10 17:20:13'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN (delayed_jobs.attempts > 0) THEN 1 ELSE 0 END) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.run_at, delayed_jobs.attempts - Sort Key: delayed_jobs.priority, delayed_jobs.queue - -> Index Only Scan using idx_delayed_jobs_live on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.run_at, delayed_jobs.attempts + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.run_at, delayed_jobs.attempts + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name + -> Index Scan using idx_delayed_jobs_live on public.delayed_jobs (cost=...) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.run_at, delayed_jobs.attempts --- -- QUERIES FOR `future_count`: --------------------------------- @@ -65,8 +70,9 @@ SELECT SUM(claimed_count) AS claimed_count, MIN(run_at) AS run_at, TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, @@ -75,25 +81,25 @@ SELECT SUM(claimed_count) AS claimed_count, FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) - Group Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.locked_at, delayed_jobs.run_at - Sort Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Index Scan using idx_delayed_jobs_live on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.locked_at, delayed_jobs.run_at + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at Index Cond: (delayed_jobs.run_at <= '2025-11-10 17:20:13'::timestamp without time zone) --- -- QUERIES FOR `erroring_count`: @@ -117,107 +123,69 @@ GroupAggregate (cost=...) -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) --- QUERIES FOR `max_age_by_name`: +SNAP + +snapshots["[legacy index] runs the expected postgresql queries with the expected plans 1"] = <<-SNAP +-- QUERIES FOR `count`: --------------------------------- -SELECT SUM(claimed_count) AS claimed_count, - SUM(claimable_count) AS claimable_count, - MIN(locked_at) AS locked_at, - MIN(run_at) AS run_at, - TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, +SELECT SUM(count) AS count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, queue AS queue, name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, - SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, - MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, - MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count FROM \"delayed_jobs\" - WHERE \"delayed_jobs\".\"failed_at\" IS NULL - AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' + WHERE \"delayed_jobs\".\"failed_at\" IS NOT NULL GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Output: sum(subquery.count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.count Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.count -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, count(*) Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name - -> Index Scan using idx_delayed_jobs_live on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at - Index Cond: (delayed_jobs.run_at <= '2025-11-10 17:20:13'::timestamp without time zone) ---- -SNAP - -snapshots["[legacy index] runs the expected postgresql queries with the expected plans 1"] = <<-SNAP --- QUERIES FOR `count`: ---------------------------------- -SELECT SUM(count) AS count, - CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count - FROM \"delayed_jobs\" - WHERE \"delayed_jobs\".\"failed_at\" IS NOT NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" - -GroupAggregate (cost=...) - Output: sum(subquery.count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.count - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.count - -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, count(*) - Group Key: delayed_jobs.priority, delayed_jobs.queue - -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue - Sort Key: delayed_jobs.priority, delayed_jobs.queue -> Index Scan using delayed_jobs_priority on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name Filter: (delayed_jobs.failed_at IS NOT NULL) --- SELECT SUM(count) AS count, SUM(future_count) AS future_count, SUM(erroring_count) AS erroring_count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count, SUM(CASE WHEN \"delayed_jobs\".\"run_at\" > '2025-11-10 17:20:13' THEN 1 ELSE 0 END) AS future_count, SUM(CASE WHEN \"delayed_jobs\".\"attempts\" > 0 THEN 1 ELSE 0 END) AS erroring_count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: sum(subquery.count), sum(subquery.future_count), sum(subquery.erroring_count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: sum(subquery.count), sum(subquery.future_count), sum(subquery.erroring_count), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.count, subquery.future_count, subquery.erroring_count - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.count, subquery.future_count, subquery.erroring_count + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.count, subquery.future_count, subquery.erroring_count + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.count, subquery.future_count, subquery.erroring_count -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, count(*), sum(CASE WHEN (delayed_jobs.run_at > '2025-11-10 17:20:13'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN (delayed_jobs.attempts > 0) THEN 1 ELSE 0 END) - Group Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, count(*), sum(CASE WHEN (delayed_jobs.run_at > '2025-11-10 17:20:13'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN (delayed_jobs.attempts > 0) THEN 1 ELSE 0 END) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.run_at, delayed_jobs.attempts - Sort Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.run_at, delayed_jobs.attempts + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Index Scan using delayed_jobs_priority on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.run_at, delayed_jobs.attempts + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.run_at, delayed_jobs.attempts Filter: (delayed_jobs.failed_at IS NULL) --- -- QUERIES FOR `future_count`: @@ -231,8 +199,9 @@ SELECT SUM(claimed_count) AS claimed_count, MIN(run_at) AS run_at, TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, @@ -241,25 +210,25 @@ SELECT SUM(claimed_count) AS claimed_count, FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" GroupAggregate (cost=...) - Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name + Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue + Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at + Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at + Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) - Group Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) + Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.locked_at, delayed_jobs.run_at - Sort Key: delayed_jobs.priority, delayed_jobs.queue + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at + Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name -> Index Scan using delayed_jobs_priority on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.locked_at, delayed_jobs.run_at + Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at Index Cond: (delayed_jobs.run_at <= '2025-11-10 17:20:13'::timestamp without time zone) Filter: (delayed_jobs.failed_at IS NULL) --- @@ -284,47 +253,6 @@ GroupAggregate (cost=...) -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) --- QUERIES FOR `max_age_by_name`: ---------------------------------- -SELECT SUM(claimed_count) AS claimed_count, - SUM(claimable_count) AS claimable_count, - MIN(locked_at) AS locked_at, - MIN(run_at) AS run_at, - TIMEZONE('UTC', STATEMENT_TIMESTAMP()) AS db_now_utc, - CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue, - name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, - SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, - MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, - MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at - FROM \"delayed_jobs\" - WHERE \"delayed_jobs\".\"failed_at\" IS NULL - AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" - -GroupAggregate (cost=...) - Output: sum(subquery.claimed_count), sum(subquery.claimable_count), min(subquery.locked_at), min(subquery.run_at), timezone('UTC'::text, statement_timestamp()), (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name - Group Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name - -> Sort (cost=...) - Output: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at - Sort Key: (CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END), subquery.queue, subquery.name - -> Subquery Scan on subquery (cost=...) - Output: CASE WHEN (subquery.priority < 10) THEN 0 WHEN (subquery.priority < 20) THEN 10 WHEN (subquery.priority < 30) THEN 20 WHEN (subquery.priority >= 30) THEN 30 ELSE NULL::integer END, subquery.queue, subquery.name, subquery.claimed_count, subquery.claimable_count, subquery.locked_at, subquery.run_at - -> GroupAggregate (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, sum(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN 1 ELSE 0 END), sum(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN 1 ELSE 0 END), min(CASE WHEN (delayed_jobs.locked_at >= '2025-11-10 16:59:43'::timestamp without time zone) THEN delayed_jobs.locked_at ELSE NULL::timestamp without time zone END), min(CASE WHEN ((delayed_jobs.locked_at IS NULL) OR (delayed_jobs.locked_at < '2025-11-10 16:59:43'::timestamp without time zone)) THEN delayed_jobs.run_at ELSE NULL::timestamp without time zone END) - Group Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name - -> Sort (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at - Sort Key: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name - -> Index Scan using delayed_jobs_priority on public.delayed_jobs (cost=...) - Output: delayed_jobs.priority, delayed_jobs.queue, delayed_jobs.name, delayed_jobs.locked_at, delayed_jobs.run_at - Index Cond: (delayed_jobs.run_at <= '2025-11-10 17:20:13'::timestamp without time zone) - Filter: (delayed_jobs.failed_at IS NULL) ---- SNAP snapshots["runs the expected sqlite3 queries with the expected plans 1"] = <<-SNAP @@ -332,15 +260,17 @@ snapshots["runs the expected sqlite3 queries with the expected plans 1"] = <<-SN --------------------------------- SELECT SUM(count) AS count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NOT NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" CO-ROUTINE subquery SCAN delayed_jobs USING INDEX idx_delayed_jobs_failed +USE TEMP B-TREE FOR GROUP BY SCAN subquery USE TEMP B-TREE FOR GROUP BY --- @@ -348,14 +278,15 @@ SELECT SUM(count) AS count, SUM(future_count) AS future_count, SUM(erroring_count) AS erroring_count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count, SUM(CASE WHEN \"delayed_jobs\".\"run_at\" > '2025-11-10 17:20:13' THEN 1 ELSE 0 END) AS future_count, SUM(CASE WHEN \"delayed_jobs\".\"attempts\" > 0 THEN 1 ELSE 0 END) AS erroring_count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" CO-ROUTINE subquery SCAN delayed_jobs USING INDEX idx_delayed_jobs_live @@ -374,8 +305,9 @@ SELECT SUM(claimed_count) AS claimed_count, MIN(run_at) AS run_at, CURRENT_TIMESTAMP AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, @@ -384,8 +316,8 @@ SELECT SUM(claimed_count) AS claimed_count, FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" CO-ROUTINE subquery SCAN delayed_jobs USING INDEX idx_delayed_jobs_live @@ -414,34 +346,6 @@ USE TEMP B-TREE FOR GROUP BY -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) --- QUERIES FOR `max_age_by_name`: ---------------------------------- -SELECT SUM(claimed_count) AS claimed_count, - SUM(claimable_count) AS claimable_count, - MIN(locked_at) AS locked_at, - MIN(run_at) AS run_at, - CURRENT_TIMESTAMP AS db_now_utc, - CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue, - name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, - SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, - MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, - MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at - FROM \"delayed_jobs\" - WHERE \"delayed_jobs\".\"failed_at\" IS NULL - AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" - -CO-ROUTINE subquery -SCAN delayed_jobs USING INDEX idx_delayed_jobs_live -USE TEMP B-TREE FOR GROUP BY -SCAN subquery -USE TEMP B-TREE FOR GROUP BY ---- SNAP snapshots["[legacy index] runs the expected sqlite3 queries with the expected plans 1"] = <<-SNAP @@ -449,12 +353,13 @@ snapshots["[legacy index] runs the expected sqlite3 queries with the expected pl --------------------------------- SELECT SUM(count) AS count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NOT NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" CO-ROUTINE subquery SCAN delayed_jobs USING INDEX delayed_jobs_priority @@ -466,14 +371,15 @@ SELECT SUM(count) AS count, SUM(future_count) AS future_count, SUM(erroring_count) AS erroring_count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", COUNT(*) AS count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", COUNT(*) AS count, SUM(CASE WHEN \"delayed_jobs\".\"run_at\" > '2025-11-10 17:20:13' THEN 1 ELSE 0 END) AS future_count, SUM(CASE WHEN \"delayed_jobs\".\"attempts\" > 0 THEN 1 ELSE 0 END) AS erroring_count FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" CO-ROUTINE subquery SCAN delayed_jobs USING INDEX delayed_jobs_priority @@ -492,8 +398,9 @@ SELECT SUM(claimed_count) AS claimed_count, MIN(run_at) AS run_at, CURRENT_TIMESTAMP AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + queue AS queue, + name AS name + FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, @@ -502,8 +409,8 @@ SELECT SUM(claimed_count) AS claimed_count, FROM \"delayed_jobs\" WHERE \"delayed_jobs\".\"failed_at\" IS NULL AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\" + GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" CO-ROUTINE subquery SCAN delayed_jobs USING INDEX delayed_jobs_priority @@ -532,34 +439,6 @@ USE TEMP B-TREE FOR GROUP BY -- QUERIES FOR `alert_age_percent`: --------------------------------- -- (no new queries) --- QUERIES FOR `max_age_by_name`: ---------------------------------- -SELECT SUM(claimed_count) AS claimed_count, - SUM(claimable_count) AS claimable_count, - MIN(locked_at) AS locked_at, - MIN(run_at) AS run_at, - CURRENT_TIMESTAMP AS db_now_utc, - CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue, - name AS name - FROM (SELECT \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\", SUM(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, - SUM(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, - MIN(CASE WHEN \"delayed_jobs\".\"locked_at\" >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, - MIN(CASE WHEN (\"delayed_jobs\".\"locked_at\" IS NULL - OR \"delayed_jobs\".\"locked_at\" < '2025-11-10 16:59:43') THEN run_at ELSE NULL END) AS run_at - FROM \"delayed_jobs\" - WHERE \"delayed_jobs\".\"failed_at\" IS NULL - AND \"delayed_jobs\".\"run_at\" <= '2025-11-10 17:20:13' - GROUP BY \"delayed_jobs\".\"priority\", \"delayed_jobs\".\"queue\", \"delayed_jobs\".\"name\") subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, \"queue\", \"name\" - -CO-ROUTINE subquery -SCAN delayed_jobs USING INDEX delayed_jobs_priority -USE TEMP B-TREE FOR GROUP BY -SCAN subquery -USE TEMP B-TREE FOR GROUP BY ---- SNAP snapshots["runs the expected mysql2 queries with the expected plans 1"] = <<-SNAP diff --git a/spec/delayed/monitor_spec.rb b/spec/delayed/monitor_spec.rb index 7804338..c9ac56d 100644 --- a/spec/delayed/monitor_spec.rb +++ b/spec/delayed/monitor_spec.rb @@ -5,13 +5,6 @@ described_class.sleep_delay = 0 end - def emitted_event_names(pattern, &block) - events = [] - callback = ->(name, *) { events << name } - ActiveSupport::Notifications.subscribed(callback, pattern, &block) - events - end - let(:default_payload) do { table: 'delayed_jobs', @@ -90,10 +83,6 @@ def emitted_event_names(pattern, &block) .and emit_notification("delayed.job.alert_age_percent").with_payload(default_payload.merge(priority: 'reporting')).approximately.with_value(0) end - it 'does not emit max_age_by_name when no jobs are present' do - expect(emitted_event_names("delayed.job.max_age_by_name") { subject.run! }).to be_empty - end - context 'when named priorities are customized' do around do |example| Delayed::Priority.names = { high: 0, low: 7 } @@ -143,10 +132,10 @@ def emitted_event_names(pattern, &block) let(:p10_attributes) { job_attributes.merge(priority: 13, locked_at: now - 1.day) } let(:p20_attributes) { job_attributes.merge(priority: 23, attempts: 1) } let(:p30_attributes) { job_attributes.merge(priority: 999, locked_at: now - 1.day) } - let(:p0_payload) { default_payload.merge(priority: 'interactive') } - let(:p10_payload) { default_payload.merge(priority: 'user_visible') } - let(:p20_payload) { default_payload.merge(priority: 'eventual') } - let(:p30_payload) { default_payload.merge(priority: 'reporting') } + let(:p0_payload) { default_payload.merge(priority: 'interactive', name: 'SimpleJob') } + let(:p10_payload) { default_payload.merge(priority: 'user_visible', name: 'SimpleJob') } + let(:p20_payload) { default_payload.merge(priority: 'eventual', name: 'SimpleJob') } + let(:p30_payload) { default_payload.merge(priority: 'reporting', name: 'SimpleJob') } let!(:p0_workable_job) { Delayed::Job.create! p0_attributes.merge(run_at: now - 30.seconds) } let!(:p0_failed_job) { Delayed::Job.create! p0_attributes.merge(failed_attributes) } let!(:p0_future_job) { Delayed::Job.create! p0_attributes.merge(run_at: now + 1.hour) } @@ -212,45 +201,43 @@ def emitted_event_names(pattern, &block) .and emit_notification("delayed.job.max_age").with_payload(p30_payload.merge(queue: 'banana')).approximately.with_value(4.hours) end - it 'emits max_age_by_name grouped by job name' do - expect { subject.run! } - .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'SimpleJob')).approximately.with_value(30.seconds) - .and emit_notification("delayed.job.max_age_by_name").with_payload(p10_payload.merge(name: 'SimpleJob')).approximately.with_value(2.minutes) - .and emit_notification("delayed.job.max_age_by_name").with_payload(p20_payload.merge(name: 'SimpleJob')).approximately.with_value(1.hour) - .and emit_notification("delayed.job.max_age_by_name").with_payload(p30_payload.merge(name: 'SimpleJob')).approximately.with_value(6.hours) - .and emit_notification("delayed.job.max_age_by_name").with_payload(p30_payload.merge(queue: 'banana', name: 'SimpleJob')).approximately.with_value(4.hours) - end - context 'when multiple job names share a priority and queue' do let!(:other_named_job) { Delayed::Job.create! p0_attributes.merge(name: 'OtherJob', run_at: now - 10.minutes) } - it 'emits a separate max_age_by_name series per name' do + it 'emits a separate series per name' do expect { subject.run! } - .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'SimpleJob')).approximately.with_value(30.seconds) - .and emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'OtherJob')).approximately.with_value(10.minutes) + .to emit_notification("delayed.job.max_age").with_payload(p0_payload).approximately.with_value(30.seconds) + .and emit_notification("delayed.job.max_age").with_payload(p0_payload.merge(name: 'OtherJob')).approximately.with_value(10.minutes) end end - context 'when metrics_by_attribute is empty' do + context 'when tag_columns is empty' do around do |example| - described_class.metrics_by_attribute = {} + described_class.tag_columns = [] example.run ensure - described_class.metrics_by_attribute = { max_age: %i(name) } + described_class.tag_columns = %i(name) end - it 'does not emit max_age_by_name' do - expect(emitted_event_names("delayed.job.max_age_by_name") { subject.run! }).to be_empty + it 'emits metrics without name tags' do + expect { subject.run! } + .to emit_notification("delayed.job.max_age").with_payload(p0_payload.except(:name)).approximately.with_value(30.seconds) end end context 'when the delayed_jobs table has no name column' do before do + described_class.instance_variable_set(:@tag_columns, nil) allow(Delayed::Job).to receive(:column_names).and_return(Delayed::Job.column_names - ['name']) end - it 'does not emit max_age_by_name' do - expect(emitted_event_names("delayed.job.max_age_by_name") { subject.run! }).to be_empty + after do + described_class.instance_variable_set(:@tag_columns, nil) + end + + it 'defaults to emitting metrics without name tags' do + expect { subject.run! } + .to emit_notification("delayed.job.max_age").with_payload(p0_payload.except(:name)).approximately.with_value(30.seconds) end end @@ -267,20 +254,20 @@ def emitted_event_names(pattern, &block) let!(:unnamed_job) { Delayed::Job.create! p0_attributes.merge(name: nil, run_at: now - 10.minutes) } - it "emits max_age_by_name under the name 'unknown'" do + it "emits metrics under the name 'unset'" do expect { subject.run! } - .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'unknown')).approximately.with_value(10.minutes) + .to emit_notification("delayed.job.max_age").with_payload(p0_payload.merge(name: 'unset')).approximately.with_value(10.minutes) end end - context 'when metrics_by_attribute pairs metrics with a custom column' do + context 'when tag_columns includes a custom column' do around do |example| Delayed::Job.connection.add_column :delayed_jobs, :owner, :string Delayed::Job.reset_column_information - described_class.metrics_by_attribute = { max_age: %i(name owner), failed_count: %i(owner) } + described_class.tag_columns = %i(name owner) example.run ensure - described_class.metrics_by_attribute = { max_age: %i(name) } + described_class.tag_columns = %i(name) Delayed::Job.connection.remove_column :delayed_jobs, :owner Delayed::Job.reset_column_information end @@ -288,36 +275,19 @@ def emitted_event_names(pattern, &block) let!(:team_a_job) { Delayed::Job.create! p0_attributes.merge(owner: 'team_a', run_at: now - 10.minutes) } let!(:team_b_job) { Delayed::Job.create! p0_attributes.merge(owner: 'team_b', run_at: now - 20.minutes) } - it "emits a max_age_by_owner series per owner, reporting ownerless rows as 'unknown'" do - expect { subject.run! } - .to emit_notification("delayed.job.max_age_by_owner").with_payload(p0_payload.merge(owner: 'team_a')).approximately.with_value(10.minutes) - .and emit_notification("delayed.job.max_age_by_owner").with_payload(p0_payload.merge(owner: 'team_b')).approximately.with_value(20.minutes) - .and emit_notification("delayed.job.max_age_by_owner").with_payload(p0_payload.merge(owner: 'unknown')).approximately.with_value(30.seconds) - end - - it 'emits other paired metrics by the same column' do + it "tags each series with the column's value, reporting NULLs as 'unset'" do expect { subject.run! } - .to emit_notification("delayed.job.failed_count_by_owner").with_payload(p0_payload.merge(owner: 'unknown')).with_value(1) - end - - it 'emits max_age_by_name independently, without owner tags' do - expect { subject.run! } - .to emit_notification("delayed.job.max_age_by_name").with_payload(p0_payload.merge(name: 'SimpleJob')).approximately.with_value(20.minutes) + .to emit_notification("delayed.job.max_age").with_payload(p0_payload.merge(owner: 'team_a')).approximately.with_value(10.minutes) + .and emit_notification("delayed.job.max_age").with_payload(p0_payload.merge(owner: 'team_b')).approximately.with_value(20.minutes) + .and emit_notification("delayed.job.max_age").with_payload(p0_payload.merge(owner: 'unset')).approximately.with_value(30.seconds) + .and emit_notification("delayed.job.failed_count").with_payload(p0_payload.merge(owner: 'unset')).with_value(1) end end - context 'when metrics_by_attribute names a column that does not exist' do - around do |example| - described_class.metrics_by_attribute = { max_age: %i(name owner) } - example.run - ensure - described_class.metrics_by_attribute = { max_age: %i(name) } - end - - it 'skips the missing column but still emits max_age_by_name' do - events = emitted_event_names(/delayed\.job\.max_age_by_/) { subject.run! } - expect(events).to include("delayed.job.max_age_by_name") - expect(events).not_to include("delayed.job.max_age_by_owner") + context 'when tag_columns names a column that does not exist' do + it 'raises loudly rather than skipping the column' do + expect { described_class.tag_columns = %i(name owner) } + .to raise_error(ArgumentError, /tag_columns includes columns missing from delayed_jobs\. Available columns: .*\bname\b/) end end @@ -328,8 +298,8 @@ def emitted_event_names(pattern, &block) ensure Delayed::Priority.names = nil end - let(:p0_payload) { default_payload.merge(priority: 'high') } - let(:p20_payload) { default_payload.merge(priority: 'low') } + let(:p0_payload) { default_payload.merge(priority: 'high', name: 'SimpleJob') } + let(:p20_payload) { default_payload.merge(priority: 'low', name: 'SimpleJob') } it 'emits the expected results for each metric' do expect { subject.run! } @@ -343,7 +313,7 @@ def emitted_event_names(pattern, &block) .and emit_notification("delayed.job.workable_count").with_payload(p0_payload).with_value(2) .and emit_notification("delayed.job.max_age").with_payload(p0_payload).approximately.with_value(2.minutes) .and emit_notification("delayed.job.max_lock_age").with_payload(p0_payload).approximately.with_value(7.minutes) - .and emit_notification("delayed.job.alert_age_percent").with_payload(p0_payload).approximately.with_value(0) + .and emit_notification("delayed.job.alert_age_percent").with_payload(p0_payload.except(:name)).approximately.with_value(0) .and emit_notification("delayed.job.count").with_payload(p20_payload).with_value(8) .and emit_notification("delayed.job.future_count").with_payload(p20_payload).with_value(2) .and emit_notification("delayed.job.locked_count").with_payload(p20_payload).with_value(2) @@ -353,7 +323,7 @@ def emitted_event_names(pattern, &block) .and emit_notification("delayed.job.workable_count").with_payload(p20_payload).with_value(2) .and emit_notification("delayed.job.max_age").with_payload(p20_payload).approximately.with_value(6.hours) .and emit_notification("delayed.job.max_lock_age").with_payload(p20_payload).approximately.with_value(11.minutes) - .and emit_notification("delayed.job.alert_age_percent").with_payload(p20_payload).approximately.with_value(0) + .and emit_notification("delayed.job.alert_age_percent").with_payload(p20_payload.except(:name)).approximately.with_value(0) .and emit_notification("delayed.job.workable_count").with_payload(p20_payload.merge(queue: 'banana')).with_value(1) .and emit_notification("delayed.job.max_age").with_payload(p20_payload.merge(queue: 'banana')).approximately.with_value(4.hours) end @@ -384,7 +354,7 @@ def emitted_event_names(pattern, &block) Delayed::Priority.names = nil Delayed::Worker.queues = [] end - let(:banana_payload) { default_payload.merge(queue: 'banana', priority: 'interactive') } + let(:banana_payload) { default_payload.merge(queue: 'banana', priority: 'interactive', name: 'SimpleJob') } let(:gram_payload) { default_payload.merge(queue: 'gram', priority: 'interactive') } it 'emits the expected results for each queue' do @@ -394,7 +364,7 @@ def emitted_event_names(pattern, &block) .and emit_notification("delayed.job.future_count").with_payload(banana_payload).with_value(0) .and emit_notification("delayed.job.locked_count").with_payload(banana_payload).with_value(0) .and emit_notification("delayed.job.erroring_count").with_payload(banana_payload).with_value(0) - .and emit_notification("delayed.job.failed_count").with_payload(banana_payload).with_value(0) + .and emit_notification("delayed.job.failed_count").with_payload(banana_payload.except(:name)).with_value(0) .and emit_notification("delayed.job.working_count").with_payload(banana_payload).with_value(0) .and emit_notification("delayed.job.workable_count").with_payload(banana_payload).with_value(1) .and emit_notification("delayed.job.max_age").with_payload(banana_payload).approximately.with_value(4.hours) @@ -466,7 +436,7 @@ def emitted_event_names(pattern, &block) end context 'when a job is locked (in-flight)' do - let(:payload) { default_payload.merge(priority: 'interactive') } + let(:payload) { default_payload.merge(priority: 'interactive', name: 'SimpleJob') } let(:base_attributes) do { priority: 0, @@ -489,11 +459,6 @@ def emitted_event_names(pattern, &block) .and emit_notification("delayed.job.alert_age_percent").with_payload(payload).approximately.with_value(0) end - it 'excludes the locked job from max_age_by_name' do - expect { subject.run! } - .to emit_notification("delayed.job.max_age_by_name").with_payload(payload.merge(name: 'SimpleJob')).approximately.with_value(0) - end - context 'and a workable job is also present in the same group' do # The workable job's run_at is newer than the locked job's, so max_age must # track the workable job (30s), not the locked job (1 hour). @@ -509,13 +474,6 @@ def emitted_event_names(pattern, &block) end end - describe '.metrics_by_attribute=' do - it 'rejects unknown metrics at assignment time' do - expect { described_class.metrics_by_attribute = { min_age: %i(name) } } - .to raise_error(ArgumentError, /Unknown metrics in metrics_by_attribute: min_age/) - end - end - describe 'SQL' do let(:monitor) { described_class.new } let(:queries) { [] } @@ -537,11 +495,10 @@ def emitted_event_names(pattern, &block) end def query_descriptions - (described_class::METRICS + %w(max_age_by_name)).each do |metric| + described_class::METRICS.each do |metric| queries << "-- QUERIES FOR `#{metric}`:" queries << "---------------------------------" - base, attribute = metric.split('_by_') - monitor.query_for(base, *[attribute&.to_sym].compact) + monitor.query_for(metric) queries << "-- (no new queries)" unless queries.last == '---' end queries.dup.map { |query| query.try(:full_description) || query } From 34c6105752908ca5b7873df77994c9aca23cce09 Mon Sep 17 00:00:00 2001 From: James Boyer Date: Fri, 31 Jul 2026 12:21:11 -0400 Subject: [PATCH 4/4] updating monitor_spec snapshot --- .../__snapshots__/monitor_spec.rb.snap | 62 ++++++++++--------- 1 file changed, 33 insertions(+), 29 deletions(-) diff --git a/spec/delayed/__snapshots__/monitor_spec.rb.snap b/spec/delayed/__snapshots__/monitor_spec.rb.snap index b20a2cf..b403da3 100644 --- a/spec/delayed/__snapshots__/monitor_spec.rb.snap +++ b/spec/delayed/__snapshots__/monitor_spec.rb.snap @@ -446,12 +446,13 @@ snapshots["runs the expected mysql2 queries with the expected plans 1"] = <<-SNA --------------------------------- SELECT SUM(count) AS count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, COUNT(*) AS count + queue AS queue, + name AS name + FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`, COUNT(*) AS count FROM `delayed_jobs` WHERE `delayed_jobs`.`failed_at` IS NOT NULL - GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`) subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue` + GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`) subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue`, `name` -> Table scan on -> Aggregate using temporary table @@ -460,20 +461,21 @@ SELECT SUM(count) AS count, -> Table scan on -> Aggregate using temporary table -> Filter: (delayed_jobs.failed_at is not null) (cost=...) - -> Covering index range scan on delayed_jobs using idx_delayed_jobs_live over (NULL < failed_at) (cost=...) + -> Table scan on delayed_jobs (cost=...) --- SELECT SUM(count) AS count, SUM(future_count) AS future_count, SUM(erroring_count) AS erroring_count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, COUNT(*) AS count, + queue AS queue, + name AS name + FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`, COUNT(*) AS count, SUM(CASE WHEN `delayed_jobs`.`run_at` > '2025-11-10 17:20:13' THEN 1 ELSE 0 END) AS future_count, SUM(CASE WHEN `delayed_jobs`.`attempts` > 0 THEN 1 ELSE 0 END) AS erroring_count FROM `delayed_jobs` WHERE `delayed_jobs`.`failed_at` IS NULL - GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`) subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue` + GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`) subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue`, `name` -> Table scan on -> Aggregate using temporary table @@ -481,8 +483,7 @@ SELECT SUM(count) AS count, -> Materialize (cost=...) -> Table scan on -> Aggregate using temporary table - -> Filter: (delayed_jobs.failed_at is null) (cost=...) - -> Covering index lookup on delayed_jobs using idx_delayed_jobs_live (failed_at = NULL) (cost=...) + -> Index lookup on delayed_jobs using idx_delayed_jobs_live (failed_at = NULL), with index condition: (delayed_jobs.failed_at is null) (cost=...) --- -- QUERIES FOR `future_count`: --------------------------------- @@ -495,8 +496,9 @@ SELECT SUM(claimed_count) AS claimed_count, MIN(run_at) AS run_at, UTC_TIMESTAMP() AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, SUM(CASE WHEN `delayed_jobs`.`locked_at` >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + queue AS queue, + name AS name + FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`, SUM(CASE WHEN `delayed_jobs`.`locked_at` >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, SUM(CASE WHEN (`delayed_jobs`.`locked_at` IS NULL OR `delayed_jobs`.`locked_at` < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, MIN(CASE WHEN `delayed_jobs`.`locked_at` >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, @@ -505,8 +507,8 @@ SELECT SUM(claimed_count) AS claimed_count, FROM `delayed_jobs` WHERE `delayed_jobs`.`failed_at` IS NULL AND `delayed_jobs`.`run_at` <= '2025-11-10 17:20:13' - GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`) subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue` + GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`) subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue`, `name` -> Table scan on -> Aggregate using temporary table @@ -514,8 +516,7 @@ SELECT SUM(claimed_count) AS claimed_count, -> Materialize (cost=...) -> Table scan on -> Aggregate using temporary table - -> Filter: ((delayed_jobs.failed_at is null) and (delayed_jobs.run_at <= TIMESTAMP'2025-11-10 17:20:13')) (cost=...) - -> Covering index lookup on delayed_jobs using idx_delayed_jobs_live (failed_at = NULL) (cost=...) + -> Index lookup on delayed_jobs using idx_delayed_jobs_live (failed_at = NULL), with index condition: ((delayed_jobs.failed_at is null) and (delayed_jobs.run_at <= TIMESTAMP'2025-11-10 17:20:13')) (cost=...) --- -- QUERIES FOR `erroring_count`: --------------------------------- @@ -545,12 +546,13 @@ snapshots["[legacy index] runs the expected mysql2 queries with the expected pla --------------------------------- SELECT SUM(count) AS count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, COUNT(*) AS count + queue AS queue, + name AS name + FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`, COUNT(*) AS count FROM `delayed_jobs` WHERE `delayed_jobs`.`failed_at` IS NOT NULL - GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`) subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue` + GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`) subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue`, `name` -> Table scan on -> Aggregate using temporary table @@ -565,14 +567,15 @@ SELECT SUM(count) AS count, SUM(future_count) AS future_count, SUM(erroring_count) AS erroring_count, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, COUNT(*) AS count, + queue AS queue, + name AS name + FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`, COUNT(*) AS count, SUM(CASE WHEN `delayed_jobs`.`run_at` > '2025-11-10 17:20:13' THEN 1 ELSE 0 END) AS future_count, SUM(CASE WHEN `delayed_jobs`.`attempts` > 0 THEN 1 ELSE 0 END) AS erroring_count FROM `delayed_jobs` WHERE `delayed_jobs`.`failed_at` IS NULL - GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`) subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue` + GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`) subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue`, `name` -> Table scan on -> Aggregate using temporary table @@ -594,8 +597,9 @@ SELECT SUM(claimed_count) AS claimed_count, MIN(run_at) AS run_at, UTC_TIMESTAMP() AS db_now_utc, CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END AS priority, - queue AS queue - FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, SUM(CASE WHEN `delayed_jobs`.`locked_at` >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, + queue AS queue, + name AS name + FROM (SELECT `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`, SUM(CASE WHEN `delayed_jobs`.`locked_at` >= '2025-11-10 16:59:43' THEN 1 ELSE 0 END) AS claimed_count, SUM(CASE WHEN (`delayed_jobs`.`locked_at` IS NULL OR `delayed_jobs`.`locked_at` < '2025-11-10 16:59:43') THEN 1 ELSE 0 END) AS claimable_count, MIN(CASE WHEN `delayed_jobs`.`locked_at` >= '2025-11-10 16:59:43' THEN locked_at ELSE NULL END) AS locked_at, @@ -604,8 +608,8 @@ SELECT SUM(claimed_count) AS claimed_count, FROM `delayed_jobs` WHERE `delayed_jobs`.`failed_at` IS NULL AND `delayed_jobs`.`run_at` <= '2025-11-10 17:20:13' - GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`) subquery - GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue` + GROUP BY `delayed_jobs`.`priority`, `delayed_jobs`.`queue`, `delayed_jobs`.`name`) subquery + GROUP BY CASE WHEN priority < 10 THEN 0 WHEN priority < 20 THEN 10 WHEN priority < 30 THEN 20 WHEN priority >= 30 THEN 30 END, `queue`, `name` -> Table scan on -> Aggregate using temporary table