From fc76250d02399db231e5dc8ca9d53ac328fb35b1 Mon Sep 17 00:00:00 2001 From: Zahary Karadjov Date: Wed, 6 May 2020 23:12:12 +0300 Subject: [PATCH] Implement timeoutToNextByte --- faststreams/inputs.nim | 57 +++++++++++++++++++++++++++++++++------ faststreams/pipelines.nim | 11 ++++---- 2 files changed, 54 insertions(+), 14 deletions(-) diff --git a/faststreams/inputs.nim b/faststreams/inputs.nim index 25cad53..f1cdae2 100644 --- a/faststreams/inputs.nim +++ b/faststreams/inputs.nim @@ -110,9 +110,6 @@ template close*(sp: AsyncInputStream) = if s.closeFut != nil: await s.closeFut -proc closeAsync*(s: AsyncInputStream) {.async.} = - close s - template closeNoWait*(sp: AsyncInputStream|InputStream) = ## Close the stream without waiting even if's async. ## This operation will use `asyncCheck` internally to detect unhandled @@ -215,6 +212,55 @@ proc readableNow*(s: InputStream): bool = 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 + + result = await s.vtable.readAsync(s, nil, 0) + + if s.buffers.eofReached: + disconnectInputDevice(s) + + if result > 0 and s.span.len == 0: + s.buffers.nextReadableSpan(s.span) + s.spanEndPos += s.span.len + +proc timeoutToNextByteImpl(s: AsyncInputStream, + deadline: Future): Future[bool] {.async.} = + let readFut = s.readOnce + await 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: + await timeoutToNextByteImpl(s, deadline) + +template timeoutToNextByte*(sp: AsyncInputStream, timeout: Duration): bool = + let s = sp + if readableNow(s): + true + else: + await timeoutToNextByteImpl(s, sleepAsync(timeout)) + +proc closeAsync*(s: AsyncInputStream) {.async.} = + close s + +# TODO: End of purely async interface + func flipPage(s: InputStream) = fsAssert s.buffers != nil and s.buffers.len > 1 discard s.buffers.popFirst @@ -716,11 +762,6 @@ template readInto*(sp: AsyncInputStream, dst: var openarray[byte]): bool = let (dstAddr, dstLen) = openArrayToPair(dst) readIntoExImpl(s, dstAddr, dstLen, fsAwait, readAsync) == dstLen -proc readOnce*(sp: AsyncInputStream): Future[Natural] = - let s = InputStream(sp) - fsAssert s.buffers != nil and s.vtable != nil - s.vtable.readAsync(s, nil, 0) - when defined(windows): proc alloca(n: int): ptr byte {.importc, header: "".} else: diff --git a/faststreams/pipelines.nim b/faststreams/pipelines.nim index 734a5b7..4953054 100644 --- a/faststreams/pipelines.nim +++ b/faststreams/pipelines.nim @@ -6,7 +6,7 @@ export inputs, outputs, async_backend type - FsAsyncPipe* = ref object + Pipe* = ref object # TODO: Make these stream handles input*: AsyncInputStream output*: AsyncOutputStream @@ -218,16 +218,15 @@ proc pipeOutput*(buffers: PageBuffers, destination: destination) func asyncPipe*(pageSize = defaultPageSize, - maxBufferedBytes = defaultPageSize * 4): FsAsyncPipe = + maxBufferedBytes = defaultPageSize * 4): Pipe = fsAssert pageSize > 0 - FsAsyncPipe(buffers: initPageBuffers(pageSize, maxBufferedBytes)) + Pipe(buffers: initPageBuffers(pageSize, maxBufferedBytes)) -func initReader*(pipe: FsAsyncPipe): AsyncInputStream = +func initReader*(pipe: Pipe): AsyncInputStream = result = pipeInput(pipe.buffers) pipe.input = result -func initWriter*(pipe: FsAsyncPipe): AsyncOutputStream = - +func initWriter*(pipe: Pipe): AsyncOutputStream = result = pipeOutput(pipe.buffers) pipe.output = result