package broadcast import ( "context" "jiby.grabit/src/node" ) // Opcode defines the type of command. type Opcode int 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 ) // Policy defines the message delivery strategy. type Policy int 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 ) // Command is a request to modify the worker's outputs. type Command[T any] struct { Op Opcode Node node.Duplex[T] } // Worker broadcasts messages to multiple output channels. type Worker[T any] struct { Input <-chan T Output map[string]chan<- T Cmds chan Command[T] Policy Policy } // New is like Init, but also does allocation func New[T any](input <-chan T, policy Policy) *Worker[T] { var w Worker[T] (&w).Init(input, policy) return &w } // Init initializes an existing worker with given input and policy. // It also allocates the internal map and command channel. func (w *Worker[T]) Init(input <-chan T, policy Policy) { *w = Worker[T]{ Input: input, Output: make(map[string]chan<- T), Cmds: make(chan Command[T]), Policy: policy, } } // Work processes a single event from input or commands. func (w *Worker[T]) Work(ctx context.Context) error { select { case <-ctx.Done(): return ctx.Err() case msg, ok := <-w.Input: if !ok { return nil } for _, ch := range w.Output { switch w.Policy { case Block: ch <- msg case Drop: select { case ch <- msg: default: } case FireForget: select { case ch <- msg: default: go func(c chan<- T, m T) { c <- m }(ch, msg) } } } return nil case cmd := <-w.Cmds: switch cmd.Op { case AddInput: if w.Input == nil { w.Input = cmd.Node.Ch } case RemoveInput: if w.Input != nil { w.Input = nil } case AddOutput: if _, ok := w.Output[cmd.Node.ID]; !ok { w.Output[cmd.Node.ID] = cmd.Node.Ch } case RemoveOutput: if _, ok := w.Output[cmd.Node.ID]; ok { delete(w.Output, cmd.Node.ID) } default: // silently drops invalid commands } return nil } } // AddOutput registers a new output channel. func (w *Worker[T]) AddOutput(id string, ch chan T) { w.Cmds <- Command[T]{ Op: AddOutput, Node: node.Duplex[T]{id,ch}, } } // RemoveOutput unregisters an output channel. func (w *Worker[T]) RemoveOutput(id string) { w.Cmds <- Command[T]{ Op: RemoveOutput, Node: (node.Duplex[T]{id,nil}), } } // AddInput registers the input channel. func (w *Worker[T]) AddInput(id string, ch chan T) { w.Cmds <- Command[T]{ Op: AddInput, Node: node.Duplex[T]{"",ch}, } } // RemoveInput unregisters the input channel. func (w *Worker[T]) RemoveInput(id string) { w.Cmds <- Command[T]{ Op: RemoveInput, Node: (node.Duplex[T]{"",nil}), } }