diff options
Diffstat (limited to 'internal/io/fs/readfile.go')
| -rw-r--r-- | internal/io/fs/readfile.go | 238 |
1 files changed, 52 insertions, 186 deletions
diff --git a/internal/io/fs/readfile.go b/internal/io/fs/readfile.go index dc1d8ea..7969cf0 100644 --- a/internal/io/fs/readfile.go +++ b/internal/io/fs/readfile.go @@ -2,7 +2,6 @@ package fs import ( "bufio" - "bytes" "compress/gzip" "context" "errors" @@ -10,25 +9,18 @@ import ( "io" "os" "strings" - "sync" "time" - "github.com/mimecast/dtail/internal/config" "github.com/mimecast/dtail/internal/io/dlog" - "github.com/mimecast/dtail/internal/io/line" - "github.com/mimecast/dtail/internal/io/pool" - "github.com/mimecast/dtail/internal/lcontext" - "github.com/mimecast/dtail/internal/regex" - - "github.com/DataDog/zstd" ) type readStatus int const ( - nothing readStatus = iota - abortReading readStatus = iota - continueReading readStatus = iota + nothing readStatus = iota + abortReading readStatus = iota + continueReading readStatus = iota + defaultMaxLineLength = 1024 * 1024 ) // Used to tail and filter a local log file. @@ -37,6 +29,8 @@ type readFile struct { stats // Path of log file to tail. filePath string + // Rooted target used for validated server-side re-opens. + validatedTarget *ValidatedReadTarget // The glob identifier of the file. globID string // Channel to send a server message to the dtail client @@ -49,6 +43,8 @@ type readFile struct { seekEOF bool // Warned already about a long line. warnedAboutLongLine bool + // Maximum line length before a line is split. + maxLineLength int } // String returns the string representation of the readFile @@ -72,52 +68,42 @@ func (f readFile) Retry() bool { return f.retry } -// Start tailing a log file. -func (f readFile) Start(ctx context.Context, ltx lcontext.LContext, - lines chan<- *line.Line, re regex.Regex) error { - - reader, fd, err := f.makeReader() - if fd != nil { - defer fd.Close() +func (f *readFile) lineLimit() int { + if f.maxLineLength <= 0 { + return defaultMaxLineLength } - if err != nil { - return err - } - - rawLines := make(chan *bytes.Buffer, 100) - truncate := make(chan struct{}) - - readCtx, readCancel := context.WithCancel(ctx) - var filterWg sync.WaitGroup - filterWg.Add(1) + return f.maxLineLength +} - go f.periodicTruncateCheck(ctx, truncate) - go func() { - f.filter(ctx, ltx, rawLines, lines, re) - filterWg.Done() - // If the filter stopped, make the reader stop too, no need to read - // more data if there is nothing more the filter wants to filter for! - // E.g. it could be that we only want to filter N matches but not more. - readCancel() - }() +func (f *readFile) warnAboutLongLine(ctx context.Context) bool { + if f.warnedAboutLongLine { + return true + } - err = f.read(readCtx, fd, reader, rawLines, truncate) - close(rawLines) - // Filter may sends some data still. So wait until it is done here. - filterWg.Wait() + if f.serverMessages == nil { + f.warnedAboutLongLine = true + return true + } - return err + select { + case f.serverMessages <- dlog.Common.Warn(f.filePath, + "Long log line, splitting into multiple lines") + "\n": + f.warnedAboutLongLine = true + return true + case <-ctx.Done(): + return false + } } -func (f *readFile) makeReader() (*bufio.Reader, *os.File, error) { +func (f *readFile) makeReader() (*bufio.Reader, *os.File, io.Closer, error) { if f.filePath == "" && f.globID == "-" { return f.makePipeReader() } return f.makeFileReader() } -func (f *readFile) makeFileReader() (reader *bufio.Reader, fd *os.File, err error) { - if fd, err = os.Open(f.filePath); err != nil { +func (f *readFile) makeFileReader() (reader *bufio.Reader, fd *os.File, decompressor io.Closer, err error) { + if fd, err = f.openFile(); err != nil { return } @@ -127,21 +113,32 @@ func (f *readFile) makeFileReader() (reader *bufio.Reader, fd *os.File, err erro } } - reader, err = f.makeCompressedFileReader(fd) + reader, decompressor, err = f.makeCompressedFileReader(fd) return } -func (f *readFile) makePipeReader() (*bufio.Reader, *os.File, error) { - return bufio.NewReader(os.Stdin), nil, nil +func (f *readFile) openFile() (*os.File, error) { + if f.validatedTarget != nil { + return f.validatedTarget.Open() + } + return os.Open(f.filePath) } -func (f *readFile) periodicTruncateCheck(ctx context.Context, truncate chan struct{}) { +func (f *readFile) makePipeReader() (*bufio.Reader, *os.File, io.Closer, error) { + return bufio.NewReader(os.Stdin), nil, nil, nil +} + +func (f *readFile) periodicTruncateCheck(ctx context.Context, truncate chan<- struct{}) { + ticker := time.NewTicker(time.Second * 3) + defer ticker.Stop() + for { select { - case <-time.After(time.Second * 3): + case <-ticker.C: select { case truncate <- struct{}{}: case <-ctx.Done(): + return } case <-ctx.Done(): return @@ -149,7 +146,7 @@ func (f *readFile) periodicTruncateCheck(ctx context.Context, truncate chan stru } } -func (f *readFile) makeCompressedFileReader(fd *os.File) (reader *bufio.Reader, err error) { +func (f *readFile) makeCompressedFileReader(fd *os.File) (reader *bufio.Reader, decompressor io.Closer, err error) { switch { case strings.HasSuffix(f.FilePath(), ".gz"): fallthrough @@ -160,82 +157,16 @@ func (f *readFile) makeCompressedFileReader(fd *os.File) (reader *bufio.Reader, if err != nil { return } + decompressor = gzipReader reader = bufio.NewReader(gzipReader) case strings.HasSuffix(f.FilePath(), ".zst"): - dlog.Common.Info(f.FilePath(), "Detected zstd compression format") - reader = bufio.NewReader(zstd.NewReader(fd)) + return f.makeZstdReader(fd) default: reader = bufio.NewReader(fd) } return } -func (f *readFile) read(ctx context.Context, fd *os.File, reader *bufio.Reader, - rawLines chan *bytes.Buffer, truncate <-chan struct{}) error { - - var offset uint64 - message := pool.BytesBuffer.Get().(*bytes.Buffer) - - for { - b, err := reader.ReadByte() - if err != nil { - status, err := f.handleReadError(ctx, err, fd, rawLines, truncate, message) - if abortReading == status { - return err - } - time.Sleep(time.Millisecond * 100) - continue - } - - offset++ - message.WriteByte(b) - - status, newMessage := f.handleReadByte(ctx, b, rawLines, message) - if status == abortReading { - return nil - } - message = newMessage - } -} - -// Filter log lines matching a given regular expression. -func (f *readFile) filter(ctx context.Context, ltx lcontext.LContext, - rawLines <-chan *bytes.Buffer, lines chan<- *line.Line, re regex.Regex) { - - // Do we have any kind of local context settings? If so then run the more complex - // filterWithLContext method. - if ltx.Has() { - // We can not skip transmitting any lines to the client with a local - // grep context specified. - f.canSkipLines = false - f.filterWithLContext(ctx, ltx, rawLines, lines, re) - return - } - - f.filterWithoutLContext(ctx, rawLines, lines, re) -} - -func (f *readFile) transmittable(rawLine *bytes.Buffer, length, capacity int, - re regex.Regex) (*line.Line, bool) { - - newLine := line.Null() - if !re.Match(rawLine.Bytes()) { - f.updateLineNotMatched() - f.updateLineNotTransmitted() - return newLine, false - } - f.updateLineMatched() - - // Can we actually send more messages, channel capacity reached? - if f.canSkipLines && length >= capacity { - f.updateLineNotTransmitted() - return newLine, false - } - f.updateLineTransmitted() - - return line.New(rawLine, f.totalLineCount(), f.transmittedPerc(), f.globID), true -} - // Check wether log file is truncated. Returns nil if not. func (f *readFile) truncated(fd *os.File) (bool, error) { if fd == nil { @@ -250,7 +181,7 @@ func (f *readFile) truncated(fd *os.File) (bool, error) { return true, err } // Can not open file at original path. - pathFd, err := os.Open(f.filePath) + pathFd, err := f.openFile() if err != nil { return true, err } @@ -267,68 +198,3 @@ func (f *readFile) truncated(fd *os.File) (bool, error) { return false, nil } -// Deal with the scenario that nothing could be read from the fd. -func (f *readFile) handleReadError(ctx context.Context, err error, fd *os.File, - rawLines chan *bytes.Buffer, truncate <-chan struct{}, - message *bytes.Buffer) (readStatus, error) { - - if err != io.EOF { - return abortReading, err - } - - select { - case <-truncate: - if isTruncated, err := f.truncated(fd); isTruncated { - return abortReading, err - } - case <-ctx.Done(): - return abortReading, nil - default: - } - - if !f.seekEOF { - dlog.Common.Info(f.FilePath(), "End of file reached") - if len(message.Bytes()) > 0 { - select { - case rawLines <- message: - case <-ctx.Done(): - } - } - return abortReading, nil - } - - return nothing, nil -} - -// Now process the byte we just read from the fd. -func (f *readFile) handleReadByte(ctx context.Context, b byte, - rawLines chan *bytes.Buffer, message *bytes.Buffer) (readStatus, *bytes.Buffer) { - - switch b { - case '\n': - select { - case rawLines <- message: - message = pool.BytesBuffer.Get().(*bytes.Buffer) - f.warnedAboutLongLine = false - case <-ctx.Done(): - return abortReading, message - } - default: - if message.Len() >= config.Server.MaxLineLength { - if !f.warnedAboutLongLine { - f.serverMessages <- dlog.Common.Warn(f.filePath, - "Long log line, splitting into multiple lines") + "\n" - f.warnedAboutLongLine = true - } - message.WriteByte('\n') - select { - case rawLines <- message: - message = pool.BytesBuffer.Get().(*bytes.Buffer) - case <-ctx.Done(): - return abortReading, message - } - } - } - - return nothing, message -} |
