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
8 changes: 6 additions & 2 deletions internal/state/range.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,8 +25,7 @@ const (
// Range represents a pair of offsets [start, end).
type Range struct {
start, end uint64
// seq handles file truncation, when a file is truncated, we increase the seq
seq uint64
seq uint64
}

var _ encoding.TextMarshaler = (*Range)(nil)
Expand Down Expand Up @@ -74,6 +73,11 @@ func (r Range) EndOffsetInt64() int64 {
return convertInt64(r.end)
}

// SequenceNumber handles file truncation detection. When a file is truncated, the sequence number will increase.
func (r Range) SequenceNumber() uint64 {
return r.seq
}

// Shift moves the previous end to the start and sets the new end. If the new end is before the previous one, it resets
// the range to [0, newEnd) and increments the sequence number.
func (r *Range) Shift(newEnd uint64) {
Expand Down
28 changes: 14 additions & 14 deletions internal/state/range_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,9 @@ func TestRange(t *testing.T) {
r.Set(10, 20)
assert.Equal(t, uint64(10), r.StartOffset())
assert.Equal(t, uint64(20), r.EndOffset())
assert.Equal(t, uint64(0), r.seq)
assert.Equal(t, uint64(0), r.SequenceNumber())
r.Set(5, 30)
assert.Equal(t, uint64(1), r.seq)
assert.Equal(t, uint64(1), r.SequenceNumber())
})
t.Run("SetGet/Int64", func(t *testing.T) {
var r Range
Expand All @@ -37,21 +37,21 @@ func TestRange(t *testing.T) {
t.Run("Shift", func(t *testing.T) {
var r Range
r.ShiftInt64(100)
assert.Equal(t, uint64(0), r.start)
assert.Equal(t, uint64(100), r.end)
assert.Equal(t, uint64(0), r.seq)
assert.Equal(t, uint64(0), r.StartOffset())
assert.Equal(t, uint64(100), r.EndOffset())
assert.Equal(t, uint64(0), r.SequenceNumber())
r.ShiftInt64(200)
assert.Equal(t, uint64(100), r.start)
assert.Equal(t, uint64(200), r.end)
assert.Equal(t, uint64(0), r.seq)
assert.Equal(t, uint64(100), r.StartOffset())
assert.Equal(t, uint64(200), r.EndOffset())
assert.Equal(t, uint64(0), r.SequenceNumber())
r.ShiftInt64(50)
assert.Equal(t, uint64(0), r.start)
assert.Equal(t, uint64(50), r.end)
assert.Equal(t, uint64(1), r.seq)
assert.Equal(t, uint64(0), r.StartOffset())
assert.Equal(t, uint64(50), r.EndOffset())
assert.Equal(t, uint64(1), r.SequenceNumber())
r.ShiftInt64(-1)
assert.Equal(t, uint64(0), r.start)
assert.Equal(t, uint64(50), r.end)
assert.Equal(t, uint64(1), r.seq)
assert.Equal(t, uint64(0), r.StartOffset())
assert.Equal(t, uint64(50), r.EndOffset())
assert.Equal(t, uint64(1), r.SequenceNumber())
})
t.Run("Contains", func(t *testing.T) {
r1 := Range{start: 0, end: 10}
Expand Down
41 changes: 41 additions & 0 deletions internal/state/statetest/manager.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
// SPDX-License-Identifier: MIT

package statetest

import "github.com/aws/amazon-cloudwatch-agent/internal/state"

type FileRangeManagerSink struct {
base state.FileRangeManager
sink state.RangeList
}

var _ state.FileRangeManager = (*FileRangeManagerSink)(nil)

func NewFileManagerSink(base state.FileRangeManager) *FileRangeManagerSink {
return &FileRangeManagerSink{
base: base,
sink: make(state.RangeList, 0),
}
}

func (f *FileRangeManagerSink) ID() string {
return f.base.ID()
}

func (f *FileRangeManagerSink) Enqueue(r state.Range) {
f.sink = append(f.sink, r)
f.base.Enqueue(r)
}

func (f *FileRangeManagerSink) Restore() (state.RangeList, error) {
return f.base.Restore()
}

func (f *FileRangeManagerSink) Run(ch state.Notification) {
f.base.Run(ch)
}

func (f *FileRangeManagerSink) GetSink() state.RangeList {
return f.sink
}
47 changes: 47 additions & 0 deletions internal/state/statetest/manager_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved.
// SPDX-License-Identifier: MIT

package statetest

import (
"sync"
"testing"
"time"

"github.com/stretchr/testify/assert"

"github.com/aws/amazon-cloudwatch-agent/internal/state"
)

func TestNewFileManagerSink(t *testing.T) {
tmpDir := t.TempDir()
sink := NewFileManagerSink(state.NewFileRangeManager(state.ManagerConfig{
StateFileDir: tmpDir,
Name: "sink",
MaxPersistedItems: 1,
}))
done := make(chan struct{})
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
sink.Run(state.Notification{Done: done})
}()
assert.Equal(t, "sink", sink.ID())
sink.Enqueue(state.NewRange(0, 5))
sink.Enqueue(state.NewRange(5, 10))
time.Sleep(time.Millisecond)
close(done)
wg.Wait()

got, err := sink.Restore()
assert.NoError(t, err)
assert.Equal(t, state.RangeList{
state.NewRange(0, 10),
}, got)

assert.Equal(t, state.RangeList{
state.NewRange(0, 5),
state.NewRange(5, 10),
}, sink.GetSink())
}
6 changes: 5 additions & 1 deletion plugins/inputs/logfile/logfile.go
Original file line number Diff line number Diff line change
Expand Up @@ -203,8 +203,11 @@ func (t *LogFile) FindLogSrc() []logs.LogSrc {
seekFile = &tail.SeekInfo{Whence: io.SeekEnd, Offset: 0}
}

var initialStateOffset int64
var gapsToRead state.RangeList
if !restored.OnlyUseMaxOffset() {
if restored.OnlyUseMaxOffset() {
initialStateOffset = restored.Last().EndOffsetInt64()
} else {
gapsToRead = state.InvertRanges(restored)
}
isutf16 := false
Expand Down Expand Up @@ -257,6 +260,7 @@ func (t *LogFile) FindLogSrc() []logs.LogSrc {
groupName, streamName,
t.Destination,
stateManager,
initialStateOffset,
fileconfig.LogGroupClass,
fileconfig.FilePath,
tailer,
Expand Down
80 changes: 78 additions & 2 deletions plugins/inputs/logfile/logfile_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ import (
"golang.org/x/text/transform"

"github.com/aws/amazon-cloudwatch-agent/internal/state"
"github.com/aws/amazon-cloudwatch-agent/internal/state/statetest"
"github.com/aws/amazon-cloudwatch-agent/logs"
)

Expand Down Expand Up @@ -571,7 +572,7 @@ func TestLogsMultilineTimeout(t *testing.T) {
func TestLogsFileTruncate(t *testing.T) {
multilineWaitPeriod = 10 * time.Millisecond
lineBeforeFileTruncate := "lineBeforeFileTruncate"
lineAfterFileTruncate := "lineAfterFileTruncate"
lineAfterFileTruncate := "afterTruncate"

tmpfile, err := createTempFile("", "")
defer os.Remove(tmpfile.Name())
Expand All @@ -589,9 +590,17 @@ func TestLogsFileTruncate(t *testing.T) {
}

lsrc := lsrcs[0]
ts, ok := lsrc.(*tailerSrc)
assert.True(t, ok)
sink := statetest.NewFileManagerSink(ts.stateManager)
ts.stateManager = sink

evts := make(chan logs.LogEvent)
lsrc.SetOutput(func(e logs.LogEvent) {
evts <- e
if e != nil {
e.Done()
evts <- e
}
})

go func() {
Expand Down Expand Up @@ -620,6 +629,73 @@ func TestLogsFileTruncate(t *testing.T) {

lsrc.Stop()
tt.Stop()

got := sink.GetSink()
assert.Len(t, got, 2)
assert.EqualValues(t, 0, got[0].SequenceNumber())
assert.EqualValues(t, 1, got.Last().SequenceNumber())
}

func TestLogsFileTruncateRestart(t *testing.T) {
logEntryString := "postTruncateRestart"
multilineWaitPeriod = 10 * time.Millisecond

tmpfile, err := createTempFile("", "")
defer os.Remove(tmpfile.Name())
require.NoError(t, err)

stateDir, err := os.MkdirTemp("", "state")
require.NoError(t, err)
defer os.Remove(stateDir)

stateFileName := state.FilePath(stateDir, tmpfile.Name())
stateFile, err := os.OpenFile(stateFileName, os.O_RDWR|os.O_CREATE|os.O_EXCL, 0600)
require.NoError(t, err)
defer os.Remove(stateFileName)

_, err = stateFile.WriteString("1000")
require.NoError(t, err)

_, err = tmpfile.WriteString(logEntryString + "\n")
require.NoError(t, err)

tt := NewLogFile()
tt.FileStateFolder = stateDir
tt.Log = TestLogger{t}
tt.FileConfig = []FileConfig{{FilePath: tmpfile.Name(), FromBeginning: true}}
tt.FileConfig[0].init()
tt.started = true

lsrcs := tt.FindLogSrc()
if len(lsrcs) != 1 {
t.Fatalf("%v log src was returned when 1 should be available", len(lsrcs))
}

lsrc := lsrcs[0]
ts, ok := lsrc.(*tailerSrc)
assert.True(t, ok)
sink := statetest.NewFileManagerSink(ts.stateManager)
ts.stateManager = sink

evts := make(chan logs.LogEvent)
lsrc.SetOutput(func(e logs.LogEvent) {
if e != nil {
e.Done()
evts <- e
}
})

e := <-evts
if e.Message() != logEntryString {
t.Errorf("Wrong log found after offset: \n%v\nExpecting:\n%v\n", e.Message(), logEntryString)
}

lsrc.Stop()
tt.Stop()

got := sink.GetSink()
assert.Len(t, got, 1)
assert.EqualValues(t, 1, got.Last().SequenceNumber())
}

func TestLogsFileWithOffset(t *testing.T) {
Expand Down
30 changes: 18 additions & 12 deletions plugins/inputs/logfile/tailersrc.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,18 +59,19 @@ func (le LogEvent) RangeQueue() state.FileRangeQueue {
}

type tailerSrc struct {
group string
stream string
class string
fileGlobPath string
destination string
stateManager state.FileRangeManager
tailer *tail.Tail
autoRemoval bool
timestampFn func(string) (time.Time, string)
enc encoding.Encoding
maxEventSize int
retentionInDays int
group string
stream string
class string
fileGlobPath string
destination string
stateManager state.FileRangeManager
initialStateOffset int64
tailer *tail.Tail
autoRemoval bool
timestampFn func(string) (time.Time, string)
enc encoding.Encoding
maxEventSize int
retentionInDays int

outputFn func(logs.LogEvent)
isMLStart func(string) bool
Expand All @@ -89,6 +90,7 @@ var _ logs.LogSrc = (*tailerSrc)(nil)
func NewTailerSrc(
group, stream, destination string,
stateManager state.FileRangeManager,
initialStateOffset int64,
logClass, fileGlobPath string,
tailer *tail.Tail,
autoRemoval bool,
Expand All @@ -105,6 +107,7 @@ func NewTailerSrc(
stream: stream,
destination: destination,
stateManager: stateManager,
initialStateOffset: initialStateOffset,
class: logClass,
fileGlobPath: fileGlobPath,
tailer: tailer,
Expand Down Expand Up @@ -195,6 +198,9 @@ func (ts *tailerSrc) runTail() {
var msgBuf bytes.Buffer
var cnt int
fo := state.Range{}
if ts.initialStateOffset > 0 {
fo.SetInt64(0, ts.initialStateOffset)
}
ignoreUntilNextEvent := false

for {
Expand Down
12 changes: 9 additions & 3 deletions plugins/inputs/logfile/tailersrc_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -72,8 +72,11 @@ func TestTailerSrc(t *testing.T) {
})

ts := NewTailerSrc(
"groupName", "streamName",
"destination", m,
"groupName",
"streamName",
"destination",
m,
0,
util.InfrequentAccessLogGroupClass,
"tailsrctest-*.log",
tailer,
Expand Down Expand Up @@ -181,9 +184,11 @@ func TestEventDoneCallback(t *testing.T) {
})

ts := NewTailerSrc(
"groupName", "streamName",
"groupName",
"streamName",
"destination",
m,
0,
util.InfrequentAccessLogGroupClass,
"tailsrctest-*.log",
tailer,
Expand Down Expand Up @@ -414,6 +419,7 @@ func setupTailer(t *testing.T, multiLineFn func(string) bool, maxEventSize int,
t.Name(),
"destination",
m,
0,
util.InfrequentAccessLogGroupClass,
"tailsrctest-*.log",
tailer,
Expand Down
Loading
Loading