From 7da1f8e702213c804940182be3ff7d4138462135 Mon Sep 17 00:00:00 2001 From: Alan Guo Xiang Tan Date: Fri, 25 Sep 2026 14:05:06 +0800 Subject: [PATCH 1/5] FEATURE: Count Pitchfork worker timeouts in Prometheus Operators cannot graph or alert on Pitchfork soft worker timeouts. This commit exposes `discourse_pitchfork_worker_timeouts_total` when core emits `web_worker_timeout`. Reporting uses a dedicated synchronous client with a one-second budget so observations can reach the collector before worker exit. Hard kills bypassing the callback are not counted. --- README.md | 6 +++ lib/reporter/worker_timeout.rb | 32 ++++++++++++ plugin.rb | 3 ++ spec/lib/reporter/worker_timeout_spec.rb | 63 ++++++++++++++++++++++++ 4 files changed, 104 insertions(+) create mode 100644 lib/reporter/worker_timeout.rb create mode 100644 spec/lib/reporter/worker_timeout_spec.rb diff --git a/README.md b/README.md index d6859ea..990dfce 100644 --- a/README.md +++ b/README.md @@ -6,4 +6,10 @@ The Discourse Prometheus plugin collects key metrics from Discourse and exposes The global reporter can pick custom metrics added by other Discourse plugins. The metric needs to define a collect method, and the `name`, `labels`, `description`, `value`, and `type` attributes. See an example [here](https://github.com/discourse/discourse-antivirus/pull/15). +## Worker timeouts + +`discourse_pitchfork_worker_timeouts_total` counts Pitchfork soft worker timeout events. It is exposed after the first event and resets when the collector restarts. Reporting is best effort, with a one-second budget before the worker exits. Hard kills that bypass the soft timeout callback are not counted. + +This metric requires a Discourse version that emits the `web_worker_timeout` event. For example, `rate(discourse_pitchfork_worker_timeouts_total[5m])` shows timeouts per second over five minutes. + For more information, please see: https://meta.discourse.org/t/prometheus-exporter-plugin-for-discourse/72666 diff --git a/lib/reporter/worker_timeout.rb b/lib/reporter/worker_timeout.rb new file mode 100644 index 0000000..59405f4 --- /dev/null +++ b/lib/reporter/worker_timeout.rb @@ -0,0 +1,32 @@ +# frozen_string_literal: true + +require "timeout" + +module DiscoursePrometheus::Reporter + class WorkerTimeout + def report + metric = DiscoursePrometheus::InternalMetric::Custom.new + metric.type = "Counter" + metric.name = "pitchfork_worker_timeouts_total" + metric.description = "Total number of Pitchfork soft worker timeouts" + metric.value = 1 + + client = + PrometheusExporter::Client.new( + host: "localhost", + port: GlobalSetting.prometheus_collector_port, + process_queue_once_and_stop: true, + ) + + Timeout.timeout(1) do + begin + client.send_json(metric.to_h) + ensure + client.stop + end + end + rescue => error + Rails.logger.warn("Failed to report worker timeout: #{error.message}") + end + end +end diff --git a/plugin.rb b/plugin.rb index 5191655..bc6a1c9 100644 --- a/plugin.rb +++ b/plugin.rb @@ -26,6 +26,7 @@ module ::DiscoursePrometheus require_relative("lib/reporter/global") require_relative("lib/reporter/web") require_relative("lib/reporter/image_processing") +require_relative("lib/reporter/worker_timeout") require_relative("lib/collector_demon") require_relative("lib/global_reporter_demon") @@ -53,6 +54,8 @@ module ::DiscoursePrometheus on(:image_processing_finished) { |payload| image_processing_reporter.report(payload) } + on(:web_worker_timeout) { DiscoursePrometheus::Reporter::WorkerTimeout.new.report } + register_demon_process(DiscoursePrometheus::CollectorDemon) register_demon_process(DiscoursePrometheus::GlobalReporterDemon) diff --git a/spec/lib/reporter/worker_timeout_spec.rb b/spec/lib/reporter/worker_timeout_spec.rb new file mode 100644 index 0000000..5eafd17 --- /dev/null +++ b/spec/lib/reporter/worker_timeout_spec.rb @@ -0,0 +1,63 @@ +# frozen_string_literal: true + +require "prometheus_exporter/server" +require_relative "../../../lib/collector" + +RSpec.describe DiscoursePrometheus::Reporter::WorkerTimeout do + describe "#report" do + it "delivers timeout events before their reporting processes exit" do + socket = TCPServer.new("127.0.0.1", 0) + port = socket.addr[1] + socket.close + global_setting :prometheus_collector_port, port + collector = DiscoursePrometheus::Collector.new + server = + PrometheusExporter::Server::WebServer.new( + port: port, + bind: "127.0.0.1", + collector: collector, + ) + runner = server.start + + 2.times do + pid = + fork do + DiscourseEvent.trigger(:web_worker_timeout) + Process.exit!(0) + end + _, status = Process.wait2(pid) + expect(status).to be_success + end + + wait_for(timeout: 5) do + collector.prometheus_metrics_text.include?("pitchfork_worker_timeouts_total 2") + end + expect(collector.prometheus_metrics_text).to include( + "pitchfork_worker_timeouts_total counter", + "pitchfork_worker_timeouts_total 2", + ) + ensure + server&.stop + runner&.join + socket&.close unless socket&.closed? + end + + it "returns when the collector is unavailable" do + socket = TCPServer.new("127.0.0.1", 0) + port = socket.addr[1] + socket.close + global_setting :prometheus_collector_port, port + + expect { described_class.new.report }.not_to raise_error + end + + it "bounds reporting when the collector connection stalls" do + allow(TCPSocket).to receive(:new) { sleep 30 } + started_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) + + described_class.new.report + + expect(Process.clock_gettime(Process::CLOCK_MONOTONIC) - started_at).to be < 3 + end + end +end From 79cf6e5c4d4691e234515965a562abd623f0a4a1 Mon Sep 17 00:00:00 2001 From: Alan Guo Xiang Tan Date: Fri, 25 Sep 2026 14:27:42 +0800 Subject: [PATCH 2/5] DEV: Report worker timeouts directly from the event handler --- lib/reporter/worker_timeout.rb | 32 ------------------- plugin.rb | 27 ++++++++++++++-- .../lib/{reporter => }/worker_timeout_spec.rb | 10 +++--- 3 files changed, 30 insertions(+), 39 deletions(-) delete mode 100644 lib/reporter/worker_timeout.rb rename spec/lib/{reporter => }/worker_timeout_spec.rb (87%) diff --git a/lib/reporter/worker_timeout.rb b/lib/reporter/worker_timeout.rb deleted file mode 100644 index 59405f4..0000000 --- a/lib/reporter/worker_timeout.rb +++ /dev/null @@ -1,32 +0,0 @@ -# frozen_string_literal: true - -require "timeout" - -module DiscoursePrometheus::Reporter - class WorkerTimeout - def report - metric = DiscoursePrometheus::InternalMetric::Custom.new - metric.type = "Counter" - metric.name = "pitchfork_worker_timeouts_total" - metric.description = "Total number of Pitchfork soft worker timeouts" - metric.value = 1 - - client = - PrometheusExporter::Client.new( - host: "localhost", - port: GlobalSetting.prometheus_collector_port, - process_queue_once_and_stop: true, - ) - - Timeout.timeout(1) do - begin - client.send_json(metric.to_h) - ensure - client.stop - end - end - rescue => error - Rails.logger.warn("Failed to report worker timeout: #{error.message}") - end - end -end diff --git a/plugin.rb b/plugin.rb index bc6a1c9..d891735 100644 --- a/plugin.rb +++ b/plugin.rb @@ -13,6 +13,7 @@ module ::DiscoursePrometheus gem "prometheus_exporter", "2.2.0" require "prometheus_exporter/client" +require "timeout" require_relative("lib/internal_metric/base") require_relative("lib/internal_metric/global") @@ -26,7 +27,6 @@ module ::DiscoursePrometheus require_relative("lib/reporter/global") require_relative("lib/reporter/web") require_relative("lib/reporter/image_processing") -require_relative("lib/reporter/worker_timeout") require_relative("lib/collector_demon") require_relative("lib/global_reporter_demon") @@ -54,7 +54,30 @@ module ::DiscoursePrometheus on(:image_processing_finished) { |payload| image_processing_reporter.report(payload) } - on(:web_worker_timeout) { DiscoursePrometheus::Reporter::WorkerTimeout.new.report } + on(:web_worker_timeout) do + metric = DiscoursePrometheus::InternalMetric::Custom.new + metric.type = "Counter" + metric.name = "pitchfork_worker_timeouts_total" + metric.description = "Total number of Pitchfork soft worker timeouts" + metric.value = 1 + + client = + PrometheusExporter::Client.new( + host: "localhost", + port: GlobalSetting.prometheus_collector_port, + process_queue_once_and_stop: true, + ) + + Timeout.timeout(1) do + begin + client.send_json(metric.to_h) + ensure + client.stop + end + end + rescue => error + Rails.logger.warn("Failed to report worker timeout: #{error.message}") + end register_demon_process(DiscoursePrometheus::CollectorDemon) register_demon_process(DiscoursePrometheus::GlobalReporterDemon) diff --git a/spec/lib/reporter/worker_timeout_spec.rb b/spec/lib/worker_timeout_spec.rb similarity index 87% rename from spec/lib/reporter/worker_timeout_spec.rb rename to spec/lib/worker_timeout_spec.rb index 5eafd17..0e5f790 100644 --- a/spec/lib/reporter/worker_timeout_spec.rb +++ b/spec/lib/worker_timeout_spec.rb @@ -1,10 +1,10 @@ # frozen_string_literal: true require "prometheus_exporter/server" -require_relative "../../../lib/collector" +require_relative "../../lib/collector" -RSpec.describe DiscoursePrometheus::Reporter::WorkerTimeout do - describe "#report" do +RSpec.describe DiscoursePrometheus do + describe ":web_worker_timeout" do it "delivers timeout events before their reporting processes exit" do socket = TCPServer.new("127.0.0.1", 0) port = socket.addr[1] @@ -48,14 +48,14 @@ socket.close global_setting :prometheus_collector_port, port - expect { described_class.new.report }.not_to raise_error + expect { DiscourseEvent.trigger(:web_worker_timeout) }.not_to raise_error end it "bounds reporting when the collector connection stalls" do allow(TCPSocket).to receive(:new) { sleep 30 } started_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) - described_class.new.report + DiscourseEvent.trigger(:web_worker_timeout) expect(Process.clock_gettime(Process::CLOCK_MONOTONIC) - started_at).to be < 3 end From bac7841c674a8011cfcb1dca0ccf1c18597969f8 Mon Sep 17 00:00:00 2001 From: Alan Guo Xiang Tan Date: Fri, 25 Sep 2026 14:37:42 +0800 Subject: [PATCH 3/5] DEV: Use the shared client to report worker timeouts --- plugin.rb | 14 ++------------ spec/lib/worker_timeout_spec.rb | 21 +++++++++++---------- 2 files changed, 13 insertions(+), 22 deletions(-) diff --git a/plugin.rb b/plugin.rb index d891735..becec13 100644 --- a/plugin.rb +++ b/plugin.rb @@ -61,19 +61,9 @@ module ::DiscoursePrometheus metric.description = "Total number of Pitchfork soft worker timeouts" metric.value = 1 - client = - PrometheusExporter::Client.new( - host: "localhost", - port: GlobalSetting.prometheus_collector_port, - process_queue_once_and_stop: true, - ) - Timeout.timeout(1) do - begin - client.send_json(metric.to_h) - ensure - client.stop - end + $prometheus_client.send_json(metric.to_h) + $prometheus_client.stop(wait_timeout_seconds: 1) end rescue => error Rails.logger.warn("Failed to report worker timeout: #{error.message}") diff --git a/spec/lib/worker_timeout_spec.rb b/spec/lib/worker_timeout_spec.rb index 0e5f790..6d576a7 100644 --- a/spec/lib/worker_timeout_spec.rb +++ b/spec/lib/worker_timeout_spec.rb @@ -5,11 +5,18 @@ RSpec.describe DiscoursePrometheus do describe ":web_worker_timeout" do + let(:port) { TCPServer.open("127.0.0.1", 0) { |socket| socket.addr[1] } } + + around do |example| + original_client = $prometheus_client + $prometheus_client = PrometheusExporter::Client.new(host: "127.0.0.1", port: port) + example.run + ensure + $prometheus_client.stop + $prometheus_client = original_client + end + it "delivers timeout events before their reporting processes exit" do - socket = TCPServer.new("127.0.0.1", 0) - port = socket.addr[1] - socket.close - global_setting :prometheus_collector_port, port collector = DiscoursePrometheus::Collector.new server = PrometheusExporter::Server::WebServer.new( @@ -39,15 +46,9 @@ ensure server&.stop runner&.join - socket&.close unless socket&.closed? end it "returns when the collector is unavailable" do - socket = TCPServer.new("127.0.0.1", 0) - port = socket.addr[1] - socket.close - global_setting :prometheus_collector_port, port - expect { DiscourseEvent.trigger(:web_worker_timeout) }.not_to raise_error end From 829fd418655dff05251b031cc2c22bd6ac308160 Mon Sep 17 00:00:00 2001 From: Alan Guo Xiang Tan Date: Fri, 25 Sep 2026 15:03:42 +0800 Subject: [PATCH 4/5] DEV: Remove timeout metric README changes --- README.md | 6 ------ 1 file changed, 6 deletions(-) diff --git a/README.md b/README.md index 990dfce..d6859ea 100644 --- a/README.md +++ b/README.md @@ -6,10 +6,4 @@ The Discourse Prometheus plugin collects key metrics from Discourse and exposes The global reporter can pick custom metrics added by other Discourse plugins. The metric needs to define a collect method, and the `name`, `labels`, `description`, `value`, and `type` attributes. See an example [here](https://github.com/discourse/discourse-antivirus/pull/15). -## Worker timeouts - -`discourse_pitchfork_worker_timeouts_total` counts Pitchfork soft worker timeout events. It is exposed after the first event and resets when the collector restarts. Reporting is best effort, with a one-second budget before the worker exits. Hard kills that bypass the soft timeout callback are not counted. - -This metric requires a Discourse version that emits the `web_worker_timeout` event. For example, `rate(discourse_pitchfork_worker_timeouts_total[5m])` shows timeouts per second over five minutes. - For more information, please see: https://meta.discourse.org/t/prometheus-exporter-plugin-for-discourse/72666 From c4f4799573781d2f5615edf9122596a697c80103 Mon Sep 17 00:00:00 2001 From: Alan Guo Xiang Tan Date: Fri, 25 Sep 2026 15:27:38 +0800 Subject: [PATCH 5/5] DEV: Send worker timeout metrics synchronously through shared client Use the synchronous client API proposed in discourse/prometheus_exporter#384. Remove explicit queue draining and Ruby Timeout. The draft remains dependent on an exporter release and a matching gem version update. --- plugin.rb | 8 +---- spec/lib/worker_timeout_spec.rb | 52 ++++----------------------------- 2 files changed, 6 insertions(+), 54 deletions(-) diff --git a/plugin.rb b/plugin.rb index becec13..942f104 100644 --- a/plugin.rb +++ b/plugin.rb @@ -13,7 +13,6 @@ module ::DiscoursePrometheus gem "prometheus_exporter", "2.2.0" require "prometheus_exporter/client" -require "timeout" require_relative("lib/internal_metric/base") require_relative("lib/internal_metric/global") @@ -61,12 +60,7 @@ module ::DiscoursePrometheus metric.description = "Total number of Pitchfork soft worker timeouts" metric.value = 1 - Timeout.timeout(1) do - $prometheus_client.send_json(metric.to_h) - $prometheus_client.stop(wait_timeout_seconds: 1) - end - rescue => error - Rails.logger.warn("Failed to report worker timeout: #{error.message}") + $prometheus_client.send_json_sync(metric.to_h) end register_demon_process(DiscoursePrometheus::CollectorDemon) diff --git a/spec/lib/worker_timeout_spec.rb b/spec/lib/worker_timeout_spec.rb index 6d576a7..fc8d9dc 100644 --- a/spec/lib/worker_timeout_spec.rb +++ b/spec/lib/worker_timeout_spec.rb @@ -5,60 +5,18 @@ RSpec.describe DiscoursePrometheus do describe ":web_worker_timeout" do - let(:port) { TCPServer.open("127.0.0.1", 0) { |socket| socket.addr[1] } } - - around do |example| - original_client = $prometheus_client - $prometheus_client = PrometheusExporter::Client.new(host: "127.0.0.1", port: port) - example.run - ensure - $prometheus_client.stop - $prometheus_client = original_client - end - - it "delivers timeout events before their reporting processes exit" do + it "increments the timeout counter through the shared client" do collector = DiscoursePrometheus::Collector.new - server = - PrometheusExporter::Server::WebServer.new( - port: port, - bind: "127.0.0.1", - collector: collector, - ) - runner = server.start - - 2.times do - pid = - fork do - DiscourseEvent.trigger(:web_worker_timeout) - Process.exit!(0) - end - _, status = Process.wait2(pid) - expect(status).to be_success + allow($prometheus_client).to receive(:send_json_sync) do |metric| + collector.process(Oj.dump(metric, mode: :object)) end - wait_for(timeout: 5) do - collector.prometheus_metrics_text.include?("pitchfork_worker_timeouts_total 2") - end + 2.times { DiscourseEvent.trigger(:web_worker_timeout) } + expect(collector.prometheus_metrics_text).to include( "pitchfork_worker_timeouts_total counter", "pitchfork_worker_timeouts_total 2", ) - ensure - server&.stop - runner&.join - end - - it "returns when the collector is unavailable" do - expect { DiscourseEvent.trigger(:web_worker_timeout) }.not_to raise_error - end - - it "bounds reporting when the collector connection stalls" do - allow(TCPSocket).to receive(:new) { sleep 30 } - started_at = Process.clock_gettime(Process::CLOCK_MONOTONIC) - - DiscourseEvent.trigger(:web_worker_timeout) - - expect(Process.clock_gettime(Process::CLOCK_MONOTONIC) - started_at).to be < 3 end end end