Don't import/export chronos by default (#20)

Chronos support is optional and should not have to be imported in order
to use faststreams for non-async use cases, as doing so pollutes the
global namespace and slows down compilation.

`async` support must now explicitly be enabled with
-d:async_backend=chronos|asyncdispatch
This commit is contained in:
Jacek Sieka 2021-08-20 13:10:26 +02:00 • committed by GitHub
commit 3a0ab42573
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
10 changed files with 841 additions and 705 deletions

View file

@ -1,14 +1,25 @@
const
faststreams_async_backend {.strdefine.} = "chronos"
# To compile with async support, use `-d:async_backend=chronos|asyncdispatch`
async_backend {.strdefine.} = "none"
const
faststreams_async_backend {.strdefine.} = ""
when faststreams_async_backend != "":
{.fatal: "use `-d:async_backend` instead".}
type
CloseBehavior* = enum
waitAsyncClose
dontWaitAsyncClose
const debugHelpers* = defined(debugHelpers)
const
debugHelpers* = defined(debugHelpers)
fsAsyncSupport* = async_backend != "none"
when faststreams_async_backend == "chronos":
when async_backend == "none":
discard
elif async_backend == "chronos":
import
chronos
@ -18,10 +29,10 @@ when faststreams_async_backend == "chronos":
template fsAwait*(f: Future): untyped =
await f
elif faststreams_async_backend in ["std", "asyncdispatch"]:
elif async_backend in ["std", "asyncdispatch"]:
import
std/asyncdispatch
export
asyncdispatch
@ -36,7 +47,7 @@ elif faststreams_async_backend in ["std", "asyncdispatch"]:
type Duration* = int
else:
{.fatal: "Unrecognized network backend: " & faststreams_async_backend.}
{.fatal: "Unrecognized network backend: " & async_backend.}
when defined(danger):
template fsAssert*(x) = discard

View file

@ -22,8 +22,9 @@ type
maxBufferedBytes*: Natural
queue*: Deque[PageRef]
waitingReader*: Future[void]
waitingWriter*: Future[void]
when fsAsyncSupport:
waitingReader*: Future[void]
waitingWriter*: Future[void]
eofReached*: bool
fauxEofPos*: Natural

View file

@ -6,16 +6,73 @@ import
export
options, CloseBehavior
type
InputStream* = ref object of RootObj
vtable*: ptr InputStreamVTable # This is nil for unsafe memory inputs
buffers*: PageBuffers # This is nil for unsafe memory inputs
span*: PageSpan
spanEndPos*: Natural
closeFut*: Future[void] # This is nil before `close` is called
when debugHelpers:
name*: string
when fsAsyncSupport:
# Circular type refs prevent more targeted `when`
type
InputStream* = ref object of RootObj
vtable*: ptr InputStreamVTable # This is nil for unsafe memory inputs
buffers*: PageBuffers # This is nil for unsafe memory inputs
span*: PageSpan
spanEndPos*: Natural
closeFut*: Future[void] # This is nil before `close` is called
when debugHelpers:
name*: string
AsyncInputStream* {.borrow: `.`.} = distinct InputStream
ReadSyncProc* = proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].}
ReadAsyncProc* = proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseSyncProc* = proc (s: InputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseAsyncProc* = proc (s: InputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
GetLenSyncProc* = proc (s: InputStream): Option[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].}
InputStreamVTable* = object
readSync*: ReadSyncProc
closeSync*: CloseSyncProc
getLenSync*: GetLenSyncProc
readAsync*: ReadAsyncProc
closeAsync*: CloseAsyncProc
MaybeAsyncInputStream* = InputStream | AsyncInputStream
else:
type
InputStream* = ref object of RootObj
vtable*: ptr InputStreamVTable # This is nil for unsafe memory inputs
buffers*: PageBuffers # This is nil for unsafe memory inputs
span*: PageSpan
spanEndPos*: Natural
when debugHelpers:
name*: string
ReadSyncProc* = proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseSyncProc* = proc (s: InputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
GetLenSyncProc* = proc (s: InputStream): Option[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].}
InputStreamVTable* = object
readSync*: ReadSyncProc
closeSync*: CloseSyncProc
getLenSync*: GetLenSyncProc
MaybeAsyncInputStream* = InputStream
type
LayeredInputStream* = ref object of InputStream
source*: InputStream
allowWaitFor*: bool
@ -23,30 +80,6 @@ type
InputStreamHandle* = object
s*: InputStream
AsyncInputStream* {.borrow: `.`.} = distinct InputStream
ReadSyncProc* = proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].}
ReadAsyncProc* = proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseSyncProc* = proc (s: InputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseAsyncProc* = proc (s: InputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
GetLenSyncProc* = proc (s: InputStream): Option[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].}
InputStreamVTable* = object
readSync*: ReadSyncProc
readAsync*: ReadAsyncProc
closeSync*: CloseSyncProc
closeAsync*: CloseAsyncProc
getLenSync*: GetLenSyncProc
MemFileInputStream = ref object of InputStream
file: MemFile
@ -54,30 +87,38 @@ type
file: File
template Sync*(s: InputStream): InputStream = s
template Async*(s: InputStream): AsyncInputStream = AsyncInputStream(s)
template Sync*(s: AsyncInputStream): InputStream = InputStream(s)
template Async*(s: AsyncInputStream): AsyncInputStream = s
when fsAsyncSupport:
template Async*(s: InputStream): AsyncInputStream = AsyncInputStream(s)
template Sync*(s: AsyncInputStream): InputStream = InputStream(s)
template Async*(s: AsyncInputStream): AsyncInputStream = s
proc disconnectInputDevice(s: InputStream) =
# TODO
# Document the behavior that closeAsync is preferred
if s.vtable != nil:
if s.vtable.closeAsync != nil:
s.closeFut = s.vtable.closeAsync(s)
elif s.vtable.closeSync != nil:
s.vtable.closeSync(s)
when fsAsyncSupport:
if s.vtable.closeAsync != nil:
s.closeFut = s.vtable.closeAsync(s)
elif s.vtable.closeSync != nil:
s.vtable.closeSync(s)
else:
if s.vtable.closeSync != nil:
s.vtable.closeSync(s)
s.vtable = nil
template disconnectInputDevice(s: AsyncInputStream) =
disconnectInputDevice InputStream(s)
when fsAsyncSupport:
template disconnectInputDevice(s: AsyncInputStream) =
disconnectInputDevice InputStream(s)
proc preventFurtherReading(s: InputStream) =
s.vtable = nil
s.span = default(PageSpan)
template preventFurtherReading(s: AsyncInputStream) =
preventFurtherReading InputStream(s)
when fsAsyncSupport:
template preventFurtherReading(s: AsyncInputStream) =
preventFurtherReading InputStream(s)
template makeHandle*(sp: InputStream): InputStreamHandle =
let s = sp
@ -94,23 +135,25 @@ proc close*(s: InputStream,
## `waitFor` to block until the async operation completes.
s.disconnectInputDevice()
s.preventFurtherReading()
if s.closeFut != nil:
fsTranslateErrors "Stream closing failed":
if behavior == waitAsyncClose:
waitFor s.closeFut
else:
asyncCheck s.closeFut
when fsAsyncSupport:
if s.closeFut != nil:
fsTranslateErrors "Stream closing failed":
if behavior == waitAsyncClose:
waitFor s.closeFut
else:
asyncCheck s.closeFut
template close*(sp: AsyncInputStream) =
## Starts the asychronous closing of the stream and returns a future that
## tracks the closing operation.
let s = InputStream sp
disconnectInputDevice(s)
preventFurtherReading(s)
if s.closeFut != nil:
fsAwait s.closeFut
when fsAsyncSupport:
template close*(sp: AsyncInputStream) =
## Starts the asychronous closing of the stream and returns a future that
## tracks the closing operation.
let s = InputStream sp
disconnectInputDevice(s)
preventFurtherReading(s)
if s.closeFut != nil:
fsAwait s.closeFut
template closeNoWait*(sp: AsyncInputStream|InputStream) =
template closeNoWait*(sp: MaybeAsyncInputStream) =
## Close the stream without waiting even if's async.
## This operation will use `asyncCheck` internally to detect unhandled
## errors from the closing operation.
@ -224,56 +267,52 @@ proc readableNow*(s: InputStream): bool =
getNewSpan s
s.span.hasRunway
template readableNow*(s: AsyncInputStream): bool =
readableNow InputStream(s)
when fsAsyncSupport:
template readableNow*(s: AsyncInputStream): bool =
readableNow InputStream(s)
# TODO: The pure async interface should be moved in a separate module
# to make FastStreams more light-weight when the async back-end
# is not used (e.g. in Confutils)
#
# The problem is that the `async` macro will pull the entire
# event loop right now.
proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} =
let s = InputStream(sp)
fsAssert s.buffers != nil and s.vtable != nil
proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} =
let s = InputStream(sp)
fsAssert s.buffers != nil and s.vtable != nil
result = fsAwait s.vtable.readAsync(s, nil, 0)
result = fsAwait s.vtable.readAsync(s, nil, 0)
if s.buffers.eofReached:
disconnectInputDevice(s)
if s.buffers.eofReached:
disconnectInputDevice(s)
if result > 0 and s.span.len == 0:
getNewSpan s
if result > 0 and s.span.len == 0:
getNewSpan s
proc timeoutToNextByteImpl(s: AsyncInputStream,
deadline: Future): Future[bool] {.async.} =
let readFut = s.readOnce
fsAwait readFut or deadline
if not readFut.finished:
readFut.cancel()
return true
else:
return false
proc timeoutToNextByteImpl(s: AsyncInputStream,
deadline: Future): Future[bool] {.async.} =
let readFut = s.readOnce
fsAwait readFut or deadline
if not readFut.finished:
readFut.cancel()
return true
else:
return false
template timeoutToNextByte*(sp: AsyncInputStream, deadline: Future): bool =
let s = sp
if readableNow(s):
true
else:
fsAwait timeoutToNextByteImpl(s, deadline)
template timeoutToNextByte*(sp: AsyncInputStream, deadline: Future): bool =
let s = sp
if readableNow(s):
true
else:
fsAwait timeoutToNextByteImpl(s, deadline)
template timeoutToNextByte*(sp: AsyncInputStream, timeout: Duration): bool =
let s = sp
if readableNow(s):
true
else:
fsAwait timeoutToNextByteImpl(s, sleepAsync(timeout))
template timeoutToNextByte*(sp: AsyncInputStream, timeout: Duration): bool =
let s = sp
if readableNow(s):
true
else:
fsAwait timeoutToNextByteImpl(s, sleepAsync(timeout))
proc closeAsync*(s: AsyncInputStream) {.async.} =
close s
proc closeAsync*(s: AsyncInputStream) {.async.} =
close s
# TODO: End of purely async interface
template totalUnconsumedBytes*(s: AsyncInputStream): Natural =
## Alias for InputStream.totalUnconsumedBytes
totalUnconsumedBytes InputStream(s)
func getBestContiguousRunway(s: InputStream): Natural =
result = s.span.len
@ -298,10 +337,6 @@ func totalUnconsumedBytes*(s: InputStream): Natural =
localRunway + runwayInBuffers
template totalUnconsumedBytes*(s: AsyncInputStream): Natural =
## Alias for InputStream.totalUnconsumedBytes
totalUnconsumedBytes InputStream(s)
proc limitReadableRange(s: InputStream, rangeLen: Natural): Natural =
s.vtable = nil
@ -317,7 +352,7 @@ proc limitReadableRange(s: InputStream, rangeLen: Natural): Natural =
s.buffers.queue.peekFirst.consumedTo -= bytesToUnconsume
return s.buffers.setFauxEof(s.spanEndPos)
template withReadableRange*(sp: InputStream|AsyncInputStream,
template withReadableRange*(sp: MaybeAsyncInputStream,
rangeLen: Natural,
rangeStreamVarName, blk: untyped) =
let
@ -409,8 +444,9 @@ proc len*(s: InputStream): Option[Natural] {.raises: [Defect, IOError].} =
else:
none Natural
template len*(s: AsyncInputStream): Option[Natural] =
len InputStream(s)
when fsAsyncSupport:
template len*(s: AsyncInputStream): Option[Natural] =
len InputStream(s)
func memoryInput*(buffers: PageBuffers): InputStreamHandle =
var spanEndPos = Natural 0
@ -541,15 +577,16 @@ template readable*(sp: InputStream): bool =
let s = sp
hasRunway(s.span) or bufferMoreDataSync(s)
template readable*(sp: AsyncInputStream): bool =
## Async version of `readable`.
## The intended API usage is the same. Instead of blocking, an async
## stream will use `await` while waiting for more data.
let s = InputStream sp
if hasRunway(s.span):
true
else:
bufferMoreDataImpl(s, fsAwait, readAsync)
when fsAsyncSupport:
template readable*(sp: AsyncInputStream): bool =
## Async version of `readable`.
## The intended API usage is the same. Instead of blocking, an async
## stream will use `await` while waiting for more data.
let s = InputStream sp
if hasRunway(s.span):
true
else:
bufferMoreDataImpl(s, fsAwait, readAsync)
func continueAfterReadN(s: InputStream,
runwayBeforeRead, bytesRead: Natural) =
@ -615,15 +652,16 @@ proc readable*(s: InputStream, n: int): bool =
## for futher discussion of this.
readableNImpl(s, n, noAwait, readSync)
template readable*(sp: AsyncInputStream, np: int): bool =
## Async version of `readable(n)`.
## The intended API usage is the same. Instead of blocking, an async
## stream will use `await` while waiting for more data.
let
s = InputStream sp
n = np
when fsAsyncSupport:
template readable*(sp: AsyncInputStream, np: int): bool =
## Async version of `readable(n)`.
## The intended API usage is the same. Instead of blocking, an async
## stream will use `await` while waiting for more data.
let
s = InputStream sp
n = np
readableNImpl(s, n, fsAwait, readAsync)
readableNImpl(s, n, fsAwait, readAsync)
template peek*(sp: InputStream): byte =
let s = sp
@ -633,8 +671,9 @@ template peek*(sp: InputStream): byte =
getNewSpanOrDieTrying s
s.span.startAddr[]
template peek*(s: AsyncInputStream): byte =
peek InputStream(s)
when fsAsyncSupport:
template peek*(s: AsyncInputStream): byte =
peek InputStream(s)
func readFromNewSpan(s: InputStream): byte =
getNewSpanOrDieTrying s
@ -650,8 +689,9 @@ template read*(sp: InputStream): byte =
else:
readFromNewSpan s
template read*(s: AsyncInputStream): byte =
read InputStream(s)
when fsAsyncSupport:
template read*(s: AsyncInputStream): byte =
read InputStream(s)
proc peekAt*(s: InputStream, pos: int): byte {.inline.} =
# TODO implement page flipping
@ -659,8 +699,9 @@ proc peekAt*(s: InputStream, pos: int): byte {.inline.} =
fsAssert cast[uint](peekHead) < cast[uint](s.span.endAddr)
return peekHead[]
template peekAt*(s: AsyncInputStream, pos: int): byte =
peekAt InputStream(s), pos
when fsAsyncSupport:
template peekAt*(s: AsyncInputStream, pos: int): byte =
peekAt InputStream(s), pos
proc advance*(s: InputStream) =
if hasRunway(s.span):
@ -673,11 +714,12 @@ proc advance*(s: InputStream, n: Natural) =
for i in 0 ..< n:
advance s
template advance*(s: AsyncInputStream) =
advance InputStream(s)
when fsAsyncSupport:
template advance*(s: AsyncInputStream) =
advance InputStream(s)
template advance*(s: AsyncInputStream, n: Natural) =
advance InputStream(s), n
template advance*(s: AsyncInputStream, n: Natural) =
advance InputStream(s), n
proc drainBuffersInto*(s: InputStream, dstAddr: ptr byte, dstLen: Natural): Natural =
var
@ -789,29 +831,30 @@ proc readInto*(s: InputStream, target: var openarray[byte]): bool =
## regarding the number of bytes read, see `readIntoEx`.
s.readIntoEx(target) == target.len
template readIntoEx*(sp: AsyncInputStream, dst: var openarray[byte]): int =
let s = InputStream(sp)
# BEWARE! `openArrayToPair` here is needed to avoid
# double evaluation of the `dst` expression:
let (dstAddr, dstLen) = openArrayToPair(dst)
readIntoExImpl(s, dstAddr, dstLen, fsAwait, readAsync)
when fsAsyncSupport:
template readIntoEx*(sp: AsyncInputStream, dst: var openarray[byte]): int =
let s = InputStream(sp)
# BEWARE! `openArrayToPair` here is needed to avoid
# double evaluation of the `dst` expression:
let (dstAddr, dstLen) = openArrayToPair(dst)
readIntoExImpl(s, dstAddr, dstLen, fsAwait, readAsync)
template readInto*(sp: AsyncInputStream, dst: var openarray[byte]): bool =
## Asynchronously read data into the destination buffer.
##
## Returns `false` if EOF was reached before the buffer
## was fully populated. if you need precise information
## regarding the number of bytes read, see `readIntoEx`.
##
## If there are enough bytes already buffered by the stream,
## the expression will complete immediately.
## Otherwise, it will await more bytes to become available.
template readInto*(sp: AsyncInputStream, dst: var openarray[byte]): bool =
## Asynchronously read data into the destination buffer.
##
## Returns `false` if EOF was reached before the buffer
## was fully populated. if you need precise information
## regarding the number of bytes read, see `readIntoEx`.
##
## If there are enough bytes already buffered by the stream,
## the expression will complete immediately.
## Otherwise, it will await more bytes to become available.
let s = InputStream(sp)
# BEWARE! `openArrayToPair` here is needed to avoid
# double evaluation of the `dst` expression:
let (dstAddr, dstLen) = openArrayToPair(dst)
readIntoExImpl(s, dstAddr, dstLen, fsAwait, readAsync) == dstLen
let s = InputStream(sp)
# BEWARE! `openArrayToPair` here is needed to avoid
# double evaluation of the `dst` expression:
let (dstAddr, dstLen) = openArrayToPair(dst)
readIntoExImpl(s, dstAddr, dstLen, fsAwait, readAsync) == dstLen
template useHeapMem(_: Natural) =
var buffer: seq[byte]
@ -869,8 +912,9 @@ template read*(sp: InputStream, np: static Natural): openarray[byte] =
template read*(s: InputStream, n: Natural): openarray[byte] =
readNImpl(s, n, useHeapMem)
template read*(s: AsyncInputStream, n: Natural): openarray[byte] =
read InputStream(s), n
when fsAsyncSupport:
template read*(s: AsyncInputStream, n: Natural): openarray[byte] =
read InputStream(s), n
proc lookAheadMatch*(s: InputStream, data: openarray[byte]): bool =
for i in 0 ..< data.len:
@ -879,23 +923,26 @@ proc lookAheadMatch*(s: InputStream, data: openarray[byte]): bool =
return true
template lookAheadMatch*(s: AsyncInputStream, data: openarray[byte]): bool =
lookAheadMatch InputStream(s)
when fsAsyncSupport:
template lookAheadMatch*(s: AsyncInputStream, data: openarray[byte]): bool =
lookAheadMatch InputStream(s)
proc next*(s: InputStream): Option[byte] =
if readable(s):
result = some read(s)
template next*(sp: AsyncInputStream): Option[byte] =
let s = sp
if readable(s):
some read(s)
else:
none byte
when fsAsyncSupport:
template next*(sp: AsyncInputStream): Option[byte] =
let s = sp
if readable(s):
some read(s)
else:
none byte
proc pos*(s: InputStream): int {.inline.} =
s.spanEndPos - s.span.len
template pos*(s: AsyncInputStream): int =
pos InputStream(s)
when fsAsyncSupport:
template pos*(s: AsyncInputStream): int =
pos InputStream(s)

View file

@ -1,35 +1,40 @@
import
stew/shims/macros,
async_backend, inputs, outputs
async_backend
export
async_backend
macro fsMultiSync*(body: untyped) =
# We will produce an identical copy of the annotated proc,
# but taking async parameters and having the async pragma.
var
asyncProcBody = copy body
asyncProcParams = asyncProcBody[3]
when fsAsyncSupport:
import
stew/shims/macros,
"."/[inputs, outputs]
asyncProcBody.addPragma(bindSym"async")
macro fsMultiSync*(body: untyped) =
# We will produce an identical copy of the annotated proc,
# but taking async parameters and having the async pragma.
var
asyncProcBody = copy body
asyncProcParams = asyncProcBody[3]
# The return types becomes Future[T]
if asyncProcParams[0].kind == nnkEmpty:
asyncProcParams[0] = newTree(nnkBracketExpr, bindSym"Future", ident"void")
else:
asyncProcParams[0] = newTree(nnkBracketExpr, bindSym"Future", asyncProcParams[0])
asyncProcBody.addPragma(bindSym"async")
# We replace all stream inputs with their async counterparts
for i in 1 ..< asyncProcParams.len:
let paramsDef = asyncProcParams[i]
let typ = paramsDef[^2]
if eqIdent(typ, "InputStream"):
paramsDef[^2] = bindSym "AsyncInputStream"
elif eqIdent(typ, "OutputStream"):
paramsDef[^2] = bindSym "AsyncOutputStream"
# The return types becomes Future[T]
if asyncProcParams[0].kind == nnkEmpty:
asyncProcParams[0] = newTree(nnkBracketExpr, bindSym"Future", ident"void")
else:
asyncProcParams[0] = newTree(nnkBracketExpr, bindSym"Future", asyncProcParams[0])
result = newStmtList(body, asyncProcBody)
when defined(debugSupportAsync):
echo result.repr
# We replace all stream inputs with their async counterparts
for i in 1 ..< asyncProcParams.len:
let paramsDef = asyncProcParams[i]
let typ = paramsDef[^2]
if eqIdent(typ, "InputStream"):
paramsDef[^2] = bindSym "AsyncInputStream"
elif eqIdent(typ, "OutputStream"):
paramsDef[^2] = bindSym "AsyncOutputStream"
result = newStmtList(body, asyncProcBody)
when defined(debugSupportAsync):
echo result.repr
else:
macro fsMultiSync*(body: untyped) = body

View file

@ -14,17 +14,78 @@ import
export
initPageBuffers, CloseBehavior
type
OutputStream* = ref object of RootObj
vtable*: ptr OutputStreamVTable # This is nil for any memory output
buffers*: PageBuffers # This is nil for unsafe memory outputs
span*: PageSpan
spanEndPos*: Natural
extCursorsCount: int
closeFut: Future[void] # This is nil before `close` is called
when debugHelpers:
name*: string
when fsAsyncSupport:
# Circular type refs prevent more targeted `when`
type
OutputStream* = ref object of RootObj
vtable*: ptr OutputStreamVTable # This is nil for any memory output
buffers*: PageBuffers # This is nil for unsafe memory outputs
span*: PageSpan
spanEndPos*: Natural
extCursorsCount: int
closeFut: Future[void] # This is nil before `close` is called
when debugHelpers:
name*: string
AsyncOutputStream* {.borrow: `.`.} = distinct OutputStream
WriteSyncProc* = proc (s: OutputStream, src: pointer, srcLen: Natural)
{.nimcall, gcsafe, raises: [IOError, Defect].}
WriteAsyncProc* = proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
FlushSyncProc* = proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
FlushAsyncProc* = proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseSyncProc* = proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseAsyncProc* = proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
OutputStreamVTable* = object
writeSync*: WriteSyncProc
writeAsync*: WriteAsyncProc
flushSync*: FlushSyncProc
flushAsync*: FlushAsyncProc
closeSync*: CloseSyncProc
closeAsync*: CloseAsyncProc
MaybeAsyncOutputStream* = OutputStream | AsyncOutputStream
else:
type
OutputStream* = ref object of RootObj
vtable*: ptr OutputStreamVTable # This is nil for any memory output
buffers*: PageBuffers # This is nil for unsafe memory outputs
span*: PageSpan
spanEndPos*: Natural
extCursorsCount: int
when debugHelpers:
name*: string
WriteSyncProc* = proc (s: OutputStream, src: pointer, srcLen: Natural)
{.nimcall, gcsafe, raises: [IOError, Defect].}
FlushSyncProc* = proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseSyncProc* = proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
OutputStreamVTable* = object
writeSync*: WriteSyncProc
flushSync*: FlushSyncProc
closeSync*: CloseSyncProc
MaybeAsyncOutputStream* = OutputStream
type
WriteCursor* = object
span: PageSpan
stream: OutputStream
@ -36,55 +97,35 @@ type
OutputStreamHandle* = object
s*: OutputStream
AsyncOutputStream* {.borrow: `.`.} = distinct OutputStream
WriteSyncProc* = proc (s: OutputStream, src: pointer, srcLen: Natural)
{.nimcall, gcsafe, raises: [IOError, Defect].}
WriteAsyncProc* = proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
FlushSyncProc* = proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
FlushAsyncProc* = proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseSyncProc* = proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].}
CloseAsyncProc* = proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].}
OutputStreamVTable* = object
writeSync*: WriteSyncProc
writeAsync*: WriteAsyncProc
flushSync*: FlushSyncProc
flushAsync*: FlushAsyncProc
closeSync*: CloseSyncProc
closeAsync*: CloseAsyncProc
VarSizeWriteCursor* = distinct WriteCursor
FileOutputStream = ref object of OutputStream
file: File
template Sync*(s: OutputStream): OutputStream = s
template Async*(s: OutputStream): AsyncOutputStream = AsyncOutputStream(s)
template Sync*(s: AsyncOutputStream): OutputStream = OutputStream(s)
template Async*(s: AsyncOutputStream): AsyncOutputStream = s
when fsAsyncSupport:
template Async*(s: OutputStream): AsyncOutputStream = AsyncOutputStream(s)
template Sync*(s: AsyncOutputStream): OutputStream = OutputStream(s)
template Async*(s: AsyncOutputStream): AsyncOutputStream = s
proc disconnectOutputDevice(s: OutputStream) =
if s.vtable != nil:
if s.vtable.closeAsync != nil:
s.closeFut = s.vtable.closeAsync(s)
elif s.vtable.closeSync != nil:
s.vtable.closeSync(s)
when fsAsyncSupport:
if s.vtable.closeAsync != nil:
s.closeFut = s.vtable.closeAsync(s)
elif s.vtable.closeSync != nil:
s.vtable.closeSync(s)
else:
if s.vtable.closeSync != nil:
s.vtable.closeSync(s)
s.vtable = nil
template disconnectOutputDevice(s: AsyncOutputStream) =
disconnectOutputDevice OutputStream(s)
when fsAsyncSupport:
template disconnectOutputDevice(s: AsyncOutputStream) =
disconnectOutputDevice OutputStream(s)
template flushImpl(s: OutputStream, awaiter, writeOp, flushOp: untyped) =
fsAssert s.extCursorsCount == 0
@ -99,36 +140,20 @@ template flushImpl(s: OutputStream, awaiter, writeOp, flushOp: untyped) =
proc flush*(s: OutputStream) =
flushImpl(s, noAwait, writeSync, flushSync)
template flush*(sp: AsyncOutputStream) =
let s = OutputStream sp
flushImpl(s, fsAwait, writeAsync, flushAsync)
proc flushAsync*(s: AsyncOutputStream) {.async.} =
flush s
proc close*(s: OutputStream,
behavior = dontWaitAsyncClose)
{.raises: [IOError, Defect].} =
flush s
disconnectOutputDevice(s)
if s.closeFut != nil:
fsTranslateErrors "Stream closing failed":
if behavior == waitAsyncClose:
waitFor s.closeFut
else:
asyncCheck s.closeFut
when fsAsyncSupport:
if s.closeFut != nil:
fsTranslateErrors "Stream closing failed":
if behavior == waitAsyncClose:
waitFor s.closeFut
else:
asyncCheck s.closeFut
template close*(sp: AsyncOutputStream) =
let s = OutputStream sp
flush(Async s)
disconnectOutputDevice(s)
if s.closeFut != nil:
fsAwait s.closeFut
proc closeAsync*(s: AsyncOutputStream) {.async.} =
close s
template closeNoWait*(sp: AsyncOutputStream|OutputStream) =
template closeNoWait*(sp: MaybeAsyncOutputStream) =
## Close the stream without waiting even if's async.
## This operation will use `asyncCheck` internally to detect unhandled
## errors from the closing operation.
@ -196,8 +221,9 @@ proc ensureRunway*(s: OutputStream, neededRunway: Natural) =
s.buffers.ensureRunway(s.span, neededRunway)
s.spanEndPos += (s.span.len - runway)
template ensureRunway*(s: AsyncOutputStream, neededRunway: Natural) =
ensureRunway OutputStream(s), neededRunway
when fsAsyncSupport:
template ensureRunway*(s: AsyncOutputStream, neededRunway: Natural) =
ensureRunway OutputStream(s), neededRunway
template implementWrites*(buffersParam: PageBuffers,
srcParam: pointer,
@ -266,8 +292,9 @@ proc fileOutput*(filename: string,
proc pos*(s: OutputStream): int =
s.spanEndPos - s.span.len
template pos*(s: AsyncOutputStream): int =
pos OutputStream(s)
when fsAsyncSupport:
template pos*(s: AsyncOutputStream): int =
pos OutputStream(s)
proc getBuffers*(s: OutputStream): PageBuffers =
fsAssert s.buffers != nil
@ -313,10 +340,11 @@ proc drainAllBuffersSync(s: OutputStream, buf: pointer, bufSize: Natural) =
s.span = s.buffers.getWritableSpan()
s.spanEndPos += s.span.len
proc drainAllBuffersAsync(s: OutputStream, buf: pointer, bufSize: Natural) {.async.} =
fsAwait s.vtable.writeAsync(s, buf, bufSize)
s.span = s.buffers.getWritableSpan()
s.spanEndPos += s.span.len
when fsAsyncSupport:
proc drainAllBuffersAsync(s: OutputStream, buf: pointer, bufSize: Natural) {.async.} =
fsAwait s.vtable.writeAsync(s, buf, bufSize)
s.span = s.buffers.getWritableSpan()
s.spanEndPos += s.span.len
proc createCursor(s: OutputStream, size: int): WriteCursor =
inc s.extCursorsCount
@ -533,21 +561,22 @@ template write*(sp: OutputStream, b: byte) =
else:
writeToNewSpan(s, b)
proc write*(sp: AsyncOutputStream, b: byte) =
let s = OutputStream sp
if atEnd(s.span):
addPage(s)
writeByte(s.span, b)
template writeAndWait*(sp: AsyncOutputStream, b: byte) =
let s = OutputStream sp
if hasRunway(s.span):
when fsAsyncSupport:
proc write*(sp: AsyncOutputStream, b: byte) =
let s = OutputStream sp
if atEnd(s.span):
addPage(s)
writeByte(s.span, b)
else:
writeToNewSpanImpl(s, b, fsAwait, writeAsync, drainAllBuffersAsync)
template write*(s: AsyncOutputStream, x: char) =
write s, byte(x)
template writeAndWait*(sp: AsyncOutputStream, b: byte) =
let s = OutputStream sp
if hasRunway(s.span):
writeByte(s.span, b)
else:
writeToNewSpanImpl(s, b, fsAwait, writeAsync, drainAllBuffersAsync)
template write*(s: AsyncOutputStream, x: char) =
write s, byte(x)
template write*(s: OutputStream|var WriteCursor, x: char) =
bind write
@ -606,7 +635,7 @@ proc write*(s: OutputStream, bytes: openArray[byte]) =
proc write*(s: OutputStream, chars: openArray[char]) =
write s, charsToBytes(chars)
proc write*(s: OutputStream|AsyncOutputStream, value: string) {.inline.} =
proc write*(s: MaybeAsyncOutputStream, value: string) {.inline.} =
write s, value.toOpenArrayByte(0, value.len - 1)
template memCopyToBytes(value: auto): untyped =
@ -618,37 +647,39 @@ template memCopyToBytes(value: auto): untyped =
proc writeMemCopy*(s: OutputStream, value: auto) =
write s, memCopyToBytes(value)
proc writeBytesAsyncImpl(sp: OutputStream,
bytes: openarray[byte]): Future[void] =
let s = sp
writeBytesImpl(s, bytes):
return s.vtable.writeAsync(s, unsafeAddr bytes[0], bytes.len)
when fsAsyncSupport:
proc writeBytesAsyncImpl(sp: OutputStream,
bytes: openarray[byte]): Future[void] =
let s = sp
writeBytesImpl(s, bytes):
return s.vtable.writeAsync(s, unsafeAddr bytes[0], bytes.len)
proc writeBytesAsyncImpl(s: OutputStream,
chars: openarray[char]): Future[void] =
writeBytesAsyncImpl s, charsToBytes(chars)
proc writeBytesAsyncImpl(s: OutputStream,
chars: openarray[char]): Future[void] =
writeBytesAsyncImpl s, charsToBytes(chars)
proc writeBytesAsyncImpl(s: OutputStream,
str: string): Future[void] =
writeBytesAsyncImpl s, toOpenArray(str, 0, str.len - 1)
proc writeBytesAsyncImpl(s: OutputStream,
str: string): Future[void] =
writeBytesAsyncImpl s, toOpenArray(str, 0, str.len - 1)
template writeAndWait*(s: OutputStream, value: untyped) =
write s, value
template writeAndWait*(sp: AsyncOutputStream, value: untyped) =
bind writeBytesAsyncImpl
when fsAsyncSupport:
template writeAndWait*(sp: AsyncOutputStream, value: untyped) =
bind writeBytesAsyncImpl
let
s = OutputStream sp
f = writeBytesAsyncImpl(s, value)
let
s = OutputStream sp
f = writeBytesAsyncImpl(s, value)
if f != nil:
fsAwait(f)
s.span = getWritableSpan s.buffers
s.spanEndPos += s.span.len
if f != nil:
fsAwait(f)
s.span = getWritableSpan s.buffers
s.spanEndPos += s.span.len
template writeMemCopyAndWait*(sp: AsyncOutputStream, value: auto) =
writeAndWait(sp, memCopyToBytes(value))
template writeMemCopyAndWait*(sp: AsyncOutputStream, value: auto) =
writeAndWait(sp, memCopyToBytes(value))
proc writeBytesToCursor(c: var WriteCursor, bytes: openarray[byte]) =
var
@ -795,9 +826,28 @@ template getOutput*(s: OutputStream, T: type seq[byte]): seq[byte] =
template getOutput*(s: OutputStream): seq[byte] =
cast[seq[byte]](s.getOutput(string))
template getOutput*(s: AsyncOutputStream): seq[byte] =
getOutput OutputStream(s)
when fsAsyncSupport:
template getOutput*(s: AsyncOutputStream): seq[byte] =
getOutput OutputStream(s)
template getOutput*(s: AsyncOutputStream, T: type): untyped =
getOutput OutputStream(s), T
template getOutput*(s: AsyncOutputStream, T: type): untyped =
getOutput OutputStream(s), T
when fsAsyncSupport:
template flush*(sp: AsyncOutputStream) =
let s = OutputStream sp
flushImpl(s, fsAwait, writeAsync, flushAsync)
proc flushAsync*(s: AsyncOutputStream) {.async.} =
flush s
template close*(sp: AsyncOutputStream) =
let s = OutputStream sp
flush(Async s)
disconnectOutputDevice(s)
if s.closeFut != nil:
fsAwait s.closeFut
proc closeAsync*(s: AsyncOutputStream) {.async.} =
close s

View file

@ -1,337 +1,341 @@
import
macros,
inputs, outputs, buffers, async_backend
"."/[inputs, outputs, async_backend]
export
inputs, outputs, async_backend
type
Pipe* = ref object
# TODO: Make these stream handles
input*: AsyncInputStream
output*: AsyncOutputStream
buffers*: PageBuffers
when fsAsyncSupport:
import
std/macros,
./buffers
template enterWait(fut: var Future, context: static string) =
let wait = newFuture[void](context)
fut = wait
try: fsAwait wait
finally: fut = nil
type
Pipe* = ref object
# TODO: Make these stream handles
input*: AsyncInputStream
output*: AsyncOutputStream
buffers*: PageBuffers
template awake(fp: Future) =
let f = fp
if f != nil and not finished(f):
complete f
template enterWait(fut: var Future, context: static string) =
let wait = newFuture[void](context)
fut = wait
try: fsAwait wait
finally: fut = nil
proc pipeRead(s: LayeredInputStream,
dst: pointer, dstLen: Natural): Future[Natural] {.async.} =
let buffers = s.buffers
if buffers.eofReached: return 0
template awake(fp: Future) =
let f = fp
if f != nil and not finished(f):
complete f
var
bytesInBuffersAtStart = buffers.totalBufferedBytes
minBytesExpected = max(1, dstLen)
bytesInBuffersNow = bytesInBuffersAtStart
proc pipeRead(s: LayeredInputStream,
dst: pointer, dstLen: Natural): Future[Natural] {.async.} =
let buffers = s.buffers
if buffers.eofReached: return 0
var
bytesInBuffersAtStart = buffers.totalBufferedBytes
minBytesExpected = max(1, dstLen)
bytesInBuffersNow = bytesInBuffersAtStart
while bytesInBuffersNow < minBytesExpected:
awake buffers.waitingWriter
buffers.waitingReader.enterWait "waiting for writer to buffer more data"
bytesInBuffersNow = buffers.totalBufferedBytes
if buffers.eofReached:
return bytesInBuffersNow - bytesInBuffersAtStart
if dst != nil:
let drained {.used.} = drainBuffersInto(s, cast[ptr byte](dst), dstLen)
fsAssert drained == dstLen
while bytesInBuffersNow < minBytesExpected:
awake buffers.waitingWriter
buffers.waitingReader.enterWait "waiting for writer to buffer more data"
bytesInBuffersNow = buffers.totalBufferedBytes
if buffers.eofReached:
return bytesInBuffersNow - bytesInBuffersAtStart
return bytesInBuffersNow - bytesInBuffersAtStart
if dst != nil:
let drained {.used.} = drainBuffersInto(s, cast[ptr byte](dst), dstLen)
fsAssert drained == dstLen
proc pipeWrite(s: LayeredOutputStream, src: pointer, srcLen: Natural) {.async.} =
let buffers = s.buffers
while buffers.canAcceptWrite(srcLen) == false:
buffers.waitingWriter.enterWait "waiting for reader to drain the buffers"
awake buffers.waitingWriter
if src != nil:
buffers.appendUnbufferedWrite(src, srcLen)
return bytesInBuffersNow - bytesInBuffersAtStart
awake buffers.waitingReader
describeBuffers "pipeWrite", buffers
proc pipeWrite(s: LayeredOutputStream, src: pointer, srcLen: Natural) {.async.} =
let buffers = s.buffers
while buffers.canAcceptWrite(srcLen) == false:
buffers.waitingWriter.enterWait "waiting for reader to drain the buffers"
template completedFuture(name: static string): untyped =
let fut = newFuture[void](name)
complete fut
fut
if src != nil:
buffers.appendUnbufferedWrite(src, srcLen)
awake buffers.waitingReader
describeBuffers "pipeWrite", buffers
template completedFuture(name: static string): untyped =
let fut = newFuture[void](name)
complete fut
fut
let pipeInputVTable = InputStreamVTable(
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to read from pipe":
let ls = LayeredInputStream(s)
fsAssert ls.allowWaitFor
return waitFor pipeRead(ls, dst, dstLen)
,
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
let pipeInputVTable = InputStreamVTable(
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Unexpected error from the async macro":
let ls = LayeredInputStream(s)
return pipeRead(ls, dst, dstLen)
,
getLenSync: proc (s: InputStream): Option[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
let source = LayeredInputStream(s).source
if source != nil:
return source.len
,
closeSync: proc (s: InputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
let source = LayeredInputStream(s).source
if source != nil:
close source
,
closeAsync: proc (s: InputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Unexpected error from the async macro":
fsTranslateErrors "Failed to read from pipe":
let ls = LayeredInputStream(s)
fsAssert ls.allowWaitFor
return waitFor pipeRead(ls, dst, dstLen)
,
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Unexpected error from the async macro":
let ls = LayeredInputStream(s)
return pipeRead(ls, dst, dstLen)
,
getLenSync: proc (s: InputStream): Option[Natural]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
let source = LayeredInputStream(s).source
if source != nil:
return closeAsync(Async source)
else:
return completedFuture("pipeInput.closeAsync")
)
return source.len
,
closeSync: proc (s: InputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
let source = LayeredInputStream(s).source
if source != nil:
close source
,
closeAsync: proc (s: InputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Unexpected error from the async macro":
let source = LayeredInputStream(s).source
if source != nil:
return closeAsync(Async source)
else:
return completedFuture("pipeInput.closeAsync")
)
let pipeOutputVTable = OutputStreamVTable(
writeSync: proc (s: OutputStream, src: pointer, srcLen: Natural)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to write all bytes to pipe":
var ls = LayeredOutputStream(s)
fsAssert ls.allowWaitFor
waitFor pipeWrite(ls, src, srcLen)
,
writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
# TODO: The async macro is raising exceptions even when
# merely forwarding a future:
fsTranslateErrors "Unexpected error from the async macro":
return pipeWrite(LayeredOutputStream s, src, srcLen)
,
flushSync: proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
let destination = LayeredOutputStream(s).destination
if destination != nil:
flush destination
,
flushAsync: proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Unexpected error from the async macro":
let pipeOutputVTable = OutputStreamVTable(
writeSync: proc (s: OutputStream, src: pointer, srcLen: Natural)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Failed to write all bytes to pipe":
var ls = LayeredOutputStream(s)
fsAssert ls.allowWaitFor
waitFor pipeWrite(ls, src, srcLen)
,
writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
# TODO: The async macro is raising exceptions even when
# merely forwarding a future:
fsTranslateErrors "Unexpected error from the async macro":
return pipeWrite(LayeredOutputStream s, src, srcLen)
,
flushSync: proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
let destination = LayeredOutputStream(s).destination
if destination != nil:
return flushAsync(Async destination)
else:
return completedFuture("pipeOutput.flushAsync")
,
closeSync: proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
flush destination
,
flushAsync: proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
fsTranslateErrors "Unexpected error from the async macro":
let destination = LayeredOutputStream(s).destination
if destination != nil:
return flushAsync(Async destination)
else:
return completedFuture("pipeOutput.flushAsync")
,
closeSync: proc (s: OutputStream)
{.nimcall, gcsafe, raises: [IOError, Defect].} =
s.buffers.eofReached = true
s.buffers.eofReached = true
fsTranslateErrors "Unexpected error from Future.complete":
awake s.buffers.waitingReader
fsTranslateErrors "Unexpected error from Future.complete":
awake s.buffers.waitingReader
let destination = LayeredOutputStream(s).destination
if destination != nil:
close destination
,
closeAsync: proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
s.buffers.eofReached = true
fsTranslateErrors "Unexpected error from Future.complete":
awake s.buffers.waitingReader
fsTranslateErrors "Unexpected error from the async macro":
let destination = LayeredOutputStream(s).destination
if destination != nil:
return closeAsync(Async destination)
else:
return completedFuture("pipeOutput.closeAsync")
)
close destination
,
closeAsync: proc (s: OutputStream): Future[void]
{.nimcall, gcsafe, raises: [IOError, Defect].} =
s.buffers.eofReached = true
func pipeInput*(source: InputStream,
pageSize = defaultPageSize,
allowWaitFor = false): AsyncInputStream =
fsAssert pageSize > 0
fsTranslateErrors "Unexpected error from Future.complete":
awake s.buffers.waitingReader
AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable,
buffers: initPageBuffers pageSize,
allowWaitFor: allowWaitFor,
source: source)
fsTranslateErrors "Unexpected error from the async macro":
let destination = LayeredOutputStream(s).destination
if destination != nil:
return closeAsync(Async destination)
else:
return completedFuture("pipeOutput.closeAsync")
)
func pipeInput*(buffers: PageBuffers,
allowWaitFor = false,
source: InputStream = nil): AsyncInputStream =
var spanEndPos = Natural 0
var span = if buffers.len == 0: default(PageSpan)
else: buffers.obtainReadableSpan(spanEndPos)
func pipeInput*(source: InputStream,
pageSize = defaultPageSize,
allowWaitFor = false): AsyncInputStream =
fsAssert pageSize > 0
AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable,
buffers: buffers,
span: span,
spanEndPos: span.len,
allowWaitFor: allowWaitFor,
source: source)
AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable,
buffers: initPageBuffers pageSize,
allowWaitFor: allowWaitFor,
source: source)
proc pipeOutput*(destination: OutputStream,
pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4,
allowWaitFor = false): AsyncOutputStream =
fsAssert pageSize > 0
func pipeInput*(buffers: PageBuffers,
allowWaitFor = false,
source: InputStream = nil): AsyncInputStream =
var spanEndPos = Natural 0
var span = if buffers.len == 0: default(PageSpan)
else: buffers.obtainReadableSpan(spanEndPos)
var
buffers = initPageBuffers pageSize
span = buffers.getWritableSpan()
AsyncInputStream LayeredInputStream(
vtable: vtableAddr pipeInputVTable,
buffers: buffers,
span: span,
spanEndPos: span.len,
allowWaitFor: allowWaitFor,
source: source)
AsyncOutputStream LayeredOutputStream(
vtable: vtableAddr pipeOutputVTable,
buffers: buffers,
span: span,
spanEndPos: span.len,
allowWaitFor: allowWaitFor,
destination: destination)
proc pipeOutput*(destination: OutputStream,
pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4,
allowWaitFor = false): AsyncOutputStream =
fsAssert pageSize > 0
proc pipeOutput*(buffers: PageBuffers,
allowWaitFor = false,
destination: OutputStream = nil): AsyncOutputStream =
var span = buffers.getWritableSpan()
AsyncOutputStream LayeredOutputStream(
vtable: vtableAddr pipeOutputVTable,
buffers: buffers,
span: span,
# TODO What if the buffers are partially populated?
# Should we adjust the spanEndPos? This would
# need the old buffers.totalBytesWritten var.
spanEndPos: span.len,
allowWaitFor: allowWaitFor,
destination: destination)
func asyncPipe*(pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4): Pipe =
fsAssert pageSize > 0
Pipe(buffers: initPageBuffers(pageSize, maxBufferedBytes))
func initReader*(pipe: Pipe): AsyncInputStream =
result = pipeInput(pipe.buffers)
pipe.input = result
func initWriter*(pipe: Pipe): AsyncOutputStream =
result = pipeOutput(pipe.buffers)
pipe.output = result
proc exchangeBuffersAfterPipilineStep(input: InputStream, output: OutputStream) =
let formerInputBuffers = input.buffers
let formerOutputBuffers = output.getBuffers
input.resetBuffers formerOutputBuffers
output.recycleBuffers formerInputBuffers
macro executePipeline*(start: InputStream, steps: varargs[untyped]): untyped =
result = newTree(nnkStmtListExpr)
var
inputVal = start
outputVal = newCall(bindSym"memoryOutput")
inputVar = genSym(nskVar, "input")
outputVar = genSym(nskVar, "output")
step0 = steps[0]
result.add quote do:
var
`inputVar` = `inputVal`
`outputVar` = OutputStream `outputVal`
buffers = initPageBuffers pageSize
span = buffers.getWritableSpan()
`step0`(`inputVar`, `outputVar`)
AsyncOutputStream LayeredOutputStream(
vtable: vtableAddr pipeOutputVTable,
buffers: buffers,
span: span,
spanEndPos: span.len,
allowWaitFor: allowWaitFor,
destination: destination)
proc pipeOutput*(buffers: PageBuffers,
allowWaitFor = false,
destination: OutputStream = nil): AsyncOutputStream =
var span = buffers.getWritableSpan()
AsyncOutputStream LayeredOutputStream(
vtable: vtableAddr pipeOutputVTable,
buffers: buffers,
span: span,
# TODO What if the buffers are partially populated?
# Should we adjust the spanEndPos? This would
# need the old buffers.totalBytesWritten var.
spanEndPos: span.len,
allowWaitFor: allowWaitFor,
destination: destination)
func asyncPipe*(pageSize = defaultPageSize,
maxBufferedBytes = defaultPageSize * 4): Pipe =
fsAssert pageSize > 0
Pipe(buffers: initPageBuffers(pageSize, maxBufferedBytes))
func initReader*(pipe: Pipe): AsyncInputStream =
result = pipeInput(pipe.buffers)
pipe.input = result
func initWriter*(pipe: Pipe): AsyncOutputStream =
result = pipeOutput(pipe.buffers)
pipe.output = result
proc exchangeBuffersAfterPipilineStep(input: InputStream, output: OutputStream) =
let formerInputBuffers = input.buffers
let formerOutputBuffers = output.getBuffers
input.resetBuffers formerOutputBuffers
output.recycleBuffers formerInputBuffers
macro executePipeline*(start: InputStream, steps: varargs[untyped]): untyped =
result = newTree(nnkStmtListExpr)
var
inputVal = start
outputVal = newCall(bindSym"memoryOutput")
inputVar = genSym(nskVar, "input")
outputVar = genSym(nskVar, "output")
step0 = steps[0]
if steps.len > 2:
let step1 = steps[1]
result.add quote do:
let formerInputBuffers = `inputVar`.buffers
`inputVar` = memoryInput(getBuffers `outputVar`)
recycleBuffers(`outputVar`, formerInputBuffers)
`step1`(`inputVar`, `outputVar`)
var
`inputVar` = `inputVal`
`outputVar` = OutputStream `outputVal`
for i in 2 .. steps.len - 2:
let step = steps[i]
result.add quote do:
exchangeBuffersAfterPipilineStep(`inputVar`, `outputVar`)
`step`(`inputVar`, `outputVar`)
`step0`(`inputVar`, `outputVar`)
var closingCall = steps[^1]
closingCall.insert(1, outputVar)
result.add closingCall
if steps.len > 2:
let step1 = steps[1]
result.add quote do:
let formerInputBuffers = `inputVar`.buffers
`inputVar` = memoryInput(getBuffers `outputVar`)
recycleBuffers(`outputVar`, formerInputBuffers)
`step1`(`inputVar`, `outputVar`)
if defined(debugMacros) or defined(debugPipelines):
echo result.repr
for i in 2 .. steps.len - 2:
let step = steps[i]
result.add quote do:
exchangeBuffersAfterPipilineStep(`inputVar`, `outputVar`)
`step`(`inputVar`, `outputVar`)
macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untyped =
var
stream = ident "stream"
pipelineSteps = ident "pipelineSteps"
pipelineBody = newTree(nnkStmtList)
var closingCall = steps[^1]
closingCall.insert(1, outputVar)
result.add closingCall
step0 = steps[0]
stepOutput = genSym(nskVar, "pipe")
if defined(debugMacros) or defined(debugPipelines):
echo result.repr
pipelineBody.add quote do:
var `pipelineSteps` = newSeq[Future[void]]()
var `stepOutput` = asyncPipe()
add `pipelineSteps`, `step0`(`stream`, initWriter(`stepOutput`))
macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untyped =
var
stream = ident "stream"
pipelineSteps = ident "pipelineSteps"
pipelineBody = newTree(nnkStmtList)
var
stepInput = stepOutput
for i in 1 .. steps.len - 2:
var step = steps[i]
stepOutput = genSym(nskVar, "pipe")
step0 = steps[0]
stepOutput = genSym(nskVar, "pipe")
pipelineBody.add quote do:
var `pipelineSteps` = newSeq[Future[void]]()
var `stepOutput` = asyncPipe()
add `pipelineSteps`, `step`(initReader(`stepInput`), initWriter(`stepOutput`))
add `pipelineSteps`, `step0`(`stream`, initWriter(`stepOutput`))
stepInput = stepOutput
var
stepInput = stepOutput
var RetTypeExpr = copy steps[^1]
RetTypeExpr.insert(1, newCall("default", ident"AsyncInputStream"))
for i in 1 .. steps.len - 2:
var step = steps[i]
stepOutput = genSym(nskVar, "pipe")
var closingCall = steps[^1]
closingCall.insert(1, newCall(bindSym"initReader", stepInput))
pipelineBody.add quote do:
var `stepOutput` = asyncPipe()
add `pipelineSteps`, `step`(initReader(`stepInput`), initWriter(`stepOutput`))
pipelineBody.add quote do:
fsAwait allFutures(`pipelineSteps`)
`closingCall`
stepInput = stepOutput
result = quote do:
type UserOpRetType = type(`RetTypeExpr`)
var RetTypeExpr = copy steps[^1]
RetTypeExpr.insert(1, newCall("default", ident"AsyncInputStream"))
when UserOpRetType is Future:
type RetType = type(default(UserOpRetType).read)
else:
type RetType = UserOpRetType
var closingCall = steps[^1]
closingCall.insert(1, newCall(bindSym"initReader", stepInput))
pipelineBody.add quote do:
fsAwait allFutures(`pipelineSteps`)
`closingCall`
result = quote do:
type UserOpRetType = type(`RetTypeExpr`)
proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} =
when UserOpRetType is Future:
var f = `pipelineBody`
return fsAwait(f)
type RetType = type(default(UserOpRetType).read)
else:
return `pipelineBody`
type RetType = UserOpRetType
pipelineProc(`start`)
proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} =
when UserOpRetType is Future:
var f = `pipelineBody`
return fsAwait(f)
else:
return `pipelineBody`
when defined(debugMacros):
echo result.repr
pipelineProc(`start`)
when defined(debugMacros):
echo result.repr