apache/beam · warning
internal: log message buffer closed
Error message
internal: log message buffer closed
What it means
The harness's remoteWriter forwards buffered log entries to the log sender over the w.buffer channel. errBuffClosed is the sentinel returned when that channel is closed while connect is still draining it — i.e. the log buffer was shut down mid-stream, an internal lifecycle condition rather than a transport failure.
Source
Thrown at sdks/go/pkg/beam/core/runtime/harness/logging.go:156
client, err := fnpb.NewBeamFnLoggingClient(conn).Logging(ctx)
if err != nil {
conn.Close()
return nil, func() {}, err
}
toDefer := func() {
client.CloseSend()
conn.Close()
}
return client, toDefer, nil
}
type logSender interface {
Send(*fnpb.LogEntry_List) error
}
var errBuffClosed = errors.New("internal: log message buffer closed")
func (w *remoteWriter) connect(ctx context.Context, makeClient func(ctx context.Context) (logSender, func(), error)) error {
client, toDefer, err := makeClient(ctx)
if err != nil {
return err
}
defer toDefer()
for {
const batchSize = 64
msgs := make([]*fnpb.LogEntry, 0, batchSize)
var flush bool
select {
case <-ctx.Done():
return nil
case newMsg, ok := <-w.buffer:
if !ok {
return errBuffClosedView on GitHub (pinned to 12126d8942)
Solutions
- Treat errBuffClosed as a normal shutdown signal: compare with errors.Is and stop retrying.
- Ensure Flush/shutdown ordering closes the buffer only after connect has finished.
- If seen during normal operation, check for early logger Close calls in harness setup/teardown.
Example fix
// before
if err := w.connect(ctx, makeClient); err != nil { log.Fatalf("log connect: %v", err) }
// after
if err := w.connect(ctx, makeClient); err != nil && !errors.Is(err, errBuffClosed) { log.Fatalf("log connect: %v", err) } Defensive patterns
Strategy: try-catch
Validate before calling
// Only start the log forwarder when the writer is open; skip if already flushed/closed
if w.closed { return nil } Type guard
func forwardable(w *remoteWriter) bool { return w != nil && !w.closed } Try / catch
if err := w.connect(ctx, makeClient); err != nil && !errors.Is(err, errBuffClosed) {
return fmt.Errorf("log forwarding failed: %w", err)
} // errBuffClosed during shutdown is expected Prevention
- Close the logger/buffer only after forwarding goroutines finish
- Use errors.Is against errBuffClosed to treat it as normal shutdown
- Test harness shutdown ordering to flush logs before teardown
When it happens
Trigger: The remoteWriter's buffer channel is closed (writer shutdown/Flush completing) while connect's select loop reads from w.buffer; the ok-receive returns false and connect returns errBuffClosed.
Common situations: Harness shutdown racing with an in-flight log flush; logger closed before the forwarding goroutine finishes; tests exercising connect against a closed buffer.
Related errors
- side input closed
- capacity of cache cannot be negative, got %v
- Processing of an element in transform %v has exceeded the sp
- empty port
- unknown ShutdownMode: %v
AI-assisted analysis of apache/beam@12126d8942 (2026-09-13).
Data as JSON: /api/errors/7159ce9b392e2535.
Report an issue: GitHub.