manager
API
manager
packageAPI reference for the manager
package.
Imports
(9)
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
ctx
context.Context
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 |
Uses
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.
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
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
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
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)
}
Uses
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
t
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
t
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
t
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
t
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
t
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
t
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
t
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
t
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
t
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
b
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
t
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
t
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
Returns
error
func (*reentrantBroker) Publish(context.Context, string, []byte) error
{
return nil
}
Subscribe
Method
Parameters
string
Returns
error
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
t
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")
}
}