From 4d8d406a8ea8bebca5a9f559f73d4f8d82277f95 Mon Sep 17 00:00:00 2001 From: simonyang08 Date: Tue, 1 Sep 2026 11:38:06 +0000 Subject: [PATCH] Fix buffer stage size accounting race Signed-off-by: simonyang08 --- lib/fluent/plugin/buffer.rb | 21 ++++---- test/plugin/test_stage_size_race.rb | 74 +++++++++++++++++++++++++++++ 2 files changed, 84 insertions(+), 11 deletions(-) create mode 100644 test/plugin/test_stage_size_race.rb diff --git a/lib/fluent/plugin/buffer.rb b/lib/fluent/plugin/buffer.rb index 59022648b0..e8cb8bb935 100644 --- a/lib/fluent/plugin/buffer.rb +++ b/lib/fluent/plugin/buffer.rb @@ -383,6 +383,12 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false) if enqueue || first_chunk.unstaged? || chunk_size_full?(first_chunk) chunks_to_enqueue << first_chunk end + if (bytesize = staged_bytesizes_by_chunk[first_chunk]) + # Account while the chunk lock is still held so enqueue_chunk + # cannot remove the chunk and subtract its bytes first. + @stage_size_metrics.add(bytesize) + log.on_trace { log.trace { "chunk #{first_chunk.path} size_added: #{bytesize} new_size: #{first_chunk.bytesize}" } } + end first_chunk.mon_exit rescue operated_chunks.unshift(first_chunk) @@ -397,6 +403,10 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false) if enqueue || chunk.unstaged? || chunk_size_full?(chunk) chunks_to_enqueue << chunk end + if (bytesize = staged_bytesizes_by_chunk[chunk]) + @stage_size_metrics.add(bytesize) + log.on_trace { log.trace { "chunk #{chunk.path} size_added: #{bytesize} new_size: #{chunk.bytesize}" } } + end chunk.mon_exit rescue => e chunk.rollback @@ -407,17 +417,6 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false) # All locks about chunks are released. - # - # Now update the stage, stage_size with proper locking - # FIX FOR stage_size miscomputation - https://github.com/fluent/fluentd/issues/2712 - # - staged_bytesizes_by_chunk.each do |chunk, bytesize| - chunk.synchronize do - synchronize { @stage_size_metrics.add(bytesize) } - log.on_trace { log.trace { "chunk #{chunk.path} size_added: #{bytesize} new_size: #{chunk.bytesize}" } } - end - end - chunks_to_enqueue.each do |c| if c.staged? && (enqueue || chunk_size_full?(c)) m = c.metadata diff --git a/test/plugin/test_stage_size_race.rb b/test/plugin/test_stage_size_race.rb new file mode 100644 index 0000000000..fec046091e --- /dev/null +++ b/test/plugin/test_stage_size_race.rb @@ -0,0 +1,74 @@ +require_relative '../helper' +require 'fluent/plugin/buffer' +require 'fluent/plugin/buffer/memory_chunk' +require 'fluent/plugin_id' +require 'fluent/log' + +module StageSizeRace + class Owner < Fluent::Plugin::Base + include Fluent::PluginId + include Fluent::PluginLoggerMixin + end + + class Buf < Fluent::Plugin::Buffer + def create_metadata(timekey = nil, tag = nil, variables = nil) + Fluent::Plugin::Buffer::Metadata.new(timekey, tag, variables) + end + + def resume + return {}, [] + end + + def generate_chunk(metadata) + Fluent::Plugin::Buffer::MemoryChunk.new(metadata) + end + end +end + +class StageSizeRaceTest < ::Test::Unit::TestCase + test 'stage_size does not go negative while a write is pending' do + b = StageSizeRace::Buf.new + b.owner = StageSizeRace::Owner.new + b.configure(config_element('buffer', '', { 'total_limit_size' => 1024, 'chunk_limit_size' => 4096 })) + b.start + + m = b.create_metadata + b.write({ m => ['a' * 400] }) + chunk = b.stage[m] + assert_equal 400, b.stage_size + + reached = Queue.new + resume = Queue.new + armed = false + + chunk.define_singleton_method(:mon_exit) do + result = super() + if armed + armed = false + reached << true + resume.pop + end + result + end + + armed = true + writer = Thread.new { b.write({ m => ['b' * 400] }) } + reached.pop + + b.enqueue_chunk(m) + assert_equal 0, b.stage_size + + resume << true + writer.join + + assert_equal 0, b.stage.size + assert_equal 0, b.stage_size + assert_equal 800, b.queue_size + ensure + if writer&.alive? + resume << true + writer.join + end + b&.stop + end +end