pub fn barrier_to_message_stream( rx: UnboundedReceiver<Barrier>, ) -> impl Stream<Item = Result<Message, StreamExecutorError>>
Receive barriers from barrier manager with the channel, error on channel close.