Documentation
¶
Index ¶
Constants ¶
This section is empty.
Variables ¶
View Source
var ErrClosed = errors.New("publisher is closed")
Functions ¶
This section is empty.
Types ¶
type ChannelPublisher ¶
type ChannelPublisher struct {
// contains filtered or unexported fields
}
func NewChannelPublisher ¶
func NewChannelPublisher(c chan *message.Message) *ChannelPublisher
func (*ChannelPublisher) Close ¶
func (p *ChannelPublisher) Close() error
func (*ChannelPublisher) GetMessageID ¶
func (p *ChannelPublisher) GetMessageID(msg *message.Message) string
type Destination ¶
type Destination string
Destination is analogous to a topic in google pubsub, a stream in redis or kafka or a queue in RabbitMQ.
type MessageBoolOptFunc ¶
func ConstantBoolMsgFn ¶
func ConstantBoolMsgFn(value bool) MessageBoolOptFunc
type MessageIntOptFunc ¶
func ConstantIntMsgFn ¶
func ConstantIntMsgFn(i int) MessageIntOptFunc
type MessagePublisher ¶
type MessagePublisher[T any] struct { RoutingFunc RoutingFunc MessageMarshaller func(msg *message.Message) (T, error) GetDestinationPublisher func(d Destination) (MessagesPublisherImpl[T], error) // contains filtered or unexported fields }
MessagePublisher is a helper struct to help with the implementation of a Publisher implementation. T is the generic type of a message in the implementation It keeps a list of internal publishers, one for each destination.
func (*MessagePublisher[T]) Close ¶
func (p *MessagePublisher[T]) Close() error
Close closes every destination publisher. All of them are closed even if some fail, and the errors are returned joined together.
type MessageStringOptFunc ¶
func ConstantStringMsgFn ¶
func ConstantStringMsgFn(s string) MessageStringOptFunc
type MessagesPublisherImpl ¶
type PublishingEngine ¶
type PublishingEngine struct {
// contains filtered or unexported fields
}
func NewPublishingEngine ¶
func NewPublishingEngine(pub Publisher) *PublishingEngine
func (*PublishingEngine) AddMiddleware ¶
func (p *PublishingEngine) AddMiddleware(m middleware.Middleware) *PublishingEngine
type RoutingFunc ¶
type RoutingFunc func(msg *message.Message) (Destination, error)
func ConstantDestination ¶
func ConstantDestination(d Destination) RoutingFunc
ConstantDestination is a routing function implementation that always returns the same destination.
Click to show internal directories.
Click to hide internal directories.