broadcast

package
v0.0.0 Latest Latest
Warning

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

Go to latest
Published: Jul 20, 2026 License: GPL-3.0 Imports: 0 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Command

type Command[T any] struct {
	Op   Opcode
	Node node.Duplex[T]
}

Command is a request to modify the worker's outputs.

type Opcode

type Opcode int

Opcode defines the type of command.

const (
	// AddOutput adds a new output channel.
	AddOutput Opcode = iota
	// RemoveOutput removes an existing output channel.
	RemoveOutput
	// AddInput adds a new input channel.
	AddInput Opcode = iota
	// RemoveInput removes an existing input channel.
	RemoveInput
)

type Policy

type Policy int

Policy defines the message delivery strategy.

const (
	// Block waits for the consumer to be ready.
	Block Policy = iota
	// Drop discards the message if the consumer is full.
	Drop
	// FireForget sends the message without waiting.
	FireForget
)

type Worker

type Worker[T any] struct {
	Input  <-chan T
	Output map[string]chan<- T
	Cmds   chan Command[T]
	Policy Policy
}

Worker broadcasts messages to multiple output channels.

func New

func New[T any](input <-chan T, policy Policy) *Worker[T]

New is like Init, but also does allocation

func (*Worker[T]) AddInput

func (w *Worker[T]) AddInput(id string, ch chan T)

AddInput registers the input channel.

func (*Worker[T]) AddOutput

func (w *Worker[T]) AddOutput(id string, ch chan T)

AddOutput registers a new output channel.

func (*Worker[T]) Init

func (w *Worker[T]) Init(input <-chan T, policy Policy)

Init initializes an existing worker with given input and policy. It also allocates the internal map and command channel.

func (*Worker[T]) RemoveInput

func (w *Worker[T]) RemoveInput(id string)

RemoveInput unregisters the input channel.

func (*Worker[T]) RemoveOutput

func (w *Worker[T]) RemoveOutput(id string)

RemoveOutput unregisters an output channel.

func (*Worker[T]) Work

func (w *Worker[T]) Work(ctx context.Context) error

Work processes a single event from input or commands.

Jump to

Keyboard shortcuts

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