From 8b863f7798f7e1a6a9266b88e7bca5aa0cb5ed8a Mon Sep 17 00:00:00 2001 From: Zahary Karadjov Date: Thu, 7 May 2020 01:18:58 +0300 Subject: [PATCH] Slightly more flexible pipelines (to be documented) --- faststreams/inputs.nim | 2 +- faststreams/pipelines.nim | 19 ++++++++++++++----- 2 files changed, 15 insertions(+), 6 deletions(-) diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index 7a511e4..a773960 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -12,7 +12,7 @@ type buffers*: PageBuffers # This is nil for unsafe memory inputs span*: PageSpan spanEndPos*: Natural - closeFut: Future[void] # This is nil before `close` is called + closeFut*: Future[void] # This is nil before `close` is called when debugHelpers: name*: string diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim index 4953054..9f5a805 100644 --- a/faststreams/pipelines.nim +++ b/faststreams/pipelines.nim @@ -305,20 +305,29 @@ macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untype stepInput = stepOutput var RetTypeExpr = copy steps[^1] - RetTypeExpr.insert(1, newCall("default", ident"AsyncOutputStream")) + RetTypeExpr.insert(1, newCall("default", ident"AsyncInputStream")) var closingCall = steps[^1] - closingCall.insert(1, newDotExpr(stepInput, ident"output")) + closingCall.insert(1, newCall(bindSym"initReader", stepInput)) pipelineBody.add quote do: await allFutures(`pipelineSteps`) - return `closingCall` + `closingCall` result = quote do: - type RetType = type(`RetTypeExpr`) + type UserOpRetType = type(`RetTypeExpr`) + + when UserOpRetType is Future: + type RetType = type(default(UserOpRetType).read) + else: + type RetType = UserOpRetType proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} = - `pipelineBody` + when UserOpRetType is Future: + var f = `pipelineBody` + return await(f) + else: + return `pipelineBody` pipelineProc(`start`)