Go

Implement a pub/sub broadcaster where one slow subscriber cannot block the others.

Question 258HardGo 1.22 to 1.25

A channel delivers each value to one receiver, so broadcasting needs a channel per subscriber. The hard parts are slow consumers, unsubscribing without a send-on-closed-channel panic, and not leaking goroutines.

type Broker[T any] struct {
	mu   sync.Mutex
	subs map[chan T]struct{}
}

func NewBroker[T any]() *Broker[T] {
	return &Broker[T]{subs: make(map[chan T]struct{})}
}

// Subscribe returns a receive-only channel and an idempotent cancel func.
func (b *Broker[T]) Subscribe(buf int) (<-chan T, func()) {
	ch := make(chan T, buf)
	b.mu.Lock()
	b.subs[ch] = struct{}{}
	b.mu.Unlock()

	var once sync.Once
	return ch, func() {
		once.Do(func() {
			b.mu.Lock()
			delete(b.subs, ch) // no Publish can see ch after this
			b.mu.Unlock()
			close(ch) // the broker is the only sender, so closing is safe
		})
	}
}

// Publish never blocks: a subscriber whose buffer is full misses the value.
func (b *Broker[T]) Publish(v T) (dropped int) {
	b.mu.Lock()
	defer b.mu.Unlock()
	for ch := range b.subs {
		select {
		case ch <- v:
		default:
			dropped++
		}
	}
	return dropped
}

Why it is correct:

  • Sends happen only while holding mu, and cancel removes the channel under the same lock before closing it, so nothing ever sends on a closed channel.
  • The non-blocking select with default means holding the lock while sending cannot stall publishers.
  • Subscribers get a <-chan T, so they cannot close or send on it.

Design choices to discuss:

  • Slow-consumer policy: drop newest (as above), drop oldest (a ring buffer per subscriber), block with a timeout, or disconnect the subscriber. Always export a dropped-messages metric.
  • Latest-value-only subscribers, such as config updates, can use a buffer of 1 and replace the pending value.
  • Tie cancel to the subscriber's context, e.g. stop := context.AfterFunc(ctx, cancel), so an abandoned subscriber cannot leak.
  • With many subscribers, publishing under one mutex becomes the bottleneck. Copy-on-write subscriber lists in an atomic.Pointer let publishes run without the lock, but cancel then needs its own coordination with in-flight sends.

More on Concurrency Patterns & sync

All 38 Concurrency Patterns & sync questions