Versions in this module Expand all Collapse all v6 v6.2.13 Sep 22, 2026 v6.2.12 Sep 21, 2026 Changes in this version + const DataStream + const ReprocessStream + const TransactionsStream + const TransactionsSubject + type Config struct + Nats NatsConfig + Verbosity string + func DefaultConfig() Config + type Conn interface + Close func() + JetStream func(opts ...nats.JSOpt) (nats.JetStreamContext, error) + type ConnectionPool interface + Acquire func(ctx context.Context) (Conn, nats.JetStreamContext, error) + Shutdown func() + type Event interface + GetStream func(streamName string) Stream + Pool func() ConnectionPool + func NewManager() Event + func NewTestManager(t *testing.T) Event + type JetStreamContext interface + type MockConn struct + func NewMockConn(ctrl *gomock.Controller) *MockConn + func (m *MockConn) Close() + func (m *MockConn) EXPECT() *MockConnMockRecorder + func (m *MockConn) JetStream(opts ...nats.JSOpt) (nats.JetStreamContext, error) + type MockConnMockRecorder struct + func (mr *MockConnMockRecorder) Close() *gomock.Call + func (mr *MockConnMockRecorder) JetStream(opts ...any) *gomock.Call + type MockConnectionPool struct + func NewMockConnectionPool(ctrl *gomock.Controller) *MockConnectionPool + func (m *MockConnectionPool) Acquire(ctx context.Context) (Conn, nats.JetStreamContext, error) + func (m *MockConnectionPool) EXPECT() *MockConnectionPoolMockRecorder + func (m *MockConnectionPool) Shutdown() + type MockConnectionPoolMockRecorder struct + func (mr *MockConnectionPoolMockRecorder) Acquire(ctx any) *gomock.Call + func (mr *MockConnectionPoolMockRecorder) Shutdown() *gomock.Call + type MockEvent struct + func NewMockEvent(ctrl *gomock.Controller) *MockEvent + func (m *MockEvent) EXPECT() *MockEventMockRecorder + func (m *MockEvent) GetStream(streamName string) Stream + func (m *MockEvent) Pool() ConnectionPool + type MockEventMockRecorder struct + func (mr *MockEventMockRecorder) GetStream(streamName any) *gomock.Call + func (mr *MockEventMockRecorder) Pool() *gomock.Call + type MockJetStreamContext struct + func NewMockJetStreamContext(ctrl *gomock.Controller) *MockJetStreamContext + func (m *MockJetStreamContext) AccountInfo(opts ...nats.JSOpt) (*nats.AccountInfo, error) + func (m *MockJetStreamContext) AddConsumer(stream string, cfg *nats.ConsumerConfig, opts ...nats.JSOpt) (*nats.ConsumerInfo, error) + func (m *MockJetStreamContext) AddStream(cfg *nats.StreamConfig, opts ...nats.JSOpt) (*nats.StreamInfo, error) + func (m *MockJetStreamContext) ChanQueueSubscribe(subj, queue string, ch chan *nats.Msg, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) ChanSubscribe(subj string, ch chan *nats.Msg, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) CleanupPublisher() + func (m *MockJetStreamContext) ConsumerInfo(stream, name string, opts ...nats.JSOpt) (*nats.ConsumerInfo, error) + func (m *MockJetStreamContext) ConsumerNames(stream string, opts ...nats.JSOpt) <-chan string + func (m *MockJetStreamContext) Consumers(stream string, opts ...nats.JSOpt) <-chan *nats.ConsumerInfo + func (m *MockJetStreamContext) ConsumersInfo(stream string, opts ...nats.JSOpt) <-chan *nats.ConsumerInfo + func (m *MockJetStreamContext) CreateKeyValue(cfg *nats.KeyValueConfig) (nats.KeyValue, error) + func (m *MockJetStreamContext) CreateObjectStore(cfg *nats.ObjectStoreConfig) (nats.ObjectStore, error) + func (m *MockJetStreamContext) DeleteConsumer(stream, consumer string, opts ...nats.JSOpt) error + func (m *MockJetStreamContext) DeleteKeyValue(bucket string) error + func (m *MockJetStreamContext) DeleteMsg(name string, seq uint64, opts ...nats.JSOpt) error + func (m *MockJetStreamContext) DeleteObjectStore(bucket string) error + func (m *MockJetStreamContext) DeleteStream(name string, opts ...nats.JSOpt) error + func (m *MockJetStreamContext) EXPECT() *MockJetStreamContextMockRecorder + func (m *MockJetStreamContext) GetLastMsg(name, subject string, opts ...nats.JSOpt) (*nats.RawStreamMsg, error) + func (m *MockJetStreamContext) GetMsg(name string, seq uint64, opts ...nats.JSOpt) (*nats.RawStreamMsg, error) + func (m *MockJetStreamContext) KeyValue(bucket string) (nats.KeyValue, error) + func (m *MockJetStreamContext) KeyValueStoreNames() <-chan string + func (m *MockJetStreamContext) KeyValueStores() <-chan nats.KeyValueStatus + func (m *MockJetStreamContext) ObjectStore(bucket string) (nats.ObjectStore, error) + func (m *MockJetStreamContext) ObjectStoreNames(opts ...nats.ObjectOpt) <-chan string + func (m *MockJetStreamContext) ObjectStores(opts ...nats.ObjectOpt) <-chan nats.ObjectStoreStatus + func (m *MockJetStreamContext) Publish(subj string, data []byte, opts ...nats.PubOpt) (*nats.PubAck, error) + func (m *MockJetStreamContext) PublishAsync(subj string, data []byte, opts ...nats.PubOpt) (nats.PubAckFuture, error) + func (m *MockJetStreamContext) PublishAsyncComplete() <-chan struct{} + func (m *MockJetStreamContext) PublishAsyncPending() int + func (m *MockJetStreamContext) PullSubscribe(subj, durable string, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) PurgeStream(name string, opts ...nats.JSOpt) error + func (m *MockJetStreamContext) QueueSubscribe(subj, queue string, cb nats.MsgHandler, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) QueueSubscribeSync(subj, queue string, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) SecureDeleteMsg(name string, seq uint64, opts ...nats.JSOpt) error + func (m *MockJetStreamContext) StreamInfo(stream string, opts ...nats.JSOpt) (*nats.StreamInfo, error) + func (m *MockJetStreamContext) StreamNameBySubject(arg0 string, arg1 ...nats.JSOpt) (string, error) + func (m *MockJetStreamContext) StreamNames(opts ...nats.JSOpt) <-chan string + func (m *MockJetStreamContext) Streams(opts ...nats.JSOpt) <-chan *nats.StreamInfo + func (m *MockJetStreamContext) StreamsInfo(opts ...nats.JSOpt) <-chan *nats.StreamInfo + func (m *MockJetStreamContext) Subscribe(subj string, cb nats.MsgHandler, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) SubscribeSync(subj string, opts ...nats.SubOpt) (*nats.Subscription, error) + func (m *MockJetStreamContext) UpdateConsumer(stream string, cfg *nats.ConsumerConfig, opts ...nats.JSOpt) (*nats.ConsumerInfo, error) + func (m *MockJetStreamContext) UpdateStream(cfg *nats.StreamConfig, opts ...nats.JSOpt) (*nats.StreamInfo, error) + func (m_2 *MockJetStreamContext) PublishMsg(m *nats.Msg, opts ...nats.PubOpt) (*nats.PubAck, error) + func (m_2 *MockJetStreamContext) PublishMsgAsync(m *nats.Msg, opts ...nats.PubOpt) (nats.PubAckFuture, error) + type MockJetStreamContextMockRecorder struct + func (mr *MockJetStreamContextMockRecorder) AccountInfo(opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) AddConsumer(stream, cfg any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) AddStream(cfg any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ChanQueueSubscribe(subj, queue, ch any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ChanSubscribe(subj, ch any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) CleanupPublisher() *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ConsumerInfo(stream, name any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ConsumerNames(stream any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) Consumers(stream any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ConsumersInfo(stream any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) CreateKeyValue(cfg any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) CreateObjectStore(cfg any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) DeleteConsumer(stream, consumer any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) DeleteKeyValue(bucket any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) DeleteMsg(name, seq any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) DeleteObjectStore(bucket any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) DeleteStream(name any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) GetLastMsg(name, subject any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) GetMsg(name, seq any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) KeyValue(bucket any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) KeyValueStoreNames() *gomock.Call + func (mr *MockJetStreamContextMockRecorder) KeyValueStores() *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ObjectStore(bucket any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ObjectStoreNames(opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) ObjectStores(opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) Publish(subj, data any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PublishAsync(subj, data any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PublishAsyncComplete() *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PublishAsyncPending() *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PublishMsg(m any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PublishMsgAsync(m any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PullSubscribe(subj, durable any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) PurgeStream(name any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) QueueSubscribe(subj, queue, cb any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) QueueSubscribeSync(subj, queue any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) SecureDeleteMsg(name, seq any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) StreamInfo(stream any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) StreamNameBySubject(arg0 any, arg1 ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) StreamNames(opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) Streams(opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) StreamsInfo(opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) Subscribe(subj, cb any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) SubscribeSync(subj any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) UpdateConsumer(stream, cfg any, opts ...any) *gomock.Call + func (mr *MockJetStreamContextMockRecorder) UpdateStream(cfg any, opts ...any) *gomock.Call + type NATSConnectFunc func(url string, options ...nats.Option) (Conn, error) + type NATSConnectionPool struct + func NewNATSConnectionPool(config Config) *NATSConnectionPool + func (pool *NATSConnectionPool) Acquire(ctx context.Context) (Conn, nats.JetStreamContext, error) + func (pool *NATSConnectionPool) Shutdown() + type NatsConfig struct + Hostname string + Port int + StorageDir string + Timeout int + type Stream interface + Config func() *nats.StreamConfig + Subscribe func(conn Conn, consumerName string, subjectFilter string, handler nats.MsgHandler) error + func NewDisposableStream(name string, subjects []string, maxMessages int64) Stream + type TransactionWithPayload struct + Payload []byte + Transaction dag.Transaction + func (t *TransactionWithPayload) UnmarshalJSON(bytes []byte) error + func (t TransactionWithPayload) MarshalJSON() ([]byte, error) Other modules containing this package github.com/nuts-foundation/nuts-node github.com/nuts-foundation/nuts-node/v5