From 209a3be229b2326a7962ea6167b0c8e06654f9cd Mon Sep 17 00:00:00 2001 From: Kentaro Hayashi Date: Fri, 31 Jul 2026 14:06:33 +0900 Subject: [PATCH 1/2] in_prometheus: do not disclosure error details for client It is reasonable to logging error, but no need to disclose detail for client. Signed-off-by: Kentaro Hayashi --- lib/fluent/plugin/in_prometheus.rb | 6 ++- spec/fluent/plugin/in_prometheus_spec.rb | 67 ++++++++++++++++++++++++ 2 files changed, 71 insertions(+), 2 deletions(-) diff --git a/lib/fluent/plugin/in_prometheus.rb b/lib/fluent/plugin/in_prometheus.rb index 22b5d71..c8249f8 100644 --- a/lib/fluent/plugin/in_prometheus.rb +++ b/lib/fluent/plugin/in_prometheus.rb @@ -210,7 +210,8 @@ def start_webrick def all_metrics response(::Prometheus::Client::Formats::Text.marshal(@registry)) rescue => e - [500, { 'Content-Type' => 'text/plain' }, e.to_s] + log.error "in_prometheus: failed to render metrics", error_class: e.class, error: e + [500, { 'Content-Type' => 'text/plain' }, "in_prometheus server error: <#{e.class}>"] end def all_workers_metrics @@ -223,7 +224,8 @@ def all_workers_metrics end response(full_result.get_metrics) rescue => e - [500, { 'Content-Type' => 'text/plain' }, e.to_s] + log.error "in_prometheus: failed to render workers metrics", error_class: e.class, error: e + [500, { 'Content-Type' => 'text/plain' }, "in_prometheus server error: <#{e.class}>"] end def send_request_to_each_worker diff --git a/spec/fluent/plugin/in_prometheus_spec.rb b/spec/fluent/plugin/in_prometheus_spec.rb index a343e58..fd212da 100644 --- a/spec/fluent/plugin/in_prometheus_spec.rb +++ b/spec/fluent/plugin/in_prometheus_spec.rb @@ -338,4 +338,71 @@ include_examples 'IPv6 server binding', '[::1]', '::1', 'handles pre-bracketed address correctly' end end + + describe 'error handling (information disclosure)' do + let(:config) { LOCAL_CONFIG } + let(:secret_message) { 'dummy secret detail: password=deadbeef' } + + shared_examples 'suppressed exception response' do + it 'returns 500 with text/plain' do + status, header, _body = subject + expect(status).to eq(500) + expect(header['Content-Type']).to eq('text/plain') + end + + it 'exposes the exception class only' do + _status, _header, body = subject + expect(body).to eq('in_prometheus server error: ') + expect(body).not_to include(secret_message) + end + + it 'logs the detail on the server side' do + subject + expect(driver.logs.any? { |log| log.include?(log_message) }).to be true + expect(driver.logs.any? { |log| log.include?(secret_message) }).to be true + end + end + + context '#all_metrics' do + subject { driver.instance.send(:all_metrics) } + + let(:log_message) { 'in_prometheus: failed to render metrics' } + + before do + allow(::Prometheus::Client::Formats::Text).to receive(:marshal).and_raise(RuntimeError, secret_message) + end + + include_examples 'suppressed exception response' + end + + context '#all_workers_metrics' do + subject { driver.instance.send(:all_workers_metrics) } + + let(:log_message) { 'in_prometheus: failed to render workers metrics' } + + before do + allow(driver.instance).to receive(:send_request_to_each_worker).and_raise(RuntimeError, secret_message) + end + + include_examples 'suppressed exception response' + end + + context 'over HTTP' do + before do + allow(::Prometheus::Client::Formats::Text).to receive(:marshal).and_raise(RuntimeError, secret_message) + end + + it 'does not leak the exception message to the client' do + driver.run(timeout: 1) do + Net::HTTP.start('127.0.0.1', port) do |http| + req = Net::HTTP::Get.new('/metrics') + res = http.request(req) + expect(res.code).to eq('500') + expect(res.body).to eq('in_prometheus server error: ') + expect(res.body).not_to include(secret_message) + end + end + end + end + end end From 3aa154752ca1e9d38a062d93ca525efcece43061 Mon Sep 17 00:00:00 2001 From: Kentaro Hayashi Date: Fri, 31 Jul 2026 18:10:13 +0900 Subject: [PATCH 2/2] Introduce error log throttling /metrics and /aggregated_metrics are assumed that these API will be called periodically. If internal server error occurs continuously, it means that same error log will be recorded. That is incompatible behavior before. To record detailed logs as often as necessary, introduced `ignore_error_log_interval`. Fluentd itself have ignore_repeated_log_interval and ignore_same_log_interval, but it must be system wide configuration. The scope of error handling should be limited to this plugin, so do not escalate system wide configuration. Signed-off-by: Kentaro Hayashi --- README.md | 1 + lib/fluent/plugin/in_prometheus.rb | 36 ++++- spec/fluent/plugin/in_prometheus_spec.rb | 171 +++++++++++++++++++++++ 3 files changed, 206 insertions(+), 2 deletions(-) diff --git a/README.md b/README.md index 4b1dcee..f353f09 100644 --- a/README.md +++ b/README.md @@ -64,6 +64,7 @@ More configuration parameters: - `metrics_path`: metrics HTTP endpoint (default: /metrics) - `aggregated_metrics_path`: metrics HTTP endpoint (default: /aggregated_metrics) - `content_encoding`: encoding format for the exposed metrics (default: identity). Supported formats are {identity, gzip} +- `ignore_error_log_interval`: Suppress repeated error logs in a certain period of time or until message was changed (default: 1h) When using multiple workers, each worker binds to port + `fluent_worker_id`. To scrape metrics from all workers at once, you can access http://localhost:24231/aggregated_metrics. diff --git a/lib/fluent/plugin/in_prometheus.rb b/lib/fluent/plugin/in_prometheus.rb index c8249f8..82d4907 100644 --- a/lib/fluent/plugin/in_prometheus.rb +++ b/lib/fluent/plugin/in_prometheus.rb @@ -36,10 +36,15 @@ class PrometheusInput < Fluent::Plugin::Input desc 'Content encoding of the exposed metrics, Currently supported encoding is identity, gzip. Ref: https://prometheus.io/docs/instrumenting/exposition_formats/#basic-info' config_param :content_encoding, :enum, list: [:identity, :gzip], default: :identity + desc 'Suppress repeated error logs in a certain period of time (1h) or until message was changed' + config_param :ignore_error_log_interval, :time, default: 3600 + def initialize super @registry = ::Prometheus::Client.registry @secure = nil + @error_log_mutex = Mutex.new + @last_error_logs = {} # scope => [logged_at, fingerprint, suppressed_count] end def configure(conf) @@ -210,7 +215,7 @@ def start_webrick def all_metrics response(::Prometheus::Client::Formats::Text.marshal(@registry)) rescue => e - log.error "in_prometheus: failed to render metrics", error_class: e.class, error: e + log_error_throttled(:metrics, "in_prometheus: failed to render metrics", error: e) [500, { 'Content-Type' => 'text/plain' }, "in_prometheus server error: <#{e.class}>"] end @@ -224,7 +229,7 @@ def all_workers_metrics end response(full_result.get_metrics) rescue => e - log.error "in_prometheus: failed to render workers metrics", error_class: e.class, error: e + log_error_throttled(:workers_metrics, "in_prometheus: failed to render workers metrics", error: e) [500, { 'Content-Type' => 'text/plain' }, "in_prometheus server error: <#{e.class}>"] end @@ -273,5 +278,32 @@ def response(metrics) end [200, { 'Content-Type' => ::Prometheus::Client::Formats::Text::CONTENT_TYPE, 'Content-Encoding' => @content_encoding.to_s }, body] end + + def log_error_throttled(scope, message, error:) + fingerprint = [error.class, error.message] + suppressed = 0 + + emit = @error_log_mutex.synchronize do + last = @last_error_logs[scope] + now = Fluent::Clock.now + if last.nil? || + last[1] != fingerprint || + (now - last[0]) >= @ignore_error_log_interval + suppressed = last && last[1] == fingerprint ? last[2] : 0 + @last_error_logs[scope] = [now, fingerprint, 0] + true + else + last[2] += 1 + false + end + end + return unless emit + + if suppressed > 0 + log.error message, error_class: error.class, error: error, suppressed_log_count: suppressed + else + log.error message, error_class: error.class, error: error + end + end end end diff --git a/spec/fluent/plugin/in_prometheus_spec.rb b/spec/fluent/plugin/in_prometheus_spec.rb index fd212da..9d05506 100644 --- a/spec/fluent/plugin/in_prometheus_spec.rb +++ b/spec/fluent/plugin/in_prometheus_spec.rb @@ -71,6 +71,22 @@ expect(driver.instance.content_encoding).to eq(:gzip) end end + + describe 'default ignore_error_log_interval' do + let(:config) { CONFIG } + it 'should be 3600 seconds by default' do + expect(driver.instance.ignore_error_log_interval).to eq(3600) + end + end + + describe 'error_log_interval' do + let(:config) { CONFIG + %[ + ignore_error_log_interval 60 +] } + it 'should be configurable' do + expect(driver.instance.ignore_error_log_interval).to eq(60) + end + end end describe '#start' do @@ -405,4 +421,159 @@ end end end + + describe 'error log throttling' do + let(:config) { LOCAL_CONFIG } + let(:secret_message) { 'dummy secret detail: password=deadbeef' } + let(:log_message) { 'in_prometheus: failed to render metrics' } + let(:workers_log_message) { 'in_prometheus: failed to render workers metrics' } + + # Fluent::Clock.now is monotonic, so a plain Hash is enough to drive it + let(:clock) { { now: 1000.0 } } + + def error_logs(message) + driver.logs.select { |log| log.include?(message) } + end + + context 'when rendering metrics keeps failing' do + before do + allow(Fluent::Clock).to receive(:now) { clock[:now] } + allow(::Prometheus::Client::Formats::Text).to receive(:marshal).and_raise(RuntimeError, secret_message) + end + + # every iteration raises from the same line, so the exceptions are equal + # to each other and only ignore_error_log_interval can let a log through + it 'logs the repeated same failure only once within ignore_error_log_interval' do + 5.times { driver.instance.send(:all_metrics) } + expect(error_logs(log_message).size).to eq(1) + end + + it 'keeps returning 500 to the client even while the log is suppressed' do + responses = 5.times.map { driver.instance.send(:all_metrics) } + expect(error_logs(log_message).size).to eq(1) + responses.each do |status, _header, body| + expect(status).to eq(500) + expect(body).to eq('in_prometheus server error: ') + end + end + + it 'logs the repeated same failure again after ignore_error_log_interval has elapsed' do + 2.times do + driver.instance.send(:all_metrics) + clock[:now] += driver.instance.ignore_error_log_interval + end + expect(error_logs(log_message).size).to eq(2) + end + + it 'does not log the repeated same failure just before ignore_error_log_interval has elapsed' do + 2.times do + driver.instance.send(:all_metrics) + clock[:now] += driver.instance.ignore_error_log_interval - 0.1 + end + expect(error_logs(log_message).size).to eq(1) + end + + it 'reports how many logs were suppressed in the meantime' do + 3.times { driver.instance.send(:all_metrics) } + clock[:now] += driver.instance.ignore_error_log_interval + driver.instance.send(:all_metrics) + logs = error_logs(log_message) + expect(logs.size).to eq(2) + expect(logs.first).not_to include('suppressed_log_count') + expect(logs.last).to include('suppressed_log_count=2') + end + end + + describe 'telling the errors apart' do + before do + allow(Fluent::Clock).to receive(:now) { clock[:now] } + end + + def log_error(scope, message, error) + driver.instance.send(:log_error_throttled, scope, message, error: error) + end + + # the plugin raises a fresh exception object per failure, so the errors + # have to be compared by value, not by identity + it 'suppresses an equal error given as a different object' do + log_error(:metrics, log_message, RuntimeError.new(secret_message)) + log_error(:metrics, log_message, RuntimeError.new(secret_message)) + expect(error_logs(log_message).size).to eq(1) + end + + it 'logs immediately when the error class differs' do + log_error(:metrics, log_message, RuntimeError.new(secret_message)) + log_error(:metrics, log_message, ArgumentError.new(secret_message)) + expect(error_logs(log_message).size).to eq(2) + end + + it 'logs immediately when the error differs' do + log_error(:metrics, log_message, RuntimeError.new(secret_message)) + log_error(:metrics, log_message, RuntimeError.new('another failure')) + expect(error_logs(log_message).size).to eq(2) + end + + # the scope, not the log message, picks the slot to throttle on + it 'keeps a separate state per scope' do + error = RuntimeError.new(secret_message) + log_error(:metrics, log_message, error) + log_error(:workers_metrics, workers_log_message, error) + expect(error_logs(log_message).size).to eq(1) + expect(error_logs(workers_log_message).size).to eq(1) + end + + it 'suppresses an equal error within a scope even when the log message differs' do + error = RuntimeError.new(secret_message) + log_error(:metrics, log_message, error) + log_error(:metrics, workers_log_message, error) + expect(error_logs(workers_log_message)).to be_empty + end + + context 'with ignore_error_log_interval 0' do + let(:config) { LOCAL_CONFIG + %[ + ignore_error_log_interval 0 +] } + + it 'logs every occurrence of the same error' do + 3.times { log_error(:metrics, log_message, RuntimeError.new(secret_message)) } + expect(error_logs(log_message).size).to eq(3) + end + end + end + + # /metrics and /aggregated_metrics are usually scraped in turn, so both of + # them must be throttled on their own slot + context 'when both endpoints keep failing alternately' do + before do + allow(Fluent::Clock).to receive(:now) { clock[:now] } + allow(::Prometheus::Client::Formats::Text).to receive(:marshal).and_raise(RuntimeError, secret_message) + allow(driver.instance).to receive(:send_request_to_each_worker).and_raise(ArgumentError, 'another failure') + end + + it 'logs each failure only once within ignore_error_log_interval' do + 5.times do + driver.instance.send(:all_metrics) + driver.instance.send(:all_workers_metrics) + end + expect(error_logs(log_message).size).to eq(1) + expect(error_logs(workers_log_message).size).to eq(1) + end + end + + context 'when errors occur concurrently' do + # long enough to keep every call within the same interval + let(:config) { LOCAL_CONFIG + %[ + ignore_error_log_interval 3600 +] } + + it 'logs the error only once' do + instance = driver.instance + error = RuntimeError.new(secret_message) + 10.times.map { + Thread.new { instance.send(:log_error_throttled, :metrics, log_message, error: error) } + }.each(&:join) + expect(error_logs(log_message).size).to eq(1) + end + end + end end