From 3a0ab42573e566ce52625760f6bbf7e0bbb6ebc4 Mon Sep 17 00:00:00 2001 From: Jacek Sieka Date: Fri, 20 Aug 2021 13:10:26 +0200 Subject: [PATCH] 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 --- faststreams.nimble | 9 +- faststreams/async_backend.nim | 23 +- faststreams/buffers.nim | 5 +- faststreams/inputs.nim | 381 +++++++++++++---------- faststreams/multisync.nim | 55 ++-- faststreams/outputs.nim | 288 ++++++++++-------- faststreams/pipelines.nim | 554 +++++++++++++++++----------------- tests/test_inputs.nim | 8 +- tests/test_outputs.nim | 9 + tests/test_pipelines.nim | 214 ++++++------- 10 files changed, 841 insertions(+), 705 deletions(-) diff --git a/faststreams.nimble b/faststreams.nimble index e088d5d..3770107 100644 --- a/faststreams.nimble +++ b/faststreams.nimble @@ -1,7 +1,7 @@ mode = ScriptMode.Verbose packageName = "faststreams" -version = "0.2.0" +version = "0.3.0" author = "Status Research & Development GmbH" description = "Nearly zero-overhead input/output streams for Nim" license = "Apache License 2.0" @@ -21,7 +21,12 @@ proc test(env, path: string) = lang = getEnv"TEST_LANG" exec "nim " & lang & " " & env & - " -r --hints:off --skipParentCfg " & path + " -d:async_backend=none -r --hints:off --skipParentCfg " & path + exec "nim " & lang & " " & env & + " -d:async_backend=chronos -r --hints:off --skipParentCfg " & path + # TODO std backend is broken / untested + # exec "nim " & lang & " " & env & + # " -d:async_backend=chronos -r --hints:off --skipParentCfg " & path task test, "Run all tests": test "-d:debug --threads:on", "tests/all_tests" diff --git a/faststreams/async_backend.nim b/faststreams/async_backend.nim index e2256e5..1a385f6 100644 --- a/faststreams/async_backend.nim +++ b/faststreams/async_backend.nim @@ -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 diff --git a/faststreams/buffers.nim b/faststreams/buffers.nim index 75111fb..99b49f2 100644 --- a/faststreams/buffers.nim +++ b/faststreams/buffers.nim @@ -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 diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index ebea223..87242b5 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -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) diff --git a/faststreams/multisync.nim b/faststreams/multisync.nim index 5f5e644..ebfbfce 100644 --- a/faststreams/multisync.nim +++ b/faststreams/multisync.nim @@ -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 diff --git a/faststreams/outputs.nim b/faststreams/outputs.nim index 758e33a..d48fce8 100644 --- a/faststreams/outputs.nim +++ b/faststreams/outputs.nim @@ -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 diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim index cf1f510..82b3bd5 100644 --- a/faststreams/pipelines.nim +++ b/faststreams/pipelines.nim @@ -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 diff --git a/tests/test_inputs.nim b/tests/test_inputs.nim index 95ef39e..ee1a103 100644 --- a/tests/test_inputs.nim +++ b/tests/test_inputs.nim @@ -2,15 +2,11 @@ import os, unittest2, strutils, random, - stew/ranges/ptr_arith, testutils, + testutils, ../faststreams, ../faststreams/textio setCurrentDir getAppDir() -proc bytes(s: string): seq[byte] = - result = newSeqOfCap[byte](s.len) - for c in s: result.add byte(c) - proc str(bytes: openarray[byte]): string = result = newStringOfCap(bytes.len) for b in items(bytes): @@ -28,7 +24,7 @@ const asciiTableFile = "files" / "ascii_table.txt" asciiTableContents = slurp(asciiTableFile) -procSuite "input stream": +suite "input stream": template emptyInputTests(suiteName, setupCode: untyped) = suite suiteName & " empty inputs": setup setupCode diff --git a/tests/test_outputs.nim b/tests/test_outputs.nim index 82ce17a..a2fa06a 100644 --- a/tests/test_outputs.nim +++ b/tests/test_outputs.nim @@ -154,6 +154,15 @@ suite "output stream": checkOutputsMatch() + test "memcpy": + var x = 0x42'u8 + + nimSeq.add x + memStream.writeMemCopy x + let memStreamRes = memStream.getOutput + + check memStreamRes == nimSeq + template undelayedOutput(content: seq[byte]) {.dirty.} = nimSeq.add content streamWritingToExistingBuffer.write content diff --git a/tests/test_pipelines.nim b/tests/test_pipelines.nim index f597540..de4861f 100644 --- a/tests/test_pipelines.nim +++ b/tests/test_pipelines.nim @@ -1,128 +1,136 @@ {.used.} import - # Std lib: - std/[strutils, random, base64, terminal], - # Other packages: testutils/unittests, + # FastStreams modules: - ../faststreams/[pipelines, multisync], - # Testing modules: - ./base64 as fsBase64 + ../faststreams/[pipelines, multisync] -include system/timers +when fsAsyncSupport: + import + # Std lib: + std/[strutils, random, base64, terminal], + # FastStreams modules: + ../faststreams/[pipelines, multisync], + # Testing modules: + ./base64 as fsBase64 -type - TestTimes = object - fsPipeline: Nanos - fsAsyncPipeline: Nanos - stdFunctionCalls: Nanos + include system/timers -proc upcaseAllCharacters(i: InputStream, o: OutputStream) {.fsMultiSync.} = - let inputLen = i.len - if inputLen.isSome: - o.ensureRunway inputLen.get + type + TestTimes = object + fsPipeline: Nanos + fsAsyncPipeline: Nanos + stdFunctionCalls: Nanos - while i.readable: - o.write toUpperAscii(i.read.char) + proc upcaseAllCharacters(i: InputStream, o: OutputStream) {.fsMultiSync.} = + let inputLen = i.len + if inputLen.isSome: + o.ensureRunway inputLen.get - close o + while i.readable: + o.write toUpperAscii(i.read.char) -proc printTimes(t: TestTimes) = - styledEcho " cpu time [FS Sync ]: ", styleBright, $t.fsPipeline, "ms" - styledEcho " cpu time [FS Async ]: ", styleBright, $t.fsAsyncPipeline, "ms" - styledEcho " cpu time [Std Lib ]: ", styleBright, $t.stdFunctionCalls, "ms" + close o -template timeit(timerVar: var Nanos, code: untyped) = - let t0 = getTicks() - code - timerVar = int(getTicks() - t0) div 1000000 + proc printTimes(t: TestTimes) = + styledEcho " cpu time [FS Sync ]: ", styleBright, $t.fsPipeline, "ms" + styledEcho " cpu time [FS Async ]: ", styleBright, $t.fsAsyncPipeline, "ms" + styledEcho " cpu time [Std Lib ]: ", styleBright, $t.stdFunctionCalls, "ms" -proc getOutput(sp: AsyncInputStream, T: type string): Future[string] {.async.} = - # this proc is a quick hack to let the test pass - # do not use it in production code - let size = sp.totalUnconsumedBytes() - if size > 0: - var data = newSeq[byte](size) - discard sp.readinto(data) - result = cast[string](data) + template timeit(timerVar: var Nanos, code: untyped) = + let t0 = getTicks() + code + timerVar = int(getTicks() - t0) div 1000000 -procSuite "pipelines": - let loremIpsum = """ - Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod - tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim - veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex - ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate - velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat - cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id - est laborum. + proc getOutput(sp: AsyncInputStream, T: type string): Future[string] {.async.} = + # this proc is a quick hack to let the test pass + # do not use it in production code + let size = sp.totalUnconsumedBytes() + if size > 0: + var data = newSeq[byte](size) + discard sp.readinto(data) + result = cast[string](data) - """ + suite "pipelines": + const loremIpsum = """ + Lorem ipsum dolor sit amet, consectetur adipiscing elit, sed do eiusmod + tempor incididunt ut labore et dolore magna aliqua. Ut enim ad minim + veniam, quis nostrud exercitation ullamco laboris nisi ut aliquip ex + ea commodo consequat. Duis aute irure dolor in reprehenderit in voluptate + velit esse cillum dolore eu fugiat nulla pariatur. Excepteur sint occaecat + cupidatat non proident, sunt in culpa qui officia deserunt mollit anim id + est laborum. - test "upper-case/base64 pipeline benchmark": - var - times: TestTimes - stdRes: string - fsRes: string - fsAsyncRes: string + """ - let inputText = loremIpsum.repeat(5000) + test "upper-case/base64 pipeline benchmark": + var + times: TestTimes + stdRes: string + fsRes: string + fsAsyncRes: string - timeIt times.stdFunctionCalls: - stdRes = base64.decode(base64.encode(toUpperAscii(inputText))) + let inputText = loremIpsum.repeat(5000) - timeIt times.fsPipeline: - fsRes = executePipeline(unsafeMemoryInput(inputText), + timeIt times.stdFunctionCalls: + stdRes = base64.decode(base64.encode(toUpperAscii(inputText))) + + timeIt times.fsPipeline: + fsRes = executePipeline(unsafeMemoryInput(inputText), + upcaseAllCharacters, + base64encode, + base64decode, + getOutput string) + + timeIt times.fsAsyncPipeline: + fsAsyncRes = waitFor executePipeline(Async unsafeMemoryInput(inputText), + upcaseAllCharacters, + base64encode, + base64decode, + getOutput string) + + check fsAsyncRes == stdRes + check fsRes == stdRes + + printTimes times + + asyncTest "upper-case/base64 async pipeline": + let pipe = asyncPipe() + let inputText = repeat(loremIpsum, 100) + + proc pipeFeeder(s: AsyncOutputStream) {.gcsafe, async.} = + randomize 1234 + var pos = 0 + + while pos != inputText.len: + let bytesToWrite = rand(15) + + if bytesToWrite == 0: + s.write inputText[pos] + inc pos + else: + let endPos = min(pos + bytesToWrite, inputText.len) + s.writeAndWait inputText[pos ..< endPos] + pos = endPos + + let sleep = rand(50) - 45 + if sleep > 0: + await sleepAsync(sleep.milliseconds) + + close s + + asyncCheck pipeFeeder(pipe.initWriter) + + let f = executePipeline(pipe.initReader, upcaseAllCharacters, base64encode, base64decode, getOutput string) - timeIt times.fsAsyncPipeline: - fsAsyncRes = waitFor executePipeline(Async unsafeMemoryInput(inputText), - upcaseAllCharacters, - base64encode, - base64decode, - getOutput string) + let fsAsyncres = await f - check fsAsyncRes == stdRes - check fsRes == stdRes - - printTimes times - - asyncTest "upper-case/base64 async pipeline": - let pipe = asyncPipe() - let inputText = repeat(loremIpsum, 100) - - proc pipeFeeder(s: AsyncOutputStream) {.gcsafe, async.} = - randomize 1234 - var pos = 0 - - while pos != inputText.len: - let bytesToWrite = rand(15) - - if bytesToWrite == 0: - s.write inputText[pos] - inc pos - else: - let endPos = min(pos + bytesToWrite, inputText.len) - s.writeAndWait inputText[pos ..< endPos] - pos = endPos - - let sleep = rand(50) - 45 - if sleep > 0: - await sleepAsync(sleep.milliseconds) - - close s - - asyncCheck pipeFeeder(pipe.initWriter) - - let f = executePipeline(pipe.initReader, - upcaseAllCharacters, - base64encode, - base64decode, - getOutput string) - - let fsAsyncres = await f - - check fsAsyncRes == toUpperAscii(inputText) + check fsAsyncRes == toUpperAscii(inputText) +else: + test "pipelines": + skip