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 errBuffClosed

View on GitHub (pinned to 12126d8942)

Solutions

  1. Treat errBuffClosed as a normal shutdown signal: compare with errors.Is and stop retrying.
  2. Ensure Flush/shutdown ordering closes the buffer only after connect has finished.
  3. 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

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


AI-assisted analysis of apache/beam@12126d8942 (2026-09-13). Data as JSON: /api/errors/7159ce9b392e2535. Report an issue: GitHub.