Skip to content

Commit 9e3bbbe

Browse files
Watson1978claude
andcommitted
buffer: fix an enqueued unstaged file chunk being re-staged and flushed twice
FileChunk#enqueued! and FileSingleChunk#enqueued! did nothing for an unstaged chunk, so a chunk pushed by Buffer#enqueue_unstaged_chunk stayed :unstaged with its staged file name while it was in the queue. When write_step_by_step unstages chunks on a single-record overflow, Buffer#write re-staged that queued chunk because the `u.unstaged?` guard could not tell it was already enqueued, and the same chunk object lived in both @Queue and @stage. The second flush of that object raised "closed stream" at FileChunk#open and ENOENT at FileChunk#purge, and the buffer size gauges leaked. Treat an unstaged chunk like a staged one in enqueued!: mark it as :queued, write its metadata and rename its files to the queued path, so resume also restores it as a queued chunk. The state changes before the rename and the gauges are updated before enqueued!, so a chunk whose rename fails stays a consistent queued chunk with its staged file name. Buffer#write only warns about such a failure and keeps enqueueing the other chunks: the records are already queued, and a chunk left unstaged there would be purged with its committed records. FileSingleBuffer#resume now enqueues a staged chunk file whose metadata is already staged, as FileBuffer#resume does, so that two chunk files with the staged name are both restored. file_rename reopens the old file when the rename fails on Windows, so that the chunk stays readable. Fixes #4662 Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com> Signed-off-by: Shizuo Fujita <fujita@clear-code.com>
1 parent 7580e2b commit 9e3bbbe

9 files changed

Lines changed: 476 additions & 44 deletions

File tree

‎lib/fluent/plugin/buf_file.rb‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -174,7 +174,7 @@ def resume
174174
when :staged
175175
# unstaged chunk created at Buffer#write_step_by_step is identified as the staged chunk here because FileChunk#assume_chunk_state checks only the file name.
176176
# https://github.com/fluent/fluentd/blob/9d113029d4550ce576d8825bfa9612aa3e55bff0/lib/fluent/plugin/buffer.rb#L663
177-
# This case can happen when fluentd process is killed by signal or other reasons between creating unstaged chunks and changing them to staged mode in Buffer#write
177+
# This case can happen when fluentd process is killed by signal or other reasons between creating unstaged chunks and enqueueing them (which renames them) in Buffer#write
178178
# these chunks(unstaged chunks) has shared the same metadata
179179
# So perform enqueue step again https://github.com/fluent/fluentd/blob/9d113029d4550ce576d8825bfa9612aa3e55bff0/lib/fluent/plugin/buffer.rb#L364
180180
if chunk_size_full?(chunk) || stage.key?(chunk.metadata)

‎lib/fluent/plugin/buf_file_single.rb‎

Lines changed: 5 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -193,7 +193,11 @@ def resume
193193

194194
case chunk.state
195195
when :staged
196-
stage[chunk.metadata] = chunk
196+
if chunk_size_full?(chunk) || stage.key?(chunk.metadata)
197+
queue << chunk.enqueued!
198+
else
199+
stage[chunk.metadata] = chunk
200+
end
197201
when :queued
198202
queue << chunk
199203
end

‎lib/fluent/plugin/buffer.rb‎

Lines changed: 32 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -418,31 +418,41 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false)
418418
end
419419
end
420420

421+
enqueue_errors = []
421422
chunks_to_enqueue.each do |c|
422-
if c.staged? && (enqueue || chunk_size_full?(c))
423-
m = c.metadata
424-
enqueue_chunk(m)
425-
if unstaged_chunks[m] && !unstaged_chunks[m].empty?
426-
u = unstaged_chunks[m].pop
427-
u.synchronize do
428-
if u.unstaged? && !chunk_size_full?(u)
429-
# `u.metadata.seq` and `m.seq` can be different but Buffer#enqueue_chunk expect them to be the same value
430-
u.metadata.seq = 0
431-
synchronize {
432-
@stage[m] = u.staged!
433-
@stage_size_metrics.add(u.bytesize)
434-
}
423+
begin
424+
if c.staged? && (enqueue || chunk_size_full?(c))
425+
m = c.metadata
426+
enqueue_chunk(m)
427+
if unstaged_chunks[m] && !unstaged_chunks[m].empty?
428+
u = unstaged_chunks[m].pop
429+
u.synchronize do
430+
if u.unstaged? && !chunk_size_full?(u)
431+
# `u.metadata.seq` and `m.seq` can be different but Buffer#enqueue_chunk expect them to be the same value
432+
u.metadata.seq = 0
433+
synchronize {
434+
@stage[m] = u.staged!
435+
@stage_size_metrics.add(u.bytesize)
436+
}
437+
end
435438
end
436439
end
440+
elsif c.unstaged?
441+
enqueue_unstaged_chunk(c)
442+
else
443+
# previously staged chunk is already enqueued, closed or purged.
444+
# no problem.
437445
end
438-
elsif c.unstaged?
439-
enqueue_unstaged_chunk(c)
440-
else
441-
# previously staged chunk is already enqueued, closed or purged.
442-
# no problem.
446+
rescue => e
447+
enqueue_errors << e
443448
end
444449
end
445450

451+
enqueue_errors.each do |e|
452+
log.warn "error occurs in enqueueing a chunk", error: e
453+
log.warn_backtrace e.backtrace
454+
end
455+
446456
operated_chunks.clear if errors.empty?
447457

448458
if errors.size > 0
@@ -492,6 +502,9 @@ def enqueue_chunk(metadata)
492502

493503
chunk.synchronize do
494504
synchronize do
505+
bytesize = chunk.bytesize
506+
@stage_size_metrics.sub(bytesize)
507+
@queue_size_metrics.add(bytesize)
495508
if chunk.empty?
496509
chunk.close
497510
else
@@ -500,9 +513,6 @@ def enqueue_chunk(metadata)
500513
@queued_num[metadata] = @queued_num.fetch(metadata, 0) + 1
501514
chunk.enqueued!
502515
end
503-
bytesize = chunk.bytesize
504-
@stage_size_metrics.sub(bytesize)
505-
@queue_size_metrics.add(bytesize)
506516
end
507517
end
508518
nil
@@ -517,9 +527,9 @@ def enqueue_unstaged_chunk(chunk)
517527
metadata.seq = 0 # metadata.seq should be 0 for counting @queued_num
518528
@queue << chunk
519529
@queued_num[metadata] = @queued_num.fetch(metadata, 0) + 1
530+
@queue_size_metrics.add(chunk.bytesize)
520531
chunk.enqueued!
521532
end
522-
@queue_size_metrics.add(chunk.bytesize)
523533
end
524534
end
525535

‎lib/fluent/plugin/buffer/file_chunk.rb‎

Lines changed: 23 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -100,8 +100,14 @@ def empty?
100100
end
101101

102102
def enqueued!
103-
return unless self.staged?
103+
return self unless self.writable?
104104

105+
super
106+
rename_to_queued_path
107+
self
108+
end
109+
110+
def rename_to_queued_path
105111
new_chunk_path = self.class.generate_queued_chunk_path(@path, @unique_id)
106112
new_meta_path = new_chunk_path + '.meta'
107113

@@ -139,9 +145,8 @@ def enqueued!
139145

140146
@path = new_chunk_path
141147
@meta_path = new_meta_path
142-
143-
super
144148
end
149+
private :rename_to_queued_path
145150

146151
def close
147152
super
@@ -266,19 +271,27 @@ def write_metadata(update: true)
266271

267272
def file_rename(file, old_path, new_path, callback=nil)
268273
pos = file.pos
274+
setup = ->(f) {
275+
f.set_encoding(Encoding::ASCII_8BIT)
276+
f.sync = true
277+
f.binmode
278+
f.pos = pos
279+
callback.call(f) if callback
280+
}
269281
if Fluent.windows?
270282
file.close
271-
File.rename(old_path, new_path)
272-
file = File.open(new_path, 'rb', @permission)
283+
renamed = false
284+
begin
285+
File.rename(old_path, new_path)
286+
renamed = true
287+
ensure
288+
setup.call(File.open(renamed ? new_path : old_path, 'rb', @permission))
289+
end
273290
else
274291
File.rename(old_path, new_path)
275292
file.reopen(new_path, 'rb')
293+
setup.call(file)
276294
end
277-
file.set_encoding(Encoding::ASCII_8BIT)
278-
file.sync = true
279-
file.binmode
280-
file.pos = pos
281-
callback.call(file) if callback
282295
end
283296

284297
def create_new_chunk(path, perm)

‎lib/fluent/plugin/buffer/file_single_chunk.rb‎

Lines changed: 23 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -93,8 +93,14 @@ def empty?
9393
end
9494

9595
def enqueued!
96-
return unless self.staged?
96+
return self unless self.writable?
9797

98+
super
99+
rename_to_queued_path
100+
self
101+
end
102+
103+
def rename_to_queued_path
98104
new_chunk_path = self.class.generate_queued_chunk_path(@path, @unique_id)
99105

100106
begin
@@ -114,9 +120,8 @@ def enqueued!
114120
end
115121

116122
@path = new_chunk_path
117-
118-
super
119123
end
124+
private :rename_to_queued_path
120125

121126
def close
122127
super
@@ -225,19 +230,27 @@ def restore_size(chunk_format)
225230

226231
def file_rename(file, old_path, new_path, callback = nil)
227232
pos = file.pos
233+
setup = ->(f) {
234+
f.set_encoding(Encoding::ASCII_8BIT)
235+
f.sync = true
236+
f.binmode
237+
f.pos = pos
238+
callback.call(f) if callback
239+
}
228240
if Fluent.windows?
229241
file.close
230-
File.rename(old_path, new_path)
231-
file = File.open(new_path, 'rb', @permission)
242+
renamed = false
243+
begin
244+
File.rename(old_path, new_path)
245+
renamed = true
246+
ensure
247+
setup.call(File.open(renamed ? new_path : old_path, 'rb', @permission))
248+
end
232249
else
233250
File.rename(old_path, new_path)
234251
file.reopen(new_path, 'rb')
252+
setup.call(file)
235253
end
236-
file.set_encoding(Encoding::ASCII_8BIT)
237-
file.sync = true
238-
file.binmode
239-
file.pos = pos
240-
callback.call(file) if callback
241254
end
242255

243256
ESCAPE_REGEXP = /[^-_.a-zA-Z0-9]/n

‎test/plugin/test_buf_file.rb‎

Lines changed: 139 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -444,6 +444,145 @@ def write_metadata(path, chunk_id, metadata, size, ctime, mtime)
444444
end
445445
end
446446

447+
sub_test_case 'unstaged chunks created by write_step_by_step' do
448+
setup do
449+
@bufdir = File.expand_path('../../tmp/buffer_file_unstaged', __FILE__)
450+
@bufpath = File.join(@bufdir, 'testbuf.*.log')
451+
FileUtils.rm_r @bufdir if File.exist?(@bufdir)
452+
453+
Fluent::Test.setup
454+
@d = FluentPluginFileBufferTest::DummyOutputPlugin.new
455+
@p = Fluent::Plugin::FileBuffer.new
456+
@p.owner = @d
457+
@p.configure(config_element('buffer', '', {'path' => @bufpath, 'chunk_limit_size' => 1024}))
458+
@p.start
459+
end
460+
461+
teardown do
462+
if @p
463+
@p.stop unless @p.stopped?
464+
@p.before_shutdown unless @p.before_shutdown?
465+
@p.shutdown unless @p.shutdown?
466+
@p.after_shutdown unless @p.after_shutdown?
467+
@p.close unless @p.closed?
468+
@p.terminate unless @p.terminated?
469+
end
470+
FileUtils.rm_r @bufdir if File.exist?(@bufdir)
471+
end
472+
473+
def write_with_single_record_overflow(m)
474+
@p.write({m => ["a" * 600]})
475+
@p.write({m => ["b" * 600, "c" * 600, "d" * 380]})
476+
end
477+
478+
test 'enqueued unstaged chunks are queued and are not re-staged' do
479+
m = @p.metadata
480+
write_with_single_record_overflow(m)
481+
482+
assert_equal [:queued, :queued, :queued], @p.queue.map(&:state)
483+
assert_equal [600, 600, 980], @p.queue.map(&:bytesize)
484+
assert_true @p.queue.all? { |c| File.basename(c.path).start_with?('testbuf.q') }
485+
assert_true @p.stage.empty?
486+
assert_equal 0, @p.stage_size
487+
assert_equal 2180, @p.queue_size
488+
end
489+
490+
test 'enqueued unstaged chunks are resumed as queued chunks' do
491+
m = @p.metadata
492+
write_with_single_record_overflow(m)
493+
@p.stop; @p.before_shutdown; @p.shutdown; @p.after_shutdown; @p.close; @p.terminate
494+
495+
@p = Fluent::Plugin::FileBuffer.new
496+
@p.owner = @d
497+
@p.configure(config_element('buffer', '', {'path' => @bufpath, 'chunk_limit_size' => 1024}))
498+
@p.start
499+
500+
assert_true @p.stage.empty?
501+
assert_equal [600, 600, 980], @p.queue.map(&:bytesize).sort
502+
assert_equal 0, @p.stage_size
503+
assert_equal 2180, @p.queue_size
504+
end
505+
506+
test 'committed data is kept when renaming an enqueued unstaged chunk fails' do
507+
m = @p.metadata
508+
@p.write({m => ["a" * 600]})
509+
first = @p.stage[m]
510+
first.define_singleton_method(:file_rename) { |*| raise Errno::EACCES, "injected" }
511+
512+
assert_nothing_raised { @p.write({m => ["b" * 600, "c" * 600, "d" * 380]}) }
513+
assert_equal 1, $log.out.logs.count { |l| l.include?("error occurs in enqueueing a chunk") }
514+
515+
assert_true @p.queue.any? { |c| c.equal?(first) }
516+
assert_equal :queued, first.state
517+
assert_equal 600, first.bytesize
518+
assert_true File.exist?(first.path)
519+
assert_equal @p.queue.sum(&:bytesize), @p.queue_size
520+
521+
read = +""
522+
while (chunk = @p.dequeue_chunk)
523+
assert_nothing_raised { read << chunk.read }
524+
@p.purge_chunk(chunk.unique_id)
525+
end
526+
assert_true read.include?("a" * 600)
527+
assert_true read.include?("b" * 600)
528+
assert_equal 0, @p.queue_size
529+
end
530+
531+
test 'all data is resumed after a rename failure and a restart' do
532+
m = @p.metadata
533+
@p.write({m => ["a" * 600]})
534+
first = @p.stage[m]
535+
first.define_singleton_method(:file_rename) { |*| raise Errno::EACCES, "injected" }
536+
@p.write({m => ["b" * 600, "c" * 600, "d" * 380]})
537+
@p.write({m => ["e" * 100]})
538+
@p.stop; @p.before_shutdown; @p.shutdown; @p.after_shutdown; @p.close; @p.terminate
539+
540+
@p = Fluent::Plugin::FileBuffer.new
541+
@p.owner = @d
542+
@p.configure(config_element('buffer', '', {'path' => @bufpath, 'chunk_limit_size' => 1024}))
543+
@p.start
544+
545+
chunks = @p.stage.values + @p.queue
546+
assert_equal 4, chunks.size
547+
assert_equal 600 + 1580 + 100, chunks.sum(&:bytesize)
548+
assert_equal chunks.sum(&:bytesize), @p.stage_size + @p.queue_size
549+
end
550+
551+
test 'a staged chunk stays queued with consistent gauges when its rename fails' do
552+
m = @p.metadata
553+
@p.write({m => ["a" * 600]})
554+
chunk = @p.stage[m]
555+
chunk.define_singleton_method(:file_rename) { |*| raise Errno::EACCES, "injected" }
556+
557+
assert_raise(RuntimeError) { @p.enqueue_chunk(m) }
558+
559+
assert_true @p.stage.empty?
560+
assert_equal [chunk], @p.queue
561+
assert_equal :queued, chunk.state
562+
assert_equal 0, @p.stage_size
563+
assert_equal 600, @p.queue_size
564+
565+
assert_nothing_raised { @p.dequeue_chunk.read }
566+
@p.purge_chunk(chunk.unique_id)
567+
assert_equal 0, @p.queue_size
568+
end
569+
570+
test 'each chunk is flushed once and the gauges return to zero' do
571+
m = @p.metadata
572+
write_with_single_record_overflow(m)
573+
@p.write({m => ["e" * 600, "f" * 300]})
574+
assert_equal @p.queue.size, @p.queue.map(&:object_id).uniq.size
575+
576+
while (chunk = @p.dequeue_chunk)
577+
assert_nothing_raised { chunk.read }
578+
@p.purge_chunk(chunk.unique_id)
579+
end
580+
581+
assert_equal 0, @p.queue_size
582+
assert_equal @p.stage.values.sum(&:bytesize), @p.stage_size
583+
end
584+
end
585+
447586
sub_test_case 'configured with system root directory and plugin @id' do
448587
setup do
449588
@root_dir = File.expand_path('../../tmp/buffer_file_root', __FILE__)

0 commit comments

Comments
 (0)