package core type Emitter[OUT any] interface { Emit(record StreamRecord[OUT]) error } type Operator[IN, OUT any] interface { Open() error Process(record StreamRecord[IN], out Emitter[OUT]) error Close() error }