amqp

package module
v0.0.0-...-7593c98 Latest Latest
Warning

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

Go to latest
Published: Aug 27, 2025 License: MIT Imports: 12 Imported by: 0

Documentation

Overview

Original license: ----------------------------------------------------------------------------------- MIT License

Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the "Software"), to deal in the Software without restriction, including without limitation the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to permit persons to whom the Software is furnished to do so, subject to the following conditions:

The above copyright notice and this permission notice shall be included in all copies or substantial portions of the Software.

THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. -----------------------------------------------------------------------------------

Original license: ----------------------------------------------------------------------------------- MIT License

Permission is hereby granted, free of charge, to any person obtaining a copy of this software and associated documentation files (the "Software"), to deal in the Software without restriction, including without limitation the rights to use, copy, modify, merge, publish, distribute, sublicense, and/or sell copies of the Software, and to permit persons to whom the Software is furnished to do so, subject to the following conditions:

The above copyright notice and this permission notice shall be included in all copies or substantial portions of the Software.

THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY, FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM, OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE SOFTWARE. -----------------------------------------------------------------------------------

Index

Constants

This section is empty.

Variables

View Source
var DefaultRetryStrategy = func(iterationCount int) time.Duration {
	return 1 * time.Second
}

DefaultRetryStrategy is just to wait 1 sec between reconnection attempts

Functions

func ExtractPrefixedMetadata

func ExtractPrefixedMetadata(headers amqp091.Table, prefix string) message.Metadata

ExtractPrefixedMetadata extracts headers with a specific prefix and returns them as metadata with the prefix removed. This is useful for extracting namespaced headers from AMQP deliveries.

func HeadersAsMetadata

func HeadersAsMetadata(headers amqp091.Table) message.Metadata

HeadersAsMetadata converts AMQP headers (amqp.Table) to message metadata. This is useful when you need to extract metadata from AMQP delivery headers.

func MarshallPayloadToJson

func MarshallPayloadToJson() middleware.Middleware

MarshallPayloadToJson marshalls the message payload into JSON bytes. This middleware converts any payload type into a JSON byte array.

func MergeHeadersToMetadata

func MergeHeadersToMetadata(msg *message.Message, headers amqp091.Table)

MergeHeadersToMetadata merges AMQP headers into existing message metadata. If there are key conflicts, the header values will override the metadata values.

func MetadataAsHeaders

func MetadataAsHeaders(msg *message.Message) map[string]interface{}

func UnmarshallPayloadFromJson

func UnmarshallPayloadFromJson[T any](payloadType T) middleware.Middleware

UnmarshallPayloadFromJson unmarshalls a JSON byte payload into the specified type. This middleware expects the message payload to be a []byte containing JSON data.

Types

type AMQPMessage

type AMQPMessage struct {
	Publishing *amqp.Publishing
	Immediate  bool
	Mandatory  bool
	RoutingKey string
}

type ConnectionOption

type ConnectionOption func(*ConnectionOptions)

func WithReconnectionRetry

func WithReconnectionRetry(s RetryDelayStrategy) ConnectionOption

type ConnectionOptions

type ConnectionOptions struct {
	// contains filtered or unexported fields
}

type DefaultMessageOptions

type DefaultMessageOptions struct {
	ContentType  string
	Priority     uint8
	DeliveryMode uint8
}

type ExchangeDestination

type ExchangeDestination struct {
	// contains filtered or unexported fields
}

func (ExchangeDestination) Close

func (d ExchangeDestination) Close() error

func (ExchangeDestination) GetMessageID

func (d ExchangeDestination) GetMessageID(message *AMQPMessage) string

func (ExchangeDestination) Publish

func (d ExchangeDestination) Publish(ctx context.Context, message *AMQPMessage) error

type HeadersProvider

type HeadersProvider func(message *message.Message) map[string]interface{}

func ExcludePrefixHeadersProvider

func ExcludePrefixHeadersProvider(prefix string) HeadersProvider

ExcludePrefixHeadersProvider creates a HeadersProvider that excludes metadata keys that start with the specified prefix. This is useful for filtering out internal metadata.

func FilteredHeadersProvider

func FilteredHeadersProvider(filter func(key string) bool) HeadersProvider

FilteredHeadersProvider creates a HeadersProvider that only includes metadata keys that match the provided filter function.

func OnlyPrefixHeadersProvider

func OnlyPrefixHeadersProvider(prefix string) HeadersProvider

OnlyPrefixHeadersProvider creates a HeadersProvider that only includes metadata keys that start with the specified prefix. This is useful for including only specific metadata.

func PrefixedHeadersProvider

func PrefixedHeadersProvider(prefix string) HeadersProvider

PrefixedHeadersProvider creates a HeadersProvider that adds a prefix to all metadata keys when converting them to AMQP headers. This is useful for namespacing custom headers.

type MessageModifierFn

type MessageModifierFn func(msg *amqp.Publishing) error

type Publisher

type Publisher struct {
	// contains filtered or unexported fields
}

func NewAMQPPublisher

func NewAMQPPublisher(
	channel *ReconnectingChannel,
	exchangeRoutingFn publisher.RoutingFunc,
	routingKeyFn RoutingKeyFunc,
	defaultMsgOptions DefaultMessageOptions,
	opts ...PublisherOption) (*Publisher, error)

func (*Publisher) Close

func (p *Publisher) Close() error

func (*Publisher) Publish

func (p *Publisher) Publish(message *message.Message) error

type PublisherOption

type PublisherOption func(*PublisherOptions)

func WithHeadersProvider

func WithHeadersProvider(provider HeadersProvider) PublisherOption

func WithImmediateMsgFn

func WithImmediateMsgFn(immediateOpt publisher.MessageBoolOptFunc) PublisherOption

func WithMandatoryMsgFn

func WithMandatoryMsgFn(mandatoryOpt publisher.MessageBoolOptFunc) PublisherOption

func WithMessageModifier

func WithMessageModifier(modifier MessageModifierFn) PublisherOption

type PublisherOptions

type PublisherOptions struct {
	// contains filtered or unexported fields
}

type ReconnectingChannel

type ReconnectingChannel struct {
	*amqp.Channel
	// contains filtered or unexported fields
}

func (*ReconnectingChannel) Close

func (ch *ReconnectingChannel) Close() error

Close ensure closed flag set

func (*ReconnectingChannel) Consume

func (ch *ReconnectingChannel) Consume(queue, consumer string, autoAck, exclusive, noLocal, noWait bool, args amqp.Table) (<-chan amqp.Delivery, error)

Consume wrap amqp.Channel.Consume, the returned delivery will end only when channel closed by developer

func (*ReconnectingChannel) IsClosed

func (ch *ReconnectingChannel) IsClosed() bool

IsClosed indicate closed by developer

type ReconnectingConnection

type ReconnectingConnection struct {
	*amqp.Connection
	// contains filtered or unexported fields
}

func Dial

func Dial(url string, opts ...ConnectionOption) (*ReconnectingConnection, error)

func DialConfig

func DialConfig(url string, config amqp.Config, opts ...ConnectionOption) (*ReconnectingConnection, error)

func (*ReconnectingConnection) Channel

Channel wrap amqp.Connection.Channel, get a auto reconnect channel

type RetryDelayStrategy

type RetryDelayStrategy func(iterationCount int) time.Duration

type RoutingKeyFunc

type RoutingKeyFunc func(message *message.Message) (string, error)

func ConstantRoutingKey

func ConstantRoutingKey(routingKey string) RoutingKeyFunc

type Subscriber

type Subscriber struct {
	// contains filtered or unexported fields
}

func NewAMQPSubscriber

func NewAMQPSubscriber(channel *ReconnectingChannel, queue string, opts ...SubscriberOption) (*Subscriber, error)

func (*Subscriber) Close

func (s *Subscriber) Close() error

func (*Subscriber) Subscribe

func (s *Subscriber) Subscribe(ctx context.Context) (<-chan *message.Message, error)

type SubscriberOption

type SubscriberOption func(*SubscriberOptions)

func WithAutoAck

func WithAutoAck() SubscriberOption

WithAutoAck will automatically ack the message when it's received.

func WithExclusive

func WithExclusive() SubscriberOption

WithExclusive will make this subscriber exclusive to the target queue.

func WithNoRequeueOnNack

func WithNoRequeueOnNack() SubscriberOption

WithNoRequeueOnNack will not requeue the message when it's nacked.

func WithProcessingTimeout

func WithProcessingTimeout(timeout time.Duration) SubscriberOption

WithProcessingTimeout will dictate how long a message will be processed before it is nacked. 0 means no timeout, wait forever.

func WithProcessingTimeoutHandler

func WithProcessingTimeoutHandler(handler func(ctx context.Context, msg *amqp.Delivery)) SubscriberOption

WithProcessingTimeoutHandler is a function that is called when a message processing times out.

type SubscriberOptions

type SubscriberOptions struct {
	// contains filtered or unexported fields
}

Directories

Path Synopsis

Jump to

Keyboard shortcuts

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