Module std::stream

Stream and Sink — the one pipe family.

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:

Example

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);
    }
}

Contents

Functions

Function pipe

pub fn pipe(capacity: i64) -> Result<(Sink<T>, Stream<T>), string>

Create 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>) = ....

Function open

pub fn open(path: string) -> Result<Stream<bytes>, string>

Open 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.

Function forward

pub fn forward(from: Stream<T>, to: Sink<T>)

Drain 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.

Examples

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()),
    }
}

Function pair_sink

pub fn pair_sink(pair: StreamPair) -> Sink<T>

Extract the write half of a pair. Used by std.net to split a socket.

Function pair_stream

pub fn pair_stream(pair: StreamPair) -> Stream<T>

Extract the read half of a pair. Used by std.net to split a socket.

Types

Struct Stream

The read half of a pipe, parameterised by element type.

Struct Sink

The write half of a pipe, parameterised by element type.

Enum CodecError

Why a frame could not be decoded.

Variants

Malformed(string)

The bytes do not form a frame of this codec.

Enum Decoded

One decode step: a complete frame with the bytes after it, or the buffer handed back unchanged because more bytes are needed.

Variants

Frame(T, bytes)
Incomplete(bytes)

Struct Lines

Newline-terminated text frames. The newline is not part of the item.

Struct LengthPrefixed

Frames prefixed by a big-endian length of width bytes (1, 2 or 4).

Fields

width: i64

Struct StreamPair

Internal paired pipe handle: both halves before extraction.

Traits

Trait Codec

Frames 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.

Methods

fn decode(self: Self, buf: bytes) -> Result<Decoded<T>, CodecError>
fn encode(self: Self, item: T) -> bytes