Slightly more flexible pipelines (to be documented)
This commit is contained in:
parent
ec5c3af313
commit
8b863f7798
2 changed files with 15 additions and 6 deletions
|
|
@ -12,7 +12,7 @@ type
|
||||||
buffers*: PageBuffers # This is nil for unsafe memory inputs
|
buffers*: PageBuffers # This is nil for unsafe memory inputs
|
||||||
span*: PageSpan
|
span*: PageSpan
|
||||||
spanEndPos*: Natural
|
spanEndPos*: Natural
|
||||||
closeFut: Future[void] # This is nil before `close` is called
|
closeFut*: Future[void] # This is nil before `close` is called
|
||||||
when debugHelpers:
|
when debugHelpers:
|
||||||
name*: string
|
name*: string
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -305,20 +305,29 @@ macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untype
|
||||||
stepInput = stepOutput
|
stepInput = stepOutput
|
||||||
|
|
||||||
var RetTypeExpr = copy steps[^1]
|
var RetTypeExpr = copy steps[^1]
|
||||||
RetTypeExpr.insert(1, newCall("default", ident"AsyncOutputStream"))
|
RetTypeExpr.insert(1, newCall("default", ident"AsyncInputStream"))
|
||||||
|
|
||||||
var closingCall = steps[^1]
|
var closingCall = steps[^1]
|
||||||
closingCall.insert(1, newDotExpr(stepInput, ident"output"))
|
closingCall.insert(1, newCall(bindSym"initReader", stepInput))
|
||||||
|
|
||||||
pipelineBody.add quote do:
|
pipelineBody.add quote do:
|
||||||
await allFutures(`pipelineSteps`)
|
await allFutures(`pipelineSteps`)
|
||||||
return `closingCall`
|
`closingCall`
|
||||||
|
|
||||||
result = quote do:
|
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.} =
|
proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} =
|
||||||
`pipelineBody`
|
when UserOpRetType is Future:
|
||||||
|
var f = `pipelineBody`
|
||||||
|
return await(f)
|
||||||
|
else:
|
||||||
|
return `pipelineBody`
|
||||||
|
|
||||||
pipelineProc(`start`)
|
pipelineProc(`start`)
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue