Documentation
¶
Index ¶
- Constants
- Variables
- func ProvideModuleConfig(config *ModuleConfig) gioc.IProvider
- type Args
- type Binding
- type ChannelProvider
- type Consumer
- func (consumer *Consumer) SetArgs(args Args) *Consumer
- func (consumer *Consumer) SetAutoAck(autoAck bool) *Consumer
- func (consumer *Consumer) SetExclusive(exclusive bool) *Consumer
- func (consumer *Consumer) SetHandler(handler MessageHandler) *Consumer
- func (consumer *Consumer) SetNoLocal(noLocal bool) *Consumer
- func (consumer *Consumer) SetNoWait(noWait bool) *Consumer
- func (consumer *Consumer) SetQueue(queue *Queue) *Consumer
- func (consumer *Consumer) SetTag(tag string) *Consumer
- type ConsumerService
- type Error
- type Exchange
- func (exchange *Exchange) SetArgs(args Args) *Exchange
- func (exchange *Exchange) SetAutoDelete(autoDelete bool) *Exchange
- func (exchange *Exchange) SetDurable(durable bool) *Exchange
- func (exchange *Exchange) SetInternal(internal bool) *Exchange
- func (exchange *Exchange) SetKind(kind string) *Exchange
- func (exchange *Exchange) SetName(name string) *Exchange
- func (exchange *Exchange) SetNoWait(noWait bool) *Exchange
- type IChannel
- type IConnection
- type IncomeMessage
- type MessageHandler
- type ModuleConfig
- type OutcomeMessage
- type ProducerService
- type Publication
- func (publication *Publication) SetDestination(destination string) *Publication
- func (publication *Publication) SetImmediate(immediate bool) *Publication
- func (publication *Publication) SetMandatary(mandatary bool) *Publication
- func (publication *Publication) SetMandatory(mandatory bool) *Publication
- func (publication *Publication) SetMessage(message OutcomeMessage) *Publication
- type Queue
- func (queue *Queue) AddBindings(bindings ...*Binding) *Queue
- func (queue *Queue) DeriveDLQ(suffix ...string) *Queue
- func (queue *Queue) SetArgs(args Args) *Queue
- func (queue *Queue) SetAutoDelete(autoDelete bool) *Queue
- func (queue *Queue) SetBindings(bindings ...*Binding) *Queue
- func (queue *Queue) SetDLQ(exchange *Exchange, key string) *Queue
- func (queue *Queue) SetDurable(durable bool) *Queue
- func (queue *Queue) SetExclusive(exclusive bool) *Queue
- func (queue *Queue) SetMessageTTL(ttl time.Duration) *Queue
- func (queue *Queue) SetName(name string) *Queue
- func (queue *Queue) SetNoWait(noWait bool) *Queue
- type Retry
- type Service
- func (service *Service) Channel() (channel IChannel, err error)
- func (service *Service) ChannelContext(ctx context.Context) (channel IChannel, err error)
- func (service *Service) CreateBinding(queue *Queue, binding *Binding) (err error)
- func (service *Service) CreateExchange(exchange *Exchange) (err error)
- func (service *Service) CreateQueues(queues ...*Queue) (result []*Queue, err error)
- func (service *Service) Restarting() bool
- func (service *Service) WaitReady() error
- func (service *Service) WaitReadyContext(ctx context.Context) error
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 Module = gioc.NewModule("RmqModule"). Provide( gioc.FactoryProvider( ServiceToken, gioc.NewFactory( ServiceInjections, gioc.Singleton, NewRmqService, ), true, ), gioc.FactoryProvider( ConsumerServiceToken, gioc.NewFactory( ConsumerServiceInjections, gioc.Singleton, NewConsumerService, ), true, ), gioc.FactoryProvider( ProducerServiceToken, gioc.NewFactory( ProducerServiceInjections, gioc.Prototype, NewProducerService, ), true, ), )
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 Binding ¶
func (*Binding) SetExchange ¶
type ChannelProvider ¶
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) SetAutoAck ¶
func (*Consumer) SetExclusive ¶
func (*Consumer) SetHandler ¶
func (consumer *Consumer) SetHandler(handler MessageHandler) *Consumer
func (*Consumer) SetNoLocal ¶
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 Exchange ¶
type Exchange struct {
Name string
Kind string
Durable bool
AutoDelete bool
Internal bool
NoWait bool
Args Args
}
func NewExchange ¶
func (*Exchange) SetAutoDelete ¶
func (*Exchange) SetDurable ¶
func (*Exchange) SetInternal ¶
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 IncomeMessage ¶
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 (*Queue) AddBindings ¶
func (*Queue) SetAutoDelete ¶
func (*Queue) SetBindings ¶
func (*Queue) SetDurable ¶
func (*Queue) SetExclusive ¶
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) ChannelContext ¶ added in v0.2.0
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) CreateExchange ¶
func (*Service) CreateQueues ¶
func (*Service) Restarting ¶
Click to show internal directories.
Click to hide internal directories.