std::streamStream
A pipe moves data between actors, sockets and files. Stream<T> is the
read half and Sink<T> the write half; the same two handle types cover an
in-memory pipe (pipe), a file (open), a socket (Connection.split)
and an actor's receive gen fn.
The contract, stated once:
input.recv() returns Option<T>: Some(item) is data, None means
every sink finished. A transport failure traps the reading actor with a
typed I/O fault; the loop never sees an error value. A zero-length
bytes or empty string is a valid item.sink.send(item) waits for capacity and returns Result<(), SendError>:
Err(SendError.Closed) once the reader is gone or the sink finished.
try_send never waits and adds Err(SendError.Full).sink.clone() adds a producer; the pipe reaches EOF when the last one
finishes. sink.finish() publishes EOF and keeps the handle (on a
socket: FIN to the peer). close() on either half consumes it; scope
exit releases a live half the same way.string, bytes, or a value record, enum or tuple. Containers and
handles are refused at check time.Stream<T> is a select source: item from input.recv() => ....import std.stream;
fn main() {
let (sink, input): (stream.Sink<string>, stream.Stream<string>) =
match stream.pipe(16) { .Ok(pair) => pair, .Err(error) => panic(error), };
sink.send("hello").expect("send");
sink.send("world").expect("send");
sink.close(); // optional — would auto-close on drop
for item in input {
println(item);
}
}
pipeCreate a bounded in-memory pipe with the given capacity.
The element type is inferred from use; write it at the binding when
nothing else fixes it: let (tx, rx): (stream.Sink<i64>, stream.Stream<i64>) = ....
openOpen a file as a chunked Stream<bytes>.
Each item is one chunk of the file; .lines() re-frames the bytes into
text lines. A read failure after open traps the reading actor.
forwardDrain from into to, then finish to and release both halves.
This is the streaming equivalent of io.Copy: it reads until EOF and
writes each item downstream, waiting on the destination's capacity.
import std.net.http;
import std.stream;
fn main() {
match http.listen(":8080") {
.Ok(server) => {
let req = server.accept();
match stream.open("data.txt") {
.Ok(file) => {
let body = match req.respond_stream(200, "text/plain") {
.Ok(body) => body,
.Err(error) => panic(to_string(error)),
};
stream.forward(file.lines(), body); // streams file to the HTTP response
}
.Err(e) => println(e),
}
req.close();
server.close();
}
.Err(_err) => println(http.listen_error()),
}
}
pair_sinkExtract the write half of a pair. Used by std.net to split a socket.
pair_streamExtract the read half of a pair. Used by std.net to split a socket.
StreamThe read half of a pipe, parameterised by element type.
SinkThe write half of a pipe, parameterised by element type.
CodecErrorWhy a frame could not be decoded.
Malformed(string)The bytes do not form a frame of this codec.
DecodedOne decode step: a complete frame with the bytes after it, or the
buffer handed back unchanged because more bytes are needed.
Frame(T, bytes)Incomplete(bytes)LinesNewline-terminated text frames. The newline is not part of the item.
LengthPrefixedFrames prefixed by a big-endian length of width bytes (1, 2 or 4).
StreamPairInternal paired pipe handle: both halves before extraction.
CodecFrames bytes into typed items and back.
decode takes the buffered bytes and returns the first complete frame
with the bytes that follow it, or Incomplete with the buffer unchanged
when more bytes are needed. encode returns one item's frame. Values
are copy-on-write, so the caller keeps what decode hands back: a var
parameter is the callee's own copy and cannot advance the caller's
buffer. Lines and LengthPrefixed are the shipped codecs;
Stream<bytes>.lines() is the Lines case as a method.