Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 9 additions & 2 deletions lib/fluent/plugin/buffer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -917,8 +917,15 @@ def write_step_by_step(metadata, data, format, splits_count, &block)
]

def statistics
stage_size, queue_size = @stage_size_metrics.get, @queue_size_metrics.get
buffer_space = 1.0 - ((stage_size + queue_size * 1.0) / @total_limit_size)
# Export-only clamp: internal gauges may go transiently negative during
# the deferred stage_size add vs enqueue_chunk sub race (#5303, #2712).
# Clamping the gauge store itself would turn that into a permanent
# over-count and break Buffer#storable? -- keep raw gauge semantics.
stage_size = [@stage_size_metrics.get, 0].max
queue_size = [@queue_size_metrics.get, 0].max
denom = @total_limit_size.to_f
# denom > 0 already excludes 0/0 NaN; stage/queue are floored above.
buffer_space = denom > 0.0 ? (1.0 - (stage_size + queue_size).to_f / denom).clamp(0.0, 1.0) : 0.0
@stage_length_metrics.set(@stage.size)
@queue_length_metrics.set(@queue.size)
@available_buffer_space_ratios_metrics.set(buffer_space * 100)
Expand Down
42 changes: 42 additions & 0 deletions test/plugin/test_buffer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -1536,5 +1536,47 @@ def create_chunk_es(metadata, es)
test 'returns available_buffer_space_ratios' do
assert_equal 10.0, @p.statistics['buffer']['available_buffer_space_ratios']
end

# Export-only clamp: internal gauges may be negative (self-correcting race
# or set_gauge); statistics must still publish non-negative sizes (#5303).
test 'exports non-negative stage/queue byte sizes when gauges are negative' do
@p.stage_size_metrics.set(-50)
@p.queue_size_metrics.set(-100)

stats = @p.statistics['buffer']
assert_equal 0, stats['stage_byte_size']
assert_equal 0, stats['queue_byte_size']
assert_equal 0, stats['total_queued_size']
assert stats['total_queued_size'] >= 0
# Negative gauges floor to 0 usage => full free space (100.0% with total_limit_size=1024)
assert_equal 100.0, stats['available_buffer_space_ratios']
end

test 'clamps available_buffer_space_ratios when usage exceeds total_limit_size' do
# Simulate counter overshoot past configured limit
@p.stage_size_metrics.set(2000)
@p.queue_size_metrics.set(2000)

stats = @p.statistics['buffer']
assert_equal 2000, stats['stage_byte_size']
assert_equal 2000, stats['queue_byte_size']
assert_equal 4000, stats['total_queued_size']
assert stats['available_buffer_space_ratios'] >= 0.0
assert stats['available_buffer_space_ratios'] <= 100.0
assert_equal 0.0, stats['available_buffer_space_ratios']
end

test 'available_buffer_space_ratios is safe when total_limit_size is zero' do
# 0/0 would be NaN; Array#max/min on NaN raises — must not crash export.
@p.instance_variable_set(:@total_limit_size, 0)
@p.stage_size_metrics.set(0)
@p.queue_size_metrics.set(0)
stats = @p.statistics['buffer']
assert_equal 0, stats['stage_byte_size']
assert_equal 0, stats['queue_byte_size']
assert stats['available_buffer_space_ratios'] >= 0.0
assert stats['available_buffer_space_ratios'] <= 100.0
refute stats['available_buffer_space_ratios'].to_f.nan?
end
end
end