Add asynctools adapter
This commit is contained in:
parent
81c24860e2
commit
5df69fc696
7 changed files with 137 additions and 31 deletions
|
|
@ -20,12 +20,12 @@ when faststreams_async_backend == "chronos":
|
||||||
|
|
||||||
elif faststreams_async_backend in ["std", "asyncdispatch"]:
|
elif faststreams_async_backend in ["std", "asyncdispatch"]:
|
||||||
import
|
import
|
||||||
std/[asyncfutures, asyncmacro]
|
std/asyncdispatch
|
||||||
|
|
||||||
export
|
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
|
# TODO revisit after https://github.com/nim-lang/Nim/pull/12085/ is merged
|
||||||
let f = awaited
|
let f = awaited
|
||||||
yield f
|
yield f
|
||||||
|
|
@ -33,6 +33,8 @@ elif faststreams_async_backend in ["std", "asyncdispatch"]:
|
||||||
raise f.error
|
raise f.error
|
||||||
f.read
|
f.read
|
||||||
|
|
||||||
|
type Duration* = int
|
||||||
|
|
||||||
else:
|
else:
|
||||||
{.fatal: "Unrecognized network backend: " & faststreams_async_backend.}
|
{.fatal: "Unrecognized network backend: " & faststreams_async_backend.}
|
||||||
|
|
||||||
|
|
|
||||||
103
faststreams/asynctools_adapters.nim
Normal file
103
faststreams/asynctools_adapters.nim
Normal file
|
|
@ -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)
|
||||||
|
|
||||||
|
|
@ -26,17 +26,15 @@ proc chronosCloseWait(t: StreamTransport)
|
||||||
await t.closeWait()
|
await t.closeWait()
|
||||||
|
|
||||||
proc chronosReadOnce(s: ChronosInputStream,
|
proc chronosReadOnce(s: ChronosInputStream,
|
||||||
dst: pointer, dstLen: Natural): Future[Natural]
|
dst: pointer, dstLen: Natural): Future[Natural] {.async.} =
|
||||||
{.async, raises: [IOError, Defect].} =
|
|
||||||
fsTranslateErrors readingErrMsg:
|
fsTranslateErrors readingErrMsg:
|
||||||
return implementSingleRead(s.buffers, dst, dstLen, {},
|
return implementSingleRead(s.buffers, dst, dstLen, ReadFlags {},
|
||||||
readStartAddr, readLen):
|
readStartAddr, readLen):
|
||||||
await s.transport.readOnce(readStartAddr, readLen)
|
await s.transport.readOnce(readStartAddr, readLen)
|
||||||
|
|
||||||
proc chronosWrites(s: ChronosOutputStream, src: pointer, srcLen: Natural)
|
proc chronosWrites(s: ChronosOutputStream, src: pointer, srcLen: Natural) {.async.} =
|
||||||
{.async, raises: [IOError, Defect].} =
|
|
||||||
fsTranslateErrors writeIncompleteErrMsg:
|
fsTranslateErrors writeIncompleteErrMsg:
|
||||||
implementWrites(s.buffers, src, srcLen, "StreamTransport"
|
implementWrites(s.buffers, src, srcLen, "StreamTransport",
|
||||||
writeStartAddr, writeLen):
|
writeStartAddr, writeLen):
|
||||||
await s.transport.write(writeStartAddr, writeLen)
|
await s.transport.write(writeStartAddr, writeLen)
|
||||||
|
|
||||||
|
|
@ -44,13 +42,15 @@ proc chronosWrites(s: ChronosOutputStream, src: pointer, srcLen: Natural)
|
||||||
let chronosInputVTable = InputStreamVTable(
|
let chronosInputVTable = InputStreamVTable(
|
||||||
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
|
readSync: proc (s: InputStream, dst: pointer, dstLen: Natural): Natural
|
||||||
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
||||||
var cs = ChronosInputStream(s)
|
fsTranslateErrors "Unexpected exception from Chronos async macro":
|
||||||
fsAssert cs.allowWaitFor
|
var cs = ChronosInputStream(s)
|
||||||
waitFor chronosReadOnce(cs, dst, dstLen)
|
fsAssert cs.allowWaitFor
|
||||||
|
return waitFor chronosReadOnce(cs, dst, dstLen)
|
||||||
,
|
,
|
||||||
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
|
readAsync: proc (s: InputStream, dst: pointer, dstLen: Natural): Future[Natural]
|
||||||
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
{.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)
|
closeSync: proc (s: InputStream)
|
||||||
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
||||||
|
|
@ -64,10 +64,11 @@ let chronosInputVTable = InputStreamVTable(
|
||||||
|
|
||||||
func chronosInput*(s: StreamTransport,
|
func chronosInput*(s: StreamTransport,
|
||||||
pageSize = defaultPageSize,
|
pageSize = defaultPageSize,
|
||||||
allowWaitFor = false): InputStreamHandle =
|
allowWaitFor = false): AsyncInputStream =
|
||||||
makeHandle ChronosInputStream(
|
AsyncInputStream ChronosInputStream(
|
||||||
vtable: vtableAddr chronosInputVTable,
|
vtable: vtableAddr chronosInputVTable,
|
||||||
buffers: initPageBuffers(pageSize),
|
buffers: initPageBuffers(pageSize),
|
||||||
|
transport: s,
|
||||||
allowWaitFor: allowWaitFor)
|
allowWaitFor: allowWaitFor)
|
||||||
|
|
||||||
let chronosOutputVTable = OutputStreamVTable(
|
let chronosOutputVTable = OutputStreamVTable(
|
||||||
|
|
@ -88,15 +89,15 @@ let chronosOutputVTable = OutputStreamVTable(
|
||||||
,
|
,
|
||||||
closeAsync: proc (s: OutputStream): Future[void]
|
closeAsync: proc (s: OutputStream): Future[void]
|
||||||
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
{.nimcall, gcsafe, raises: [IOError, Defect].} =
|
||||||
chronosCloseWait ChronosOutputtream(s).transport
|
chronosCloseWait ChronosOutputStream(s).transport
|
||||||
)
|
)
|
||||||
|
|
||||||
func chronosOutput*(s: StreamTransport,
|
func chronosOutput*(s: StreamTransport,
|
||||||
pageSize = defaultPageSize,
|
pageSize = defaultPageSize,
|
||||||
allowWaitFor = false): OutputStreamHandle =
|
allowWaitFor = false): AsyncOutputStream =
|
||||||
makeHandle ChronosOutputStream(
|
AsyncOutputStream ChronosOutputStream(
|
||||||
vtable: vtableAddr(chronosOuputVTable),
|
vtable: vtableAddr(chronosOutputVTable),
|
||||||
buffers: initPageBuffers(pageSize)
|
buffers: initPageBuffers(pageSize),
|
||||||
transport: s,
|
transport: s,
|
||||||
allowWaitFor: allowWaitFor)
|
allowWaitFor: allowWaitFor)
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -108,7 +108,7 @@ template close*(sp: AsyncInputStream) =
|
||||||
disconnectInputDevice(s)
|
disconnectInputDevice(s)
|
||||||
preventFurtherReading(s)
|
preventFurtherReading(s)
|
||||||
if s.closeFut != nil:
|
if s.closeFut != nil:
|
||||||
await s.closeFut
|
fsAwait s.closeFut
|
||||||
|
|
||||||
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.
|
||||||
|
|
@ -234,7 +234,7 @@ proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} =
|
||||||
let s = InputStream(sp)
|
let s = InputStream(sp)
|
||||||
fsAssert s.buffers != nil and s.vtable != nil
|
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:
|
if s.buffers.eofReached:
|
||||||
disconnectInputDevice(s)
|
disconnectInputDevice(s)
|
||||||
|
|
@ -245,7 +245,7 @@ proc readOnce*(sp: AsyncInputStream): Future[Natural] {.async.} =
|
||||||
proc timeoutToNextByteImpl(s: AsyncInputStream,
|
proc timeoutToNextByteImpl(s: AsyncInputStream,
|
||||||
deadline: Future): Future[bool] {.async.} =
|
deadline: Future): Future[bool] {.async.} =
|
||||||
let readFut = s.readOnce
|
let readFut = s.readOnce
|
||||||
await readFut or deadline
|
fsAwait readFut or deadline
|
||||||
if not readFut.finished:
|
if not readFut.finished:
|
||||||
readFut.cancel()
|
readFut.cancel()
|
||||||
return true
|
return true
|
||||||
|
|
@ -257,14 +257,14 @@ template timeoutToNextByte*(sp: AsyncInputStream, deadline: Future): bool =
|
||||||
if readableNow(s):
|
if readableNow(s):
|
||||||
true
|
true
|
||||||
else:
|
else:
|
||||||
await timeoutToNextByteImpl(s, deadline)
|
fsAwait timeoutToNextByteImpl(s, deadline)
|
||||||
|
|
||||||
template timeoutToNextByte*(sp: AsyncInputStream, timeout: Duration): bool =
|
template timeoutToNextByte*(sp: AsyncInputStream, timeout: Duration): bool =
|
||||||
let s = sp
|
let s = sp
|
||||||
if readableNow(s):
|
if readableNow(s):
|
||||||
true
|
true
|
||||||
else:
|
else:
|
||||||
await timeoutToNextByteImpl(s, sleepAsync(timeout))
|
fsAwait timeoutToNextByteImpl(s, sleepAsync(timeout))
|
||||||
|
|
||||||
proc closeAsync*(s: AsyncInputStream) {.async.} =
|
proc closeAsync*(s: AsyncInputStream) {.async.} =
|
||||||
close s
|
close s
|
||||||
|
|
|
||||||
|
|
@ -123,7 +123,7 @@ template close*(sp: AsyncOutputStream) =
|
||||||
flush(Async s)
|
flush(Async s)
|
||||||
disconnectOutputDevice(s)
|
disconnectOutputDevice(s)
|
||||||
if s.closeFut != nil:
|
if s.closeFut != nil:
|
||||||
await s.closeFut
|
fsAwait s.closeFut
|
||||||
|
|
||||||
proc closeAsync*(s: AsyncOutputStream) {.async.} =
|
proc closeAsync*(s: AsyncOutputStream) {.async.} =
|
||||||
close s
|
close s
|
||||||
|
|
@ -311,7 +311,7 @@ proc drainAllBuffersSync(s: OutputStream, buf: pointer, bufSize: Natural) =
|
||||||
s.spanEndPos += s.span.len
|
s.spanEndPos += s.span.len
|
||||||
|
|
||||||
proc drainAllBuffersAsync(s: OutputStream, buf: pointer, bufSize: Natural) {.async.} =
|
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.span = s.buffers.getWritableSpan()
|
||||||
s.spanEndPos += s.span.len
|
s.spanEndPos += s.span.len
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -15,7 +15,7 @@ type
|
||||||
template enterWait(fut: var Future, context: static string) =
|
template enterWait(fut: var Future, context: static string) =
|
||||||
let wait = newFuture[void](context)
|
let wait = newFuture[void](context)
|
||||||
fut = wait
|
fut = wait
|
||||||
try: await wait
|
try: fsAwait wait
|
||||||
finally: fut = nil
|
finally: fut = nil
|
||||||
|
|
||||||
template awake(fp: Future) =
|
template awake(fp: Future) =
|
||||||
|
|
@ -312,7 +312,7 @@ macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untype
|
||||||
closingCall.insert(1, newCall(bindSym"initReader", stepInput))
|
closingCall.insert(1, newCall(bindSym"initReader", stepInput))
|
||||||
|
|
||||||
pipelineBody.add quote do:
|
pipelineBody.add quote do:
|
||||||
await allFutures(`pipelineSteps`)
|
fsAwait allFutures(`pipelineSteps`)
|
||||||
`closingCall`
|
`closingCall`
|
||||||
|
|
||||||
result = quote do:
|
result = quote do:
|
||||||
|
|
@ -326,7 +326,7 @@ macro executePipeline*(start: AsyncInputStream, steps: varargs[untyped]): untype
|
||||||
proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} =
|
proc pipelineProc(`stream`: AsyncInputStream): Future[RetType] {.async.} =
|
||||||
when UserOpRetType is Future:
|
when UserOpRetType is Future:
|
||||||
var f = `pipelineBody`
|
var f = `pipelineBody`
|
||||||
return await(f)
|
return fsAwait(f)
|
||||||
else:
|
else:
|
||||||
return `pipelineBody`
|
return `pipelineBody`
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -126,7 +126,7 @@ const
|
||||||
NewLines* = {'\r', '\n'}
|
NewLines* = {'\r', '\n'}
|
||||||
Digits* = {'0'..'9'}
|
Digits* = {'0'..'9'}
|
||||||
|
|
||||||
proc readLine*(s: InputStream, keepEol = false): TaintedString =
|
proc readLine*(s: InputStream, keepEol = false): TaintedString {.fsMultiSync.} =
|
||||||
fsAssert readableNow(s)
|
fsAssert readableNow(s)
|
||||||
|
|
||||||
while s.readable:
|
while s.readable:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue