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
60 changes: 38 additions & 22 deletions lib/fluent/plugin/in_tail.rb
Original file line number Diff line number Diff line change
Expand Up @@ -307,7 +307,7 @@ def start_watchers(paths)
pe = @pf[path]
if @read_from_head && pe.read_inode.zero?
begin
pe.update(Fluent::FileWrapper.stat(path).ino, 0)
pe.update(Fluent::FileWrapper.stat(path).ino, 0, 0)
rescue Errno::ENOENT
$log.warn "#{path} not found. Continuing without tailing it."
end
Expand Down Expand Up @@ -372,7 +372,7 @@ def detach_watcher(tw, close_io = true)
tw.close if close_io
flush_buffer(tw)
if tw.unwatched && @pf
@pf[tw.path].update_pos(PositionFile::UNWATCHED_POSITION)
@pf[tw.path].update_pos(PositionFile::UNWATCHED_POSITION, 0)
end
end

Expand Down Expand Up @@ -580,19 +580,19 @@ def on_rotate(stat)
# in either case of a and b, seek to the saved position
# c) file was once renamed, truncated and then backed
# in this case, consider it truncated
@pe.update(inode, 0) if fsize < @pe.read_pos
@pe.update(inode, 0, fsize) if fsize < @pe.read_pos
elsif last_inode != 0
# this is FilePositionEntry and fluentd once started.
# read data from the head of the rotated file.
# logs never duplicate because this file is a rotated new file.
@pe.update(inode, 0)
@pe.update(inode, 0, fize)
else
# this is MemoryPositionEntry or this is the first time fluentd started.
# seek to the end of the any files.
# logs may duplicate without this seek because it's not sure the file is
# existent file or rotated new file.
pos = @read_from_head ? 0 : fsize
@pe.update(inode, pos)
@pe.update(inode, pos, fsize)
end
@io_handler = IOHandler.new(self, &method(:wrap_receive_lines))
else
Expand All @@ -604,10 +604,10 @@ def on_rotate(stat)
if stat
inode = stat.ino
if inode == @pe.read_inode # truncated
@pe.update_pos(0)
@pe.update_pos(0, stat.size)
@io_handler.close
elsif !@io_handler.opened? # There is no previous file. Reuse TailWatcher
@pe.update(inode, 0)
@pe.update(inode, 0, stat.size)
else # file is rotated and new file found
watcher_needs_update = true
# Handle the old log file before renewing TailWatcher [fluentd#1055]
Expand All @@ -634,7 +634,7 @@ def on_rotate(stat)
def swap_state(pe)
# Use MemoryPositionEntry for rotated file temporary
mpe = MemoryPositionEntry.new
mpe.update(pe.read_inode, pe.read_pos)
mpe.update(pe.read_inode, pe.read_pos, pe.read_size)
@pe = mpe
pe # This pe will be updated in on_rotate after TailWatcher is initialized
end
Expand Down Expand Up @@ -766,7 +766,7 @@ def handle_notify

unless @lines.empty?
if @receive_lines.call(@lines)
@watcher.pe.update_pos(io.pos - @fifo.bytesize)
@watcher.pe.update_pos(io.pos - @fifo.bytesize, io.size)
@lines.clear
else
read_more = false
Expand Down Expand Up @@ -912,10 +912,10 @@ def [](path)

@file_mutex.synchronize {
@file.pos = @last_pos
@file.write "#{path}\t0000000000000000\t0000000000000000\n"
@file.write "#{path}\t0000000000000000\t0000000000000000\t0000000000000000\n"
seek = @last_pos + path.bytesize + 1
@last_pos = @file.pos
@map[path] = FilePositionEntry.new(@file, @file_mutex, seek, 0, 0)
@map[path] = FilePositionEntry.new(@file, @file_mutex, seek, 0, 0, 0)
}
end

Expand All @@ -926,16 +926,17 @@ def self.parse(file)
map = {}
file.pos = 0
file.each_line {|line|
m = /^([^\t]+)\t([0-9a-fA-F]+)\t([0-9a-fA-F]+)/.match(line)
m = /^([^\t]+)\t([0-9a-fA-F]+)\t([0-9a-fA-F]+)[\t]?([0-9a-fA-F]+)?/.match(line)
unless m
$log.warn "Unparsable line in pos_file: #{line}"
next
end
path = m[1]
pos = m[2].to_i(16)
ino = m[3].to_i(16)
size = m[4].to_i(16)
seek = file.pos - line.bytesize + path.bytesize + 1
map[path] = FilePositionEntry.new(file, file_mutex, seek, pos, ino)
map[path] = FilePositionEntry.new(file, file_mutex, seek, pos, size, ino)
}
new(file, file_mutex, map, file.pos)
end
Expand All @@ -944,16 +945,17 @@ def self.parse(file)
def self.compact(file)
file.pos = 0
existent_entries = file.each_line.map { |line|
m = /^([^\t]+)\t([0-9a-fA-F]+)\t([0-9a-fA-F]+)/.match(line)
m = /^([^\t]+)\t([0-9a-fA-F]+)\t([0-9a-fA-F]+)[\t]?([0-9a-fA-F]+)?/.match(line)
unless m
$log.warn "Unparsable line in pos_file: #{line}"
next
end
path = m[1]
pos = m[2].to_i(16)
ino = m[3].to_i(16)
m[4].nil? ? size = 0 : size = m[4].to_i(16)
# 32bit inode converted to 64bit at this phase
pos == UNWATCHED_POSITION ? nil : ("%s\t%016x\t%016x\n" % [path, pos, ino])
pos == UNWATCHED_POSITION ? nil : ("%s\t%016x\t%016x\t%016x\n" % [path, pos, ino, size])
}.compact

file.pos = 0
Expand All @@ -971,35 +973,42 @@ class FilePositionEntry
LN_OFFSET = 33
SIZE = 34

def initialize(file, file_mutex, seek, pos, inode)
def initialize(file, file_mutex, seek, pos, size, inode)
@file = file
@file_mutex = file_mutex
@seek = seek
@pos = pos
@size = size
@inode = inode
end

def update(ino, pos)
def update(ino, pos, size)
@file_mutex.synchronize {
@file.pos = @seek
@file.write "%016x\t%016x" % [pos, ino]
@file.write "%016x\t%016x\t%016x" % [pos, ino, size]
}
@pos = pos
@inode = ino
@size = size
end

def update_pos(pos)
def update_pos(pos, size)
@file_mutex.synchronize {
@file.pos = @seek
@file.write "%016x" % pos
@file.write "%016x\t%016x\t%016x" % [pos, @inode, size]
}
@pos = pos
@size = size
end

def read_inode
@inode
end

def read_size
@size
end

def read_pos
@pos
end
Expand All @@ -1009,21 +1018,28 @@ class MemoryPositionEntry
def initialize
@pos = 0
@inode = 0
@size = 0
end

def update(ino, pos)
def update(ino, pos, size)
@inode = ino
@pos = pos
@size = size
end

def update_pos(pos)
def update_pos(pos, size)
@pos = pos
@size = size
end

def read_pos
@pos
end

def read_size
@size
end

def read_inode
@inode
end
Expand Down
26 changes: 26 additions & 0 deletions test/plugin/test_in_tail.rb
Original file line number Diff line number Diff line change
Expand Up @@ -1050,6 +1050,32 @@ def test_pos_file_dir_creation
d.instance_shutdown
end

def test_pos_file_legacy_file_format_migrate_to_new_pos_file_format
config = config_element("", "", {
"tag" => "tail",
"path" => "#{TMP_DIR}/*.txt",
"format" => "none",
"pos_file" => "#{TMP_DIR}/pos/tail.pos",
"read_from_head" => true,
"refresh_interval" => 1
})
d = create_driver(config, false)
d.run(expect_emits: 1, shutdown: false) do
File.open("#{TMP_DIR}/pos/tail.pos", "wb") { |f|
f.puts "#{TMP_DIR}/tail.txt\t0000000000000000\t0000000000000000\n"
}
File.open("#{TMP_DIR}/tail.txt", "ab") { |f| f.puts "test3\n" }
end
assert_path_exist("#{TMP_DIR}/pos/tail.pos")

first_line_pos_file = ''
File.open("#{TMP_DIR}/pos/tail.pos", "r") { |f| first_line_pos_file = f.readline }
assert_equal(5, /^([^\t]+)\t([0-9a-fA-F]+)\t([0-9a-fA-F]+)\t([0-9a-fA-F]+)/.match(first_line_pos_file).length)

cleanup_directory(TMP_DIR)
d.instance_shutdown
end

def test_z_refresh_watchers
plugin = create_driver(EX_CONFIG, false).instance
sio = StringIO.new
Expand Down