diff --git a/internal/tailer/logstream/filestream.go b/internal/tailer/logstream/filestream.go index f77b1adbf..3deb856b4 100644 --- a/internal/tailer/logstream/filestream.go +++ b/internal/tailer/logstream/filestream.go @@ -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 @@ -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, @@ -52,13 +57,13 @@ 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) @@ -66,6 +71,17 @@ func (fs *fileStream) stream(ctx context.Context, wg *sync.WaitGroup, waker wake } 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. @@ -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 @@ -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 @@ -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. diff --git a/internal/tailer/logstream/filestream_internal_test.go b/internal/tailer/logstream/filestream_internal_test.go new file mode 100644 index 000000000..41cc815ad --- /dev/null +++ b/internal/tailer/logstream/filestream_internal_test.go @@ -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") + } +} diff --git a/internal/tailer/logstream/logstream.go b/internal/tailer/logstream/logstream.go index 7dbd5314a..cd42caa34 100644 --- a/internal/tailer/logstream/logstream.go +++ b/internal/tailer/logstream/logstream.go @@ -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