Skip to content

Streams

Chris Michael edited this page Oct 6, 2026 · 1 revision

Pudu

Streams

A stream handler emits items by calling a sink, which answers whether it wants another. Stream behaviors wrap the pipeline and may hand next a sink of their own — to map, filter, or count items. A consumer takes items three ways: Mediator.stream with a sink, Mediator.collect for all of them, or Mediator.reader, which runs the stream on its own thread behind a bounded channel and lets you pull one item at a time:

module Ex06Streams

import Std.Io as Io
import PuduLangMediator.Context as Context
import PuduLangMediator.Mediator as Mediator
import PuduLangMediator as Messaging
import PuduLangMediator.Reader as Reader
import PuduLangMediator.Registration as Registration
import PuduLangMediator.Stream as Stream

export fn main() -> Int {
  let lines: Stream.Kind[Int, Str, Str] = Stream.kind("log.lines")
  let mediator = match Mediator.build([
      Registration.streamHandler(&lines, fn(count: Int, context: Context.Context, sink: Stream.Sink[Str]) -> Messaging.Outcome[(), Str] {
          var index = 1
          while index <= count {
            Context.check(&context) ?
            if !sink("line " + show(index)) { return Ok(()) }
            index = index + 1
          }
          Ok(())
        }),
      Registration.streamBehavior(&lines, fn(_count: Int, context: Context.Context, sink: Stream.Sink[Str], next: Stream.Next[Str, Str]) -> Messaging.Outcome[(), Str] {
          next(context, |line: Str| sink(line.toUpper()))
        })
    ]) {
    case Ok(built) => built
    case Err(invalid) => panic(Mediator.explain(&invalid))
  }
  let all = Mediator.collect(&mediator, &lines, 3)
  let _all = Io.writeLine("collected: " + Messaging.summarize(&all))
  let firstTwo = Mediator.stream(&mediator, &lines, 100, fn(line: Str) -> Bool {
      let _said = Io.writeLine("sink got " + line)
      line != "LINE 2"
    })
  let reader = Mediator.reader(&mediator, &lines, 1000000, &Context.create(), 8)
  let pulled = [Reader.next(&reader), Reader.next(&reader)]
  Reader.close(&reader)
  let _pulled = Io.writeLine("pulled: " + show(pulled) + ", then " + show(Reader.next(&reader)))
  if all == Ok(["LINE 1", "LINE 2", "LINE 3"]) && firstTwo == Ok(()) { 0 } else { 1 }
}

Check it and run it:

pudu check src/Ex06Streams.pudu
pudu run src/Ex06Streams.pudu

Output:

collected: Ok(["LINE 1", "LINE 2", "LINE 3"])
sink got LINE 1
sink got LINE 2
pulled: [Some(Ok("LINE 1")), Some(Ok("LINE 2"))], then None

The sink declined after LINE 2, so the handler stopped there and the stream still succeeded. The reader held at most eight items ahead of the consumer: the producer waits while it is full. Reader.close stops the producer through its token, wakes it if it is waiting, and joins its thread; a closed reader answers None. A producer that fails ends the reader with that failure once; one whose thread stops ends it with Crashed.

Related

Clone this wiki locally