std/stream
Streams: one flow in one direction, with a reading end and a writing end.
"Stream" is the word for the flow, not a type. What a signature names is one of its two ends - Source<Item, Failure>
to read from, Sink<Item, Failure> to write into - and they are the asynchronous siblings of Iterator and
Accumulator, with the same verbs (next, add, finish). Everything in between, every map, filter and
framer, is a Stage of std/iteration and is written once for both worlds.
const all = response.body.through(Json.items<User>()).checked().toList().await()?
File.write(path, response.body.mapFailure(IoError.from)).await()?
See docs/design/STREAMS.md for the whole model: the contracts of both ends, what is synchronous and what has to wait, the
drivers, and what was taken from Rust, C#, Swift, Scala, Java, Node, Web Streams and Bun.
Modules
std/stream/bytesBytes, and the two stages between bytes and text.std/stream/sinkThe writing end of a stream: the asynchronous sibling ofAccumulator, with the same two verbs and the samevar self.std/stream/sourceThe reading end of a stream:Source, the stages built from it (map,take,through, ...), the terminal operations that read it to a value (toList,fold,into, ...), and its three ways to make one (from,pulling,produce).
Everything
- type
BufferedA sink in front of another one that hands items on in groups. - alias
BytesA chunk of bytes. - fn
decodedTextBytes to text, one chunk at a time and cut at character borders: what arrives is whatever was complete, and a character split across two chunks is held back until the rest is there. - fn
encodedTextText to bytes. - fn
linesBytes to lines, split at\n, with a\rbefore it dropped so that CRLF files read the same as LF files. - type
PullingSource.pulling { ... }: one closure fornext, and the thing above it that has to be released. - type
PushingSink.pushing(accept, complete): one closure per verb, and the thing below it that has to be released. - extend
Sink<Bytes, Failure>Writing text into a sink of bytes. - trait
Sinkusing file = File.create(path)? file.addLine("first").await()? file.end().await(). - trait
SourceThe reading end of a stream: the asynchronous sibling ofIterator, with the same verb and the samevar self. - extend
Source<Result<Item, Problem>, Failure>A stage that can fail answersResultitems (Json.items<User>(),lines()), because aStageis synchronous and knows nothing about the stream's failure. - type
Stagedsource.through(stage). - fn
textOfA whole chunk of bytes as text - the short form for everything that has all of its bytes already (a file that was read in one piece, a body below its limit). - type
Utf8ErrorBytes that are not UTF-8, with the offset of the byte that broke it.