mqs

package
v1.6.0 Latest Latest
Warning

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

Go to latest
Published: Oct 2, 2026 License: Apache-2.0 Imports: 34 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type AWSQueue

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

func NewAWSQueue

func NewAWSQueue(config *AWSSQSConfig) *AWSQueue

func (*AWSQueue) Init

func (q *AWSQueue) Init(ctx context.Context) (func(), error)

func (*AWSQueue) InitSDK

func (q *AWSQueue) InitSDK(ctx context.Context) error

func (*AWSQueue) Publish

func (q *AWSQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error

func (*AWSQueue) Subscribe

func (q *AWSQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type AWSSQSConfig

type AWSSQSConfig struct {
	Endpoint                  string // optional - dev-focused
	Region                    string
	ServiceAccountCredentials string
	Topic                     string
	WaitTime                  time.Duration // optional - defaults to 20s if not set
}

func (*AWSSQSConfig) ToCredentials

ToCredentials returns a static credentials provider, or (nil, nil) when no static credentials are set so the caller falls back to the AWS SDK default credential chain.

type AzureServiceBusConfig

type AzureServiceBusConfig struct {
	Topic        string
	Subscription string
	DLQ          bool // Set to true to subscribe to the dead letter queue

	// Credentials
	ConnectionString string
	// or
	TenantID       string
	ClientID       string
	ClientSecret   string
	SubscriptionID string
	ResourceGroup  string
	Namespace      string
}

type AzureServiceBusQueue added in v0.3.0

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

func NewAzureServiceBusQueue added in v0.3.0

func NewAzureServiceBusQueue(config *AzureServiceBusConfig) *AzureServiceBusQueue

func (*AzureServiceBusQueue) Init added in v0.3.0

func (q *AzureServiceBusQueue) Init(ctx context.Context) (func(), error)

func (*AzureServiceBusQueue) InitClient added in v0.3.0

func (q *AzureServiceBusQueue) InitClient(ctx context.Context) error

func (*AzureServiceBusQueue) Publish added in v0.3.0

func (q *AzureServiceBusQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error

func (*AzureServiceBusQueue) Subscribe added in v0.3.0

func (q *AzureServiceBusQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type ConcurrentSubscription added in v0.14.0

type ConcurrentSubscription interface {
	SupportsConcurrency() bool
}

ConcurrentSubscription indicates a subscription that manages its own concurrency internally (e.g. via SDK flow control). When true, the consumer should skip its own semaphore-based concurrency limiting.

type GCPPubSubConfig

type GCPPubSubConfig struct {
	ProjectID                 string
	TopicID                   string
	SubscriptionID            string
	ServiceAccountCredentials string // JSON key file content
}

type GCPPubSubQueue

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

func NewGCPPubSubQueue

func NewGCPPubSubQueue(config *GCPPubSubConfig, visibilityTimeout time.Duration) *GCPPubSubQueue

func (*GCPPubSubQueue) Init

func (q *GCPPubSubQueue) Init(ctx context.Context) (func(), error)

func (*GCPPubSubQueue) Publish

func (q *GCPPubSubQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error

func (*GCPPubSubQueue) Subscribe

func (q *GCPPubSubQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type InMemoryConfig

type InMemoryConfig struct {
	Name string
}

type InMemoryQueue

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

func NewInMemoryQueue

func NewInMemoryQueue(config *InMemoryConfig) *InMemoryQueue

func (*InMemoryQueue) Init

func (q *InMemoryQueue) Init(ctx context.Context) (func(), error)

func (*InMemoryQueue) Publish

func (q *InMemoryQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error

func (*InMemoryQueue) Subscribe

func (q *InMemoryQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type IncomingMessage

type IncomingMessage interface {
	ToMessage() (*Message, error)
	FromMessage(msg *Message) error
}

type Message

type Message struct {
	QueueMessage
	LoggableID string
	// ID is the broker's message ID, stable across redeliveries. Empty when
	// the message has none.
	ID   string
	Body []byte
}

func (*Message) Reject added in v1.6.0

func (m *Message) Reject()

Reject stops the broker from redelivering the message. RabbitMQ nacks without requeue, which dead-letters the message if the queue has a dead-letter exchange and discards it otherwise. Azure Service Bus dead-letters it. Other brokers have no such option, so the message is nacked and their own redelivery policy applies.

func (*Message) Rejectable added in v1.6.0

func (m *Message) Rejectable() bool

Rejectable reports whether Reject stops redelivery.

type NATSConfig added in v1.6.0

type NATSConfig struct {
	ServerURL  string
	Stream     string
	Subject    string
	DLQSubject string // optional; exhausted messages are moved here instead of being lost
}

NATSConfig configures a JetStream-backed queue. Unlike core NATS pub/sub, JetStream persists messages and supports at-least-once delivery via explicit ack/nak — the same durability guarantee RabbitMQ provides here.

NATS servers default max_payload to 1MiB. A larger event passes Outpost's own API validation and then fails at publish time here, so a server handling larger events needs max_payload raised in its own config.

type NATSQueue added in v1.6.0

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

func NewNATSQueue added in v1.6.0

func NewNATSQueue(config *NATSConfig) *NATSQueue

func (*NATSQueue) Init added in v1.6.0

func (q *NATSQueue) Init(ctx context.Context) (func(), error)

func (*NATSQueue) Publish added in v1.6.0

func (q *NATSQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error

func (*NATSQueue) Subscribe added in v1.6.0

func (q *NATSQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type NATSSubscription added in v1.6.0

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

func (*NATSSubscription) Receive added in v1.6.0

func (s *NATSSubscription) Receive(ctx context.Context) (*Message, error)

Receive uses the Messages() iterator rather than repeated single-message Fetch calls: Fetch doesn't take a context, so a Receive blocked waiting on it can hold up shutdown for its full FetchMaxWait (30s here); Next, via NextContext, returns as soon as ctx is done. Messages also lets JetStream overlap pull requests instead of paying a full round-trip per message, which Fetch(1, ...) caps regardless of consumer concurrency settings.

Every delivery is handed back here, including the last one JetStream will ever make for this message — unlike RabbitMQ/SQS, nothing here decides a message is "exhausted" before the caller gets to see it. That decision instead happens in the returned message's own Nack (see natsQueueMessage), so the handler gets its full MaxDeliver attempts, the same as every other provider, rather than losing the last one to this provider's own client-side DLQ handling.

func (*NATSSubscription) Shutdown added in v1.6.0

func (s *NATSSubscription) Shutdown(ctx context.Context) error

type Queue

type Queue interface {
	Init(ctx context.Context) (func(), error)
	Publish(ctx context.Context, msg IncomingMessage) error
	Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)
}

func NewQueue

func NewQueue(config *QueueConfig) Queue

type QueueConfig

type QueueConfig struct {
	AWSSQS          *AWSSQSConfig
	AzureServiceBus *AzureServiceBusConfig
	GCPPubSub       *GCPPubSubConfig
	RabbitMQ        *RabbitMQConfig
	NATS            *NATSConfig
	InMemory        *InMemoryConfig // mainly for testing purposes

	VisibilityTimeout time.Duration
}

type QueueMessage

type QueueMessage interface {
	Ack()
	Nack()
}

type RabbitMQConfig

type RabbitMQConfig struct {
	ServerURL string
	Exchange  string // optional
	Queue     string
	// Dial, when set, opens the broker connection in place of a direct TCP
	// dial, e.g. through a proxy chain.
	Dial proxychain.DialFunc
}

type RabbitMQQueue

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

func NewRabbitMQQueue

func NewRabbitMQQueue(config *RabbitMQConfig) *RabbitMQQueue

func (*RabbitMQQueue) Init

func (q *RabbitMQQueue) Init(ctx context.Context) (func(), error)

func (*RabbitMQQueue) Publish

func (q *RabbitMQQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error

func (*RabbitMQQueue) Subscribe

func (q *RabbitMQQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type Rejecter added in v1.6.0

type Rejecter interface {
	Reject()
}

Rejecter is implemented by queue messages whose broker can stop redelivering a message in a way its owner can see.

type SubscribeOption added in v0.14.0

type SubscribeOption func(*SubscribeOptions)

SubscribeOption configures subscription behavior.

func WithConcurrency added in v0.14.0

func WithConcurrency(n int) SubscribeOption

WithConcurrency sets the max in-flight messages for the subscription.

type SubscribeOptions added in v0.14.0

type SubscribeOptions struct {
	Concurrency int
}

SubscribeOptions holds options for Subscribe.

func ApplySubscribeOptions added in v0.14.0

func ApplySubscribeOptions(opts []SubscribeOption) SubscribeOptions

ApplySubscribeOptions applies all options and returns the result.

type Subscription

type Subscription interface {
	Receive(ctx context.Context) (*Message, error)
	Shutdown(ctx context.Context) error
}

type UnimplementedQueue

type UnimplementedQueue struct{}

func (*UnimplementedQueue) Init

func (q *UnimplementedQueue) Init(ctx context.Context) (func(), error)

func (*UnimplementedQueue) Publish

func (*UnimplementedQueue) Subscribe

func (q *UnimplementedQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)

type WrappedSubscription

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

func (*WrappedSubscription) Receive

func (s *WrappedSubscription) Receive(ctx context.Context) (*Message, error)

func (*WrappedSubscription) Shutdown

func (s *WrappedSubscription) Shutdown(ctx context.Context) error

Jump to

Keyboard shortcuts

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