Documentation
¶
Overview ¶
Package amqpadapter provides a production-oriented RabbitMQ adapter for Go.
It provides publishing and consuming of RabbitMQ messages with automatic connection recovery, publisher confirms, concurrent consumers, delayed retries using RabbitMQ delay queues or dead-letter exchanges, graceful shutdown, context propagation, correlation IDs and OpenTelemetry support.
The package is intended for Go services that need reliable RabbitMQ message processing without implementing connection recovery and retry infrastructure in application code.
Index ¶
- Constants
- Variables
- func IsRetryable(err error) bool
- func Retry(err error) error
- type Config
- type ConsumeHeadersExtractor
- type ConsumerHandler
- type FailJobHandler
- type FailedJob
- type Logger
- type Option
- type PublishHeadersBuilder
- type PublishMessage
- type PublishMessagePriority
- type Queue
- type QueueItem
- type QueueName
- type RetryConfig
- type RetryableError
Constants ¶
const DefaultContentType = "text/plain"
DefaultContentType is the default content type for published messages.
Variables ¶
var ErrBindQueueToExchange = errors.New("error binding queue to exchange")
ErrBindQueueToExchange is returned when binding a queue to an exchange fails.
var ErrClientClosedBeforeReady = errors.New("client closed before ready")
ErrClientClosedBeforeReady is returned when the client is closed before becoming ready.
var ErrClosingQueueChannel = errors.New("error closing queue channel")
ErrClosingQueueChannel is returned when closing the queue channel fails.
var ErrClosingQueueConnection = errors.New("error closing queue connection")
ErrClosingQueueConnection is returned when closing the queue connection fails.
var ErrConnectToQueue = errors.New("failed to connect to queue at address")
ErrConnectToQueue is returned when dialing the broker fails.
var ErrConsumeQueueConnectionClosed = errors.New("queue connection closed during consume")
ErrConsumeQueueConnectionClosed is returned when consuming on a closed connection.
var ErrConsumersAlreadyStarted = errors.New("consumers already started")
ErrConsumersAlreadyStarted is returned when InitConsumer was already called.
var ErrContextCanceledBeforeAck = errors.New("context canceled before publish confirm")
ErrContextCanceledBeforeAck is returned when the context is canceled before publish confirm.
var ErrContextCanceledOnConsumerStart = errors.New("context canceled while starting consumers")
ErrContextCanceledOnConsumerStart is returned when the context is canceled during InitConsumer.
var ErrDoneSignalBeforeAck = errors.New("done signal received before publish confirm")
ErrDoneSignalBeforeAck is returned when the client is closed before publish confirm.
var ErrEmptyURL = errors.New("RabbitMQ URL is not set")
ErrEmptyURL is returned when the RabbitMQ URL is empty.
var ErrEnableConfirms = errors.New("failed to enable publisher confirms")
ErrEnableConfirms is returned when enabling publisher confirms fails.
var ErrFailHandlerRequired = errors.New("retry enabled but WithFailHandler was not provided")
ErrFailHandlerRequired is returned when retry is enabled but WithFailHandler is missing.
var ErrInvalidDelay = errors.New("config delay must be greater than zero")
ErrInvalidDelay is returned when a configured delay is not positive.
var ErrInvalidPriority = errors.New("message priority must be in range 0–9")
ErrInvalidPriority is returned when message priority is outside 0–9.
var ErrInvalidRetryConfig = errors.New("invalid retry configuration")
ErrInvalidRetryConfig is returned when retry parameters are invalid.
var ErrNoConsumersRegistered = errors.New("no consumers registered")
ErrNoConsumersRegistered is returned when InitConsumer is called with no consumers.
var ErrOpenChannel = errors.New("failed to open queue channel")
ErrOpenChannel is returned when opening an AMQP channel fails.
var ErrPublishFailed = errors.New("error publishing message to queue")
ErrPublishFailed is returned when publishing a message fails.
var ErrPublishNack = errors.New("broker nacked the published message")
ErrPublishNack is returned when the broker nacks a publish (publisher confirm).
var ErrPushQueueConnectionClosed = errors.New("queue connection closed during publish")
ErrPushQueueConnectionClosed is returned when publishing on a closed connection.
var ErrQueueConnectionClosed = errors.New("queue connection closed")
ErrQueueConnectionClosed is returned when the queue connection is closed.
var ErrQueueConsumerOnly = errors.New("queue is configured for consume only (ConsumerOnly)")
ErrQueueConsumerOnly is returned when publishing to a ConsumerOnly queue.
var ErrQueueNotFound = errors.New("queue not found in configuration")
ErrQueueNotFound is returned when the queue is missing from configuration.
var ErrStartingQueueConsumption = errors.New("error starting queue consumption")
ErrStartingQueueConsumption is returned when starting consume fails.
Functions ¶
func IsRetryable ¶
IsRetryable reports whether err is marked as retryable.
Types ¶
type Config ¶
type Config struct {
URL string
ReconnectDelay time.Duration
ReInitDelay time.Duration
ResendDelay time.Duration
QueueParams map[QueueName]QueueItem
}
Config holds RabbitMQ connection settings and per-queue parameters.
type ConsumeHeadersExtractor ¶
ConsumeHeadersExtractor restores context from AMQP headers on consume. parent is the consumer-loop context (from InitConsumer); the returned context must be derived from it.
type ConsumerHandler ¶
ConsumerHandler processes a message body. Return mq.Retry(err) to request delayed reprocessing.
type FailJobHandler ¶
FailJobHandler is called on permanent failure in retry mode.
type Logger ¶
type Logger interface {
Debug(msg string, args ...any)
Info(msg string, args ...any)
Warn(msg string, args ...any)
Error(msg string, args ...any)
DebugContext(ctx context.Context, msg string, args ...any)
InfoContext(ctx context.Context, msg string, args ...any)
WarnContext(ctx context.Context, msg string, args ...any)
ErrorContext(ctx context.Context, msg string, args ...any)
}
Logger is a minimal logging interface with slog-style method names.
type Option ¶
type Option func(*options)
Option configures Queue at New time.
func WithConsumeHeadersExtractor ¶
func WithConsumeHeadersExtractor(e ConsumeHeadersExtractor) Option
WithConsumeHeadersExtractor sets the function that restores context from AMQP headers on consume.
func WithFailHandler ¶
func WithFailHandler(h FailJobHandler) Option
WithFailHandler sets the permanent-failure handler (required when QueueItem.Retry is set).
func WithLogger ¶
WithLogger sets the logger; the default is a noop logger.
func WithPublishHeadersBuilder ¶
func WithPublishHeadersBuilder(b PublishHeadersBuilder) Option
WithPublishHeadersBuilder sets the function that builds AMQP headers on publish.
type PublishHeadersBuilder ¶
PublishHeadersBuilder builds AMQP headers from context at publish time.
type PublishMessage ¶
type PublishMessage struct {
Body []byte
ContentType string
MaxRetryDuration *time.Duration
Priority *PublishMessagePriority
}
PublishMessage describes the body and options of a published message.
type PublishMessagePriority ¶
type PublishMessagePriority uint8
PublishMessagePriority is the RabbitMQ message priority (0–9).
type Queue ¶
type Queue interface {
Publish(ctx context.Context, queue QueueName, msg PublishMessage) error
AddConsumer(ctx context.Context, queue QueueName, handler ConsumerHandler) error
AddConsumerN(ctx context.Context, queue QueueName, parallelism int, handler ConsumerHandler) error
InitConsumer(ctx context.Context) error
Shutdown(ctx context.Context) error
}
Queue is the public contract for publish, consume registration, and shutdown.
type QueueItem ¶
type QueueItem struct {
// Retry enables retry mode; nil means no retry.
Retry *RetryConfig
// ConsumerOnly means consume-only (no publisher connection).
ConsumerOnly bool
}
QueueItem describes processing options for a single queue.
type RetryConfig ¶
RetryConfig describes retry parameters for a queue.
type RetryableError ¶
type RetryableError struct {
Err error
}
RetryableError wraps an error to signal that the consumer loop should retry.
func (*RetryableError) Error ¶
func (e *RetryableError) Error() string
func (*RetryableError) Unwrap ¶
func (e *RetryableError) Unwrap() error