Documentation
¶
Overview ¶
Package nats provides tools for comfort work during event processing with NATS broker.
Core mode — CommonPublisher and CommonConsumer — is an at-most-once broadcast: it has neither memory, nor acknowledgements. JetStream mode — JetStreamPublisher and JetStreamConsumer — is an at-least-once queue: the stream keeps messages on disk, a durable consumer holds the reading position on the server, and an unacknowledged message is redelivered.
Index ¶
- type CommonConsumer
- type CommonPublisher
- type Consumer
- type ConsumerAlreadyRunningError
- type ConsumerAlreadyStoppedError
- type ConsumerOption
- func WithCloseHandler(handler func(connection *natsbroker.Conn)) ConsumerOption
- func WithDisconnectErrorHandler(handler func(connection *natsbroker.Conn, err error)) ConsumerOption
- func WithErrorHandler(...) ConsumerOption
- func WithGoroutinesPoolSize(size int) ConsumerOption
- func WithMessageChannelBufferSize(size int) ConsumerOption
- func WithMessageHandler(handler func(message *natsbroker.Msg)) ConsumerOption
- func WithNatsOptions(opts ...natsbroker.Option) ConsumerOption
- type IDPublisher
- type JetStreamConsumer
- type JetStreamConsumerOption
- func WithJetStreamAckWait(ackWait time.Duration) JetStreamConsumerOption
- func WithJetStreamDurable(name string) JetStreamConsumerOption
- func WithJetStreamGoroutinesPoolSize(size int) JetStreamConsumerOption
- func WithJetStreamMaxAckPending(maxAckPending int) JetStreamConsumerOption
- func WithJetStreamMaxDeliver(maxDeliver int) JetStreamConsumerOption
- func WithJetStreamMessageChannelBufferSize(size int) JetStreamConsumerOption
- func WithJetStreamMessageHandler(handler func(message *natsbroker.Msg)) JetStreamConsumerOption
- func WithJetStreamNatsOptions(opts ...natsbroker.Option) JetStreamConsumerOption
- func WithJetStreamStream(name string) JetStreamConsumerOption
- type JetStreamPublisher
- type Publisher
- type StreamConfig
- type StreamNotConfiguredError
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type CommonConsumer ¶ added in v1.11.0
type CommonConsumer struct {
// contains filtered or unexported fields
}
CommonConsumer is a base consumer for processing NATS messages.
func NewConsumer ¶ added in v1.11.0
func NewConsumer( url string, subject string, opts ...ConsumerOption, ) (*CommonConsumer, error)
NewConsumer creates *CommonConsumer with provided options.
func (*CommonConsumer) Run ¶ added in v1.11.0
func (c *CommonConsumer) Run() error
Run starts goroutines for NATS messages processing.
func (*CommonConsumer) Stop ¶ added in v1.11.0
func (c *CommonConsumer) Stop() error
Stop stops launched goroutines, which processes NATS messages.
type CommonPublisher ¶
type CommonPublisher struct {
// contains filtered or unexported fields
}
CommonPublisher is a base NATS publisher.
func NewPublisher ¶ added in v1.3.2
func NewPublisher(url string, opts ...natsbroker.Option) (*CommonPublisher, error)
NewPublisher creates *CommonPublisher.
func (*CommonPublisher) Close ¶
func (p *CommonPublisher) Close() error
Close closes NATS connection.
type ConsumerAlreadyRunningError ¶ added in v1.11.0
ConsumerAlreadyRunningError is an error, which represents, that consumer was already started and can not be started again.
func (ConsumerAlreadyRunningError) Error ¶ added in v1.11.0
func (e ConsumerAlreadyRunningError) Error() string
func (ConsumerAlreadyRunningError) Unwrap ¶ added in v1.11.0
func (e ConsumerAlreadyRunningError) Unwrap() error
type ConsumerAlreadyStoppedError ¶ added in v1.11.0
ConsumerAlreadyStoppedError is an error, which represents, that consumer was not started yet or was already stopped.
func (ConsumerAlreadyStoppedError) Error ¶ added in v1.11.0
func (e ConsumerAlreadyStoppedError) Error() string
func (ConsumerAlreadyStoppedError) Unwrap ¶ added in v1.11.0
func (e ConsumerAlreadyStoppedError) Unwrap() error
type ConsumerOption ¶ added in v1.11.0
type ConsumerOption func(options *consumerOptions) error
ConsumerOption represents golang functional option pattern func for Consumer configuration.
func WithCloseHandler ¶
func WithCloseHandler(handler func(connection *natsbroker.Conn)) ConsumerOption
WithCloseHandler sets handler for connection with NATS closure.
func WithDisconnectErrorHandler ¶
func WithDisconnectErrorHandler( handler func(connection *natsbroker.Conn, err error), ) ConsumerOption
WithDisconnectErrorHandler sets handler for disconnection from server.
func WithErrorHandler ¶
func WithErrorHandler( handler func(connection *natsbroker.Conn, subscription *natsbroker.Subscription, err error), ) ConsumerOption
WithErrorHandler sets handler for processing error during message processing.
func WithGoroutinesPoolSize ¶
func WithGoroutinesPoolSize(size int) ConsumerOption
WithGoroutinesPoolSize sets number of goroutines for process messages from NATS via message channel.
func WithMessageChannelBufferSize ¶
func WithMessageChannelBufferSize(size int) ConsumerOption
WithMessageChannelBufferSize sets buffer for channel, where NATS will store messages for processing.
func WithMessageHandler ¶
func WithMessageHandler(handler func(message *natsbroker.Msg)) ConsumerOption
WithMessageHandler sets handler for received message.
func WithNatsOptions ¶
func WithNatsOptions(opts ...natsbroker.Option) ConsumerOption
WithNatsOptions sets NATS option for connection with broker configuration.
type IDPublisher ¶ added in v1.17.0
type IDPublisher interface {
Publisher
// PublishWithID sends a message with a deduplication id. A repeated publish
// with the same id inside the server deduplication window is dropped.
PublishWithID(subject, msgID string, content []byte) error
}
IDPublisher publishes messages with a deduplication id.
type JetStreamConsumer ¶ added in v1.17.0
type JetStreamConsumer struct {
// contains filtered or unexported fields
}
JetStreamConsumer processes messages of a durable JetStream consumer.
Reading position lives on the server, so a restarted worker continues from the place, where it stopped, and not from the current moment. Acknowledgement is manual: the handler acks after a successful delivery, and an unacked message is redelivered after AckWait.
func NewJetStreamConsumer ¶ added in v1.17.0
func NewJetStreamConsumer( url string, subject string, opts ...JetStreamConsumerOption, ) (*JetStreamConsumer, error)
NewJetStreamConsumer creates *JetStreamConsumer with provided options.
func (*JetStreamConsumer) Run ¶ added in v1.17.0
func (c *JetStreamConsumer) Run() error
Run starts goroutines for messages processing.
func (*JetStreamConsumer) Stop ¶ added in v1.17.0
func (c *JetStreamConsumer) Stop() error
Stop drains the subscription and waits for launched goroutines.
Drain instead of Unsubscribe — but neither of them is what keeps the durable consumer alive: both delete a consumer, created by the subscribe call. The consumer survives, because ensureConsumer creates it apart and the subscription only binds to it.
Drain is asynchronous: it returns at once, and the library keeps delivering buffered messages into the channel. Closing the channel right away would panic with a send on a closed channel, so the drain is awaited first. The timeout insures against a hung connection and loses nothing: an unacked task comes back after ack_wait.
type JetStreamConsumerOption ¶ added in v1.17.0
type JetStreamConsumerOption func(options *jetStreamConsumerOptions) error
JetStreamConsumerOption represents functional option for JetStreamConsumer.
func WithJetStreamAckWait ¶ added in v1.17.0
func WithJetStreamAckWait(ackWait time.Duration) JetStreamConsumerOption
WithJetStreamAckWait sets how long server waits for an acknowledgement before redelivering the message.
func WithJetStreamDurable ¶ added in v1.17.0
func WithJetStreamDurable(name string) JetStreamConsumerOption
WithJetStreamDurable sets durable consumer name, which holds reading position on the server.
func WithJetStreamGoroutinesPoolSize ¶ added in v1.17.0
func WithJetStreamGoroutinesPoolSize(size int) JetStreamConsumerOption
WithJetStreamGoroutinesPoolSize sets number of processing goroutines.
func WithJetStreamMaxAckPending ¶ added in v1.17.0
func WithJetStreamMaxAckPending(maxAckPending int) JetStreamConsumerOption
WithJetStreamMaxAckPending limits number of unacknowledged messages in flight.
func WithJetStreamMaxDeliver ¶ added in v1.17.0
func WithJetStreamMaxDeliver(maxDeliver int) JetStreamConsumerOption
WithJetStreamMaxDeliver limits number of delivery attempts.
func WithJetStreamMessageChannelBufferSize ¶ added in v1.17.0
func WithJetStreamMessageChannelBufferSize(size int) JetStreamConsumerOption
WithJetStreamMessageChannelBufferSize sets buffer of the messages channel.
func WithJetStreamMessageHandler ¶ added in v1.17.0
func WithJetStreamMessageHandler(handler func(message *natsbroker.Msg)) JetStreamConsumerOption
WithJetStreamMessageHandler sets handler for received message.
func WithJetStreamNatsOptions ¶ added in v1.17.0
func WithJetStreamNatsOptions(opts ...natsbroker.Option) JetStreamConsumerOption
WithJetStreamNatsOptions sets NATS connection options.
func WithJetStreamStream ¶ added in v1.17.0
func WithJetStreamStream(name string) JetStreamConsumerOption
WithJetStreamStream binds consumer to an existing stream.
type JetStreamPublisher ¶ added in v1.17.0
type JetStreamPublisher struct {
// contains filtered or unexported fields
}
JetStreamPublisher publishes messages to a JetStream stream.
Unlike CommonPublisher it waits for a server acknowledgement: a publish, which returned nil, is stored on disk.
func NewJetStreamPublisher ¶ added in v1.17.0
func NewJetStreamPublisher( url string, stream StreamConfig, opts ...natsbroker.Option, ) (*JetStreamPublisher, error)
NewJetStreamPublisher creates *JetStreamPublisher and ensures the stream.
func (*JetStreamPublisher) Close ¶ added in v1.17.0
func (p *JetStreamPublisher) Close() error
Close closes NATS connection.
func (*JetStreamPublisher) Publish ¶ added in v1.17.0
func (p *JetStreamPublisher) Publish(subject string, data []byte) error
Publish sends a message without deduplication.
It matches Publisher by signature, but it is NOT a drop-in replacement for the core one, and assuming otherwise is the trap this comment exists for: a JetStream publish waits for a stream acknowledgement, so a subject, which no stream captures, fails with `no responders`. Live WebSocket events therefore keep going through the core publisher.
func (*JetStreamPublisher) PublishWithID ¶ added in v1.17.0
func (p *JetStreamPublisher) PublishWithID(subject, msgID string, data []byte) error
PublishWithID sends a message with Nats-Msg-Id header.
A repeated publish with the same id inside the deduplication window is dropped by the server. It closes the case, when a publisher did not get an acknowledgement and sent the task for the second time.
type StreamConfig ¶ added in v1.17.0
type StreamConfig struct {
// Name is a stream name. Stream is created on publisher start, if missing.
Name string
// Subjects is a list of subject patterns, captured by the stream.
Subjects []string
// MaxAge is an upper bound of a message lifetime in the stream.
MaxAge time.Duration
// MaxBytes is an upper bound of a stream size on disk. Reaching it fails
// a publish (see DiscardNew above), and does not drop a stored task.
MaxBytes int64
// Duplicates is a deduplication window for Nats-Msg-Id header.
Duplicates time.Duration
}
StreamConfig describes a JetStream stream with file storage.
Retention is always WorkQueuePolicy: a message leaves the stream when it is acknowledged, not when it expires. MaxAge and MaxBytes are the upper bounds, which keep an undeliverable task from growing the stream forever.
Discard is always DiscardNew, and it is not a detail. With the default DiscardOld a stream, which hit MaxBytes, silently drops the OLDEST unacknowledged messages — that is, loses the very tasks the stream exists to keep. DiscardNew fails the publish instead, loudly and at the publisher.
Mind the WorkQueuePolicy constraint on the consumer side: filter subjects of consumers of such a stream must not overlap. One consumer per channel (`notifications.*.email`, `.web_push`, `.bell`) is fine; a second one on a subset of those subjects is refused with `filtered consumer not unique on workqueue stream`.
type StreamNotConfiguredError ¶ added in v1.17.0
StreamNotConfiguredError is an error, which represents, that JetStream entity was created without a stream configuration.
func (StreamNotConfiguredError) Error ¶ added in v1.17.0
func (e StreamNotConfiguredError) Error() string
func (StreamNotConfiguredError) Unwrap ¶ added in v1.17.0
func (e StreamNotConfiguredError) Unwrap() error