diff --git a/pkg/log/writer.go b/pkg/log/writer.go index 4fa7d7b22..c82f64074 100644 --- a/pkg/log/writer.go +++ b/pkg/log/writer.go @@ -1,23 +1,20 @@ package log import ( + "bytes" "io" + "sync" - "go.uber.org/zap" "go.uber.org/zap/zapcore" ) -// Writer returns an io.WriteCloser that writes each line as a log entry at the given level. -// Level uses the package constants: LevelInfo, LevelDebug, etc. +// maxPendingLine bounds how much unterminated output levelWriter buffers +// before logging it anyway, so a subprocess that never emits a newline +// can't grow the pending line without limit. +const maxPendingLine = 64 * 1024 + func Writer(level int) io.WriteCloser { - zapLevel := verbosityConstToZapLevel(level) - w, closer, _ := zap.Open("stderr") - _ = closer // stderr doesn't need closing - return &levelWriter{ - sink: w, - level: zapLevel, - core: sugar.Load().Desugar().Core(), - } + return &levelWriter{level: verbosityConstToZapLevel(level)} } func verbosityConstToZapLevel(level int) zapcore.Level { @@ -34,18 +31,56 @@ func verbosityConstToZapLevel(level int) zapcore.Level { } type levelWriter struct { - sink zapcore.WriteSyncer level zapcore.Level - core zapcore.Core + + mu sync.Mutex + buf bytes.Buffer } func (w *levelWriter) Write(p []byte) (int, error) { - if !w.core.Enabled(w.level) { + w.mu.Lock() + defer w.mu.Unlock() + + if !sugar.Load().Desugar().Core().Enabled(w.level) { return len(p), nil // discard if below current level } - return w.sink.Write(p) + + total := len(p) + for { + i := bytes.IndexByte(p, '\n') + if i < 0 { + break + } + w.buf.Write(p[:i]) + w.logLine(w.buf.String()) + w.buf.Reset() + p = p[i+1:] + } + w.buf.Write(p) + if w.buf.Len() > maxPendingLine { + w.logLine(w.buf.String()) + w.buf.Reset() + } + return total, nil } func (w *levelWriter) Close() error { + w.mu.Lock() + defer w.mu.Unlock() + if w.buf.Len() > 0 { + w.logLine(w.buf.String()) + w.buf.Reset() + } return nil } + +func (w *levelWriter) logLine(line string) { + switch w.level { + case zapcore.DebugLevel: + Debug(line) + case zapcore.ErrorLevel: + Error(line) + default: + Info(line) + } +} diff --git a/pkg/log/writer_test.go b/pkg/log/writer_test.go new file mode 100644 index 000000000..3d90ce04e --- /dev/null +++ b/pkg/log/writer_test.go @@ -0,0 +1,151 @@ +package log + +import ( + "bytes" + "encoding/json" + "strings" + "testing" +) + +const testFormatJSON = "json" + +func TestWriter_EmitsStructuredJSONLine(t *testing.T) { + Init(Config{Verbosity: 2, Format: testFormatJSON}) + + var sink bytes.Buffer + remove := AddSink(&sink) + defer remove() + + w := Writer(LevelInfo) + _, err := w.Write([]byte("Cloning into 'repo'...\n")) + if err != nil { + t.Fatalf("Write: %v", err) + } + _ = Sync() + + got := strings.TrimSpace(sink.String()) + if !strings.HasPrefix(got, "{") || !strings.HasSuffix(got, "}") { + t.Fatalf("expected a single JSON object, got %q", got) + } + if !strings.Contains(got, `"msg":"Cloning into 'repo'..."`) { + t.Errorf("missing expected msg field: %q", got) + } + if !strings.Contains(got, `"level":"info"`) { + t.Errorf("missing expected level field: %q", got) + } +} + +func TestWriter_SplitsMultipleLinesInOneWrite(t *testing.T) { + Init(Config{Verbosity: 2, Format: testFormatJSON}) + + var sink bytes.Buffer + remove := AddSink(&sink) + defer remove() + + w := Writer(LevelInfo) + _, err := w.Write([]byte("line one\nline two\n")) + if err != nil { + t.Fatalf("Write: %v", err) + } + + // A line split across two Write calls (as os/exec delivers subprocess + // output in arbitrary chunks) must still be logged as one complete + // record, not two fragments. + if _, err := w.Write([]byte("line three\nline ")); err != nil { + t.Fatalf("Write: %v", err) + } + if _, err := w.Write([]byte("four\n")); err != nil { + t.Fatalf("Write: %v", err) + } + _ = Sync() + + lines := strings.Split(strings.TrimSpace(sink.String()), "\n") + wantMsgs := []string{"line one", "line two", "line three", "line four"} + if len(lines) != len(wantMsgs) { + t.Fatalf("got %d lines, want %d: %q", len(lines), len(wantMsgs), lines) + } + for i, l := range lines { + assertJSONLineMsg(t, i, l, wantMsgs[i]) + } +} + +func assertJSONLineMsg(t *testing.T, i int, line, wantMsg string) { + t.Helper() + if !strings.HasPrefix(line, "{") || !strings.HasSuffix(line, "}") { + t.Errorf("line %d is not valid single-object JSON: %q", i, line) + return + } + var rec struct { + Msg string `json:"msg"` + } + if err := json.Unmarshal([]byte(line), &rec); err != nil { + t.Errorf("line %d: json.Unmarshal: %v", i, err) + return + } + if rec.Msg != wantMsg { + t.Errorf("line %d msg = %q, want %q", i, rec.Msg, wantMsg) + } +} + +func TestWriter_FlushesTrailingPartialLineOnClose(t *testing.T) { + Init(Config{Verbosity: 2, Format: testFormatJSON}) + + var sink bytes.Buffer + remove := AddSink(&sink) + defer remove() + + w := Writer(LevelInfo) + _, _ = w.Write([]byte("no trailing newline")) + _ = Sync() + if sink.Len() != 0 { + t.Errorf("expected nothing logged before Close, got %q", sink.String()) + } + + if err := w.Close(); err != nil { + t.Fatalf("Close: %v", err) + } + _ = Sync() + + if !strings.Contains(sink.String(), "no trailing newline") { + t.Errorf("expected trailing partial line flushed on Close, got %q", sink.String()) + } +} + +func TestWriter_PreservesBlankLines(t *testing.T) { + Init(Config{Verbosity: 2, Format: testFormatJSON}) + + var sink bytes.Buffer + remove := AddSink(&sink) + defer remove() + + w := Writer(LevelInfo) + if _, err := w.Write([]byte("one\n\ntwo\n")); err != nil { + t.Fatalf("Write: %v", err) + } + _ = Sync() + + lines := strings.Split(strings.TrimSpace(sink.String()), "\n") + if len(lines) != 3 { + t.Fatalf("got %d lines, want 3 (blank line preserved): %q", len(lines), lines) + } + if !strings.Contains(lines[1], `"msg":""`) { + t.Errorf("line 1 = %q, want an empty msg field", lines[1]) + } +} + +func TestWriter_DiscardsBelowConfiguredLevel(t *testing.T) { + Init(Config{Verbosity: 1, Format: testFormatJSON}) // info+ only, debug disabled + + var sink bytes.Buffer + remove := AddSink(&sink) + defer remove() + + w := Writer(LevelDebug) + _, _ = w.Write([]byte("should not appear\n")) + _ = w.Close() + _ = Sync() + + if sink.Len() != 0 { + t.Errorf("expected debug output discarded below configured level, got %q", sink.String()) + } +}