Documentation
¶
Index ¶
- Constants
- type AcknowledgerCmd
- type AknowledgerToken
- type Amqp10AsyncConsumer
- func (dc *Amqp10AsyncConsumer) DecrementActiveProcessors()
- func (dc *Amqp10AsyncConsumer) Start(logger *l.CapiLogger, listenerChannel chan *wfmodel.Message, ...) error
- func (dc *Amqp10AsyncConsumer) StopAcknowledger() error
- func (dc *Amqp10AsyncConsumer) StopListener() error
- func (dc *Amqp10AsyncConsumer) SupportsHeartbeat() bool
- type Amqp10Consumer
- type Amqp10Producer
- type CapimqAsyncConsumer
- func (dc *CapimqAsyncConsumer) DecrementActiveProcessors()
- func (dc *CapimqAsyncConsumer) Start(logger *l.CapiLogger, listenerChannel chan *wfmodel.Message, ...) error
- func (dc *CapimqAsyncConsumer) StopAcknowledger() error
- func (dc *CapimqAsyncConsumer) StopListener() error
- func (dc *CapimqAsyncConsumer) SupportsHeartbeat() bool
- type CapimqProducer
- type MqAsyncConsumer
- type MqClientType
- type MqProducer
- type RetryMethodType
Constants ¶
const Amqp10AcknowledgerAckTimeout time.Duration = 2000
const Amqp10AcknowledgerCloseTimeout time.Duration = 2000
const Amqp10AcknowledgerOpenTimeout time.Duration = 2000
const Amqp10AcknowledgerReconnectTimeout time.Duration = 2000
const Amqp10AcknowledgerTotalTimeout time.Duration = Amqp10AcknowledgerOpenTimeout + Amqp10AcknowledgerAckTimeout + Amqp10AcknowledgerCloseTimeout + Amqp10AcknowledgerReconnectTimeout + 1000
const Amqp10FullListenerChannelTimeout time.Duration = 50
const Amqp10ListenerAckTimeout time.Duration = 2000
const Amqp10ListenerCloseTimeout time.Duration = 2000
const Amqp10ListenerOpenTimeout time.Duration = 2000
const Amqp10ListenerReceiveTimeout time.Duration = 2000
const Amqp10ListenerReconnectTimeout time.Duration = 2000
const Amqp10ListenerTotalTimeout time.Duration = Amqp10ListenerOpenTimeout + Amqp10ListenerReceiveTimeout + Amqp10ListenerAckTimeout + Amqp10ListenerCloseTimeout + Amqp10ListenerReconnectTimeout + 1000
const Amqp10ProducerCloseTimeout time.Duration = 2000
const Amqp10ProducerOpenTimeout time.Duration = 2000
const Amqp10ProducerSendTimeout time.Duration = 2000
const CapimqAcknowledgerAckTimeout time.Duration = 2000
const CapimqAcknowledgerTotalTimeout time.Duration = CapimqAcknowledgerAckTimeout + 1000
const CapimqAllProcessorsBusyTimeout time.Duration = 100
const CapimqListenerConnErrorTimeout time.Duration = 2000
const CapimqListenerNothingToClaimTimeout time.Duration = 500
const CapimqListenerReceiveTimeout time.Duration = 500
const CapimqListenerTotalTimeout time.Duration = CapimqListenerReceiveTimeout + CapimqListenerNothingToClaimTimeout + 1000
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) 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) 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) 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 RetryMethodType ¶
type RetryMethodType string
const ( RetryMethodRelease RetryMethodType = "release" RetryMethodReject RetryMethodType = "reject" RetryMethodUnknown RetryMethodType = "unknown" )
func StringToRetryMethod ¶
func StringToRetryMethod(s string) (RetryMethodType, error)