Documentation ¶
Overview ¶
Package pipeline contains structures that implement both the Producer and Consumer interfaces. They can be used as extra pipeline components for various utilities.
Index ¶
- type Processor
- func (p *Processor) CloseAsync()
- func (p *Processor) MessageChan() <-chan types.Message
- func (p *Processor) ResponseChan() <-chan types.Response
- func (p *Processor) StartListening(responses <-chan types.Response) error
- func (p *Processor) StartReceiving(msgs <-chan types.Message) error
- func (p *Processor) WaitForClose(timeout time.Duration) error
- type Type
Constants ¶
This section is empty.
Variables ¶
This section is empty.
Functions ¶
This section is empty.
Types ¶
type Processor ¶
type Processor struct {
// contains filtered or unexported fields
}
Processor is a pipeline that supports both Consumer and Producer interfaces. The processor will read from a source, perform some processing, and then either propagate a new message or drop it.
func NewProcessor ¶
func NewProcessor( log log.Modular, stats metrics.Type, msgProcessors ...processor.Type, ) *Processor
NewProcessor returns a new message processing pipeline.
func (*Processor) CloseAsync ¶
func (p *Processor) CloseAsync()
CloseAsync shuts down the pipeline and stops processing messages.
func (*Processor) MessageChan ¶
MessageChan returns the channel used for consuming messages from this pipeline.
func (*Processor) ResponseChan ¶
ResponseChan returns the response channel from this pipeline.
func (*Processor) StartListening ¶
StartListening sets the channel that this pipeline will read responses from.
func (*Processor) StartReceiving ¶
StartReceiving assigns a messages channel for the pipeline to read.