ezbus

package module
v0.13.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 28, 2026 License: MIT Imports: 15 Imported by: 0

README

go-ezbus CI

This is a package for communication between software components. It makes sending, publishing and receiving messages super easy!

Using RabbitMQ as transport for messages. More transports can and will (hopefully) be added

install

go get github.com/zapote/go-ezbus

idea

Ezbus is great to use when working in a distrubuted system. Publish events when a software executes a command and let rest of the system know.

Plugin new software components in your architecture without touching the existing ones.

Ezbus is super easy to use and will get you started in no time.

code example

//PlaceOrder command
type PlaceOrder struct {
	ID string
}

//OrderPlaced event
type OrderPlaced struct {
	ID string
}

//create message router
r := ezbus.NewRouter()

//register handler for message PlaceOrder
r.Handle("PlaceOrder", func(message ezbus.Message) {
    var po PlaceOrder
    json.Unmarshal(m.Body, &po) 
    bus.Publish(OrderPlaced {po.ID})
})

//create a rabbitmq broker
b := rabbitmq.NewBroker("my-queue", rabbitmq.WithURL("amqp://guest:guest@localhost:5672"))

//create the bus with router and broker
bus := ezbus.NewBus(b, r)

//Go!
bus.Go()

Stopping

Stop takes no more messages and waits for the handler that is running, so its message is acked before the connection closes. Deliveries the broker had buffered go back on the queue in the order they came.

Run is Go, a wait and Stop in one call. It returns when the context has ended and the bus has stopped:

ctx, cancel := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer cancel()

err := bus.Run(ctx)

Stop waits at most the drain timeout: 20 seconds, or what the broker is given with rabbitmq.WithDrainTimeout. Keep it below the time the process is given to exit, terminationGracePeriodSeconds in Kubernetes. Shutdown(ctx) waits until ctx ends instead.

When the time runs out the context of the message is cancelled and the connection closed. A handler that passes m.Context() on to its database calls and HTTP requests returns at once. Its message goes back on the queue, not to the error queue, and is delivered again.

Traces and metrics

The bus is instrumented with the OpenTelemetry API. Install a tracer provider and a meter provider in the service and it starts reporting. Without them everything is a no-op.

Traces: SendContext and PublishContext start a producer span and write traceparent into the message headers. A handler gets the consumer span through m.Context(). Pass it on to database calls, HTTP requests and new messages, or the trace stops there.

The library's own log lines about a message, failed attempts and the move to the error queue, are logged with the message context, so a slog handler that adds trace_id from the context puts them on the trace too.

Metrics:

Name Kind What
messaging.client.sent.messages counter sends and publishes, by destination and outcome
messaging.client.operation.duration histogram, s time to send or publish
messaging.client.consumed.messages counter handled messages, by queue, message name and outcome: ok, error_queue, discarded or interrupted
messaging.process.duration histogram, s time to handle a message, all attempts included
ezbus.process.attempts histogram attempts it took; five means the message ended on the error queue

Queue lengths are not here. RabbitMQ reports them itself through rabbitmq_prometheus.

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

func IsHandlerNotFoundErr added in v0.3.0

func IsHandlerNotFoundErr(err error) bool

IsHandlerNotFoundErr checks if error is HandlerNotFoundErr

Types

type Broker

type Broker interface {
	Send(dst string, m Message) error
	Publish(m Message) error
	Start(handle MessageHandler) error
	Stop() error
	Shutdown(ctx context.Context) error
	Endpoint() string
	Subscribe(endpoint string, messageName string) error
}

Broker interface

type Bus

type Bus interface {
	StarterStopper
	Sender
	Publisher
	Subscriber
}

Bus for publishing, sending and receiving messages

func NewBus

func NewBus(b Broker, r Router) Bus

NewBus creates a bus instance for sending and receiving messages.

type HandlerNotFoundErr added in v0.2.2

type HandlerNotFoundErr struct {
	MessageName string
}

HandlerNotFoundErr is returned when handler is not found

func (HandlerNotFoundErr) Error added in v0.2.2

func (e HandlerNotFoundErr) Error() string

type Message

type Message struct {
	Headers map[string]string
	Body    []byte
	// contains filtered or unexported fields
}

Message in EzBus

func NewMessage

func NewMessage(h map[string]string, b []byte) Message

NewMessage creates a new Message instance Using h as headers and b as body

func (Message) Context added in v0.8.0

func (m Message) Context() context.Context

Context of the message. Background when the message carries none, which is the case for messages built outside of Bus.

func (Message) WithContext added in v0.8.0

func (m Message) WithContext(ctx context.Context) Message

WithContext returns a copy of the message with ctx attached, in the same way as http.Request.WithContext.

type MessageHandler

type MessageHandler = func(m Message) error

MessageHandler func for handling messsages

type Middleware

type Middleware = func(next MessageHandler) MessageHandler

Middleware for router message handling pipeline

type Publisher

type Publisher interface {
	Publish(msg interface{}) error
	PublishContext(ctx context.Context, msg interface{}) error
}

Publisher publishes a message to subscribers. PublishContext carries ctx along with the message; Publish is PublishContext with context.Background().

type Router

type Router interface {
	Handle(messageName string, h MessageHandler)
	Middleware(mw Middleware)
	Receive(n string, m Message) error
}

Router routes message to correct MessageHandler func

func NewRouter

func NewRouter() Router

NewRouter creates a new router instance.

type Sender

type Sender interface {
	Send(dst string, msg interface{}) error
	SendContext(ctx context.Context, dst string, msg interface{}) error
}

Sender sends a message to a destination. SendContext carries ctx along with the message; Send is SendContext with context.Background().

type StarterStopper added in v0.1.1

type StarterStopper interface {
	Go() error
	Run(ctx context.Context) error
	Stop() error
	Shutdown(ctx context.Context) error
}

StarterStopper interface. Stop and Shutdown both let the handler that is running finish: Stop within the broker's drain timeout, Shutdown until ctx ends. Run is Go, a wait for ctx to end, and Stop.

type Subscriber

type Subscriber interface {
	Subscribe(endpoint string)
	SubscribeMessage(endpoint string, messageName string)
}

Subscriber interface

Directories

Path Synopsis

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL