Documentation
¶
Overview ¶
Package messagebus provides a transport-agnostic asynchronous message bus with a middleware stack, routing, and a consumer worker.
Index ¶
- Constants
- func BusMustFromContainer(serviceContainer containercontract.Container) messagebuscontract.Bus
- func BusMustFromResolver(resolver containercontract.Resolver) messagebuscontract.Bus
- func ConsumeBusFromResolver(resolver containercontract.Resolver) messagebuscontract.Bus
- func DeadLetterAttemptCount(envelopeInstance messagebuscontract.Envelope) int
- func EnsureEnvelope(message any) messagebuscontract.Envelope
- func HandlerLocatorMustFromContainer(serviceContainer containercontract.Container) messagebuscontract.HandlerLocator
- func HandlerLocatorMustFromResolver(resolver containercontract.Resolver) messagebuscontract.HandlerLocator
- func LastStampOfType[T messagebuscontract.Stamp](envelopeInstance messagebuscontract.Envelope) (T, bool)
- func MessageId(envelopeInstance messagebuscontract.Envelope) (string, bool)
- func NewEnvelope(message any, stamps ...messagebuscontract.Stamp) messagebuscontract.Envelope
- func NewHandleMessageMiddleware(locator messagebuscontract.HandlerLocator) messagebuscontract.Middleware
- func NewHandleMessageMiddlewareWithOptions(locator messagebuscontract.HandlerLocator, options HandleOptions) messagebuscontract.Middleware
- func NewSendMessageMiddleware(routingByType map[reflect.Type]TransportRouting) messagebuscontract.Middleware
- func NewSendMessageMiddlewareFromRouting(routing *Routing) messagebuscontract.Middleware
- func RedeliveryCount(envelopeInstance messagebuscontract.Envelope) int
- func RegisterHandler[T any](locator *HandlerLocator, ...)
- func RegisterTransports(registrar ServiceRegistrar, transports map[string]messagebuscontract.Transport)
- func TransportsMustFromResolver(resolver containercontract.Resolver) map[string]messagebuscontract.Transport
- type BusNameStamp
- type ConsumeCommand
- func NewConsumeCommand(bus messagebuscontract.Bus, transports map[string]messagebuscontract.Transport) *ConsumeCommand
- func NewConsumeCommandFromContainer() *ConsumeCommand
- func NewConsumeCommandWithRetry(bus messagebuscontract.Bus, transports map[string]messagebuscontract.Transport, ...) *ConsumeCommand
- func (instance *ConsumeCommand) Description() string
- func (instance *ConsumeCommand) Flags() []clicontract.Flag
- func (instance *ConsumeCommand) Name() string
- func (instance *ConsumeCommand) Run(runtimeInstance runtimecontract.Runtime, ...) error
- func (instance *ConsumeCommand) WithShutdownGrace(grace time.Duration) *ConsumeCommand
- type DeadLetterAttemptStamp
- type DelayStamp
- type HandleOptions
- type HandledStamp
- type HandlerLocator
- type InMemoryTransport
- func (instance *InMemoryTransport) Ack(runtimeInstance runtimecontract.Runtime, ...) error
- func (instance *InMemoryTransport) Close(runtimeInstance runtimecontract.Runtime) error
- func (instance *InMemoryTransport) Nack(runtimeInstance runtimecontract.Runtime, ...) error
- func (instance *InMemoryTransport) Receive(runtimeInstance runtimecontract.Runtime) (<-chan messagebuscontract.Envelope, error)
- func (instance *InMemoryTransport) Send(runtimeInstance runtimecontract.Runtime, ...) error
- func (instance *InMemoryTransport) WithLogger(logger loggingcontract.Logger) *InMemoryTransport
- type Manager
- type MessageIdStamp
- type ReceivedStamp
- type RedeliveryStamp
- type RetryPolicy
- type Routing
- type SentStamp
- type ServiceRegistrar
- type TransportRouting
Constants ¶
const ( ServiceBus = "service.messagebus.bus" ServiceConsumeBus = "service.messagebus.consume_bus" ServiceHandlerLocator = "service.messagebus.handler_locator" ServiceTransports = "service.messagebus.transports" ServiceRetryPolicy = "service.messagebus.retry_policy" )
const ( StampNameBusName = "bus_name" StampNameSent = "sent" StampNameReceived = "received" StampNameHandled = "handled" StampNameRedelivery = "redelivery" StampNameDelay = "delay" StampNameDeadLetterAttempt = "dead_letter_attempt" StampNameMessageId = "message_id" )
Variables ¶
This section is empty.
Functions ¶
func BusMustFromContainer ¶
func BusMustFromContainer(serviceContainer containercontract.Container) messagebuscontract.Bus
func BusMustFromResolver ¶
func BusMustFromResolver(resolver containercontract.Resolver) messagebuscontract.Bus
func ConsumeBusFromResolver ¶ added in v3.9.0
func ConsumeBusFromResolver(resolver containercontract.Resolver) messagebuscontract.Bus
ConsumeBusFromResolver returns the bus the consumer dispatches received messages into, preferring a dedicated consume bus and falling back to the shared bus so single-bus applications need no extra wiring.
func DeadLetterAttemptCount ¶ added in v3.8.0
func DeadLetterAttemptCount(envelopeInstance messagebuscontract.Envelope) int
func EnsureEnvelope ¶
func EnsureEnvelope(message any) messagebuscontract.Envelope
func HandlerLocatorMustFromContainer ¶
func HandlerLocatorMustFromContainer(serviceContainer containercontract.Container) messagebuscontract.HandlerLocator
func HandlerLocatorMustFromResolver ¶
func HandlerLocatorMustFromResolver(resolver containercontract.Resolver) messagebuscontract.HandlerLocator
func LastStampOfType ¶
func LastStampOfType[T messagebuscontract.Stamp](envelopeInstance messagebuscontract.Envelope) (T, bool)
func MessageId ¶ added in v3.9.0
func MessageId(envelopeInstance messagebuscontract.Envelope) (string, bool)
MessageId returns the producer-assigned message id stamped on the envelope, if any, so a transport can carry it for consumer-side deduplication.
func NewEnvelope ¶
func NewEnvelope(message any, stamps ...messagebuscontract.Stamp) messagebuscontract.Envelope
func NewHandleMessageMiddleware ¶
func NewHandleMessageMiddleware(locator messagebuscontract.HandlerLocator) messagebuscontract.Middleware
func NewHandleMessageMiddlewareWithOptions ¶
func NewHandleMessageMiddlewareWithOptions( locator messagebuscontract.HandlerLocator, options HandleOptions, ) messagebuscontract.Middleware
func NewSendMessageMiddleware ¶
func NewSendMessageMiddleware(routingByType map[reflect.Type]TransportRouting) messagebuscontract.Middleware
func NewSendMessageMiddlewareFromRouting ¶
func NewSendMessageMiddlewareFromRouting(routing *Routing) messagebuscontract.Middleware
func RedeliveryCount ¶
func RedeliveryCount(envelopeInstance messagebuscontract.Envelope) int
func RegisterHandler ¶
func RegisterHandler[T any]( locator *HandlerLocator, handle func(runtimeInstance runtimecontract.Runtime, message T) error, )
func RegisterTransports ¶ added in v3.9.0
func RegisterTransports(registrar ServiceRegistrar, transports map[string]messagebuscontract.Transport)
RegisterTransports registers the named transports the consume command resolves at run time, so registering them is enough for the framework to expose melody:messagebus:consume.
func TransportsMustFromResolver ¶ added in v3.9.0
func TransportsMustFromResolver(resolver containercontract.Resolver) map[string]messagebuscontract.Transport
Types ¶
type BusNameStamp ¶
type BusNameStamp struct {
BusName string
}
func (BusNameStamp) StampName ¶
func (instance BusNameStamp) StampName() string
type ConsumeCommand ¶
type ConsumeCommand struct {
// contains filtered or unexported fields
}
func NewConsumeCommand ¶
func NewConsumeCommand( bus messagebuscontract.Bus, transports map[string]messagebuscontract.Transport, ) *ConsumeCommand
func NewConsumeCommandFromContainer ¶ added in v3.9.0
func NewConsumeCommandFromContainer() *ConsumeCommand
NewConsumeCommandFromContainer builds the command so it resolves the bus, the transport map and an optional retry policy from the service container at run time, letting the framework auto-register melody:messagebus:consume once the application registers its transports.
func NewConsumeCommandWithRetry ¶
func NewConsumeCommandWithRetry( bus messagebuscontract.Bus, transports map[string]messagebuscontract.Transport, retryPolicy RetryPolicy, ) *ConsumeCommand
func (*ConsumeCommand) Description ¶
func (instance *ConsumeCommand) Description() string
func (*ConsumeCommand) Flags ¶
func (instance *ConsumeCommand) Flags() []clicontract.Flag
func (*ConsumeCommand) Name ¶
func (instance *ConsumeCommand) Name() string
func (*ConsumeCommand) Run ¶
func (instance *ConsumeCommand) Run( runtimeInstance runtimecontract.Runtime, commandContext *clicontract.CommandContext, ) error
func (*ConsumeCommand) WithShutdownGrace ¶
func (instance *ConsumeCommand) WithShutdownGrace(grace time.Duration) *ConsumeCommand
type DeadLetterAttemptStamp ¶ added in v3.8.0
type DeadLetterAttemptStamp struct {
Count int
}
func (DeadLetterAttemptStamp) StampName ¶ added in v3.8.0
func (instance DeadLetterAttemptStamp) StampName() string
type DelayStamp ¶
func (DelayStamp) StampName ¶
func (instance DelayStamp) StampName() string
type HandleOptions ¶
type HandleOptions struct {
RequireHandler bool
}
type HandledStamp ¶
type HandledStamp struct {
HandlerName string
}
func (HandledStamp) StampName ¶
func (instance HandledStamp) StampName() string
type HandlerLocator ¶
type HandlerLocator struct {
// contains filtered or unexported fields
}
func NewHandlerLocator ¶
func NewHandlerLocator() *HandlerLocator
func (*HandlerLocator) HandlersFor ¶
func (instance *HandlerLocator) HandlersFor(message any) []messagebuscontract.MessageHandler
func (*HandlerLocator) Register ¶
func (instance *HandlerLocator) Register(messageType reflect.Type, handler messagebuscontract.MessageHandler)
type InMemoryTransport ¶
type InMemoryTransport struct {
// contains filtered or unexported fields
}
func NewInMemoryTransport ¶
func NewInMemoryTransport(bufferSize int) *InMemoryTransport
func (*InMemoryTransport) Ack ¶
func (instance *InMemoryTransport) Ack( runtimeInstance runtimecontract.Runtime, envelopeInstance messagebuscontract.Envelope, ) error
func (*InMemoryTransport) Close ¶
func (instance *InMemoryTransport) Close(runtimeInstance runtimecontract.Runtime) error
func (*InMemoryTransport) Nack ¶
func (instance *InMemoryTransport) Nack( runtimeInstance runtimecontract.Runtime, envelopeInstance messagebuscontract.Envelope, requeue bool, ) error
func (*InMemoryTransport) Receive ¶
func (instance *InMemoryTransport) Receive( runtimeInstance runtimecontract.Runtime, ) (<-chan messagebuscontract.Envelope, error)
func (*InMemoryTransport) Send ¶
func (instance *InMemoryTransport) Send( runtimeInstance runtimecontract.Runtime, envelopeInstance messagebuscontract.Envelope, ) error
func (*InMemoryTransport) WithLogger ¶
func (instance *InMemoryTransport) WithLogger(logger loggingcontract.Logger) *InMemoryTransport
type Manager ¶
type Manager struct {
// contains filtered or unexported fields
}
func NewManager ¶
func NewManager(name string, middlewares ...messagebuscontract.Middleware) *Manager
func (*Manager) Dispatch ¶
func (instance *Manager) Dispatch( runtimeInstance runtimecontract.Runtime, message any, stamps ...messagebuscontract.Stamp, ) (messagebuscontract.Envelope, error)
type MessageIdStamp ¶ added in v3.9.0
type MessageIdStamp struct {
MessageId string
}
MessageIdStamp carries a stable, producer-assigned identifier for the message so a transport can publish it (for example as the AMQP message id) and a consumer can deduplicate redeliveries. A producer with at-least-once semantics — such as the outbox relay, which may redeliver after a transport-success-then-crash — stamps it with a deterministic id per logical message.
func (MessageIdStamp) StampName ¶ added in v3.9.0
func (instance MessageIdStamp) StampName() string
type ReceivedStamp ¶
type ReceivedStamp struct {
TransportName string
}
func (ReceivedStamp) StampName ¶
func (instance ReceivedStamp) StampName() string
type RedeliveryStamp ¶
type RedeliveryStamp struct {
Count int
}
func (RedeliveryStamp) StampName ¶
func (instance RedeliveryStamp) StampName() string
type RetryPolicy ¶
type RetryPolicy struct {
MaxRetries int
BaseDelay time.Duration
FailureTransport messagebuscontract.Transport
MaxDelay time.Duration
FailureRequeueDelay time.Duration
/* @important bound on requeues of an exhausted message after the FailureTransport rejects it; 0 keeps the default no-loss behavior (requeue until it recovers), a positive value nacks without requeue after that many failed routings so a transport-native dead-letter (AMQP DLX) can claim it instead of looping forever */
MaxDeadLetterAttempts int
}
func RetryPolicyFromResolver ¶ added in v3.9.0
func RetryPolicyFromResolver(resolver containercontract.Resolver) (RetryPolicy, bool)
RetryPolicyFromResolver returns the application-provided retry policy when one is registered; the second result is false when the consumer should keep the framework defaults.
type Routing ¶
type Routing struct {
// contains filtered or unexported fields
}
func NewRouting ¶
func NewRouting() *Routing
type ServiceRegistrar ¶ added in v3.9.0
type ServiceRegistrar interface {
RegisterService(serviceName string, provider any, options ...containercontract.RegisterOption)
}
type TransportRouting ¶
type TransportRouting struct {
Name string
Transport messagebuscontract.Transport
}