Versions in this module Expand all Collapse all v0 v0.1.1 Sep 7, 2026 v0.1.0 Aug 31, 2026 Changes in this version + type AMQPHeaderCarrier amqp091.Table + func (c AMQPHeaderCarrier) Get(key string) string + func (c AMQPHeaderCarrier) Keys() []string + func (c AMQPHeaderCarrier) Set(key string, value string) + type DelayStrategy int + const DelayDLXTTL + type ExchangeName string + type Logger interface + Errorw func(msg string, keysAndValues ...any) + Infof func(template string, args ...any) + Infow func(msg string, keysAndValues ...any) + type Message struct + Body []byte + ContentType string + Expiration time.Duration + Headers amqp091.Table + MessageID string + Priority uint8 + Timestamp time.Time + func NewMessage() *Message + func (m *Message) WithBody(body []byte) *Message + func (m *Message) WithContentType(contentType string) *Message + func (m *Message) WithExpiration(d time.Duration) *Message + func (m *Message) WithHeader(key string, value any) *Message + func (m *Message) WithJSONBody(v any) *Message + func (m *Message) WithMessageID(id string) *Message + func (m *Message) WithPriority(priority uint8) *Message + type Metrics interface + IncMessagesAcked func(ctx context.Context, queue string) + IncMessagesNacked func(ctx context.Context, queue string) + IncMessagesReceived func(ctx context.Context, queue string) + IncPublishFailed func(ctx context.Context, queue string) + IncPublishSuccess func(ctx context.Context, queue string) + IncReconnects func(ctx context.Context) + ObserveProcessingDuration func(ctx context.Context, queue string, durationSeconds float64) + type MockRabbitMQKit struct + func NewMockRabbitMQKit(ctrl *gomock.Controller) *MockRabbitMQKit + func (m *MockRabbitMQKit) BatchPublish(ctx context.Context, exchange, routingKey string, messages []*Message) error + func (m *MockRabbitMQKit) DeclareQueues(ctx context.Context, queues ...QueueConfig) error + func (m *MockRabbitMQKit) EXPECT() *MockRabbitMQKitMockRecorder + func (m *MockRabbitMQKit) Publish(ctx context.Context, exchange, routingKey string, msg *Message) error + func (m *MockRabbitMQKit) PublishWithDelayDLXTTL(ctx context.Context, queueName string, msg *Message) error + func (m *MockRabbitMQKit) SetLogger(logger Logger) + func (m *MockRabbitMQKit) SetQos(ctx context.Context, prefetchCount, prefetchSize int, global bool) error + func (m *MockRabbitMQKit) Shutdown(ctx context.Context) error + func (m *MockRabbitMQKit) Start(ctx context.Context, queues []QueueConfig) error + type MockRabbitMQKitMockRecorder struct + func (mr *MockRabbitMQKitMockRecorder) BatchPublish(ctx, exchange, routingKey, messages any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) DeclareQueues(ctx any, queues ...any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) Publish(ctx, exchange, routingKey, msg any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) PublishWithDelayDLXTTL(ctx, queueName, msg any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) SetLogger(logger any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) SetQos(ctx, prefetchCount, prefetchSize, global any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) Shutdown(ctx any) *gomock.Call + func (mr *MockRabbitMQKitMockRecorder) Start(ctx, queues any) *gomock.Call + type MsgHandler func(ctx context.Context, queueName QueueName, msg RabbitMQMsg) error + type NoOpLogger struct + func (l *NoOpLogger) Errorw(_ string, _ ...any) + func (l *NoOpLogger) Infof(_ string, _ ...any) + func (l *NoOpLogger) Infow(_ string, _ ...any) + type QueueConfig struct + AutoDelete bool + ConsumerTag string + ConsumerTimeout time.Duration + DLXTTL time.Duration + DelayStrategy DelayStrategy + DummyMessageEnabled bool + DummyMessageFrequency time.Duration + Durable bool + EnableDLQ bool + Exchange ExchangeName + Exclusive bool + Handler MsgHandler + IsDelayQueue bool + MaxRedelivery int + Name QueueName + NoWait bool + NumWorkers int + ProcessTimeout time.Duration + QueueType QueueType + RoutingKey RoutingKeyName + type QueueName string + type QueueType string + const Classic + const Quorum + type RabbitMQKit interface + BatchPublish func(ctx context.Context, exchange, routingKey string, messages []*Message) error + DeclareQueues func(ctx context.Context, queues ...QueueConfig) error + Publish func(ctx context.Context, exchange, routingKey string, msg *Message) error + PublishWithDelayDLXTTL func(ctx context.Context, queueName string, msg *Message) error + SetLogger func(logger Logger) + SetQos func(ctx context.Context, prefetchCount, prefetchSize int, global bool) error + Shutdown func(ctx context.Context) error + Start func(ctx context.Context, queues []QueueConfig) error + func NewSDK(ctx context.Context, cfg SDKConfig) (RabbitMQKit, error) + type RabbitMQMsg struct + type RoutingKeyName string + type SDK struct + func (s *SDK) BatchPublish(ctx context.Context, exchange, routingKey string, messages []*Message) error + func (s *SDK) DeclareQueues(ctx context.Context, queues ...QueueConfig) error + func (s *SDK) Publish(ctx context.Context, exchange, routingKey string, msg *Message) error + func (s *SDK) PublishWithDelayDLXTTL(ctx context.Context, queueName string, msg *Message) error + func (s *SDK) SetLogger(logger Logger) + func (s *SDK) SetQos(ctx context.Context, prefetchCount, prefetchSize int, global bool) error + func (s *SDK) Shutdown(ctx context.Context) error + func (s *SDK) Start(ctx context.Context, queues []QueueConfig) error + type SDKConfig struct + Addr string + ConnectionName string + DialTimeout time.Duration + EnableMetrics bool + EnableTracing bool + ExtractContext func(ctx context.Context, d amqp091.Delivery) context.Context + Heartbeat time.Duration + InjectHeaders func(ctx context.Context, headers amqp091.Table) amqp091.Table + MetricLabelName string + MetricLabelValue func(ctx context.Context) string + MetricsNamespace string + MetricsRegisterer prometheus.Registerer + Password string + PublisherMaxPoolSize int + PublisherPoolSize int + PublisherPoolWait time.Duration + ReconnectDelay time.Duration + TLSConfig *tls.Config + Username string + VHost string