Skip to main content

stream_into

Function stream_into 

Source
pub fn stream_into<O, S, B, E>(
    runtime: &Handle,
    open: O,
    stop: impl Fn() -> bool + Clone + Send + Sync + 'static,
    write: impl FnMut(&[u8]) -> Result<()>,
) -> Result<u64, StreamError>
where O: Future<Output = Opened<S>> + Send + 'static, S: Stream<Item = Result<B, E>> + Send + 'static, B: AsRef<[u8]> + Send + 'static, E: Display,
Expand description

Run open on runtime and pass each stream chunk to write on this thread, in order, returning the bytes written. stop (and whether this side still listens) is checked between chunks and every STALL_CHECK of silence. Errors, never a short success, on open or chunk failure, a refused write, a stop, or runtime shutdown. Blocks: never call on a runtime worker.