manager API

manager

package

API reference for the manager package.

S
struct

Job

Job represents a unit of work with metadata.

core/relay/manager/manager.go:16-23
type Job struct

Fields

Name Type Description
ID string
Queue string
Topic string
Payload []byte
CreatedAt time.Time
TryCount int
T
type

Handler

Handler processes a typed payload from the relay.

core/relay/manager/manager.go:26-26
type Handler func(ctx context.Context, payload T) error
T
type

Broker

Broker publishes and subscribes to topics.

core/relay/manager/manager.go:29-29
type Broker broker.Broker
S
struct

Relay

Relay coordinates typed message handlers with a pluggable broker.

core/relay/manager/manager.go:34-46
type Relay struct

Methods

Start
Method

Start subscribes all registered handlers and returns a ready channel.

Parameters

Returns

<-chan struct{}
error
func (*Relay) Start(ctx context.Context) (<-chan struct{}, error)
{
	ready := make(chan struct{})
	if err := ctx.Err(); err != nil {
		close(ready)
		return ready, err
	}
	r.lifecycle.Lock()
	if r.started || r.starting || r.stopping {
		r.lifecycle.Unlock()
		close(ready)
		return ready, errors.New("relay: already started")
	}
	r.starting = true
	r.stopStart = false
	r.handlerMu.RLock()
	handlers := make(map[string]any, len(r.handlers))
	for topic, handler := range r.handlers {
		handlers[topic] = handler
	}
	r.handlerMu.RUnlock()
	r.lifecycle.Unlock()

	var subscriptions []broker.Subscription
	for topic, wrapperFn := range handlers {
		userHandler := wrapperFn.(func(ctx context.Context, data []byte) error)

		subscription, err := safeSubscribe(r.broker, topic, func(ctx context.Context, data []byte) (handlerErr error) {
			defer func() {
				if rec := recover(); rec != nil {
					handlerErr = fmt.Errorf("relay: panic in topic %q handler: %v", topic, rec)
				}
			}()

			return userHandler(ctx, data)
		})
		if err != nil {
			r.lifecycle.Lock()
			r.starting = false
			r.stopping = true
			r.lifecycle.Unlock()
			cleanupErr := unsubscribeAll(subscriptions)
			r.lifecycle.Lock()
			r.stopping = false
			r.lifecycle.Unlock()
			close(ready)
			return ready, errors.Join(err, cleanupErr)
		}
		if subscription == nil {
			r.lifecycle.Lock()
			r.starting = false
			r.stopping = true
			r.lifecycle.Unlock()
			cleanupErr := unsubscribeAll(subscriptions)
			r.lifecycle.Lock()
			r.stopping = false
			r.lifecycle.Unlock()
			close(ready)
			return ready, errors.Join(
				errors.New("relay: broker returned a nil subscription"),
				cleanupErr,
			)
		}
		subscriptions = append(subscriptions, subscription)
	}

	runCtx, cancel := context.WithCancel(ctx)
	r.lifecycle.Lock()
	if r.stopStart || ctx.Err() != nil {
		r.starting = false
		r.stopping = true
		r.lifecycle.Unlock()
		cancel()
		cleanupErr := unsubscribeAll(subscriptions)
		r.lifecycle.Lock()
		r.stopping = false
		r.lifecycle.Unlock()
		close(ready)
		return ready, errors.Join(errors.New("relay: start stopped"), ctx.Err(), cleanupErr)
	}
	r.runID++
	runID := r.runID
	r.cancel = cancel
	r.subs = subscriptions
	r.started = true
	r.starting = false
	r.stopStart = false
	r.lifecycle.Unlock()
	close(ready)
	go func() {
		<-runCtx.Done()
		_ = r.stopRun(runID)
	}()
	return ready, nil
}
Stop
Method

Stop removes active broker subscriptions.

Returns

error
func (*Relay) Stop() error
{
	return r.stopRun(0)
}
stopRun
Method

Parameters

runID uint64

Returns

error
func (*Relay) stopRun(runID uint64) error
{
	r.lifecycle.Lock()
	if r.starting {
		r.stopStart = true
		r.lifecycle.Unlock()
		return nil
	}
	if !r.started || runID != 0 && runID != r.runID {
		r.lifecycle.Unlock()
		return nil
	}
	r.started = false
	r.stopping = true
	cancel := r.cancel
	r.cancel = nil
	subscriptions := r.subs
	r.subs = nil
	r.lifecycle.Unlock()
	if cancel != nil {
		cancel()
	}
	err := unsubscribeAll(subscriptions)
	r.lifecycle.Lock()
	r.stopping = false
	r.lifecycle.Unlock()
	return err
}

Fields

Name Type Description
broker Broker
handlers map[string]any
handlerMu sync.RWMutex
lifecycle sync.Mutex
started bool
starting bool
stopping bool
stopStart bool
runID uint64
cancel context.CancelFunc
subs []broker.Subscription
T
type

Option

Option configures a Relay.

core/relay/manager/manager.go:49-49
type Option options.Option[Relay]
F
function

New

New creates a Relay with the given options.

Parameters

opts
...Option

Returns

core/relay/manager/manager.go:52-62
func New(opts ...Option) *Relay

{
	r := &Relay{
		broker:   broker.NewMemoryBroker(),
		handlers: make(map[string]any),
	}
	options.Apply(r, opts...)
	if r.broker == nil {
		panic("relay: broker cannot be nil")
	}
	return r
}
F
function

WithBroker

WithBroker sets the broker implementation.

Parameters

b

Returns

core/relay/manager/manager.go:65-69
func WithBroker(b Broker) Option

{
	return func(r *Relay) {
		r.broker = b
	}
}
F
function

Register

Register adds a typed handler for a topic.

Parameters

r
topic
string
fn
Handler[T]

Returns

error
core/relay/manager/manager.go:72-108
func Register[T any](r *Relay, topic string, fn Handler[T]) error

{
	if r == nil {
		return errors.New("relay: relay cannot be nil")
	}
	if topic == "" {
		return errors.New("relay: topic cannot be empty")
	}
	if fn == nil {
		return errors.New("relay: handler cannot be nil")
	}
	wrapper := func(ctx context.Context, raw []byte) error {
		if len(raw) > maxMessageSize {
			return errors.New("relay: message exceeds 4 MiB limit")
		}
		var payload T
		if err := json.Unmarshal(raw, &payload); err != nil {
			return fmt.Errorf("payload unmarshal failed: %w", err)
		}
		return fn(ctx, payload)
	}

	r.lifecycle.Lock()
	if r.started || r.starting || r.stopping {
		r.lifecycle.Unlock()
		return errors.New("relay: cannot register handlers after start")
	}
	r.handlerMu.Lock()
	if _, exists := r.handlers[topic]; exists {
		r.handlerMu.Unlock()
		r.lifecycle.Unlock()
		return fmt.Errorf("relay: topic %q already has a handler", topic)
	}
	r.handlers[topic] = wrapper
	r.handlerMu.Unlock()
	r.lifecycle.Unlock()
	return nil
}
F
function

Enqueue

Enqueue publishes a typed payload to a topic.

Parameters

r
topic
string
payload
T

Returns

error
core/relay/manager/manager.go:111-124
func Enqueue[T any](ctx context.Context, r *Relay, topic string, payload T) error

{
	if r == nil {
		return errors.New("relay: relay cannot be nil")
	}
	data, err := json.Marshal(payload)
	if err != nil {
		return fmt.Errorf("payload marshal failed: %w", err)
	}
	if len(data) > maxMessageSize {
		return errors.New("relay: message exceeds 4 MiB limit")
	}

	return r.broker.Publish(ctx, topic, data)
}
F
function

safeSubscribe

Parameters

b
topic
string
handler

Returns

subscription
err
error
core/relay/manager/manager.go:255-266
func safeSubscribe(b Broker, topic string, handler broker.Handler) (subscription broker.Subscription, err error)

{
	defer func() {
		if recovered := recover(); recovered != nil {
			err = fmt.Errorf("relay: broker subscribe panic: %v", recovered)
		}
	}()
	return b.Subscribe(topic, handler)
}
F
function

unsubscribeAll

Parameters

subscriptions

Returns

error
core/relay/manager/manager.go:268-285
func unsubscribeAll(subscriptions []broker.Subscription) error

{
	var errs []error
	for index := len(subscriptions) - 1; index >= 0; index-- {
		subscription := subscriptions[index]
		err := func() (err error) {
			defer func() {
				if recovered := recover(); recovered != nil {
					err = fmt.Errorf("relay: unsubscribe panic: %v", recovered)
				}
			}()
			return subscription.Unsubscribe()
		}()
		if err != nil {
			errs = append(errs, err)
		}
	}
	return errors.Join(errs...)
}
S
struct

testPayload

core/relay/manager/manager_test.go:13-16
type testPayload struct

Fields

Name Type Description
Message string json:"message"
Value int json:"value"
F
function

TestNew_DefaultBroker

Parameters

core/relay/manager/manager_test.go:18-23
func TestNew_DefaultBroker(t *testing.T)

{
	r := New()
	if r == nil {
		t.Fatal("New returned nil")
	}
}
F
function

TestRegisterAndEnqueue

Parameters

core/relay/manager/manager_test.go:25-59
func TestRegisterAndEnqueue(t *testing.T)

{
	r := New()
	var mu sync.Mutex
	var received testPayload
	handlerCalled := make(chan struct{})

	Register[testPayload](r, "test.topic", func(ctx context.Context, p testPayload) error {
		mu.Lock()
		received = p
		mu.Unlock()
		close(handlerCalled)
		return nil
	})

	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()
	ready, err := r.Start(ctx)
	if err != nil {
		t.Fatalf("Start failed: %v", err)
	}
	<-ready

	Enqueue[testPayload](context.Background(), r, "test.topic", testPayload{Message: "hello", Value: 42})

	select {
	case <-handlerCalled:
		mu.Lock()
		if received.Message != "hello" || received.Value != 42 {
			t.Errorf("got %+v, want {Message:hello Value:42}", received)
		}
		mu.Unlock()
	case <-time.After(time.Second):
		t.Fatal("timeout waiting for handler")
	}
}
F
function

TestRegister_TypeSafety

Parameters

core/relay/manager/manager_test.go:61-74
func TestRegister_TypeSafety(t *testing.T)

{
	r := New()

	Register[testPayload](r, "topic1", func(ctx context.Context, p testPayload) error {
		return nil
	})
	Register[int](r, "topic2", func(ctx context.Context, p int) error {
		return nil
	})

	if len(r.handlers) != 2 {
		t.Errorf("expected 2 handlers, got %d", len(r.handlers))
	}
}
F
function

TestRegisterRejectsDuplicateTopic

Parameters

core/relay/manager/manager_test.go:76-84
func TestRegisterRejectsDuplicateTopic(t *testing.T)

{
	relay := New()
	if err := Register[int](relay, "topic", func(context.Context, int) error { return nil }); err != nil {
		t.Fatal(err)
	}
	if err := Register[int](relay, "topic", func(context.Context, int) error { return nil }); err == nil {
		t.Fatal("Register() accepted a duplicate topic")
	}
}
F
function

TestEnqueue_MarshalError

Parameters

core/relay/manager/manager_test.go:86-93
func TestEnqueue_MarshalError(t *testing.T)

{
	r := New()

	err := Enqueue[testPayload](context.Background(), r, "test", testPayload{})
	if err != nil {
		t.Fatalf("Enqueue failed: %v", err)
	}
}
F
function

TestBroker_Interface

Parameters

core/relay/manager/manager_test.go:95-101
func TestBroker_Interface(t *testing.T)

{
	b := broker.NewMemoryBroker()
	r := New(WithBroker(b))
	if r.broker != b {
		t.Error("broker not set correctly")
	}
}
F
function

TestStart_ContextCancel

Parameters

core/relay/manager/manager_test.go:103-113
func TestStart_ContextCancel(t *testing.T)

{
	r := New()
	ctx, cancel := context.WithCancel(context.Background())
	cancel()

	ready, err := r.Start(ctx)
	if err == nil {
		t.Fatal("Start accepted a cancelled context")
	}
	<-ready
}
F
function

TestMultipleTopics

Parameters

core/relay/manager/manager_test.go:115-155
func TestMultipleTopics(t *testing.T)

{
	r := New()
	ch1 := make(chan int, 1)
	ch2 := make(chan int, 1)

	Register[int](r, "topic.a", func(ctx context.Context, p int) error {
		ch1 <- p
		return nil
	})
	Register[int](r, "topic.b", func(ctx context.Context, p int) error {
		ch2 <- p
		return nil
	})

	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()
	ready, err := r.Start(ctx)
	if err != nil {
		t.Fatalf("Start failed: %v", err)
	}
	<-ready

	Enqueue[int](context.Background(), r, "topic.a", 1)
	Enqueue[int](context.Background(), r, "topic.b", 2)

	var results [2]int
	select {
	case results[0] = <-ch1:
	case <-time.After(time.Second):
		t.Fatal("timeout for topic.a")
	}
	select {
	case results[1] = <-ch2:
	case <-time.After(time.Second):
		t.Fatal("timeout for topic.b")
	}

	if results[0] != 1 || results[1] != 2 {
		t.Errorf("got %v, want [1 2]", results[:])
	}
}
F
function

TestEnqueueAndRegisterAreThreadSafe

Parameters

core/relay/manager/manager_test.go:157-178
func TestEnqueueAndRegisterAreThreadSafe(t *testing.T)

{
	r := New()

	var wg sync.WaitGroup
	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func(i int) {
			defer wg.Done()
			handler := func(ctx context.Context, p int) error { return nil }
			Register[int](r, "", handler)
		}(i)
	}

	for i := 0; i < 10; i++ {
		wg.Add(1)
		go func() {
			defer wg.Done()
			Enqueue[int](context.Background(), r, "", 0)
		}()
	}
	wg.Wait()
}
F
function

BenchmarkEnqueue

Parameters

core/relay/manager/manager_test.go:180-191
func BenchmarkEnqueue(b *testing.B)

{
	r := New()
	Register[int](r, "bench", func(ctx context.Context, p int) error { return nil })
	ctx, cancel := context.WithCancel(context.Background())
	defer cancel()
	go r.Start(ctx)

	b.ResetTimer()
	for i := 0; i < b.N; i++ {
		Enqueue[int](context.Background(), r, "bench", i)
	}
}
F
function

TestStartRejectsDuplicateAndContextCancellationUnsubscribes

Parameters

core/relay/manager/manager_test.go:193-231
func TestStartRejectsDuplicateAndContextCancellationUnsubscribes(t *testing.T)

{
	r := New()
	var calls int
	if err := Register[int](r, "topic", func(context.Context, int) error {
		calls++
		return nil
	}); err != nil {
		t.Fatal(err)
	}
	ctx, cancel := context.WithCancel(context.Background())
	ready, err := r.Start(ctx)
	if err != nil {
		t.Fatal(err)
	}
	<-ready
	if _, err := r.Start(context.Background()); err == nil {
		t.Fatal("second Start() succeeded")
	}
	cancel()
	deadline := time.Now().Add(time.Second)
	for {
		r.lifecycle.Lock()
		started := r.started
		r.lifecycle.Unlock()
		if !started {
			break
		}
		if time.Now().After(deadline) {
			t.Fatal("relay did not stop after context cancellation")
		}
		time.Sleep(time.Millisecond)
	}
	if err := Enqueue(context.Background(), r, "topic", 1); err != nil {
		t.Fatal(err)
	}
	if calls != 0 {
		t.Fatalf("handler called %d times after stop", calls)
	}
}
F
function

TestRegisterRejectsOversizedMessage

Parameters

core/relay/manager/manager_test.go:233-241
func TestRegisterRejectsOversizedMessage(t *testing.T)

{
	r := New()
	if err := Register[string](r, "topic", func(context.Context, string) error { return nil }); err != nil {
		t.Fatal(err)
	}
	if err := Enqueue(context.Background(), r, "topic", string(make([]byte, maxMessageSize))); err == nil {
		t.Fatal("Enqueue() accepted an oversized message")
	}
}
S
struct

reentrantSubscription

core/relay/manager/manager_test.go:243-245
type reentrantSubscription struct

Methods

Unsubscribe
Method

Returns

error
func (*reentrantSubscription) Unsubscribe() error
{
	_, err := s.relay.Start(context.Background())
	if err == nil {
		return errors.New("relay restarted during unsubscribe")
	}
	return nil
}

Fields

Name Type Description
relay *Relay
S
struct

reentrantBroker

core/relay/manager/manager_test.go:255-257
type reentrantBroker struct

Methods

Publish
Method

Parameters

string
[]byte

Returns

error
func (*reentrantBroker) Publish(context.Context, string, []byte) error
{
	return nil
}
Subscribe
Method

Parameters

Returns

func (*reentrantBroker) Subscribe(string, broker.Handler) (broker.Subscription, error)
{
	if err := b.relay.Stop(); err != nil {
		return nil, err
	}
	return &reentrantSubscription{relay: b.relay}, nil
}

Fields

Name Type Description
relay *Relay
F
function

TestRelayBrokerCallbacksCanReenterLifecycle

Parameters

core/relay/manager/manager_test.go:270-293
func TestRelayBrokerCallbacksCanReenterLifecycle(t *testing.T)

{
	customBroker := &reentrantBroker{}
	relay := New(WithBroker(customBroker))
	customBroker.relay = relay
	if err := Register[int](relay, "topic", func(context.Context, int) error {
		return nil
	}); err != nil {
		t.Fatal(err)
	}

	done := make(chan error, 1)
	go func() {
		_, err := relay.Start(context.Background())
		done <- err
	}()
	select {
	case err := <-done:
		if err == nil {
			t.Fatal("Start() ignored the reentrant Stop()")
		}
	case <-time.After(time.Second):
		t.Fatal("Start() deadlocked in a reentrant broker callback")
	}
}