
Reader.pudu
Pudu88 lines3.6 KB
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}2526272829export 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}51525354export 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}686970export 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}777879export 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}858687fn finish[T, E](reader: &Reader[T, E]) -> () { let _marked = Sync.set(&reader.done, true) }88