Documentation
¶
Index ¶
- Constants
- Variables
- func ProvideModuleConfig(config *ModuleConfig) gioc.IProvider
- type Args
- type Binding
- type Channel
- type ChannelProvider
- type Connection
- 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 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 Channel, 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
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 Channel ¶
type Channel 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 ChannelProvider ¶
type Connection ¶
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 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) CreateBinding ¶
func (*Service) CreateExchange ¶
func (*Service) CreateQueues ¶
func (*Service) Restarting ¶
Click to show internal directories.
Click to hide internal directories.