diff --git a/src/tink/io/Sink.hx b/src/tink/io/Sink.hx index 089462a..06769da 100644 --- a/src/tink/io/Sink.hx +++ b/src/tink/io/Sink.hx @@ -4,6 +4,7 @@ import tink.Chunk; import tink.io.PipeOptions; import tink.streams.Stream; +using tink.io.PipeResult; using tink.io.Source; using tink.CoreApi; @@ -138,4 +139,49 @@ class SinkBase implements SinkObject { //override public function endSafely():Future { //return target.end().recover(function (_) return Future.sync(false)); //} -//} \ No newline at end of file +//} + +class CollectSink extends SinkBase { + var ended = false; + var result:PromiseTrigger = Promise.trigger(); + var collected:Chunk = Chunk.EMPTY; + var process:Chunk->Promise; + + public function new(process) { + this.process = process; + } + + override function get_sealed() return ended; + + override function consume(source:Stream, options:PipeOptions):Future> { + return + if(ended) + result.asPromise().map(function(o) return switch o { + case Success(result): SinkEnded(result, (cast source:Source)); + case Failure(e): SinkFailed(e, (cast source:Source)); + }); + else + source.forEach(function(chunk) { + collected = collected & chunk; + return Resume; + }).flatMap(function(o):Future> return switch o { + case Depleted: + if(options.end) { + ended = true; + process(collected).map(function(o) { + result.trigger(o); + return switch o { + case Success(result): SinkEnded(result, (cast Source.EMPTY:Source)); + case Failure(e): SinkFailed(e, (cast Source.EMPTY:Source)); + } + }); + } else { + Future.sync(AllWritten); + } + case Failed(e): + Future.sync(SourceFailed(e)); + case Halted(rest): + throw 'unreachable'; + }); + } +} \ No newline at end of file