nats

package
v1.18.2 Latest Latest
Warning

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

Go to latest
Published: Oct 6, 2026 License: GPL-3.0 Imports: 5 Imported by: 6

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

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.

func (*CommonPublisher) Publish

func (p *CommonPublisher) Publish(topic string, data []byte) error

Publish sends message to provided topic (subject).

type Consumer added in v1.11.0

type Consumer interface {
	Run() error
	Stop() error
}

Consumer asynchronously processes NATS messages in goroutines.

type ConsumerAlreadyRunningError added in v1.11.0

type ConsumerAlreadyRunningError struct {
	Message string
	BaseErr error
}

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 (ConsumerAlreadyRunningError) Unwrap added in v1.11.0

type ConsumerAlreadyStoppedError added in v1.11.0

type ConsumerAlreadyStoppedError struct {
	Message string
	BaseErr error
}

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 (ConsumerAlreadyStoppedError) Unwrap added in v1.11.0

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 Publisher

type Publisher interface {
	Publish(subject string, content []byte) error
	Close() error
}

Publisher publishes messages to NATS broker.

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

type StreamNotConfiguredError struct {
	Message string
	BaseErr error
}

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

Directories

Path Synopsis
Package mocks is a generated GoMock package.
Package mocks is a generated GoMock package.

Jump to

Keyboard shortcuts

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