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
42 changes: 36 additions & 6 deletions internal/tailer/logstream/filestream.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,11 @@ import (
// fileTruncates counts the truncations of a file stream.
var fileTruncates = expvar.NewMap("file_truncates_total")

// newLogFileReader returns the reader a file stream reads log data from.
// Tests replace it to inject read errors that a real filesystem won't produce
// on demand, like ESTALE.
var newLogFileReader = func(f *os.File) io.Reader { return f }

// fileStream streams log lines from a regular file on the file system. These
// log files are appended to by another process, and are either rotated or
// truncated by that (or yet another) process. Rotation implies that a new
Expand All @@ -40,7 +45,7 @@ type fileStream struct {
}

// newFileStream creates a new log stream from a regular file.
func newFileStream(ctx context.Context, wg *sync.WaitGroup, waker waker.Waker, pathname string, fi os.FileInfo, oneShot OneShotMode) (LogStream, error) {
func newFileStream(ctx context.Context, wg *sync.WaitGroup, waker waker.Waker, pathname string, oneShot OneShotMode) (LogStream, error) {
ctx, cancel := context.WithCancel(ctx)
fs := &fileStream{
cancel: cancel,
Expand All @@ -52,20 +57,31 @@ func newFileStream(ctx context.Context, wg *sync.WaitGroup, waker waker.Waker, p
}
// Stream from the start of the file when in one shot mode.
streamFromStart := oneShot == OneShotEnabled
if err := fs.stream(ctx, wg, waker, fi, oneShot, streamFromStart); err != nil {
if err := fs.stream(ctx, wg, waker, oneShot, streamFromStart); err != nil {
return nil, err
}
return fs, nil
}

func (fs *fileStream) stream(ctx context.Context, wg *sync.WaitGroup, waker waker.Waker, fi os.FileInfo, oneShot OneShotMode, streamFromStart bool) error {
func (fs *fileStream) stream(ctx context.Context, wg *sync.WaitGroup, waker waker.Waker, oneShot OneShotMode, streamFromStart bool) error {
fd, err := os.OpenFile(fs.pathname, os.O_RDONLY, 0o600)
if err != nil {
logErrors.Add(fs.sourcename, 1)
return err
}
logOpens.Add(fs.sourcename, 1)
glog.V(2).Infof("stream(%s): opened new file", fs.sourcename)
// Rotation is detected by comparing against the file this stream is
// reading, not against the file the caller last saw.
fi, err := fd.Stat()
if err != nil {
logErrors.Add(fs.sourcename, 1)
if err := fd.Close(); err != nil {
logErrors.Add(fs.sourcename, 1)
glog.Infof("stream(%s): closing file: %v", fs.sourcename, err)
}
return err
}
if !streamFromStart {
// Normal operation for first stream is to ignore the past, and seek to
// EOF immediately to start tailing.
Expand All @@ -80,7 +96,7 @@ func (fs *fileStream) stream(ctx context.Context, wg *sync.WaitGroup, waker wake
glog.V(2).Infof("stream(%s): seeked to end", fs.sourcename)
}

lr := NewLineReader(fs.sourcename, fs.lines, fd, defaultReadBufferSize, fs.cancel)
lr := NewLineReader(fs.sourcename, fs.lines, newLogFileReader(fd), defaultReadBufferSize, fs.cancel)

started := make(chan struct{})
var total int
Expand Down Expand Up @@ -118,9 +134,23 @@ func (fs *fileStream) stream(ctx context.Context, wg *sync.WaitGroup, waker wake
// retryable.
if errors.Is(err, syscall.ESTALE) {
glog.Infof("stream(%s): reopening stream due to %s", fs.sourcename, err)
// The reopened file can be stale as well, so wait for the
// next wakeup instead of reopening in a tight loop.
select {
case <-ctx.Done():
lr.Finish(ctx)
close(fs.lines)
return
case <-waker.Wake():
}
// streamFromStart always true on a stream reopen
if nerr := fs.stream(ctx, wg, waker, fi, oneShot, true); nerr != nil {
if nerr := fs.stream(ctx, wg, waker, oneShot, true); nerr != nil {
glog.Infof("stream(%s): new stream: %v", fs.sourcename, nerr)
// Nothing owns the lines channel now, so close it to
// tell the Tailer this source is finished and can be
// tailed again.
lr.Finish(ctx)
close(fs.lines)
}
// Close this stream.
return
Expand Down Expand Up @@ -155,7 +185,7 @@ func (fs *fileStream) stream(ctx context.Context, wg *sync.WaitGroup, waker wake
if !os.SameFile(fi, newfi) {
glog.V(2).Infof("stream(%s): adding a new file routine", fs.sourcename)
// Stream from start always true on a stream reopen
if err := fs.stream(ctx, wg, waker, newfi, oneShot, true); err != nil {
if err := fs.stream(ctx, wg, waker, oneShot, true); err != nil {
glog.Info("stream(%s): new stream: %v", fs.sourcename, err)
}
// We're at EOF so there's nothing left to read here.
Expand Down
136 changes: 136 additions & 0 deletions internal/tailer/logstream/filestream_internal_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
// Copyright 2026 Google Inc. All Rights Reserved.
// This file is available under the Apache license.

package logstream

import (
"context"
"io"
"os"
"path/filepath"
"sync"
"syscall"
"testing"

"github.com/google/mtail/internal/testutil"
"github.com/google/mtail/internal/waker"
)

// staleReader returns ESTALE on every read, counting them, until `limit` reads
// have been made, after which it returns EOF. The limit stops a stream that
// reopens without waiting from spinning until the test deadline. Every read
// announces itself on `entered` and then waits for `release` to be closed, so
// a test can act while the stream is blocked in a read.
type staleReader struct {
limit int
entered chan struct{}
release chan struct{}

mu sync.Mutex
count int
}

func (r *staleReader) Read(_ []byte) (int, error) {
if r.entered != nil {
r.entered <- struct{}{}
}
if r.release != nil {
<-r.release
}
r.mu.Lock()
defer r.mu.Unlock()
r.count++
if r.count > r.limit {
return 0, io.EOF
}
return 0, &os.PathError{Op: "read", Path: "log", Err: syscall.ESTALE}
}

func (r *staleReader) reads() int {
r.mu.Lock()
defer r.mu.Unlock()
return r.count
}

// injectReader makes every file stream read from `r` instead of the file it
// opened, for the duration of the test.
func injectReader(t *testing.T, r io.Reader) {
t.Helper()
saved := newLogFileReader
newLogFileReader = func(_ *os.File) io.Reader { return r }
t.Cleanup(func() { newLogFileReader = saved })
}

// testLogFile creates an empty log to stream from, returning its pathname.
func testLogFile(t *testing.T) string {
t.Helper()
name := filepath.Join(testutil.TestTempDir(t), "log")
f := testutil.TestOpenFile(t, name)
testutil.FatalIfErr(t, f.Close())
return name
}

func TestFileStreamStaleWaitsBetweenReopens(t *testing.T) {
var wg sync.WaitGroup

r := &staleReader{limit: 100}
injectReader(t, r)
name := testLogFile(t)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
w, awaken := waker.NewTest(ctx, 1, "stream")

fs, err := New(ctx, &wg, w, name, OneShotDisabled)
testutil.FatalIfErr(t, err)
go func() {
for range fs.Lines() {
}
}()

// The first read happens before the stream waits, so there is one more
// read than there are wakeups.
for i := 0; i < 3; i++ {
awaken(1, 1)
}

if got, want := r.reads(), 4; got != want {
t.Errorf("stale reads: got %d, want %d (the stream is reopening without waiting)", got, want)
}

cancel()
wg.Wait()
}

func TestFileStreamStaleReopenFailureEndsStream(t *testing.T) {
var wg sync.WaitGroup

r := &staleReader{limit: 100, entered: make(chan struct{}), release: make(chan struct{})}
injectReader(t, r)
name := testLogFile(t)

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
w, awaken := waker.NewTest(ctx, 1, "stream")

fs, err := New(ctx, &wg, w, name, OneShotDisabled)
testutil.FatalIfErr(t, err)

// Take the log away while the stream is blocked in its first read, so
// that the reopen after the stale read is certain to fail.
<-r.entered
testutil.FatalIfErr(t, os.Remove(name))
close(r.release)

go awaken(1, 0)
wg.Wait()

select {
case line, ok := <-fs.Lines():
if ok {
t.Errorf("expecting the stream to be complete because the reopen failed, got %v", line)
}
default:
t.Error("expecting the lines channel to be closed after a failed reopen")
}
}
2 changes: 1 addition & 1 deletion internal/tailer/logstream/logstream.go
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,7 @@ func New(ctx context.Context, wg *sync.WaitGroup, waker waker.Waker, pathname st
}
switch m := fi.Mode(); {
case m.IsRegular():
return newFileStream(ctx, wg, waker, path, fi, oneShot)
return newFileStream(ctx, wg, waker, path, oneShot)
case m&os.ModeType == os.ModeNamedPipe:
return newFifoStream(ctx, wg, waker, path, fi)
// TODO(jaq): in order to listen on an existing socket filepath, we must unlink and recreate it
Expand Down
Loading