Fixes incorrect async exception handling. Adds sleepAsync.
The tasyncexceptions test has been added which tests for this incorrect exception handling behaviour. The problem was that the exception was raised inside a callback which was called from a previously finished async procedure. This caused a "Future already finished" error. The fix was to simply reraise the exception if the retFutureSym is already finished. sleepAsync was added to help with the reproduction of this test. It should also be useful for users however. Finally some debug information was added to futures to help with future bugs.
This commit is contained in:
parent
fd086abb43
commit
4f5f98f0b1
4 changed files with 141 additions and 50 deletions
|
|
@ -9,7 +9,7 @@
|
||||||
|
|
||||||
include "system/inclrtl"
|
include "system/inclrtl"
|
||||||
|
|
||||||
import os, oids, tables, strutils, macros
|
import os, oids, tables, strutils, macros, times
|
||||||
|
|
||||||
import rawsockets, net
|
import rawsockets, net
|
||||||
|
|
||||||
|
|
@ -41,27 +41,40 @@ type
|
||||||
cb: proc () {.closure,gcsafe.}
|
cb: proc () {.closure,gcsafe.}
|
||||||
finished: bool
|
finished: bool
|
||||||
error*: ref EBase
|
error*: ref EBase
|
||||||
stackTrace: string ## For debugging purposes only.
|
when defined(debug):
|
||||||
|
stackTrace: string ## For debugging purposes only.
|
||||||
|
id: int
|
||||||
|
fromProc: string
|
||||||
|
|
||||||
PFuture*[T] = ref object of PFutureBase
|
PFuture*[T] = ref object of PFutureBase
|
||||||
value: T
|
value: T
|
||||||
|
|
||||||
proc newFuture*[T](): PFuture[T] =
|
var currentID* = 0
|
||||||
|
proc newFuture*[T](fromProc: string = "unspecified"): PFuture[T] =
|
||||||
## Creates a new future.
|
## Creates a new future.
|
||||||
|
##
|
||||||
|
## Specifying ``fromProc``, which is a string specifying the name of the proc
|
||||||
|
## that this future belongs to, is a good habit as it helps with debugging.
|
||||||
new(result)
|
new(result)
|
||||||
result.finished = false
|
result.finished = false
|
||||||
result.stackTrace = getStackTrace()
|
when defined(debug):
|
||||||
|
result.stackTrace = getStackTrace()
|
||||||
|
result.id = currentID
|
||||||
|
result.fromProc = fromProc
|
||||||
|
currentID.inc()
|
||||||
|
|
||||||
proc checkFinished[T](future: PFuture[T]) =
|
proc checkFinished[T](future: PFuture[T]) =
|
||||||
if future.finished:
|
when defined(debug):
|
||||||
echo("<----->")
|
if future.finished:
|
||||||
echo(future.stackTrace)
|
echo("<-----> ", future.id, " ", future.fromProc)
|
||||||
echo("-----")
|
echo(future.stackTrace)
|
||||||
when T is string:
|
echo("-----")
|
||||||
echo("Contents: ", future.value.repr)
|
when T is string:
|
||||||
echo("<----->")
|
echo("Contents: ", future.value.repr)
|
||||||
echo("Future already finished, cannot finish twice.")
|
echo("<----->")
|
||||||
assert false
|
echo("Future already finished, cannot finish twice.")
|
||||||
|
echo getStackTrace()
|
||||||
|
assert false
|
||||||
|
|
||||||
proc complete*[T](future: PFuture[T], val: T) =
|
proc complete*[T](future: PFuture[T], val: T) =
|
||||||
## Completes ``future`` with value ``val``.
|
## Completes ``future`` with value ``val``.
|
||||||
|
|
@ -121,7 +134,8 @@ proc read*[T](future: PFuture[T]): T =
|
||||||
##
|
##
|
||||||
## If the result of the future is an error then that error will be raised.
|
## If the result of the future is an error then that error will be raised.
|
||||||
if future.finished:
|
if future.finished:
|
||||||
if future.error != nil: raise future.error
|
if future.error != nil:
|
||||||
|
raise future.error
|
||||||
when T isnot void:
|
when T isnot void:
|
||||||
return future.value
|
return future.value
|
||||||
else:
|
else:
|
||||||
|
|
@ -150,7 +164,21 @@ proc asyncCheck*[T](future: PFuture[T]) =
|
||||||
## This should be used instead of ``discard`` to discard void futures.
|
## This should be used instead of ``discard`` to discard void futures.
|
||||||
future.callback =
|
future.callback =
|
||||||
proc () =
|
proc () =
|
||||||
if future.failed: raise future.error
|
if future.failed:
|
||||||
|
raise future.error
|
||||||
|
|
||||||
|
type
|
||||||
|
PDispatcherBase = ref object of PObject
|
||||||
|
timers: seq[tuple[finishAt: float, fut: PFuture[void]]]
|
||||||
|
|
||||||
|
proc processTimers(p: PDispatcherBase) =
|
||||||
|
var oldTimers = p.timers
|
||||||
|
p.timers = @[]
|
||||||
|
for t in oldTimers:
|
||||||
|
if epochTime() >= t.finishAt:
|
||||||
|
t.fut.complete()
|
||||||
|
else:
|
||||||
|
p.timers.add(t)
|
||||||
|
|
||||||
when defined(windows) or defined(nimdoc):
|
when defined(windows) or defined(nimdoc):
|
||||||
import winlean, sets, hashes
|
import winlean, sets, hashes
|
||||||
|
|
@ -162,7 +190,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
cb: proc (sock: TAsyncFD, bytesTransferred: DWORD,
|
cb: proc (sock: TAsyncFD, bytesTransferred: DWORD,
|
||||||
errcode: TOSErrorCode) {.closure,gcsafe.}
|
errcode: TOSErrorCode) {.closure,gcsafe.}
|
||||||
|
|
||||||
PDispatcher* = ref object
|
PDispatcher* = ref object of PDispatcherBase
|
||||||
ioPort: THandle
|
ioPort: THandle
|
||||||
handles: TSet[TAsyncFD]
|
handles: TSet[TAsyncFD]
|
||||||
|
|
||||||
|
|
@ -181,6 +209,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
new result
|
new result
|
||||||
result.ioPort = CreateIOCompletionPort(INVALID_HANDLE_VALUE, 0, 0, 1)
|
result.ioPort = CreateIOCompletionPort(INVALID_HANDLE_VALUE, 0, 0, 1)
|
||||||
result.handles = initSet[TAsyncFD]()
|
result.handles = initSet[TAsyncFD]()
|
||||||
|
result.timers = @[]
|
||||||
|
|
||||||
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
||||||
proc getGlobalDispatcher*(): PDispatcher =
|
proc getGlobalDispatcher*(): PDispatcher =
|
||||||
|
|
@ -207,8 +236,9 @@ when defined(windows) or defined(nimdoc):
|
||||||
proc poll*(timeout = 500) =
|
proc poll*(timeout = 500) =
|
||||||
## Waits for completion events and processes them.
|
## Waits for completion events and processes them.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
if p.handles.len == 0:
|
if p.handles.len == 0 and p.timers.len == 0:
|
||||||
raise newException(EInvalidValue, "No handles registered in dispatcher.")
|
raise newException(EInvalidValue,
|
||||||
|
"No handles or timers registered in dispatcher.")
|
||||||
|
|
||||||
let llTimeout =
|
let llTimeout =
|
||||||
if timeout == -1: winlean.INFINITE
|
if timeout == -1: winlean.INFINITE
|
||||||
|
|
@ -242,6 +272,9 @@ when defined(windows) or defined(nimdoc):
|
||||||
discard
|
discard
|
||||||
else: osError(errCode)
|
else: osError(errCode)
|
||||||
|
|
||||||
|
# Timer processing.
|
||||||
|
processTimers(p)
|
||||||
|
|
||||||
var connectExPtr: pointer = nil
|
var connectExPtr: pointer = nil
|
||||||
var acceptExPtr: pointer = nil
|
var acceptExPtr: pointer = nil
|
||||||
var getAcceptExSockAddrsPtr: pointer = nil
|
var getAcceptExSockAddrsPtr: pointer = nil
|
||||||
|
|
@ -314,7 +347,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
## Returns a ``PFuture`` which will complete when the connection succeeds
|
## Returns a ``PFuture`` which will complete when the connection succeeds
|
||||||
## or an error occurs.
|
## or an error occurs.
|
||||||
verifyPresence(socket)
|
verifyPresence(socket)
|
||||||
var retFuture = newFuture[void]()
|
var retFuture = newFuture[void]("connect")
|
||||||
# Apparently ``ConnectEx`` expects the socket to be initially bound:
|
# Apparently ``ConnectEx`` expects the socket to be initially bound:
|
||||||
var saddr: Tsockaddr_in
|
var saddr: Tsockaddr_in
|
||||||
saddr.sin_family = int16(toInt(af))
|
saddr.sin_family = int16(toInt(af))
|
||||||
|
|
@ -384,7 +417,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
# '\0' in the message currently signifies a socket disconnect. Who
|
# '\0' in the message currently signifies a socket disconnect. Who
|
||||||
# knows what will happen when someone sends that to our socket.
|
# knows what will happen when someone sends that to our socket.
|
||||||
verifyPresence(socket)
|
verifyPresence(socket)
|
||||||
var retFuture = newFuture[string]()
|
var retFuture = newFuture[string]("recv")
|
||||||
var dataBuf: TWSABuf
|
var dataBuf: TWSABuf
|
||||||
dataBuf.buf = cast[cstring](alloc0(size))
|
dataBuf.buf = cast[cstring](alloc0(size))
|
||||||
dataBuf.len = size
|
dataBuf.len = size
|
||||||
|
|
@ -459,7 +492,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
## Sends ``data`` to ``socket``. The returned future will complete once all
|
## Sends ``data`` to ``socket``. The returned future will complete once all
|
||||||
## data has been sent.
|
## data has been sent.
|
||||||
verifyPresence(socket)
|
verifyPresence(socket)
|
||||||
var retFuture = newFuture[void]()
|
var retFuture = newFuture[void]("send")
|
||||||
|
|
||||||
var dataBuf: TWSABuf
|
var dataBuf: TWSABuf
|
||||||
dataBuf.buf = data # since this is not used in a callback, this is fine
|
dataBuf.buf = data # since this is not used in a callback, this is fine
|
||||||
|
|
@ -502,7 +535,7 @@ when defined(windows) or defined(nimdoc):
|
||||||
##
|
##
|
||||||
## The resulting client socket is automatically registered to dispatcher.
|
## The resulting client socket is automatically registered to dispatcher.
|
||||||
verifyPresence(socket)
|
verifyPresence(socket)
|
||||||
var retFuture = newFuture[tuple[address: string, client: TAsyncFD]]()
|
var retFuture = newFuture[tuple[address: string, client: TAsyncFD]]("acceptAddr")
|
||||||
|
|
||||||
var clientSock = newRawSocket()
|
var clientSock = newRawSocket()
|
||||||
if clientSock == osInvalidSocket: osError(osLastError())
|
if clientSock == osInvalidSocket: osError(osLastError())
|
||||||
|
|
@ -614,6 +647,7 @@ else:
|
||||||
proc newDispatcher*(): PDispatcher =
|
proc newDispatcher*(): PDispatcher =
|
||||||
new result
|
new result
|
||||||
result.selector = newSelector()
|
result.selector = newSelector()
|
||||||
|
result.timers = @[]
|
||||||
|
|
||||||
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
var gDisp{.threadvar.}: PDispatcher ## Global dispatcher
|
||||||
proc getGlobalDispatcher*(): PDispatcher =
|
proc getGlobalDispatcher*(): PDispatcher =
|
||||||
|
|
@ -694,6 +728,8 @@ else:
|
||||||
# FD no longer a part of the selector. Likely been closed
|
# FD no longer a part of the selector. Likely been closed
|
||||||
# (e.g. socket disconnected).
|
# (e.g. socket disconnected).
|
||||||
|
|
||||||
|
processTimers(p)
|
||||||
|
|
||||||
proc connect*(socket: TAsyncFD, address: string, port: TPort,
|
proc connect*(socket: TAsyncFD, address: string, port: TPort,
|
||||||
af = AF_INET): PFuture[void] =
|
af = AF_INET): PFuture[void] =
|
||||||
var retFuture = newFuture[void]()
|
var retFuture = newFuture[void]()
|
||||||
|
|
@ -814,11 +850,19 @@ else:
|
||||||
addRead(socket, cb)
|
addRead(socket, cb)
|
||||||
return retFuture
|
return retFuture
|
||||||
|
|
||||||
|
proc sleepAsync*(ms: int): PFuture[void] =
|
||||||
|
## Suspends the execution of the current async procedure for the next
|
||||||
|
## ``ms`` miliseconds.
|
||||||
|
var retFuture = newFuture[void]("sleepAsync")
|
||||||
|
let p = getGlobalDispatcher()
|
||||||
|
p.timers.add((epochTime() + (ms / 1000), retFuture))
|
||||||
|
return retFuture
|
||||||
|
|
||||||
proc accept*(socket: TAsyncFD): PFuture[TAsyncFD] =
|
proc accept*(socket: TAsyncFD): PFuture[TAsyncFD] =
|
||||||
## Accepts a new connection. Returns a future containing the client socket
|
## Accepts a new connection. Returns a future containing the client socket
|
||||||
## corresponding to that connection.
|
## corresponding to that connection.
|
||||||
## The future will complete when the connection is successfully accepted.
|
## The future will complete when the connection is successfully accepted.
|
||||||
var retFut = newFuture[TAsyncFD]()
|
var retFut = newFuture[TAsyncFD]("accept")
|
||||||
var fut = acceptAddr(socket)
|
var fut = acceptAddr(socket)
|
||||||
fut.callback =
|
fut.callback =
|
||||||
proc (future: PFuture[tuple[address: string, client: TAsyncFD]]) =
|
proc (future: PFuture[tuple[address: string, client: TAsyncFD]]) =
|
||||||
|
|
@ -845,11 +889,16 @@ template createCb*(retFutureSym, iteratorNameSym,
|
||||||
else:
|
else:
|
||||||
next.callback = cb
|
next.callback = cb
|
||||||
except:
|
except:
|
||||||
retFutureSym.fail(getCurrentException())
|
if retFutureSym.finished:
|
||||||
|
# Take a look at tasyncexceptions for the bug which this fixes.
|
||||||
|
# That test explains it better than I can here.
|
||||||
|
raise
|
||||||
|
else:
|
||||||
|
retFutureSym.fail(getCurrentException())
|
||||||
cb()
|
cb()
|
||||||
#{.pop.}
|
#{.pop.}
|
||||||
proc generateExceptionCheck(futSym,
|
proc generateExceptionCheck(futSym,
|
||||||
exceptBranch, rootReceiver: PNimrodNode): PNimrodNode {.compileTime.} =
|
exceptBranch, rootReceiver, fromNode: PNimrodNode): PNimrodNode {.compileTime.} =
|
||||||
if exceptBranch == nil:
|
if exceptBranch == nil:
|
||||||
result = rootReceiver
|
result = rootReceiver
|
||||||
else:
|
else:
|
||||||
|
|
@ -869,20 +918,21 @@ proc generateExceptionCheck(futSym,
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
let elseNode = newNimNode(nnkElse)
|
let elseNode = newNimNode(nnkElse, fromNode)
|
||||||
elseNode.add newNimNode(nnkStmtList)
|
elseNode.add newNimNode(nnkStmtList, fromNode)
|
||||||
elseNode[0].add rootReceiver
|
elseNode[0].add rootReceiver
|
||||||
result.add elseNode
|
result.add elseNode
|
||||||
|
|
||||||
template createVar(result: var PNimrodNode, futSymName: string,
|
template createVar(result: var PNimrodNode, futSymName: string,
|
||||||
asyncProc: PNimrodNode,
|
asyncProc: PNimrodNode,
|
||||||
valueReceiver, rootReceiver: expr) =
|
valueReceiver, rootReceiver: expr,
|
||||||
result = newNimNode(nnkStmtList)
|
fromNode: PNimrodNode) =
|
||||||
|
result = newNimNode(nnkStmtList, fromNode)
|
||||||
var futSym = genSym(nskVar, "future")
|
var futSym = genSym(nskVar, "future")
|
||||||
result.add newVarStmt(futSym, asyncProc) # -> var future<x> = y
|
result.add newVarStmt(futSym, asyncProc) # -> var future<x> = y
|
||||||
result.add newNimNode(nnkYieldStmt).add(futSym) # -> yield future<x>
|
result.add newNimNode(nnkYieldStmt, fromNode).add(futSym) # -> yield future<x>
|
||||||
valueReceiver = newDotExpr(futSym, newIdentNode("read")) # -> future<x>.read
|
valueReceiver = newDotExpr(futSym, newIdentNode("read")) # -> future<x>.read
|
||||||
result.add generateExceptionCheck(futSym, exceptBranch, rootReceiver)
|
result.add generateExceptionCheck(futSym, exceptBranch, rootReceiver, fromNode)
|
||||||
|
|
||||||
proc processBody(node, retFutureSym: PNimrodNode,
|
proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
subTypeIsVoid: bool,
|
subTypeIsVoid: bool,
|
||||||
|
|
@ -891,7 +941,7 @@ proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
result = node
|
result = node
|
||||||
case node.kind
|
case node.kind
|
||||||
of nnkReturnStmt:
|
of nnkReturnStmt:
|
||||||
result = newNimNode(nnkStmtList)
|
result = newNimNode(nnkStmtList, node)
|
||||||
if node[0].kind == nnkEmpty:
|
if node[0].kind == nnkEmpty:
|
||||||
if not subtypeIsVoid:
|
if not subtypeIsVoid:
|
||||||
result.add newCall(newIdentNode("complete"), retFutureSym,
|
result.add newCall(newIdentNode("complete"), retFutureSym,
|
||||||
|
|
@ -902,19 +952,19 @@ proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
result.add newCall(newIdentNode("complete"), retFutureSym,
|
result.add newCall(newIdentNode("complete"), retFutureSym,
|
||||||
node[0].processBody(retFutureSym, subtypeIsVoid, exceptBranch))
|
node[0].processBody(retFutureSym, subtypeIsVoid, exceptBranch))
|
||||||
|
|
||||||
result.add newNimNode(nnkReturnStmt).add(newNilLit())
|
result.add newNimNode(nnkReturnStmt, node).add(newNilLit())
|
||||||
return # Don't process the children of this return stmt
|
return # Don't process the children of this return stmt
|
||||||
of nnkCommand:
|
of nnkCommand:
|
||||||
if node[0].kind == nnkIdent and node[0].ident == !"await":
|
if node[0].kind == nnkIdent and node[0].ident == !"await":
|
||||||
case node[1].kind
|
case node[1].kind
|
||||||
of nnkIdent:
|
of nnkIdent:
|
||||||
# await x
|
# await x
|
||||||
result = newNimNode(nnkYieldStmt).add(node[1]) # -> yield x
|
result = newNimNode(nnkYieldStmt, node).add(node[1]) # -> yield x
|
||||||
of nnkCall:
|
of nnkCall:
|
||||||
# await foo(p, x)
|
# await foo(p, x)
|
||||||
var futureValue: PNimrodNode
|
var futureValue: PNimrodNode
|
||||||
result.createVar("future" & $node[1][0].toStrLit, node[1], futureValue,
|
result.createVar("future" & $node[1][0].toStrLit, node[1], futureValue,
|
||||||
futureValue)
|
futureValue, node)
|
||||||
else:
|
else:
|
||||||
error("Invalid node kind in 'await', got: " & $node[1].kind)
|
error("Invalid node kind in 'await', got: " & $node[1].kind)
|
||||||
elif node[1].kind == nnkCommand and node[1][0].kind == nnkIdent and
|
elif node[1].kind == nnkCommand and node[1][0].kind == nnkIdent and
|
||||||
|
|
@ -922,7 +972,7 @@ proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
# foo await x
|
# foo await x
|
||||||
var newCommand = node
|
var newCommand = node
|
||||||
result.createVar("future" & $node[0].toStrLit, node[1][1], newCommand[1],
|
result.createVar("future" & $node[0].toStrLit, node[1][1], newCommand[1],
|
||||||
newCommand)
|
newCommand, node)
|
||||||
|
|
||||||
of nnkVarSection, nnkLetSection:
|
of nnkVarSection, nnkLetSection:
|
||||||
case node[0][2].kind
|
case node[0][2].kind
|
||||||
|
|
@ -931,7 +981,7 @@ proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
# var x = await y
|
# var x = await y
|
||||||
var newVarSection = node # TODO: Should this use copyNimNode?
|
var newVarSection = node # TODO: Should this use copyNimNode?
|
||||||
result.createVar("future" & $node[0][0].ident, node[0][2][1],
|
result.createVar("future" & $node[0][0].ident, node[0][2][1],
|
||||||
newVarSection[0][2], newVarSection)
|
newVarSection[0][2], newVarSection, node)
|
||||||
else: discard
|
else: discard
|
||||||
of nnkAsgn:
|
of nnkAsgn:
|
||||||
case node[1].kind
|
case node[1].kind
|
||||||
|
|
@ -939,7 +989,7 @@ proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
if node[1][0].ident == !"await":
|
if node[1][0].ident == !"await":
|
||||||
# x = await y
|
# x = await y
|
||||||
var newAsgn = node
|
var newAsgn = node
|
||||||
result.createVar("future" & $node[0].toStrLit, node[1][1], newAsgn[1], newAsgn)
|
result.createVar("future" & $node[0].toStrLit, node[1][1], newAsgn[1], newAsgn, node)
|
||||||
else: discard
|
else: discard
|
||||||
of nnkDiscardStmt:
|
of nnkDiscardStmt:
|
||||||
# discard await x
|
# discard await x
|
||||||
|
|
@ -947,10 +997,10 @@ proc processBody(node, retFutureSym: PNimrodNode,
|
||||||
node[0][0].ident == !"await":
|
node[0][0].ident == !"await":
|
||||||
var newDiscard = node
|
var newDiscard = node
|
||||||
result.createVar("futureDiscard_" & $toStrLit(node[0][1]), node[0][1],
|
result.createVar("futureDiscard_" & $toStrLit(node[0][1]), node[0][1],
|
||||||
newDiscard[0], newDiscard)
|
newDiscard[0], newDiscard, node)
|
||||||
of nnkTryStmt:
|
of nnkTryStmt:
|
||||||
# try: await x; except: ...
|
# try: await x; except: ...
|
||||||
result = newNimNode(nnkStmtList)
|
result = newNimNode(nnkStmtList, node)
|
||||||
proc processForTry(n: PNimrodNode, i: var int,
|
proc processForTry(n: PNimrodNode, i: var int,
|
||||||
res: PNimrodNode): bool {.compileTime.} =
|
res: PNimrodNode): bool {.compileTime.} =
|
||||||
result = false
|
result = false
|
||||||
|
|
@ -1009,7 +1059,7 @@ macro async*(prc: stmt): stmt {.immediate.} =
|
||||||
(returnType.kind == nnkBracketExpr and
|
(returnType.kind == nnkBracketExpr and
|
||||||
returnType[1].kind == nnkIdent and returnType[1].ident == !"void")
|
returnType[1].kind == nnkIdent and returnType[1].ident == !"void")
|
||||||
|
|
||||||
var outerProcBody = newNimNode(nnkStmtList)
|
var outerProcBody = newNimNode(nnkStmtList, prc[6])
|
||||||
|
|
||||||
# -> var retFuture = newFuture[T]()
|
# -> var retFuture = newFuture[T]()
|
||||||
var retFutureSym = genSym(nskVar, "retFuture")
|
var retFutureSym = genSym(nskVar, "retFuture")
|
||||||
|
|
@ -1019,9 +1069,10 @@ macro async*(prc: stmt): stmt {.immediate.} =
|
||||||
outerProcBody.add(
|
outerProcBody.add(
|
||||||
newVarStmt(retFutureSym,
|
newVarStmt(retFutureSym,
|
||||||
newCall(
|
newCall(
|
||||||
newNimNode(nnkBracketExpr).add(
|
newNimNode(nnkBracketExpr, prc[6]).add(
|
||||||
newIdentNode(!"newFuture"), # TODO: Strange bug here? Remove the `!`.
|
newIdentNode(!"newFuture"), # TODO: Strange bug here? Remove the `!`.
|
||||||
subRetType)))) # Get type from return type of this proc
|
subRetType),
|
||||||
|
newLit(prc[0].getName)))) # Get type from return type of this proc
|
||||||
|
|
||||||
# -> iterator nameIter(): PFutureBase {.closure.} =
|
# -> iterator nameIter(): PFutureBase {.closure.} =
|
||||||
# -> var result: T
|
# -> var result: T
|
||||||
|
|
@ -1030,7 +1081,7 @@ macro async*(prc: stmt): stmt {.immediate.} =
|
||||||
var iteratorNameSym = genSym(nskIterator, $prc[0].getName & "Iter")
|
var iteratorNameSym = genSym(nskIterator, $prc[0].getName & "Iter")
|
||||||
var procBody = prc[6].processBody(retFutureSym, subtypeIsVoid, nil)
|
var procBody = prc[6].processBody(retFutureSym, subtypeIsVoid, nil)
|
||||||
if not subtypeIsVoid:
|
if not subtypeIsVoid:
|
||||||
procBody.insert(0, newNimNode(nnkVarSection).add(
|
procBody.insert(0, newNimNode(nnkVarSection, prc[6]).add(
|
||||||
newIdentDefs(newIdentNode("result"), returnType[1]))) # -> var result: T
|
newIdentDefs(newIdentNode("result"), returnType[1]))) # -> var result: T
|
||||||
procBody.add(
|
procBody.add(
|
||||||
newCall(newIdentNode("complete"),
|
newCall(newIdentNode("complete"),
|
||||||
|
|
@ -1041,7 +1092,7 @@ macro async*(prc: stmt): stmt {.immediate.} =
|
||||||
|
|
||||||
var closureIterator = newProc(iteratorNameSym, [newIdentNode("PFutureBase")],
|
var closureIterator = newProc(iteratorNameSym, [newIdentNode("PFutureBase")],
|
||||||
procBody, nnkIteratorDef)
|
procBody, nnkIteratorDef)
|
||||||
closureIterator[4] = newNimNode(nnkPragma).add(newIdentNode("closure"))
|
closureIterator[4] = newNimNode(nnkPragma, prc[6]).add(newIdentNode("closure"))
|
||||||
outerProcBody.add(closureIterator)
|
outerProcBody.add(closureIterator)
|
||||||
|
|
||||||
# -> createCb(retFuture)
|
# -> createCb(retFuture)
|
||||||
|
|
@ -1051,7 +1102,7 @@ macro async*(prc: stmt): stmt {.immediate.} =
|
||||||
outerProcBody.add procCb
|
outerProcBody.add procCb
|
||||||
|
|
||||||
# -> return retFuture
|
# -> return retFuture
|
||||||
outerProcBody.add newNimNode(nnkReturnStmt).add(retFutureSym)
|
outerProcBody.add newNimNode(nnkReturnStmt, prc[6][prc[6].len-1]).add(retFutureSym)
|
||||||
|
|
||||||
result = prc
|
result = prc
|
||||||
|
|
||||||
|
|
@ -1068,8 +1119,8 @@ macro async*(prc: stmt): stmt {.immediate.} =
|
||||||
result[6] = outerProcBody
|
result[6] = outerProcBody
|
||||||
|
|
||||||
#echo(treeRepr(result))
|
#echo(treeRepr(result))
|
||||||
#if prc[0].getName == "routeReq":
|
#if prc[0].getName == "processClient":
|
||||||
#echo(toStrLit(result))
|
# echo(toStrLit(result))
|
||||||
|
|
||||||
proc recvLine*(socket: TAsyncFD): PFuture[string] {.async.} =
|
proc recvLine*(socket: TAsyncFD): PFuture[string] {.async.} =
|
||||||
## Reads a line of data from ``socket``. Returned future will complete once
|
## Reads a line of data from ``socket``. Returned future will complete once
|
||||||
|
|
|
||||||
|
|
@ -199,6 +199,8 @@ proc serve*(server: PAsyncHttpServer, port: TPort,
|
||||||
#var (address, client) = await server.socket.acceptAddr()
|
#var (address, client) = await server.socket.acceptAddr()
|
||||||
var fut = await server.socket.acceptAddr()
|
var fut = await server.socket.acceptAddr()
|
||||||
asyncCheck processClient(fut.client, fut.address, callback)
|
asyncCheck processClient(fut.client, fut.address, callback)
|
||||||
|
#echo(f.isNil)
|
||||||
|
#echo(f.repr)
|
||||||
|
|
||||||
proc close*(server: PAsyncHttpServer) =
|
proc close*(server: PAsyncHttpServer) =
|
||||||
## Terminates the async http server instance.
|
## Terminates the async http server instance.
|
||||||
|
|
|
||||||
|
|
@ -140,7 +140,7 @@ proc acceptAddr*(socket: PAsyncSocket):
|
||||||
## Accepts a new connection. Returns a future containing the client socket
|
## Accepts a new connection. Returns a future containing the client socket
|
||||||
## corresponding to that connection and the remote address of the client.
|
## corresponding to that connection and the remote address of the client.
|
||||||
## The future will complete when the connection is successfully accepted.
|
## The future will complete when the connection is successfully accepted.
|
||||||
var retFuture = newFuture[tuple[address: string, client: PAsyncSocket]]()
|
var retFuture = newFuture[tuple[address: string, client: PAsyncSocket]]("asyncnet.acceptAddr")
|
||||||
var fut = acceptAddr(socket.fd.TAsyncFD)
|
var fut = acceptAddr(socket.fd.TAsyncFD)
|
||||||
fut.callback =
|
fut.callback =
|
||||||
proc (future: PFuture[tuple[address: string, client: TAsyncFD]]) =
|
proc (future: PFuture[tuple[address: string, client: TAsyncFD]]) =
|
||||||
|
|
@ -157,7 +157,7 @@ proc accept*(socket: PAsyncSocket): PFuture[PAsyncSocket] =
|
||||||
## Accepts a new connection. Returns a future containing the client socket
|
## Accepts a new connection. Returns a future containing the client socket
|
||||||
## corresponding to that connection.
|
## corresponding to that connection.
|
||||||
## The future will complete when the connection is successfully accepted.
|
## The future will complete when the connection is successfully accepted.
|
||||||
var retFut = newFuture[PAsyncSocket]()
|
var retFut = newFuture[PAsyncSocket]("asyncnet.accept")
|
||||||
var fut = acceptAddr(socket)
|
var fut = acceptAddr(socket)
|
||||||
fut.callback =
|
fut.callback =
|
||||||
proc (future: PFuture[tuple[address: string, client: PAsyncSocket]]) =
|
proc (future: PFuture[tuple[address: string, client: PAsyncSocket]]) =
|
||||||
|
|
|
||||||
38
tests/async/tasyncexceptions.nim
Normal file
38
tests/async/tasyncexceptions.nim
Normal file
|
|
@ -0,0 +1,38 @@
|
||||||
|
discard """
|
||||||
|
file: "tasyncexceptions.nim"
|
||||||
|
exitcode: 1
|
||||||
|
outputsub: "Error: unhandled exception: foobar [E_Base]"
|
||||||
|
"""
|
||||||
|
import asyncdispatch
|
||||||
|
|
||||||
|
proc accept(): PFuture[int] {.async.} =
|
||||||
|
await sleepAsync(100)
|
||||||
|
result = 4
|
||||||
|
|
||||||
|
proc recvLine(fd: int): PFuture[string] {.async.} =
|
||||||
|
await sleepAsync(100)
|
||||||
|
return "get"
|
||||||
|
|
||||||
|
proc processClient(fd: int) {.async.} =
|
||||||
|
# these finish synchronously, we need some async delay to emulate this bug.
|
||||||
|
var line = await recvLine(fd)
|
||||||
|
var foo = line[0]
|
||||||
|
if foo == 'g':
|
||||||
|
raise newException(EBase, "foobar")
|
||||||
|
|
||||||
|
|
||||||
|
proc serve() {.async.} =
|
||||||
|
|
||||||
|
while true:
|
||||||
|
var fut = await accept()
|
||||||
|
await processClient(fut)
|
||||||
|
|
||||||
|
when isMainModule:
|
||||||
|
var fut = serve()
|
||||||
|
fut.callback =
|
||||||
|
proc () =
|
||||||
|
if fut.failed:
|
||||||
|
# This test ensures that this exception crashes the application
|
||||||
|
# as it is not handled.
|
||||||
|
raise fut.error
|
||||||
|
runForever()
|
||||||
Loading…
Add table
Add a link
Reference in a new issue