Implement timeoutToNextByte
This commit is contained in:
parent
e9f50bc847
commit
fc76250d02
2 changed files with 54 additions and 14 deletions
|
|
@ -110,9 +110,6 @@ template close*(sp: AsyncInputStream) =
|
||||||
if s.closeFut != nil:
|
if s.closeFut != nil:
|
||||||
await s.closeFut
|
await s.closeFut
|
||||||
|
|
||||||
proc closeAsync*(s: AsyncInputStream) {.async.} =
|
|
||||||
close s
|
|
||||||
|
|
||||||
template closeNoWait*(sp: AsyncInputStream|InputStream) =
|
template closeNoWait*(sp: AsyncInputStream|InputStream) =
|
||||||
## Close the stream without waiting even if's async.
|
## Close the stream without waiting even if's async.
|
||||||
## This operation will use `asyncCheck` internally to detect unhandled
|
## This operation will use `asyncCheck` internally to detect unhandled
|
||||||
|
|
@ -215,6 +212,55 @@ proc readableNow*(s: InputStream): bool =
|
||||||
template readableNow*(s: AsyncInputStream): bool =
|
template readableNow*(s: AsyncInputStream): bool =
|
||||||
readableNow InputStream(s)
|
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) =
|
func flipPage(s: InputStream) =
|
||||||
fsAssert s.buffers != nil and s.buffers.len > 1
|
fsAssert s.buffers != nil and s.buffers.len > 1
|
||||||
discard s.buffers.popFirst
|
discard s.buffers.popFirst
|
||||||
|
|
@ -716,11 +762,6 @@ template readInto*(sp: AsyncInputStream, dst: var openarray[byte]): bool =
|
||||||
let (dstAddr, dstLen) = openArrayToPair(dst)
|
let (dstAddr, dstLen) = openArrayToPair(dst)
|
||||||
readIntoExImpl(s, dstAddr, dstLen, fsAwait, readAsync) == dstLen
|
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):
|
when defined(windows):
|
||||||
proc alloca(n: int): ptr byte {.importc, header: "<malloc.h>".}
|
proc alloca(n: int): ptr byte {.importc, header: "<malloc.h>".}
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -6,7 +6,7 @@ export
|
||||||
inputs, outputs, async_backend
|
inputs, outputs, async_backend
|
||||||
|
|
||||||
type
|
type
|
||||||
FsAsyncPipe* = ref object
|
Pipe* = ref object
|
||||||
# TODO: Make these stream handles
|
# TODO: Make these stream handles
|
||||||
input*: AsyncInputStream
|
input*: AsyncInputStream
|
||||||
output*: AsyncOutputStream
|
output*: AsyncOutputStream
|
||||||
|
|
@ -218,16 +218,15 @@ proc pipeOutput*(buffers: PageBuffers,
|
||||||
destination: destination)
|
destination: destination)
|
||||||
|
|
||||||
func asyncPipe*(pageSize = defaultPageSize,
|
func asyncPipe*(pageSize = defaultPageSize,
|
||||||
maxBufferedBytes = defaultPageSize * 4): FsAsyncPipe =
|
maxBufferedBytes = defaultPageSize * 4): Pipe =
|
||||||
fsAssert pageSize > 0
|
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)
|
result = pipeInput(pipe.buffers)
|
||||||
pipe.input = result
|
pipe.input = result
|
||||||
|
|
||||||
func initWriter*(pipe: FsAsyncPipe): AsyncOutputStream =
|
func initWriter*(pipe: Pipe): AsyncOutputStream =
|
||||||
|
|
||||||
result = pipeOutput(pipe.buffers)
|
result = pipeOutput(pipe.buffers)
|
||||||
pipe.output = result
|
pipe.output = result
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue