From 5df69fc6961e58205189cd92ae2477769fa8c4c0 Mon Sep 17 00:00:00 2001 From: Zahary Karadjov Date: Mon, 8 Jun 2020 22:37:22 +0300 Subject: [PATCH] Add asynctools adapter --- faststreams/async_backend.nim | 8 ++- faststreams/asynctools_adapters.nim | 103 ++++++++++++++++++++++++++++ faststreams/chronos_adapters.nim | 35 +++++----- faststreams/inputs.nim | 10 +-- faststreams/outputs.nim | 4 +- faststreams/pipelines.nim | 6 +- faststreams/textio.nim | 2 +- 7 files changed, 137 insertions(+), 31 deletions(-) create mode 100644 faststreams/asynctools_adapters.nim diff --git a/faststreams/async_backend.nim b/faststreams/async_backend.nim index 784544c..e2256e5 100644 --- a/faststreams/async_backend.nim +++ b/faststreams/async_backend.nim @@ -20,12 +20,12 @@ when faststreams_async_backend == "chronos": elif faststreams_async_backend in ["std", "asyncdispatch"]: import - std/[asyncfutures, asyncmacro] + std/asyncdispatch export - asyncfutures, asyncmacro + asyncdispatch - template fsAwait*(awaited: Future[T]): untyped = + template fsAwait*(awaited: Future): untyped = # TODO revisit after https://github.com/nim-lang/Nim/pull/12085/ is merged let f = awaited yield f @@ -33,6 +33,8 @@ elif faststreams_async_backend in ["std", "asyncdispatch"]: raise f.error f.read + type Duration* = int + else: {.fatal: "Unrecognized network backend: " & faststreams_async_backend.} diff --git a/faststreams/asynctools_adapters.nim b/faststreams/asynctools_adapters.nim new file mode 100644 index 0000000..e40a4fe --- /dev/null +++ b/faststreams/asynctools_adapters.nim @@ -0,0 +1,103 @@ +import + asynctools/asyncpipe, + inputs, outputs, buffers, multisync + +export + inputs, outputs, asyncpipe, fsMultiSync + +type + AsyncPipeInput* = ref object of InputStream + pipe: AsyncPipe + allowWaitFor: bool + + AsyncPipeOutput* = ref object of OutputStream + pipe: AsyncPipe + allowWaitFor: bool + +const + readingErrMsg = "Failed to read from AsyncPipe" + writingErrMsg = "Failed to write to AsyncPipe" + closingErrMsg = "Failed to close AsyncPipe" + writeIncompleteErrMsg = "Failed to write all bytes to AsyncPipe" + +proc closeAsyncPipe(pipe: AsyncPipe) + {.raises: [Defect, IOError].} = + fsTranslateErrors closingErrMsg: + close pipe + +proc readOnce(s: AsyncPipeInput, + dst: pointer, dstLen: Natural): Future[Natural] {.async.} = + fsTranslateErrors readingErrMsg: + return implementSingleRead(s.buffers, dst, dstLen, ReadFlags {}, + readStartAddr, readLen): + await s.pipe.readInto(readStartAddr, readLen) + +proc write(s: AsyncPipeOutput, src: pointer, srcLen: Natural) {.async.} = + fsTranslateErrors writeIncompleteErrMsg: + implementWrites(s.buffers, src, srcLen, "AsyncPipe", + writeStartAddr, writeLen): + await s.pipe.write(writeStartAddr, writeLen) + +# TODO: Use the Raising type here +let asyncPipeInputVTable = InputStreamVTable( + readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural + {.nimcall, gcsafe, raises: [IOError, Defect].} = + fsTranslateErrors "Unexpected exception from asyncdispatch": + var cs = AsyncPipeInput(s) + fsAssert cs.allowWaitFor + return waitFor readOnce(cs, dst, dstLen) + , + readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural] + {.nimcall, gcsafe, raises: [IOError, Defect].} = + fsTranslateErrors "Unexpected exception from merely forwarding a future": + return readOnce(AsyncPipeInput s, dst, dstLen) + , + closeSync: proc (s: InputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + closeAsyncPipe AsyncPipeInput(s).pipe + , + closeAsync: proc (s: InputStream): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} = + closeAsyncPipe AsyncPipeInput(s).pipe +) + +func asyncPipeInput*(pipe: AsyncPipe, + pageSize = defaultPageSize, + allowWaitFor = false): AsyncInputStream = + AsyncInputStream AsyncPipeInput( + vtable: vtableAddr asyncPipeInputVTable, + buffers: initPageBuffers(pageSize), + pipe: pipe, + allowWaitFor: allowWaitFor) + +let asyncPipeOutputVTable = OutputStreamVTable( + writeSync: proc (s: OutputStream, src: pointer, srcLen: Natural) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + fsTranslateErrors "Unexpected exception from asyncdispatch": + var cs = AsyncPipeOutput(s) + fsAssert cs.allowWaitFor + waitFor write(cs, src, srcLen) + , + writeAsync: proc (s: OutputStream, src: pointer, srcLen: Natural): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} = + fsTranslateErrors "Unexpected exception from merely forwarding a future": + return write(AsyncPipeOutput s, src, srcLen) + , + closeSync: proc (s: OutputStream) + {.nimcall, gcsafe, raises: [IOError, Defect].} = + closeAsyncPipe AsyncPipeOutput(s).pipe + , + closeAsync: proc (s: OutputStream): Future[void] + {.nimcall, gcsafe, raises: [IOError, Defect].} = + closeAsyncPipe AsyncPipeOutput(s).pipe +) + +func asyncPipeOutput*(pipe: AsyncPipe, + pageSize = defaultPageSize, + allowWaitFor = false): AsyncOutputStream = + AsyncOutputStream AsyncPipeOutput( + vtable: vtableAddr(asyncPipeOutputVTable), + buffers: initPageBuffers(pageSize), + pipe: pipe, + allowWaitFor: allowWaitFor) + diff --git a/faststreams/chronos_adapters.nim b/faststreams/chronos_adapters.nim index 4dcc503..eb0dc46 100644 --- a/faststreams/chronos_adapters.nim +++ b/faststreams/chronos_adapters.nim @@ -26,17 +26,15 @@ proc chronosCloseWait(t: StreamTransport) await t.closeWait() proc chronosReadOnce(s: ChronosInputStream, - dst: pointer, dstLen: Natural): Future[Natural] - {.async, raises: [IOError, Defect].} = + dst: pointer, dstLen: Natural): Future[Natural] {.async.} = fsTranslateErrors readingErrMsg: - return implementSingleRead(s.buffers, dst, dstLen, {}, + return implementSingleRead(s.buffers, dst, dstLen, ReadFlags {}, readStartAddr, readLen): await s.transport.readOnce(readStartAddr, readLen) -proc chronosWrites(s: ChronosOutputStream, src: pointer, srcLen: Natural) - {.async, raises: [IOError, Defect].} = +proc chronosWrites(s: ChronosOutputStream, src: pointer, srcLen: Natural) {.async.} = fsTranslateErrors writeIncompleteErrMsg: - implementWrites(s.buffers, src, srcLen, "StreamTransport" + implementWrites(s.buffers, src, srcLen, "StreamTransport", writeStartAddr, writeLen): await s.transport.write(writeStartAddr, writeLen) @@ -44,13 +42,15 @@ proc chronosWrites(s: ChronosOutputStream, src: pointer, srcLen: Natural) let chronosInputVTable = InputStreamVTable( readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural {.nimcall, gcsafe, raises: [IOError, Defect].} = - var cs = ChronosInputStream(s) - fsAssert cs.allowWaitFor - waitFor chronosReadOnce(cs, dst, dstLen) + fsTranslateErrors "Unexpected exception from Chronos async macro": + var cs = ChronosInputStream(s) + fsAssert cs.allowWaitFor + return waitFor chronosReadOnce(cs, dst, dstLen) , readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural] {.nimcall, gcsafe, raises: [IOError, Defect].} = - chronosReadOnce(ChronosInputStream s, dst, dstLen) + fsTranslateErrors "Unexpected exception from merely forwarding a future": + return chronosReadOnce(ChronosInputStream s, dst, dstLen) , closeSync: proc (s: InputStream) {.nimcall, gcsafe, raises: [IOError, Defect].} = @@ -64,10 +64,11 @@ let chronosInputVTable = InputStreamVTable( func chronosInput*(s: StreamTransport, pageSize = defaultPageSize, - allowWaitFor = false): InputStreamHandle = - makeHandle ChronosInputStream( + allowWaitFor = false): AsyncInputStream = + AsyncInputStream ChronosInputStream( vtable: vtableAddr chronosInputVTable, buffers: initPageBuffers(pageSize), + transport: s, allowWaitFor: allowWaitFor) let chronosOutputVTable = OutputStreamVTable( @@ -88,15 +89,15 @@ let chronosOutputVTable = OutputStreamVTable( , closeAsync: proc (s: OutputStream): Future[void] {.nimcall, gcsafe, raises: [IOError, Defect].} = - chronosCloseWait ChronosOutputtream(s).transport + chronosCloseWait ChronosOutputStream(s).transport ) func chronosOutput*(s: StreamTransport, pageSize = defaultPageSize, - allowWaitFor = false): OutputStreamHandle = - makeHandle ChronosOutputStream( - vtable: vtableAddr(chronosOuputVTable), - buffers: initPageBuffers(pageSize) + allowWaitFor = false): AsyncOutputStream = + AsyncOutputStream ChronosOutputStream( + vtable: vtableAddr(chronosOutputVTable), + buffers: initPageBuffers(pageSize), transport: s, allowWaitFor: allowWaitFor) diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index 10c8950..bf6cf03 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -108,7 +108,7 @@ template close*(sp: AsyncInputStream) = disconnectInputDevice(s) preventFurtherReading(s) if s.closeFut != nil: - await s.closeFut + fsAwait s.closeFut template closeNoWait*(sp: AsyncInputStream|InputStream) = ## Close the stream without waiting even if's async. @@ -234,7 +234,7 @@ proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} = let s = InputStream(sp) fsAssert s.buffers != nil and s.vtable != nil - result = await s.vtable.readAsync(s, nil, 0) + result = fsAwait s.vtable.readAsync(s, nil, 0) if s.buffers.eofReached: disconnectInputDevice(s) @@ -245,7 +245,7 @@ proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} = proc timeoutToNextByteImpl(s: AsyncInputStream, deadline: Future): Future[bool] {.async.} = let readFut = s.readOnce - await readFut or deadline + fsAwait readFut or deadline if not readFut.finished: readFut.cancel() return true @@ -257,14 +257,14 @@ template timeoutToNextByte*(sp: AsyncInputStream, deadline: Future): bool = if readableNow(s): true else: - await timeoutToNextByteImpl(s, deadline) + fsAwait timeoutToNextByteImpl(s, deadline) template timeoutToNextByte*(sp: AsyncInputStream, timeout: Duration): bool = let s = sp if readableNow(s): true else: - await timeoutToNextByteImpl(s, sleepAsync(timeout)) + fsAwait timeoutToNextByteImpl(s, sleepAsync(timeout)) proc closeAsync*(s: AsyncInputStream) {.async.} = close s diff --git a/faststreams/outputs.nim b/faststreams/outputs.nim index 53adc68..dae8e13 100644 --- a/faststreams/outputs.nim +++ b/faststreams/outputs.nim @@ -123,7 +123,7 @@ template close*(sp: AsyncOutputStream) = flush(Async s) disconnectOutputDevice(s) if s.closeFut != nil: - await s.closeFut + fsAwait s.closeFut proc closeAsync*(s: AsyncOutputStream) {.async.} = close s @@ -311,7 +311,7 @@ proc drainAllBuffersSync(s: OutputStream, buf: pointer, bufSize: Natural) = s.spanEndPos += s.span.len proc drainAllBuffersAsync(s: OutputStream, buf: pointer, bufSize: Natural) {.async.} = - await s.vtable.writeAsync(s, buf, bufSize) + fsAwait s.vtable.writeAsync(s, buf, bufSize) s.span = s.buffers.getWritableSpan() s.spanEndPos += s.span.len diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim index 6e5c3ac..cf1f510 100644 --- a/faststreams/pipelines.nim +++ b/faststreams/pipelines.nim @@ -15,7 +15,7 @@ type template enterWait(fut: var Future, context: static string) = let wait = newFuture[void](context) fut = wait - try: await wait + try: fsAwait wait finally: fut = nil template awake(fp: Future) = @@ -312,7 +312,7 @@ macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untype closingCall.insert(1, newCall(bindSym"initReader", stepInput)) pipelineBody.add quote do: - await allFutures(`pipelineSteps`) + fsAwait allFutures(`pipelineSteps`) `closingCall` result = quote do: @@ -326,7 +326,7 @@ macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untype proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} = when UserOpRetType is Future: var f = `pipelineBody` - return await(f) + return fsAwait(f) else: return `pipelineBody` diff --git a/faststreams/textio.nim b/faststreams/textio.nim index 74f69bf..4041d08 100644 --- a/faststreams/textio.nim +++ b/faststreams/textio.nim @@ -126,7 +126,7 @@ const NewLines* = {'\r', '\n'} Digits* = {'0'..'9'} -proc readLine*(s: InputStream, keepEol = false): TaintedString = +proc readLine*(s: InputStream, keepEol = false): TaintedString {.fsMultiSync.} = fsAssert readableNow(s) while s.readable: