Skip to content

Commit 7b9380b

Browse files
authored
Buffer BufferedLogger by newline to avoid log splitting (#39288)
* Buffer BufferedLogger by newline to avoid log splitting * Update BufferedLogger.Write to search for newlines and accumulate partial log lines in the builder instead of immediately emitting them as separate entries. * Update FlushAtDebug and FlushAtError to flush any remaining trailing text in the builder on exit or flush events. * Fix bug in buffered_logging_test.go where log list assertions only verified the first element of logCatcher.msgs instead of checking all gathered log messages. * Add TestBufferedLogger/partial_write_splitting to verify correct chunked write buffering and line assembly behavior. * Call BuferredLogger.Printf so the deps will not be separate lines. * Address comments * Refactor BufferedLogger flushing and container command execution - Add Flush(ctx, err) to BufferedLogger to automatically flush at ERROR on failure or DEBUG on success. - Add executeWithLogger and executeWithOutput helpers to eliminate repetitive flush boilerplate. - Unify runtime dependency log outputs into single log entries. - Add unit tests and docstrings for BufferedLogger. * Address comments.
1 parent f5feab8 commit 7b9380b

4 files changed

Lines changed: 240 additions & 54 deletions

File tree

sdks/go/container/tools/buffered_logging.go

Lines changed: 55 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,7 @@
1616
package tools
1717

1818
import (
19+
"bytes"
1920
"context"
2021
"log"
2122
"os"
@@ -52,23 +53,59 @@ func NewBufferedLoggerWithFlushInterval(ctx context.Context, logger *Logger, int
5253
return &BufferedLogger{logger: logger, lastFlush: time.Now(), flushInterval: interval, periodicFlushContext: ctx, now: time.Now}
5354
}
5455

55-
// Write implements the io.Writer interface, converting input to a string
56-
// and storing it in the BufferedLogger's buffer. If a logger is not provided,
57-
// the output is sent directly to os.Stderr.
56+
// Write implements the io.Writer interface. It buffers byte streams line-by-line
57+
// into memory and flushes periodically or upon calling Flush(), FlushAtError(), or
58+
// FlushAtDebug(). It is used primarily to redirect stdout/stderr of subprocesses or
59+
// standard Go log output. If a logger is not provided, the output is sent directly to os.Stderr.
5860
func (b *BufferedLogger) Write(p []byte) (int, error) {
5961
if b.logger == nil {
6062
return os.Stderr.Write(p)
6163
}
62-
n, err := b.builder.Write(p)
64+
6365
if b.logs == nil {
6466
b.logs = make([]string, 0, initialLogSize)
6567
}
66-
b.logs = append(b.logs, b.builder.String())
67-
b.builder.Reset()
68+
69+
start := 0
70+
for {
71+
// Look for the next newline in the incoming byte slice directly
72+
nl := bytes.IndexByte(p[start:], '\n')
73+
if nl == -1 {
74+
break
75+
}
76+
77+
// Write the segment up to the newline into the builder
78+
b.builder.Write(p[start : start+nl])
79+
80+
// The builder now contains any previous partial line + the current complete segment
81+
b.logs = append(b.logs, strings.TrimSuffix(b.builder.String(), "\r"))
82+
b.builder.Reset()
83+
84+
start += nl + 1
85+
}
86+
87+
// Buffer any remaining bytes that didn't end in a newline
88+
if start < len(p) {
89+
b.builder.Write(p[start:])
90+
}
91+
6892
if b.now().Sub(b.lastFlush) > b.flushInterval {
6993
b.FlushAtDebug(b.periodicFlushContext)
7094
}
71-
return n, err
95+
96+
return len(p), nil
97+
}
98+
99+
// Flush flushes the contents of the buffer to the logging service.
100+
// If err is non-nil, it flushes at Error severity; otherwise it flushes at Debug severity.
101+
// It returns the provided error.
102+
func (b *BufferedLogger) Flush(ctx context.Context, err error) error {
103+
if err != nil {
104+
b.FlushAtError(ctx)
105+
} else {
106+
b.FlushAtDebug(ctx)
107+
}
108+
return err
72109
}
73110

74111
// FlushAtError flushes the contents of the buffer to the logging
@@ -77,6 +114,10 @@ func (b *BufferedLogger) FlushAtError(ctx context.Context) {
77114
if b.logger == nil {
78115
return
79116
}
117+
if b.builder.Len() > 0 {
118+
b.logs = append(b.logs, strings.TrimSuffix(b.builder.String(), "\r"))
119+
b.builder.Reset()
120+
}
80121
for _, message := range b.logs {
81122
b.logger.Errorf(ctx, "%s", message)
82123
}
@@ -90,15 +131,20 @@ func (b *BufferedLogger) FlushAtDebug(ctx context.Context) {
90131
if b.logger == nil {
91132
return
92133
}
134+
if b.builder.Len() > 0 {
135+
b.logs = append(b.logs, strings.TrimSuffix(b.builder.String(), "\r"))
136+
b.builder.Reset()
137+
}
93138
for _, message := range b.logs {
94139
b.logger.Printf(ctx, "%s", message)
95140
}
96141
b.logs = nil
97142
b.lastFlush = time.Now()
98143
}
99144

100-
// Prints directly to the logging service. If the logger is nil, prints directly to the
101-
// console. Used for the container pre-build workflow.
145+
// Printf directly writes formatted messages to the underlying logger/service,
146+
// bypassing line buffering. If the logger is nil, it prints directly to the
147+
// console. Used for direct informational logs and the container pre-build workflow.
102148
func (b *BufferedLogger) Printf(ctx context.Context, format string, args ...any) {
103149
if b.logger == nil {
104150
log.Printf(format, args...)

sdks/go/container/tools/buffered_logging_test.go

Lines changed: 156 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -17,21 +17,48 @@ package tools
1717

1818
import (
1919
"context"
20+
"errors"
2021
"testing"
2122
"time"
2223

2324
fnpb "github.com/apache/beam/sdks/v2/go/pkg/beam/model/fnexecution_v1"
2425
)
2526

27+
func getAllLogEntries(catcher *logCatcher) []*fnpb.LogEntry {
28+
var entries []*fnpb.LogEntry
29+
for _, list := range catcher.msgs {
30+
entries = append(entries, list.GetLogEntries()...)
31+
}
32+
return entries
33+
}
34+
2635
func TestBufferedLogger(t *testing.T) {
2736
ctx := context.Background()
2837

38+
t.Run("printf", func(t *testing.T) {
39+
catcher := &logCatcher{}
40+
l := &Logger{client: catcher}
41+
bl := NewBufferedLogger(l)
42+
43+
bl.Printf(ctx, "test message")
44+
45+
received := catcher.msgs[0].GetLogEntries()[0]
46+
47+
if got, want := received.Message, "test message"; got != want {
48+
t.Errorf("got message %q, want %q", got, want)
49+
}
50+
51+
if got, want := received.Severity, fnpb.LogEntry_Severity_DEBUG; got != want {
52+
t.Errorf("got severity %v, want %v", got, want)
53+
}
54+
})
55+
2956
t.Run("write", func(t *testing.T) {
3057
catcher := &logCatcher{}
3158
l := &Logger{client: catcher}
3259
bl := NewBufferedLogger(l)
3360

34-
message := []byte("test message")
61+
message := []byte("test message\n")
3562
n, err := bl.Write(message)
3663
if err != nil {
3764
t.Errorf("got error %v", err)
@@ -77,7 +104,8 @@ func TestBufferedLogger(t *testing.T) {
77104
l := &Logger{client: catcher}
78105
bl := NewBufferedLogger(l)
79106

80-
messages := []string{"foo", "bar", "baz"}
107+
messages := []string{"foo\n", "bar\n", "baz\n"}
108+
expected := []string{"foo", "bar", "baz"}
81109

82110
for _, message := range messages {
83111
messBytes := []byte(message)
@@ -93,10 +121,14 @@ func TestBufferedLogger(t *testing.T) {
93121

94122
bl.FlushAtDebug(ctx)
95123

96-
received := catcher.msgs[0].GetLogEntries()
124+
received := getAllLogEntries(catcher)
125+
126+
if got, want := len(received), len(expected); got != want {
127+
t.Fatalf("expected %d log entries received, got %d", want, got)
128+
}
97129

98130
for i, message := range received {
99-
if got, want := message.Message, messages[i]; got != want {
131+
if got, want := message.Message, expected[i]; got != want {
100132
t.Errorf("got message %q, want %q", got, want)
101133
}
102134

@@ -139,7 +171,8 @@ func TestBufferedLogger(t *testing.T) {
139171
l := &Logger{client: catcher}
140172
bl := NewBufferedLogger(l)
141173

142-
messages := []string{"foo", "bar", "baz"}
174+
messages := []string{"foo\n", "bar\n", "baz\n"}
175+
expected := []string{"foo", "bar", "baz"}
143176

144177
for _, message := range messages {
145178
messBytes := []byte(message)
@@ -155,10 +188,14 @@ func TestBufferedLogger(t *testing.T) {
155188

156189
bl.FlushAtError(ctx)
157190

158-
received := catcher.msgs[0].GetLogEntries()
191+
received := getAllLogEntries(catcher)
192+
193+
if got, want := len(received), len(expected); got != want {
194+
t.Fatalf("expected %d log entries received, got %d", want, got)
195+
}
159196

160197
for i, message := range received {
161-
if got, want := message.Message, messages[i]; got != want {
198+
if got, want := message.Message, expected[i]; got != want {
162199
t.Errorf("got message %q, want %q", got, want)
163200
}
164201

@@ -168,6 +205,55 @@ func TestBufferedLogger(t *testing.T) {
168205
}
169206
})
170207

208+
t.Run("flush with nil error", func(t *testing.T) {
209+
catcher := &logCatcher{}
210+
l := &Logger{client: catcher}
211+
bl := NewBufferedLogger(l)
212+
213+
message := []byte("success message\n")
214+
_, err := bl.Write(message)
215+
if err != nil {
216+
t.Fatalf("unexpected write error: %v", err)
217+
}
218+
219+
if gotErr := bl.Flush(ctx, nil); gotErr != nil {
220+
t.Errorf("Flush(ctx, nil) returned error %v, want nil", gotErr)
221+
}
222+
223+
received := catcher.msgs[0].GetLogEntries()[0]
224+
if got, want := received.Message, "success message"; got != want {
225+
t.Errorf("got message %q, want %q", got, want)
226+
}
227+
if got, want := received.Severity, fnpb.LogEntry_Severity_DEBUG; got != want {
228+
t.Errorf("got severity %v, want %v", got, want)
229+
}
230+
})
231+
232+
t.Run("flush with non-nil error", func(t *testing.T) {
233+
catcher := &logCatcher{}
234+
l := &Logger{client: catcher}
235+
bl := NewBufferedLogger(l)
236+
237+
message := []byte("error message\n")
238+
_, err := bl.Write(message)
239+
if err != nil {
240+
t.Fatalf("unexpected write error: %v", err)
241+
}
242+
243+
originalErr := errors.New("command failed")
244+
if gotErr := bl.Flush(ctx, originalErr); gotErr != originalErr {
245+
t.Errorf("Flush(ctx, err) returned %v, want %v", gotErr, originalErr)
246+
}
247+
248+
received := catcher.msgs[0].GetLogEntries()[0]
249+
if got, want := received.Message, "error message"; got != want {
250+
t.Errorf("got message %q, want %q", got, want)
251+
}
252+
if got, want := received.Severity, fnpb.LogEntry_Severity_ERROR; got != want {
253+
t.Errorf("got severity %v, want %v", got, want)
254+
}
255+
})
256+
171257
t.Run("direct print", func(t *testing.T) {
172258
catcher := &logCatcher{}
173259
l := &Logger{client: catcher}
@@ -195,7 +281,8 @@ func TestBufferedLogger(t *testing.T) {
195281
startTime := time.Now()
196282
bl.now = func() time.Time { return startTime }
197283

198-
messages := []string{"foo", "bar"}
284+
messages := []string{"foo\n", "bar\n"}
285+
expected := []string{"foo", "bar"}
199286

200287
for i, message := range messages {
201288
if i > 1 {
@@ -212,7 +299,8 @@ func TestBufferedLogger(t *testing.T) {
212299
}
213300
}
214301

215-
lastMessage := "baz"
302+
lastMessage := "baz\n"
303+
expected = append(expected, "baz")
216304
bl.now = func() time.Time { return startTime.Add(6 * time.Second) }
217305
messBytes := []byte(lastMessage)
218306
n, err := bl.Write(messBytes)
@@ -225,11 +313,14 @@ func TestBufferedLogger(t *testing.T) {
225313
}
226314

227315
// Type should have auto-flushed at debug after the third message
228-
received := catcher.msgs[0].GetLogEntries()
229-
messages = append(messages, lastMessage)
316+
received := getAllLogEntries(catcher)
317+
318+
if got, want := len(received), len(expected); got != want {
319+
t.Fatalf("expected %d log entries received, got %d", want, got)
320+
}
230321

231322
for i, message := range received {
232-
if got, want := message.Message, messages[i]; got != want {
323+
if got, want := message.Message, expected[i]; got != want {
233324
t.Errorf("got message %q, want %q", got, want)
234325
}
235326

@@ -238,4 +329,57 @@ func TestBufferedLogger(t *testing.T) {
238329
}
239330
}
240331
})
332+
333+
t.Run("partial write splitting", func(t *testing.T) {
334+
catcher := &logCatcher{}
335+
l := &Logger{client: catcher}
336+
bl := NewBufferedLogger(l)
337+
338+
// Write a partial line
339+
n, err := bl.Write([]byte("hello "))
340+
if err != nil {
341+
t.Errorf("got error %v", err)
342+
}
343+
if n != 6 {
344+
t.Errorf("got %d, want 6", n)
345+
}
346+
if len(bl.logs) != 0 {
347+
t.Errorf("expected no logs buffered yet, got %d", len(bl.logs))
348+
}
349+
350+
// Write remainder and a second line
351+
n, err = bl.Write([]byte("world\nline2\npartial"))
352+
if err != nil {
353+
t.Errorf("got error %v", err)
354+
}
355+
if n != 19 {
356+
t.Errorf("got %d, want 19", n)
357+
}
358+
359+
if got, want := len(bl.logs), 2; got != want {
360+
t.Errorf("expected 2 logs buffered, got %d", got)
361+
}
362+
if got, want := bl.logs[0], "hello world"; got != want {
363+
t.Errorf("got %q, want %q", got, want)
364+
}
365+
if got, want := bl.logs[1], "line2"; got != want {
366+
t.Errorf("got %q, want %q", got, want)
367+
}
368+
369+
// Flush should flush the final partial message
370+
bl.FlushAtDebug(ctx)
371+
received := getAllLogEntries(catcher)
372+
if got, want := len(received), 3; got != want {
373+
t.Fatalf("expected 3 log entries received, got %d", got)
374+
}
375+
if got, want := received[0].Message, "hello world"; got != want {
376+
t.Errorf("got message %q, want %q", got, want)
377+
}
378+
if got, want := received[1].Message, "line2"; got != want {
379+
t.Errorf("got message %q, want %q", got, want)
380+
}
381+
if got, want := received[2].Message, "partial"; got != want {
382+
t.Errorf("got message %q, want %q", got, want)
383+
}
384+
})
241385
}

0 commit comments

Comments
 (0)