broker API

broker

package

API reference for the broker package.

T
type

Handler

Handler processes a broker message.

core/relay/broker/broker.go:6-6
type Handler func(ctx context.Context, payload []byte) error
I
interface

Subscription

Subscription controls one broker subscription.

core/relay/broker/broker.go:9-11
type Subscription interface

Methods

Unsubscribe
Method

Returns

error
func Unsubscribe(...)
I
interface

Broker

Broker publishes messages and creates subscriptions.

core/relay/broker/broker.go:14-17
type Broker interface

Methods

Publish
Method

Parameters

topic string
payload []byte

Returns

error
func Publish(...)
Subscribe
Method

Parameters

topic string
handler Handler

Returns

error
func Subscribe(...)
S
struct
Implements: Broker

MemoryBroker

MemoryBroker is an in-memory pub/sub implementation for testing and single-process use.

core/relay/broker/memory.go:17-21
type MemoryBroker struct

Methods

Publish
Method

Publish sends a payload to all topic subscribers.

Parameters

topic string
payload []byte

Returns

error
func (*MemoryBroker) Publish(ctx context.Context, topic string, payload []byte) error
{
	m.subMu.RLock()
	topicHandlers := m.subs[topic]
	handlers := make([]Handler, 0, len(topicHandlers))
	for _, handler := range topicHandlers {
		handlers = append(handlers, handler)
	}
	m.subMu.RUnlock()

	var errs []error
	for _, h := range handlers {
		if err := func() (err error) {
			defer func() {
				if recovered := recover(); recovered != nil {
					err = errors.New("relay: subscriber panic")
				}
			}()
			return h(ctx, append([]byte(nil), payload...))
		}(); err != nil {
			errs = append(errs, err)
		}
	}
	return errors.Join(errs...)
}
Subscribe
Method

Subscribe registers a handler for a topic.

Parameters

topic string
handler Handler

Returns

error
func (*MemoryBroker) Subscribe(topic string, handler Handler) (Subscription, error)
{
	if topic == "" {
		return nil, errors.New("relay: subscription topic cannot be empty")
	}
	if handler == nil {
		return nil, errors.New("relay: subscription handler cannot be nil")
	}
	m.subMu.Lock()
	m.next++
	id := m.next
	if m.subs[topic] == nil {
		m.subs[topic] = make(map[uint64]Handler)
	}
	m.subs[topic][id] = handler
	m.subMu.Unlock()
	return &memorySubscription{broker: m, topic: topic, id: id}, nil
}

Fields

Name Type Description
subs map[string]map[uint64]Handler
subMu sync.RWMutex
next uint64
F
function

NewMemoryBroker

NewMemoryBroker creates a MemoryBroker.

Returns

core/relay/broker/memory.go:24-28
func NewMemoryBroker() *MemoryBroker

{
	return &MemoryBroker{
		subs: make(map[string]map[uint64]Handler),
	}
}
S
struct
Implements: Subscription

memorySubscription

core/relay/broker/memory.go:75-80
type memorySubscription struct

Methods

Unsubscribe
Method

Returns

error
func (*memorySubscription) Unsubscribe() error
{
	s.once.Do(func() {
		s.broker.subMu.Lock()
		delete(s.broker.subs[s.topic], s.id)
		if len(s.broker.subs[s.topic]) == 0 {
			delete(s.broker.subs, s.topic)
		}
		s.broker.subMu.Unlock()
	})
	return nil
}

Fields

Name Type Description
broker *MemoryBroker
topic string
id uint64
once sync.Once
F
function

TestMemoryBroker_PublishSubscribe

Parameters

core/relay/broker/memory_test.go:11-33
func TestMemoryBroker_PublishSubscribe(t *testing.T)

{
	b := NewMemoryBroker()
	ch := make(chan []byte, 1)

	b.Subscribe("test", func(_ context.Context, data []byte) error {
		ch <- data
		return nil
	})

	err := b.Publish(context.Background(), "test", []byte("hello"))
	if err != nil {
		t.Fatalf("Publish failed: %v", err)
	}

	select {
	case result := <-ch:
		if string(result) != "hello" {
			t.Errorf("got %q, want %q", string(result), "hello")
		}
	case <-time.After(time.Second):
		t.Fatal("timeout waiting for handler")
	}
}
F
function

TestMemoryBroker_MultipleSubscribers

Parameters

core/relay/broker/memory_test.go:35-60
func TestMemoryBroker_MultipleSubscribers(t *testing.T)

{
	b := NewMemoryBroker()
	ch := make(chan int, 3)
	count := 0

	for i := 0; i < 3; i++ {
		b.Subscribe("topic", func(_ context.Context, data []byte) error {
			ch <- 1
			return nil
		})
	}

	b.Publish(context.Background(), "topic", []byte("data"))

	for i := 0; i < 3; i++ {
		select {
		case <-ch:
			count++
		case <-time.After(time.Second):
			t.Fatalf("timeout waiting for handler %d", i)
		}
	}
	if count != 3 {
		t.Errorf("expected 3 handler calls, got %d", count)
	}
}
F
function

TestMemoryBroker_DifferentTopics

Parameters

core/relay/broker/memory_test.go:62-85
func TestMemoryBroker_DifferentTopics(t *testing.T)

{
	b := NewMemoryBroker()
	ch := make(chan string, 1)

	b.Subscribe("topic1", func(_ context.Context, data []byte) error {
		ch <- "topic1:" + string(data)
		return nil
	})
	b.Subscribe("topic2", func(_ context.Context, data []byte) error {
		ch <- "topic2:" + string(data)
		return nil
	})

	b.Publish(context.Background(), "topic1", []byte("msg"))

	select {
	case result := <-ch:
		if result != "topic1:msg" {
			t.Errorf("got %q, want %q", result, "topic1:msg")
		}
	case <-time.After(time.Second):
		t.Fatal("timeout waiting for handler")
	}
}
F
function

TestMemoryBroker_PublishNoSubscribers

Parameters

core/relay/broker/memory_test.go:87-93
func TestMemoryBroker_PublishNoSubscribers(t *testing.T)

{
	b := NewMemoryBroker()
	err := b.Publish(context.Background(), "nonexistent", []byte("msg"))
	if err != nil {
		t.Fatalf("Publish to nonexistent topic: %v", err)
	}
}
F
function

TestMemoryBroker_PropagatesSubscriberError

Parameters

core/relay/broker/memory_test.go:95-116
func TestMemoryBroker_PropagatesSubscriberError(t *testing.T)

{
	broker := NewMemoryBroker()
	want := errors.New("delivery failed")
	var successfulCalls int
	if _, err := broker.Subscribe("topic", func(context.Context, []byte) error {
		return want
	}); err != nil {
		t.Fatal(err)
	}
	if _, err := broker.Subscribe("topic", func(context.Context, []byte) error {
		successfulCalls++
		return nil
	}); err != nil {
		t.Fatal(err)
	}
	if err := broker.Publish(context.Background(), "topic", nil); !errors.Is(err, want) {
		t.Fatalf("Publish() error = %v, want %v", err, want)
	}
	if successfulCalls != 1 {
		t.Fatalf("successful subscriber calls = %d, want 1", successfulCalls)
	}
}
F
function

TestMemoryBrokerCopiesPayloadPerSubscriber

Parameters

core/relay/broker/memory_test.go:118-139
func TestMemoryBrokerCopiesPayloadPerSubscriber(t *testing.T)

{
	broker := NewMemoryBroker()
	if _, err := broker.Subscribe("topic", func(_ context.Context, payload []byte) error {
		payload[0] = 'X'
		return nil
	}); err != nil {
		t.Fatal(err)
	}
	var received string
	if _, err := broker.Subscribe("topic", func(_ context.Context, payload []byte) error {
		received = string(payload)
		return nil
	}); err != nil {
		t.Fatal(err)
	}
	if err := broker.Publish(context.Background(), "topic", []byte("safe")); err != nil {
		t.Fatal(err)
	}
	if received != "safe" {
		t.Fatalf("second subscriber received %q", received)
	}
}
F
function

TestMemoryBrokerUnsubscribe

Parameters

core/relay/broker/memory_test.go:141-160
func TestMemoryBrokerUnsubscribe(t *testing.T)

{
	broker := NewMemoryBroker()
	calls := 0
	subscription, err := broker.Subscribe("topic", func(context.Context, []byte) error {
		calls++
		return nil
	})
	if err != nil {
		t.Fatal(err)
	}
	if err := subscription.Unsubscribe(); err != nil {
		t.Fatal(err)
	}
	if err := broker.Publish(context.Background(), "topic", nil); err != nil {
		t.Fatal(err)
	}
	if calls != 0 {
		t.Fatalf("handler called %d times after unsubscribe", calls)
	}
}
F
function

TestMemoryBroker_ConcurrentPublish

Parameters

core/relay/broker/memory_test.go:162-207
func TestMemoryBroker_ConcurrentPublish(t *testing.T)

{
	b := NewMemoryBroker()
	ch := make(chan int, 10)
	received := 0

	b.Subscribe("test", func(_ context.Context, data []byte) error {
		ch <- 1
		return nil
	})

	var wg sync.WaitGroup
	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func() {
			defer wg.Done()
			b.Publish(context.Background(), "test", []byte("data"))
		}()
	}
	wg.Wait()

	done := make(chan struct{})
	go func() {
		for {
			select {
			case <-ch:
				received++
				if received >= 10 {
					close(done)
					return
				}
			case <-time.After(100 * time.Millisecond):
				close(done)
				return
			}
		}
	}()

	select {
	case <-done:
		if received != 10 {
			t.Errorf("expected 10 handler calls, got %d", received)
		}
	case <-time.After(2 * time.Second):
		t.Fatal("timeout")
	}
}