Documentation
¶
Index ¶
- type AWSQueue
- type AWSSQSConfig
- type AzureServiceBusConfig
- type AzureServiceBusQueue
- func (q *AzureServiceBusQueue) Init(ctx context.Context) (func(), error)
- func (q *AzureServiceBusQueue) InitClient(ctx context.Context) error
- func (q *AzureServiceBusQueue) Publish(ctx context.Context, incomingMessage IncomingMessage) error
- func (q *AzureServiceBusQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)
- type ConcurrentSubscription
- type GCPPubSubConfig
- type GCPPubSubQueue
- type InMemoryConfig
- type InMemoryQueue
- type IncomingMessage
- type Message
- type NATSConfig
- type NATSQueue
- type NATSSubscription
- type Queue
- type QueueConfig
- type QueueMessage
- type RabbitMQConfig
- type RabbitMQQueue
- type Rejecter
- type SubscribeOption
- type SubscribeOptions
- type Subscription
- type UnimplementedQueue
- type WrappedSubscription
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) 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 ¶
func (c *AWSSQSConfig) ToCredentials() (*credentials.StaticCredentialsProvider, error)
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 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 GCPPubSubQueue ¶
type GCPPubSubQueue struct {
// contains filtered or unexported fields
}
func NewGCPPubSubQueue ¶
func NewGCPPubSubQueue(config *GCPPubSubConfig, visibilityTimeout time.Duration) *GCPPubSubQueue
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) 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 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
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) 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.
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) 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 UnimplementedQueue ¶
type UnimplementedQueue struct{}
func (*UnimplementedQueue) Init ¶
func (q *UnimplementedQueue) Init(ctx context.Context) (func(), error)
func (*UnimplementedQueue) Publish ¶
func (q *UnimplementedQueue) Publish(ctx context.Context, msg IncomingMessage) error
func (*UnimplementedQueue) Subscribe ¶
func (q *UnimplementedQueue) Subscribe(ctx context.Context, opts ...SubscribeOption) (Subscription, error)
type WrappedSubscription ¶
type WrappedSubscription struct {
// contains filtered or unexported fields
}