Skip to content
Draft
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
21 changes: 10 additions & 11 deletions lib/fluent/plugin/buffer.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
Expand All @@ -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
Expand Down
74 changes: 74 additions & 0 deletions test/plugin/test_stage_size_race.rb
Original file line number Diff line number Diff line change
@@ -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