Documentation
¶
Overview ¶
Package libbus is a high-level publish-subscribe abstraction over a message broker, offering fire-and-forget publish, streaming subscriptions, and request-reply (Serve/Request) on top of pluggable backends (NATS, SQLite, in-memory).
Index ¶
- Variables
- func SetupNatsInstance(ctx context.Context) (string, testcontainers.Container, func(), error)
- type Config
- type Handler
- type InMem
- func (p *InMem) Close() error
- func (p *InMem) Publish(ctx context.Context, subject string, data []byte) error
- func (p *InMem) Request(ctx context.Context, subject string, data []byte) ([]byte, error)
- func (p *InMem) Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)
- func (p *InMem) Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)
- type Messenger
- type SQLiteBus
- func (b *SQLiteBus) Close() error
- func (b *SQLiteBus) Publish(ctx context.Context, subject string, data []byte) error
- func (b *SQLiteBus) Request(ctx context.Context, subject string, data []byte) ([]byte, error)
- func (b *SQLiteBus) Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)
- func (b *SQLiteBus) Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)
- type SQLiteBusOptions
- type Subscription
Constants ¶
This section is empty.
Variables ¶
var ( // ErrConnectionClosed is returned when an operation is attempted on a closed connection. ErrConnectionClosed = errors.New("connection closed") // ErrStreamSubscriptionFail is returned when a stream subscription fails. ErrStreamSubscriptionFail = errors.New("stream subscription failed") // ErrMessagePublish is returned when publishing a message fails for reasons other than a closed connection. ErrMessagePublish = errors.New("message publishing failed") // ErrRequestTimeout is returned when a request-reply operation times out. ErrRequestTimeout = errors.New("request timed out") ErrNoResponders = fmt.Errorf("%w: no handler was subscribed to the subject, so the request was never delivered and is safe to retry", ErrRequestTimeout) )
Functions ¶
func SetupNatsInstance ¶
Types ¶
type Handler ¶
Handler processes a request and returns a response for Messenger.Serve.
type InMem ¶
type InMem struct {
// contains filtered or unexported fields
}
InMem is an in-memory Messenger for single-process use. It reproduces the NATS backend's observable contract.
func (*InMem) Publish ¶
Publish hands the message to every Stream subscriber's queue and returns. It never blocks on a consumer; a full queue drops the message.
func (*InMem) Request ¶
Request invokes the Serve handler registered for the subject, in the caller's goroutine. A missing handler fails immediately rather than at the deadline.
type Messenger ¶
type Messenger interface {
// Publish sends a fire-and-forget message to a given subject.
Publish(ctx context.Context, subject string, data []byte) error
// Stream subscribes to a subject and delivers messages to ch until ctx is
// cancelled.
Stream(ctx context.Context, subject string, ch chan<- []byte) (Subscription, error)
// Request sends a request message and waits for a reply. The context can be
// used to set a timeout or to cancel the request.
Request(ctx context.Context, subject string, data []byte) ([]byte, error)
// Serve registers a handler for a subject; the returned Subscription stops it.
Serve(ctx context.Context, subject string, handler Handler) (Subscription, error)
// Close disconnects from the messaging server and releases its resources.
Close() error
}
Messenger is a pub-sub/request-reply interface for distributing lightweight messages between services. Every backend guarantees that Publish to a subject with no subscribers never blocks, Stream delivers in publish order, a handler error still yields a reply, and every method returns ErrConnectionClosed after Close. Backends differ in durability, backpressure and latency, so always give Request a deadline.
func NewTestPubSub ¶
NewTestPubSub starts a NATS container using SetupNatsInstance, creates a new PubSub instance, and returns it along with a cleanup function.
type SQLiteBus ¶ added in v0.4.0
type SQLiteBus struct {
// contains filtered or unexported fields
}
SQLiteBus implements Messenger over a SQLite database. The bus_events, bus_requests and bus_replies tables must exist before use.
func NewSQLite ¶ added in v0.4.0
func NewSQLite(exec sqlExec) *SQLiteBus
NewSQLite creates a SQLite-backed Messenger over the result of dbManager.WithoutTransaction().
func NewSQLiteWithOptions ¶ added in v0.6.8
func NewSQLiteWithOptions(exec sqlExec, opt SQLiteBusOptions) *SQLiteBus
NewSQLiteWithOptions is like NewSQLite but allows tuning poll intervals.
func (*SQLiteBus) Close ¶ added in v0.4.0
Close stops all background goroutines. The underlying database is not closed.
func (*SQLiteBus) Publish ¶ added in v0.4.0
Publish inserts a row into bus_events so Stream subscribers can pick it up.
func (*SQLiteBus) Request ¶ added in v0.4.0
Request inserts a request row and polls for the reply until ctx deadline or 10s timeout.
type SQLiteBusOptions ¶ added in v0.6.8
SQLiteBusOptions overrides poll intervals.
type Subscription ¶
type Subscription interface {
// Unsubscribe removes the subscription, stopping the delivery of messages.
Unsubscribe() error
}
Subscription represents an active subscription to a subject.