Describe the bug
Buffer#write cannot update @stage_size_metrics while it still holds chunk locks, so it defers the add until every chunk lock has been 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 |
That leaves a window between chunk.mon_exit and the deferred add where the bytes are already committed into the chunk and visible through @stage[metadata], but have not been counted yet. If a flush thread calls enqueue_chunk inside that window it subtracts the current chunk.bytesize, including the bytes the writer has not added:
|
bytesize = chunk.bytesize |
|
@stage_size_metrics.sub(bytesize) |
|
@queue_size_metrics.add(bytesize) |
So sub runs before the matching add and stage_size drops below zero until the deferred add lands.
#5467 clamps the values at export time so exporters no longer publish negative sizes. It deliberately does not touch the gauge store, because clamping LocalMetrics#sub/#dec turns the self-correcting negative into a permanent over-count and makes Buffer#storable? reject every write. This issue tracks the underlying accounting window that #5467 masks.
To Reproduce
The window is deterministically reproducible by pausing the writer on the staged chunk's mon_exit. This test passes on master today, i.e. it asserts the negative intermediate value:
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 goes negative while the deferred add 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
# pause the writer right after it releases the chunk lock (buffer.rb L386)
# and before the deferred add (buffer.rb L414-L419)
chunk.define_singleton_method(:mon_exit) do
r = super()
if armed
armed = false
reached << true
resume.pop
end
r
end
armed = true
writer = Thread.new { b.write({ m => ['b' * 400] }) }
reached.pop
b.enqueue_chunk(m)
assert_equal(-400, b.stage_size) # nothing is staged, yet the gauge is negative
resume << true
writer.join
assert_equal 0, b.stage.size
assert_equal 0, b.stage_size # self-heals once the deferred add lands
end
end
Output at the two checkpoints:
mid: stage.size=0 stage_size=-400 queue_size=800
end: stage.size=0 stage_size=0 queue_size=800
Expected behavior
stage_size should never be negative, and it should equal the total bytesize of the chunks currently in @stage once concurrent operations have settled.
Your Environment
- Fluentd version:
- Package version:
- Operating system:
- Kernel version:
Your Configuration
Your Error Log
Additional context
No response
Describe the bug
Buffer#writecannot update@stage_size_metricswhile it still holds chunk locks, so it defers the add until every chunk lock has been released:fluentd/lib/fluent/plugin/buffer.rb
Lines 410 to 419 in 4d5527a
That leaves a window between
chunk.mon_exitand the deferredaddwhere the bytes are already committed into the chunk and visible through@stage[metadata], but have not been counted yet. If a flush thread callsenqueue_chunkinside that window it subtracts the currentchunk.bytesize, including the bytes the writer has not added:fluentd/lib/fluent/plugin/buffer.rb
Lines 503 to 505 in 4d5527a
So
subruns before the matchingaddandstage_sizedrops below zero until the deferred add lands.#5467 clamps the values at export time so exporters no longer publish negative sizes. It deliberately does not touch the gauge store, because clamping
LocalMetrics#sub/#decturns the self-correcting negative into a permanent over-count and makesBuffer#storable?reject every write. This issue tracks the underlying accounting window that #5467 masks.To Reproduce
The window is deterministically reproducible by pausing the writer on the staged chunk's
mon_exit. This test passes on master today, i.e. it asserts the negative intermediate value:Output at the two checkpoints:
Expected behavior
stage_sizeshould never be negative, and it should equal the total bytesize of the chunks currently in@stageonce concurrent operations have settled.Your Environment
Your Configuration
Your Error Log
Additional context
No response