Pudu programming language
Menu
Package

@chrismichaelps / pudu-lang-mediator

In-process messaging for Pudu: requests, notifications, streams, pipeline behaviors, processors, and exception handling

0.1.0Apache-2.01

InstallClose

Reader.pudu

Pudu88 lines3.6 KB

GitHub ↗
1/** @Mediator.Reader.Puller — a stream read one item at a time while another thread produces it */2module PuduLangMediator.Reader34import Std.Channel as Channel5import Std.Concurrent.Cancel as Cancel6import Std.Concurrent as Concurrent7import Std.Math as Math8import Std.Sync as Sync9import PuduLangMediator.Constants.Messages as Messages10import PuduLangMediator.Context as Context11import PuduLangMediator as Messaging12import PuduLangMediator.Stream as Stream13import PuduLangMediator.Utils.Threads as Threads1415/** @Mediator.Reader.Step — one item, or the end of the stream with its outcome */16export type Step[T, E] = Item(T) | Finished(Messaging.Outcome[(), E])1718/** @Mediator.Reader.Reader — the consuming end of a stream produced on its own thread */19export type Reader[T, E] = {20  channel: Channel.Channel[Step[T, E]],21  context: Context.Context,22  worker: Option[Concurrent.Task],23  done: Sync.Cell[Bool]24}2526/// A reader of what `produce` emits on a thread of its own, holding at most `capacity` items27/// (at least one) ahead of the consumer. The producer waits while the reader is full and stops28/// once the reader is closed or `context` is cancelled.29export fn open[T, E](capacity: Int, context: &Context.Context, produce: fn(Context.Context, Stream.Sink[T]) -> Messaging.Outcome[(), E]) -> Reader[T, E] {30  let channel: Channel.Channel[Step[T, E]] = Channel.channel(Math.max(capacity, 1))31  let own = Context.withToken(context, Cancel.child(&context.token))32  let started = Concurrent.start(fn() -> () {33      let result: Sync.Cell[Option[Messaging.Outcome[(), E]]] = Sync.cell(None)34      let sink = fn(item: T) -> Bool { !Context.stopped(&own) && Channel.send(&channel, Item(item)) == Ok(()) }35      let ran = Concurrent.contain(fn() -> () { let _stored = Sync.set(&result, Some(produce(own, sink))) })36      let outcome = match (Sync.get(&result), ran) {37        case (Ok(Some(found)), _) => found38        case (_, Err(problem)) => Err(Messaging.Crashed(Threads.reasonOf(&problem)))39        case _ => Err(Messaging.Crashed(Messages.READER_CLOSED))40      }41      let _ended = Channel.send(&channel, Finished(outcome))42    })43  match started {44    case Ok(worker) => Reader{channel: channel, context: own, worker: Some(worker), done: Sync.cell(false)}45    case Err(problem) => {46      let _ended = Channel.send(&channel, Finished(Err(Messaging.Crashed(Threads.reasonOf(&problem)))))47      Reader{channel: channel, context: own, worker: None, done: Sync.cell(false)}48    }49  }50}5152/// The next item, waiting for the producer when none is ready; the stream's failure once, when53/// it ended with one; and `None` after it ended or the reader was closed.54export fn next[T, E](reader: &Reader[T, E]) -> Option[Messaging.Outcome[T, E]] {55  if Sync.get(&reader.done) == Ok(true) { return None }56  match Channel.receive(&reader.channel) {57    case Ok(Some(Item(item))) => Some(Ok(item))58    case Ok(Some(Finished(Err(failure)))) => {59      finish(reader)60      Some(Err(failure))61    }62    case _ => {63      finish(reader)64      None65    }66  }67}6869/// Every item still to come, or the stream's failure.70export fn rest[T, E](reader: &Reader[T, E]) -> Messaging.Outcome[Array[T], E] {71  var items: Array[T] = []72  while let Some(step) = next(reader) {73    items = items.push(step ?)74  }75  Ok(items)76}7778/// Stops the producer, discards what it had produced, and waits for its thread to end.79export fn close[T, E](reader: &Reader[T, E]) -> () {80  Context.cancel(&reader.context, Messages.READER_CLOSED)81  let _closed = Channel.close(&reader.channel)82  if let Some(worker) = reader.worker { let _joined = Concurrent.join(&worker) }83  finish(reader)84}8586/// Marks the reader as ended.87fn finish[T, E](reader: &Reader[T, E]) -> () { let _marked = Sync.set(&reader.done, true) }88