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
19 changes: 17 additions & 2 deletions collector/log.go
Original file line number Diff line number Diff line change
Expand Up @@ -320,11 +320,17 @@ func (c *LogCollectOptions) Collect(m *Manager, cls *models.TiDBCluster) error {
return err
}
for _, f := range c.fileStats[inst.GetHost()] {
source := f.Target
if f.Attributes != nil {
if v, ok := f.Attributes["source"].(string); ok && v != "" {
source = v
}
}
// build checking tasks
t2 = t2.
// check for listening ports
CopyFile(
f.Target,
source,
filepath.Join(c.resultDir, inst.GetHost(), f.Target),
inst.GetHost(),
true,
Expand Down Expand Up @@ -479,9 +485,18 @@ func parseScraperSamples(ctx context.Context, host string) (map[string][]Collect
})
}
for k, v := range s.Log {
target := k
if s.LogTargets != nil {
if original, ok := s.LogTargets[k]; ok && original != "" {
target = original
}
}
stats[host] = append(stats[host], CollectStat{
Target: k,
Target: target,
Size: v,
Attributes: map[string]interface{}{
"source": k,
},
})
}
for k, v := range s.TSDB {
Expand Down
2 changes: 1 addition & 1 deletion collector/log/iterator/log.go
Original file line number Diff line number Diff line change
Expand Up @@ -195,7 +195,7 @@ func (iter *LogIterator) carpet(bpos, epos int64, point time.Time) error {
if item == nil {
continue
}
if item.GetTime().After(point) {
if !item.GetTime().Before(point) {
if item.GetTime().After(iter.end) {
iter.current = nil
iter.nextError = io.EOF
Expand Down
47 changes: 47 additions & 0 deletions collector/log/iterator/log_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,47 @@
package iterator

import (
"io"
"os"
"path/filepath"
"testing"
"time"

"github.com/pingcap/diag/collector/log/parser"
"github.com/stretchr/testify/require"
)

func TestLogIteratorIncludesEntryAtBeginTime(t *testing.T) {
root := t.TempDir()
logDir := filepath.Join(root, "127.0.0.1", "tidb-4000")
require.NoError(t, os.MkdirAll(logDir, 0755))
require.NoError(t, os.WriteFile(filepath.Join(logDir, "tidb.log"), []byte(
`[2026/05/26 09:59:59.000 +08:00] [INFO] [test.go:1] ["before"]`+"\n"+
`[2026/05/26 10:00:00.000 +08:00] [INFO] [test.go:2] ["at start"]`+"\n"+
`[2026/05/26 10:00:01.000 +08:00] [INFO] [test.go:3] ["inside"]`+"\n",
), 0644))

begin := mustParseIteratorTime(t, "2026/05/26 10:00:00.000 +08:00")
end := mustParseIteratorTime(t, "2026/05/26 10:00:01.000 +08:00")
iter, err := New(parser.NewFileWrapper(root, "127.0.0.1", "tidb-4000", "tidb.log"), begin, end)
require.NoError(t, err)
defer iter.Close()

first, err := iter.Next()
require.NoError(t, err)
require.Contains(t, string(first.GetContent()), "at start")

second, err := iter.Next()
require.NoError(t, err)
require.Contains(t, string(second.GetContent()), "inside")

_, err = iter.Next()
require.ErrorIs(t, err, io.EOF)
}

func mustParseIteratorTime(t *testing.T, s string) time.Time {
t.Helper()
ts, err := time.Parse("2006/01/02 15:04:05.000 -07:00", s)
require.NoError(t, err)
return ts
}
179 changes: 173 additions & 6 deletions scraper/log.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ package scraper

import (
"bufio"
"bytes"
"compress/gzip"
"fmt"
"io"
Expand Down Expand Up @@ -46,10 +47,11 @@ func IsValidLogType(logtype string) bool {

// LogScraper scraps log files of components
type LogScraper struct {
Paths []string // paths of log files
Types map[string]bool // log type
Start time.Time // start time
End time.Time // end time
Paths []string // paths of log files
Types map[string]bool // log type
Start time.Time // start time
End time.Time // end time
outputDir string
}

// Scrap implements the Scraper interface
Expand All @@ -60,6 +62,9 @@ func (s *LogScraper) Scrap(result *Sample) error {
if result.LogTypes == nil {
result.LogTypes = make(FileTypes)
}
if result.LogTargets == nil {
result.LogTargets = make(FileTypes)
}
fileList := make([]string, 0)

// extend all file paths
Expand All @@ -81,8 +86,23 @@ func (s *LogScraper) Scrap(result *Sample) error {

logtype, in, err := getLogType(fp, fi, s.Start, s.End)
if s.Types[logtype] && in {
result.Log[fp] = fi.Size()
result.LogTypes[fp] = logtype
target := fp
size := fi.Size()
if canFilterLogFile(fp, logtype) {
filtered, filteredSize, ok, err := s.filterLogFile(fp, logtype)
if err != nil {
fmt.Fprintf(os.Stderr, "error filtering %s: %s\n", fi.Name(), err)
continue
}
if !ok {
continue
}
target = filtered
size = filteredSize
}
result.Log[target] = size
result.LogTargets[target] = fp
result.LogTypes[target] = logtype
}
if err != nil {
fmt.Fprintf(os.Stderr, "error checking %s: %s\n", fi.Name(), err)
Expand Down Expand Up @@ -152,6 +172,153 @@ func getLogType(fpath string, fi fs.FileInfo, start, end time.Time) (logtype str
return LogTypeUnknown, true, nil
}

func canFilterLogFile(fpath, logtype string) bool {
if strings.Contains(filepath.Base(fpath), "stderr") {
return false
}
return logtype == LogTypeStd || logtype == LogTypeSlow
}

func (s *LogScraper) filterLogFile(fpath, logtype string) (string, int64, bool, error) {
outputDir, err := s.filteredOutputDir()
if err != nil {
return "", 0, false, err
}
outPath := filteredLogPath(outputDir, fpath)
if err := os.MkdirAll(filepath.Dir(outPath), 0755); err != nil {
return "", 0, false, err
}

in, err := os.Open(fpath)
if err != nil {
return "", 0, false, err
}
defer in.Close()

var reader io.Reader = in
var gzReader *gzip.Reader
if strings.HasSuffix(fpath, ".gz") {
gzReader, err = gzip.NewReader(in)
if err != nil {
return "", 0, false, err
}
defer gzReader.Close()
reader = gzReader
}

out, err := os.Create(outPath)
if err != nil {
return "", 0, false, err
}
var gzWriter *gzip.Writer
defer func() {
if gzWriter != nil {
_ = gzWriter.Close()
}
_ = out.Close()
}()

var writer io.Writer = out
if strings.HasSuffix(outPath, ".gz") {
gzWriter = gzip.NewWriter(out)
writer = gzWriter
}

written, err := writeFilteredLog(reader, writer, logtype, s.Start, s.End)
if err != nil {
return "", 0, false, err
}
if gzWriter != nil {
if err := gzWriter.Close(); err != nil {
return "", 0, false, err
}
gzWriter = nil
}
if err := out.Close(); err != nil {
return "", 0, false, err
}
if !written {
_ = os.Remove(outPath)
return "", 0, false, nil
}
fi, err := os.Stat(outPath)
if err != nil {
return "", 0, false, err
}
return outPath, fi.Size(), true, nil
}

func (s *LogScraper) filteredOutputDir() (string, error) {
if s.outputDir != "" {
return s.outputDir, nil
}
outputDir, err := os.MkdirTemp("", "diag-scraped-logs-*")
if err != nil {
return "", err
}
s.outputDir = outputDir
return s.outputDir, nil
}

func filteredLogPath(outputDir, src string) string {
clean := filepath.Clean(src)
if filepath.IsAbs(clean) {
clean = strings.TrimPrefix(clean, string(filepath.Separator))
}
return filepath.Join(outputDir, clean)
}

func writeFilteredLog(reader io.Reader, writer io.Writer, logtype string, start, end time.Time) (bool, error) {
bufr := bufio.NewReader(reader)
parsers := parser.ListStd()
if logtype == LogTypeSlow {
parsers = []parser.Parser{&parser.SlowQueryParser{}}
}

var current []byte
currentInRange := false
written := false
flush := func() error {
if currentInRange && len(current) > 0 {
if _, err := writer.Write(current); err != nil {
return err
}
written = true
}
current = nil
currentInRange = false
return nil
}

for {
line, err := bufr.ReadBytes('\n')
if err != nil && err != io.EOF {
return false, err
}
if len(line) > 0 {
if ts := parseLine(bytes.TrimRight(line, "\r\n"), parsers); ts != nil {
if err := flush(); err != nil {
return false, err
}
if ts.After(end) {
break
}
current = append(current, line...)
currentInRange = !ts.Before(start)
} else if current != nil {
current = append(current, line...)
}
}
if err == io.EOF {
break
}
}
if err := flush(); err != nil {
return false, err
}
return written, nil
}

func parseLine(line []byte, parsers []parser.Parser) *time.Time {
for _, p := range parsers {
if t, _ := p.ParseHead(line); t != nil {
Expand Down
85 changes: 85 additions & 0 deletions scraper/log_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,85 @@
package scraper

import (
"bytes"
"os"
"path/filepath"
"strings"
"testing"
"time"

"github.com/stretchr/testify/require"
)

func TestWriteFilteredLogIncludesOnlyTimeRange(t *testing.T) {
start := mustParseLogTime(t, "2026/05/26 10:00:00.000 +08:00")
end := mustParseLogTime(t, "2026/05/26 10:02:00.000 +08:00")
input := strings.Join([]string{
`[2026/05/26 09:59:59.999 +08:00] [INFO] [test.go:1] ["before"]`,
`before continuation`,
`[2026/05/26 10:00:00.000 +08:00] [INFO] [test.go:2] ["at start"]`,
`start continuation`,
`[2026/05/26 10:01:00.000 +08:00] [WARN] [test.go:3] ["inside"]`,
`[2026/05/26 10:02:00.000 +08:00] [ERROR] [test.go:4] ["at end"]`,
`[2026/05/26 10:02:00.001 +08:00] [INFO] [test.go:5] ["after"]`,
}, "\n") + "\n"

var out bytes.Buffer
written, err := writeFilteredLog(strings.NewReader(input), &out, LogTypeStd, start, end)
require.NoError(t, err)
require.True(t, written)

result := out.String()
require.NotContains(t, result, "before")
require.Contains(t, result, "at start")
require.Contains(t, result, "start continuation")
require.Contains(t, result, "inside")
require.Contains(t, result, "at end")
require.NotContains(t, result, "after")
}

func TestLogScraperOutputsFilteredLogAndOriginalTarget(t *testing.T) {
start := mustParseLogTime(t, "2026/05/26 10:00:00.000 +08:00")
end := mustParseLogTime(t, "2026/05/26 10:01:00.000 +08:00")
dir := t.TempDir()
logPath := filepath.Join(dir, "tidb.log")
require.NoError(t, os.WriteFile(logPath, []byte(strings.Join([]string{
`[2026/05/26 09:59:59.000 +08:00] [INFO] [test.go:1] ["before"]`,
`[2026/05/26 10:00:00.000 +08:00] [INFO] [test.go:2] ["at start"]`,
`[2026/05/26 10:01:00.000 +08:00] [INFO] [test.go:3] ["at end"]`,
`[2026/05/26 10:01:01.000 +08:00] [INFO] [test.go:4] ["after"]`,
}, "\n")+"\n"), 0644))

s := &LogScraper{
Paths: []string{logPath},
Types: map[string]bool{LogTypeStd: true},
Start: start,
End: end,
}
result := &Sample{}
require.NoError(t, s.Scrap(result))
require.Len(t, result.Log, 1)

var filteredPath string
for p := range result.Log {
filteredPath = p
}
require.NotEqual(t, logPath, filteredPath)
require.Equal(t, logPath, result.LogTargets[filteredPath])
require.Equal(t, LogTypeStd, result.LogTypes[filteredPath])

filtered, err := os.ReadFile(filteredPath)
require.NoError(t, err)
filteredContent := string(filtered)
require.NotContains(t, filteredContent, "before")
require.Contains(t, filteredContent, "at start")
require.Contains(t, filteredContent, "at end")
require.NotContains(t, filteredContent, "after")
}

func mustParseLogTime(t *testing.T, s string) time.Time {
t.Helper()
ts, err := time.Parse("2006/01/02 15:04:05.000 -07:00", s)
require.NoError(t, err)
return ts
}
Loading
Loading