Documentation
¶
Index ¶
- Variables
- func NewMessageHandler(eventHandler EventHandler, opts ...MessageHandlerOption) consumer.MessageHandler
- func WithQueue(queueConfig *mqs.QueueConfig) func(opts *PublishMQOption)
- type EventHandler
- type HandleResult
- type MessageHandlerOption
- type PublishMQ
- type PublishMQOption
- type PublishedEvent
- type RedeliveryCounter
Constants ¶
This section is empty.
Variables ¶
Functions ¶
func NewMessageHandler ¶
func NewMessageHandler(eventHandler EventHandler, opts ...MessageHandlerOption) consumer.MessageHandler
func WithQueue ¶
func WithQueue(queueConfig *mqs.QueueConfig) func(opts *PublishMQOption)
Types ¶
type EventHandler ¶
type EventHandler interface {
Handle(ctx context.Context, event *models.Event) (*HandleResult, error)
}
func NewEventHandler ¶
func NewEventHandler( logger *logging.Logger, deliveryMQ *deliverymq.DeliveryMQ, tenantStore tenantstore.TenantStore, eventTracer eventtracer.EventTracer, topics []string, topicsAllowWildcards bool, idempotence idempotence.Idempotence, ) EventHandler
type HandleResult ¶ added in v0.9.0
type MessageHandlerOption ¶ added in v1.6.0
type MessageHandlerOption func(*messageHandler)
func WithMaxRedeliveries ¶ added in v1.6.0
func WithMaxRedeliveries(maxRedeliveries int, counter RedeliveryCounter) MessageHandlerOption
WithMaxRedeliveries stops redelivering a message once it has failed maxRedeliveries+1 times: rejected where the broker supports it, acked otherwise. Without it, failed messages are redelivered without limit.
type PublishMQ ¶
type PublishMQ struct {
// contains filtered or unexported fields
}
func New ¶
func New(opts ...func(opts *PublishMQOption)) *PublishMQ
func (*PublishMQ) Subscribe ¶
func (q *PublishMQ) Subscribe(ctx context.Context, opts ...mqs.SubscribeOption) (mqs.Subscription, error)
type PublishMQOption ¶
type PublishMQOption struct {
QueueConfig *mqs.QueueConfig
}
type PublishedEvent ¶
type PublishedEvent struct {
ID string `json:"id"`
TenantID string `json:"tenant_id" binding:"required"`
DestinationID string `json:"destination_id"`
Topic string `json:"topic"`
EligibleForRetry *bool `json:"eligible_for_retry"`
Time time.Time `json:"time"`
Metadata map[string]string `json:"metadata"`
Data json.RawMessage `json:"data"`
}
type RedeliveryCounter ¶ added in v1.6.0
RedeliveryCounter counts failed receives per message.
func NewRedisRedeliveryCounter ¶ added in v1.6.0
func NewRedisRedeliveryCounter(client redis.Cmdable, deploymentID string) RedeliveryCounter
Click to show internal directories.
Click to hide internal directories.