Skip to content
Merged
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
2 changes: 1 addition & 1 deletion lib/fluent/plugin/in_tail.rb
Original file line number Diff line number Diff line change
Expand Up @@ -647,7 +647,7 @@ def detach_watcher(tw, ino, close_io = true)

tw.close if close_io

if @pf && tw.unwatched && (@follow_inode || !@tails[tw.path])
if @pf && tw.unwatched && (@follow_inodes || !@tails[tw.path])
target_info = TargetInfo.new(tw.path, ino)
@pf.unwatch(target_info)
end
Expand Down
112 changes: 112 additions & 0 deletions test/plugin/test_in_tail.rb
Original file line number Diff line number Diff line change
Expand Up @@ -2093,6 +2093,118 @@ def test_should_delete_file_pos_entry_for_non_existing_file_with_follow_inodes
end
end

# When `follow_inodes` is enabled, position entries are keyed by inode.
# The inode of a rotated file disappears from `@tails` when its TailWatcher
# is detached by `update_watcher`, so `detach_watcher` must unwatch that
# entry even if the same path is still followed by a new watcher.
# Otherwise the entry of the rotated inode is left in the position file.
def test_detach_watcher_unwatches_pos_entry_of_rotated_inode_with_follow_inodes
config = config_element("ROOT", "", {
"path" => "#{@tmp_dir}/tail.txt*",
"pos_file" => "#{@tmp_dir}/tail.pos",
"tag" => "t1",
"format" => "none",
"read_from_head" => "true",
"follow_inodes" => "true",
"rotate_wait" => "1s",
# In order to reproduce the same condition stably, ensure that
# `refresh_watchers` is not called by a timer.
"refresh_interval" => "1h",
# stat_watcher often calls `TailWatcher::on_notify` unexpectedly,
# so disable it in order to reproduce the same condition stably.
"enable_stat_watcher" => "false",
})

path = "#{@tmp_dir}/tail.txt"
Fluent::FileWrapper.open(path, "wb") {|f| f.puts "log1" }

d = create_driver(config, false)
d.run(shutdown: false) do
plugin = d.instance
pf = plugin.instance_variable_get(:@pf)
tails = plugin.instance_variable_get(:@tails)

current_tw = tails[path]
assert_not_nil(current_tw)

# The state right after a rotation: the position entry of the rotated
# inode is still in the position file, and the path is followed by a
# new watcher which has a different inode.
rotated_ino = 0xffffffffffffff01
rotated_target_info = Fluent::Plugin::TailInput::TargetInfo.new(path, rotated_ino)
rotated_pe = pf[rotated_target_info]
rotated_tw = Fluent::Plugin::TailInput::TailWatcher.new(
rotated_target_info, rotated_pe, $log, true, true, nil, nil, nil, nil
)
rotated_tw.unwatched = true

plugin.detach_watcher(rotated_tw, rotated_ino)

assert_equal(
{
map_keys: [current_tw.ino],
unwatched_pe_pos: Fluent::Plugin::TailInput::PositionFile::UNWATCHED_POSITION,
},
{
map_keys: pf.instance_variable_get(:@map).keys,
unwatched_pe_pos: rotated_pe.read_pos,
}
)
end

d.instance_shutdown
end

# With `follow_inodes false`, position entries are keyed by path, so the
# entry must not be unwatched while the path is still followed by a watcher.
# (See https://github.com/fluent/fluentd/pull/4327)
def test_detach_watcher_keeps_pos_entry_of_still_watched_path_without_follow_inodes
config = config_element("ROOT", "", {
"path" => "#{@tmp_dir}/tail.txt*",
"pos_file" => "#{@tmp_dir}/tail.pos",
"tag" => "t1",
"format" => "none",
"read_from_head" => "true",
"follow_inodes" => "false",
"rotate_wait" => "1s",
# In order to reproduce the same condition stably, ensure that
# `refresh_watchers` is not called by a timer.
"refresh_interval" => "1h",
"enable_stat_watcher" => "false",
})

path = "#{@tmp_dir}/tail.txt"
Fluent::FileWrapper.open(path, "wb") {|f| f.puts "log1" }

d = create_driver(config, false)
d.run(shutdown: false) do
plugin = d.instance
pf = plugin.instance_variable_get(:@pf)
tails = plugin.instance_variable_get(:@tails)

current_tw = tails[path]
assert_not_nil(current_tw)

# A watcher which was following the same path with a different inode.
# Since the path is still followed by `current_tw`, the position entry
# of that path must be kept.
old_ino = 0xffffffffffffff01
old_target_info = Fluent::Plugin::TailInput::TargetInfo.new(path, old_ino)
old_pe = pf[old_target_info]
old_tw = Fluent::Plugin::TailInput::TailWatcher.new(
old_target_info, old_pe, $log, true, false, nil, nil, nil, nil
)
old_tw.unwatched = true

plugin.detach_watcher(old_tw, old_ino)

assert_equal([path], pf.instance_variable_get(:@map).keys)
assert_not_equal(Fluent::Plugin::TailInput::PositionFile::UNWATCHED_POSITION, old_pe.read_pos)
end

d.instance_shutdown
end

def test_should_write_latest_offset_after_rotate_wait
config = common_follow_inode_config
Fluent::FileWrapper.open("#{@tmp_dir}/tail.txt", "wb") {|f|
Expand Down
Loading