Documentation
¶
Overview ¶
Package broker defines an asynchronous publish/subscribe abstraction with acknowledgement semantics, modeled after go-micro's Broker. It is the cross-process counterpart to pkg/event (which is process-local): a Broker moves messages between services, with at-least-once delivery when the implementation supports Ack.
Index ¶
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
func RegisterBroker ¶
RegisterBroker registers a Broker factory under name. Backend packages call this from init().
func Registered ¶
Registered reports whether a broker backend is registered under name.
Types ¶
type Broker ¶
type Broker interface {
// Publish sends a message to topic.
Publish(ctx context.Context, topic string, m *Message, opts ...PublishOption) error
// Subscribe registers a handler for topic.
Subscribe(ctx context.Context, topic string, h Handler, opts ...SubscribeOption) (Subscriber, error)
// Connect establishes the underlying connection.
Connect() error
// Disconnect closes the underlying connection.
Disconnect() error
// String names the broker implementation.
String() string
}
Broker publishes messages to topics and delivers them to subscribers.
type Event ¶
type Event interface {
// Topic returns the topic this event was published to.
Topic() string
// Message returns the underlying message.
Message() *Message
// Ack acknowledges the message. For at-least-once brokers this commits the
// message; for at-most-once brokers it is a no-op.
Ack() error
// Error returns a delivery error, if any.
Error() error
}
Event is a delivered message received by a subscriber.
type Factory ¶
Factory constructs a Broker from a backend-specific options value. Concrete backends register their factory in init() and type-assert the options.
type Message ¶
type Message struct {
Header metadata.Metadata `json:"header,omitempty"`
Body []byte `json:"body,omitempty"`
}
Message is a single broker message.
type PublishOption ¶
type PublishOption func(*PublishOptions)
PublishOption configures a Publish call.
type SubscribeOption ¶
type SubscribeOption func(*SubscribeOptions)
SubscribeOption configures a Subscribe call.
func WithAutoAck ¶
func WithAutoAck(v bool) SubscribeOption
WithAutoAck sets whether messages are auto-acknowledged on handler success.
func WithConcurrency ¶
func WithConcurrency(n int) SubscribeOption
WithConcurrency sets the number of concurrent consumers.
func WithQueue ¶
func WithQueue(queue string) SubscribeOption
WithQueue sets the consumer group name.
type SubscribeOptions ¶
type SubscribeOptions struct {
// Queue is a consumer-group name. Subscribers sharing a queue (and topic)
// load-balance messages; otherwise each subscriber gets every message.
Queue string
// AutoAck acknowledges a message automatically when the handler returns nil.
// When false, the handler must call Event.Ack explicitly.
AutoAck bool
// Concurrency is the number of goroutines consuming the subscription.
Concurrency int
}
SubscribeOptions holds Subscribe tunables.
type Subscriber ¶
type Subscriber interface {
// Topic returns the subscribed topic.
Topic() string
// Unsubscribe cancels the subscription and releases its resources.
Unsubscribe() error
}
Subscriber is an active subscription.
Directories
¶
| Path | Synopsis |
|---|---|
|
Package memory provides an in-process broker.Broker for tests and single-node deployments.
|
Package memory provides an in-process broker.Broker for tests and single-node deployments. |
|
Package redis provides a Redis Streams-backed broker.Broker with at-least-once delivery and explicit Ack.
|
Package redis provides a Redis Streams-backed broker.Broker with at-least-once delivery and explicit Ack. |