broker
API
broker
packageAPI reference for the broker
package.
Imports
(6)
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
S
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
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
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
t
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
t
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
t
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
t
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
t
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
t
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
t
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
t
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")
}
}