rmq

package
v0.2.0 Latest Latest
Warning

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

Go to latest
Published: Aug 23, 2026 License: MIT Imports: 11 Imported by: 0

Documentation

Index

Constants

View Source
const ConsumerServiceToken = gioc.Token("RMQConsumerService")
View Source
const ModuleConfigToken = gioc.Token("RmqModuleConfig")
View Source
const ProducerServiceToken = gioc.Token("RMQProducerService")
View Source
const ServiceToken = "RmqService"

Variables

View Source
var ConsumerServiceInjections = react.InjectFromBase(ServiceToken)
View Source
var ProducerServiceInjections = react.InjectFromBase(ServiceToken)
View Source
var ServiceInjections = react.InjectFromBase(
	ModuleConfigToken,
)

Functions

func ProvideModuleConfig

func ProvideModuleConfig(config *ModuleConfig) gioc.IProvider

Types

type Args

type Args = amqp091.Table

type Binding

type Binding struct {
	Key      string
	Exchange *Exchange
	NoWait   bool
	Args     Args
}

func Bind

func Bind(exchange *Exchange, keys ...string) (bindings []*Binding)

func MultiBind

func MultiBind(bindings ...[]*Binding) (result []*Binding)

func (*Binding) SetArgs

func (binding *Binding) SetArgs(args Args) *Binding

func (*Binding) SetExchange

func (binding *Binding) SetExchange(exchange *Exchange) *Binding

func (*Binding) SetKey

func (binding *Binding) SetKey(key string) *Binding

func (*Binding) SetNoWait

func (binding *Binding) SetNoWait(noWait bool) *Binding

type ChannelProvider

type ChannelProvider func() (IChannel, error)

type Consumer

type Consumer struct {
	Queue     *Queue
	Handler   MessageHandler
	Tag       string
	AutoAck   bool
	Exclusive bool
	NoLocal   bool
	NoWait    bool
	Args      Args
	// contains filtered or unexported fields
}

func (*Consumer) SetArgs

func (consumer *Consumer) SetArgs(args Args) *Consumer

func (*Consumer) SetAutoAck

func (consumer *Consumer) SetAutoAck(autoAck bool) *Consumer

func (*Consumer) SetExclusive

func (consumer *Consumer) SetExclusive(exclusive bool) *Consumer

func (*Consumer) SetHandler

func (consumer *Consumer) SetHandler(handler MessageHandler) *Consumer

func (*Consumer) SetNoLocal

func (consumer *Consumer) SetNoLocal(noLocal bool) *Consumer

func (*Consumer) SetNoWait

func (consumer *Consumer) SetNoWait(noWait bool) *Consumer

func (*Consumer) SetQueue

func (consumer *Consumer) SetQueue(queue *Queue) *Consumer

func (*Consumer) SetTag

func (consumer *Consumer) SetTag(tag string) *Consumer

type ConsumerService

type ConsumerService struct {
	react.BaseService
	// contains filtered or unexported fields
}

func NewConsumerService

func NewConsumerService(injections gioc.Injections) (service *ConsumerService, err error)

func (*ConsumerService) Consume

func (service *ConsumerService) Consume(consumer *Consumer) (err error)

type Error

type Error = *amqp091.Error

type Exchange

type Exchange struct {
	Name       string
	Kind       string
	Durable    bool
	AutoDelete bool
	Internal   bool
	NoWait     bool
	Args       Args
}

func NewExchange

func NewExchange(name string) *Exchange

func (*Exchange) SetArgs

func (exchange *Exchange) SetArgs(args Args) *Exchange

func (*Exchange) SetAutoDelete

func (exchange *Exchange) SetAutoDelete(autoDelete bool) *Exchange

func (*Exchange) SetDurable

func (exchange *Exchange) SetDurable(durable bool) *Exchange

func (*Exchange) SetInternal

func (exchange *Exchange) SetInternal(internal bool) *Exchange

func (*Exchange) SetKind

func (exchange *Exchange) SetKind(kind string) *Exchange

func (*Exchange) SetName

func (exchange *Exchange) SetName(name string) *Exchange

func (*Exchange) SetNoWait

func (exchange *Exchange) SetNoWait(noWait bool) *Exchange

type IChannel added in v0.2.0

type IChannel interface {
	NotifyClose(chan Error) chan Error
	IsClosed() bool
	Close() error
	PublishWithContext(context.Context, string, string, bool, bool, OutcomeMessage) error
	QueueDeclare(string, bool, bool, bool, bool, Args) (amqp091.Queue, error)
	ExchangeDeclare(string, string, bool, bool, bool, bool, Args) error
	QueueBind(string, string, string, bool, Args) error
	ConsumeWithContext(context.Context, string, string, bool, bool, bool, bool, Args) (<-chan amqp091.Delivery, error)
	QueueDelete(string, bool, bool, bool) (int, error)
	QueueInspect(string) (amqp091.Queue, error)
	Get(string, bool) (amqp091.Delivery, bool, error)
}

type IConnection added in v0.2.0

type IConnection interface {
	NotifyClose(chan Error) chan Error
	IsClosed() bool
	Channel() (*amqp091.Channel, error)
	Close() error
}

type IncomeMessage

type IncomeMessage amqp091.Delivery

func (IncomeMessage) Ack

func (message IncomeMessage) Ack(multiple bool) error

func (IncomeMessage) Nack

func (message IncomeMessage) Nack(multiple, requeue bool) error

func (IncomeMessage) Reject

func (message IncomeMessage) Reject(requeue bool) error

func (IncomeMessage) Retry

func (message IncomeMessage) Retry(multiple, requeue bool) (Retry, error)

func (IncomeMessage) RetryState

func (message IncomeMessage) RetryState() Retry

type MessageHandler

type MessageHandler func(message IncomeMessage) error

type ModuleConfig

type ModuleConfig struct {
	Host        string `env:"RMQ_HOST"`
	Port        int    `env:"RMQ_PORT"`
	User        string `env:"RMQ_USER"`
	Password    string `env:"RMQ_PASSWORD"`
	VirtualHost string `env:"RMQ_VIRTUAL_HOST"`

	RetryCount int           `env:"RMQ_CONNECTION_RETRY_COUNT"`
	RetryDelay time.Duration `env:"RMQ_CONNECTION_RETRY_DELAY"`
}

type OutcomeMessage

type OutcomeMessage = amqp091.Publishing

type ProducerService

type ProducerService struct {
	react.BaseService
	// contains filtered or unexported fields
}

func NewProducerService

func NewProducerService(injections gioc.Injections) (service *ProducerService, err error)

func (*ProducerService) Bind

func (producer *ProducerService) Bind(exchange *Exchange)

func (*ProducerService) Produce

func (producer *ProducerService) Produce(publication *Publication) (err error)

func (*ProducerService) WithTimeout

func (producer *ProducerService) WithTimeout(timeout time.Duration)

type Publication

type Publication struct {
	Destination string
	Mandatary   bool
	Immediate   bool
	Message     OutcomeMessage
}

func (*Publication) SetDestination

func (publication *Publication) SetDestination(destination string) *Publication

func (*Publication) SetImmediate

func (publication *Publication) SetImmediate(immediate bool) *Publication

func (*Publication) SetMandatary

func (publication *Publication) SetMandatary(mandatary bool) *Publication

func (*Publication) SetMandatory

func (publication *Publication) SetMandatory(mandatory bool) *Publication

func (*Publication) SetMessage

func (publication *Publication) SetMessage(message OutcomeMessage) *Publication

type Queue

type Queue struct {
	Name       string
	Durable    bool
	AutoDelete bool
	Exclusive  bool
	NoWait     bool
	Args       Args

	Bindings []*Binding
	// contains filtered or unexported fields
}

func NewQueue

func NewQueue(name string) *Queue

func (*Queue) AddBindings

func (queue *Queue) AddBindings(bindings ...*Binding) *Queue

func (*Queue) DeriveDLQ

func (queue *Queue) DeriveDLQ(suffix ...string) *Queue

func (*Queue) SetArgs

func (queue *Queue) SetArgs(args Args) *Queue

func (*Queue) SetAutoDelete

func (queue *Queue) SetAutoDelete(autoDelete bool) *Queue

func (*Queue) SetBindings

func (queue *Queue) SetBindings(bindings ...*Binding) *Queue

func (*Queue) SetDLQ

func (queue *Queue) SetDLQ(exchange *Exchange, key string) *Queue

func (*Queue) SetDurable

func (queue *Queue) SetDurable(durable bool) *Queue

func (*Queue) SetExclusive

func (queue *Queue) SetExclusive(exclusive bool) *Queue

func (*Queue) SetMessageTTL

func (queue *Queue) SetMessageTTL(ttl time.Duration) *Queue

func (*Queue) SetName

func (queue *Queue) SetName(name string) *Queue

func (*Queue) SetNoWait

func (queue *Queue) SetNoWait(noWait bool) *Queue

type Retry

type Retry struct {
	Count  int
	Reason string
}

type Service

type Service struct {
	react.BaseConfigurableService[*ModuleConfig]
	// contains filtered or unexported fields
}

func NewRmqService

func NewRmqService(injections gioc.Injections) (service *Service, err error)

func (*Service) Channel

func (service *Service) Channel() (channel IChannel, err error)

func (*Service) ChannelContext added in v0.2.0

func (service *Service) ChannelContext(ctx context.Context) (channel IChannel, err error)

ChannelContext opens a channel after waiting for connection readiness while honoring the caller's deadline. Long-running infrastructure adapters should prefer it to Channel so an individual operation cannot wait for reconnect beyond its own context.

func (*Service) CreateBinding

func (service *Service) CreateBinding(queue *Queue, binding *Binding) (err error)

func (*Service) CreateExchange

func (service *Service) CreateExchange(exchange *Exchange) (err error)

func (*Service) CreateQueues

func (service *Service) CreateQueues(queues ...*Queue) (result []*Queue, err error)

func (*Service) Restarting

func (service *Service) Restarting() bool

func (*Service) WaitReady

func (service *Service) WaitReady() error

func (*Service) WaitReadyContext added in v0.2.0

func (service *Service) WaitReadyContext(ctx context.Context) error

WaitReadyContext waits for a usable connection until either the caller or the RMQ service lifecycle is cancelled.

Jump to

Keyboard shortcuts

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