mq

package
v1.2.0 Latest Latest
Warning

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

Go to latest
Published: Nov 23, 2025 License: MIT Imports: 17 Imported by: 0

Documentation

Index

Constants

View Source
const Amqp10AcknowledgerAckTimeout time.Duration = 2000
View Source
const Amqp10AcknowledgerCloseTimeout time.Duration = 2000
View Source
const Amqp10AcknowledgerOpenTimeout time.Duration = 2000
View Source
const Amqp10AcknowledgerReconnectTimeout time.Duration = 2000
View Source
const Amqp10FullListenerChannelTimeout time.Duration = 50
View Source
const Amqp10ListenerAckTimeout time.Duration = 2000
View Source
const Amqp10ListenerCloseTimeout time.Duration = 2000
View Source
const Amqp10ListenerOpenTimeout time.Duration = 2000
View Source
const Amqp10ListenerReceiveTimeout time.Duration = 2000
View Source
const Amqp10ListenerReconnectTimeout time.Duration = 2000
View Source
const Amqp10ProducerCloseTimeout time.Duration = 2000
View Source
const Amqp10ProducerOpenTimeout time.Duration = 2000
View Source
const Amqp10ProducerSendTimeout time.Duration = 2000
View Source
const CapimqAcknowledgerAckTimeout time.Duration = 2000
View Source
const CapimqAcknowledgerTotalTimeout time.Duration = CapimqAcknowledgerAckTimeout + 1000
View Source
const CapimqAllProcessorsBusyTimeout time.Duration = 100
View Source
const CapimqListenerConnErrorTimeout time.Duration = 2000
View Source
const CapimqListenerNothingToClaimTimeout time.Duration = 500
View Source
const CapimqListenerReceiveTimeout time.Duration = 500
View Source
const CapimqProducerSendTimeout time.Duration = 2000

Variables

This section is empty.

Functions

This section is empty.

Types

type AcknowledgerCmd

type AcknowledgerCmd int
const (
	AcknowledgerCmdAck AcknowledgerCmd = iota
	AcknowledgerCmdRetry
	AcknowledgerCmdHeartbeat
)

type AknowledgerToken

type AknowledgerToken struct {
	MsgId string
	Cmd   AcknowledgerCmd
}

type Amqp10AsyncConsumer

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

The idea behind this async consumer is to Receive AMQP messages with one go-amqp receiver (we call it listener), and Ack/Retry AMQP messages with another go-amqp receiver (we call it acknoledger). This way, we do not need to introduce any multi-thread protection for the Open/Close code. It uses amqpMessagesInHandlingMutex for protecting the helper map, which is cheap. To work properly, and gracefully close async consumer, the caller should do this: asyncConsumer.Start(listenerChannel, acknowledgerChannel) ... caller reads messages from listenerChannel, passes them to processors that write to acknowledgerChannel asyncConsumer.StopListener() close(listenerChannel) waitForAllProcessorsToCompleteSoTheyDoNotWriteToAcknowledgerChannel asyncConsumer.StopAcknowledger() close(acknowledgerChannel) The size of listenerChannel is 1 because we want zero prefetch The size of acknowledgerChannel is between 1 and any reasonale value

func NewAmqp10Consumer

func NewAmqp10Consumer(brokerUrl string, address string, ackMethod RetryMethodType, maxProcessors int) *Amqp10AsyncConsumer

func (*Amqp10AsyncConsumer) DecrementActiveProcessors

func (dc *Amqp10AsyncConsumer) DecrementActiveProcessors()

func (*Amqp10AsyncConsumer) Start

func (dc *Amqp10AsyncConsumer) Start(logger *l.CapiLogger, listenerChannel chan *wfmodel.Message, acknowledgerChannel chan AknowledgerToken) error

func (*Amqp10AsyncConsumer) StopAcknowledger

func (dc *Amqp10AsyncConsumer) StopAcknowledger() error

func (*Amqp10AsyncConsumer) StopListener

func (dc *Amqp10AsyncConsumer) StopListener() error

func (*Amqp10AsyncConsumer) SupportsHeartbeat

func (dc *Amqp10AsyncConsumer) SupportsHeartbeat() bool

type Amqp10Consumer

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

func (*Amqp10Consumer) Open

func (c *Amqp10Consumer) Open(ctx context.Context, url string, address string) error

func (*Amqp10Consumer) Receiver

func (c *Amqp10Consumer) Receiver() *amqp10.Receiver

type Amqp10Producer

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

func NewAmqp10Producer

func NewAmqp10Producer(url string, address string) *Amqp10Producer

func (*Amqp10Producer) Close

func (p *Amqp10Producer) Close() error

func (*Amqp10Producer) Open

func (p *Amqp10Producer) Open() error

func (*Amqp10Producer) Send

func (p *Amqp10Producer) Send(msg *wfmodel.Message) error

func (*Amqp10Producer) SendBulk

func (p *Amqp10Producer) SendBulk([]*wfmodel.Message) error

func (*Amqp10Producer) SupportsSendBulk

func (p *Amqp10Producer) SupportsSendBulk() bool

type CapimqAsyncConsumer

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

func NewCapimqConsumer

func NewCapimqConsumer(brokerUrl string, clientName string, maxProcessors int) *CapimqAsyncConsumer

func (*CapimqAsyncConsumer) DecrementActiveProcessors

func (dc *CapimqAsyncConsumer) DecrementActiveProcessors()

func (*CapimqAsyncConsumer) Start

func (dc *CapimqAsyncConsumer) Start(logger *l.CapiLogger, listenerChannel chan *wfmodel.Message, acknowledgerChannel chan AknowledgerToken) error

func (*CapimqAsyncConsumer) StopAcknowledger

func (dc *CapimqAsyncConsumer) StopAcknowledger() error

func (*CapimqAsyncConsumer) StopListener

func (dc *CapimqAsyncConsumer) StopListener() error

func (*CapimqAsyncConsumer) SupportsHeartbeat

func (dc *CapimqAsyncConsumer) SupportsHeartbeat() bool

type CapimqProducer

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

func NewCapimqProducer

func NewCapimqProducer(url string) *CapimqProducer

func (*CapimqProducer) Close

func (p *CapimqProducer) Close() error

func (*CapimqProducer) Open

func (p *CapimqProducer) Open() error

func (*CapimqProducer) Send

func (p *CapimqProducer) Send(wfmodelMsg *wfmodel.Message) error

func (*CapimqProducer) SendBulk

func (p *CapimqProducer) SendBulk(wfmodelMsgs []*wfmodel.Message) error

func (*CapimqProducer) SupportsSendBulk

func (p *CapimqProducer) SupportsSendBulk() bool

type MqAsyncConsumer

type MqAsyncConsumer interface {
	Start(logger *l.CapiLogger, listenerChannel chan *wfmodel.Message, acknowledgerChannel chan AknowledgerToken) error
	StopListener() error
	StopAcknowledger() error
	SupportsHeartbeat() bool
	DecrementActiveProcessors()
}

type MqClientType

type MqClientType string
const (
	MqClientAmqp10 MqClientType = "amqp10"
	MqClientCapimq MqClientType = "capimq"
)

type MqProducer

type MqProducer interface {
	Open() error
	Close() error
	Send(msg *wfmodel.Message) error
	SendBulk(msgs []*wfmodel.Message) error
	SupportsSendBulk() bool
}

type RetryMethodType

type RetryMethodType string
const (
	RetryMethodRelease RetryMethodType = "release"
	RetryMethodReject  RetryMethodType = "reject"
	RetryMethodUnknown RetryMethodType = "unknown"
)

func StringToRetryMethod

func StringToRetryMethod(s string) (RetryMethodType, error)

Jump to

Keyboard shortcuts

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