publisher

package
v0.0.0-...-f414a8c Latest Latest
Warning

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

Go to latest
Published: Oct 1, 2026 License: MIT Imports: 6 Imported by: 0

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

func (*ChannelPublisher) Publish

func (p *ChannelPublisher) Publish(msg *message.Message) error

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

type MessageBoolOptFunc func(msg *message.Message) (bool, error)

func ConstantBoolMsgFn

func ConstantBoolMsgFn(value bool) MessageBoolOptFunc

type MessageIntOptFunc

type MessageIntOptFunc func(msg *message.Message) (int, error)

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.

func (*MessagePublisher[T]) Publish

func (p *MessagePublisher[T]) Publish(message *message.Message) error

type MessageStringOptFunc

type MessageStringOptFunc func(msg *message.Message) (string, error)

func ConstantStringMsgFn

func ConstantStringMsgFn(s string) MessageStringOptFunc

type MessagesPublisherImpl

type MessagesPublisherImpl[T any] interface {
	io.Closer
	Publish(ctx context.Context, message T) error
	GetMessageID(message T) string
}

type Publisher

type Publisher interface {
	io.Closer
	Publish(message *message.Message) error
}

type PublishingEngine

type PublishingEngine struct {
	// contains filtered or unexported fields
}

func NewPublishingEngine

func NewPublishingEngine(pub Publisher) *PublishingEngine

func (*PublishingEngine) AddMiddleware

func (*PublishingEngine) Publish

func (p *PublishingEngine) Publish(msg *message.Message) error

Publish publishes the messages to the destination topic calculated by the routing function.

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.

Jump to

Keyboard shortcuts

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