Compare commits

...
Sign in to create a new pull request.

18 commits

Author SHA1 Message Date
Dominik Picheta
b3469c43a6 Implemented 'await connect' fully. Fixed a buffering and 'reqReg' bug.
* reqReq did not reset lastReq to reqNil, this caused problems in the
new tasyncitermacro test. (await accept stopped accepting requests after
some arbitrary amount have been accepted).
* Rewritten and moved tasyncitermacro to tests/run/ so that the tester
actually tests whether it works at runtime.
* tasynciterraw is now a copy of tasyncitermacro after macro expansion.
It will be used as a control test, useful if the macro starts generating
incorrect code.
2013-08-03 14:19:17 +01:00
Dominik Picheta
ab534804c6 Merge branch 'master' into asyncmacro 2013-07-31 23:44:47 +01:00
Dominik Picheta
e42f877664 Optimised `await` on custom async procs.
* The async proc is now invoked immediately after the 'await' request
is processed. This means that the iterator is executed until it yields
with a request which cannot be immediately satisfied (e.g. readLine)
after which point a
select() call will need to be made. Previously the async proc would only
be added to the list of requests, the event loop would need to start
over and thus would need to go
through each request again before invoking the iterator and reaching a
request which could not be immediately satisfied.
* This was implemented by recursively calling the processWorker
procedure in asyncio.nim. In the future it may be of benefit to also
perform this for reqReg.
* When the awaited user-defined async proc finishes, the parent worker
will be executed, the same optimisation has been used for that case.
* The way that this optimisation is implemented means that even multiple
levels of awaited user-defined procs will be quickly processed. (Look at
the tasyncitermacro test, the auth procs are arranged this way to test
this)
2013-07-31 23:27:07 +01:00
Dominik Picheta
70d67daa14 Introduced 'reg' keyword, and changed the await behaviour for custom
procs.
2013-07-31 19:28:37 +01:00
Dominik Picheta
ae64089e53 Fixed crash when worker has just been registered and didn't get executed
yet.
2013-07-28 18:28:50 +01:00
Dominik Picheta
426c2d6c01 Implemented default arguments for async procs. 2013-07-28 14:41:09 +01:00
Dominik Picheta
2330c95b22 Implemented genSym for async function calls. 2013-07-28 13:50:32 +01:00
Dominik Picheta
adc74f8ab0 Merge branch 'master' into asyncmacro 2013-07-25 23:24:30 +01:00
Dominik Picheta
3ead23cb26 Merge branch 'master' into asyncmacro 2013-07-25 22:16:06 +01:00
Dominik Picheta
c5ff9650f7 Changed some asyncio macro code to use the new genSym. 2013-07-25 22:13:37 +01:00
Dominik Picheta
2cedc1d1af Merge branch 'master' into asyncmacro 2013-07-25 20:31:47 +01:00
Dominik Picheta
3e7ea859ec Merge branch 'master' into asyncmacro 2013-07-25 19:33:08 +01:00
Dominik Picheta
d8659eaaa2 Another docs fix for asyncio. 2013-07-21 15:18:44 +01:00
Dominik Picheta
1e1a583903 Fixes problems with async docs. 2013-07-21 14:51:25 +01:00
Dominik Picheta
a1a8925692 Calling an async proc with invalid arguments will now fail at compile-time. 2013-07-21 14:28:02 +01:00
Dominik Picheta
d15e67e020 Better param passing. Implemented send, readLine and support for
documentation stubs.
2013-07-21 13:58:54 +01:00
Dominik Picheta
84ea209fe1 Merge branch 'master' into asyncmacro 2013-07-21 11:09:20 +01:00
Dominik Picheta
75893d1b36 First implementation of the async/await-like system with macros. 2013-07-20 23:42:10 +01:00
4 changed files with 948 additions and 57 deletions

View file

@ -6,7 +6,7 @@
# distribution, for details about the copyright.
#
import sockets, os
import sockets, os, macros, strutils
## This module implements an asynchronous event loop together with asynchronous sockets
## which use this event loop.
@ -31,10 +31,10 @@ import sockets, os
## on with the events. The type that you set userArg to must be inheriting from
## TObject!
##
## **Note:** If you want to provide async ability to your module please do not
## use the ``TDelegate`` object, instead use ``PAsyncSocket``. It is possible
## that in the future this type's fields will not be exported therefore breaking
## your code.
## **Note:** If you are after using this module to provide async functionality
## for one of your modules then it is best to use PAsyncSocket if your module
## only requires sockets. For non-socket objects with a select-like interface
## a TDelegate implementation should be created, if one doesn't already exist.
##
## **Warning:** The API of this module is unstable, and therefore is subject
## to change.
@ -109,10 +109,6 @@ type
PDelegate* = ref TDelegate
PDispatcher* = ref TDispatcher
TDispatcher = object
delegates: seq[PDelegate]
PAsyncSocket* = ref TAsyncSocket
TAsyncSocket* = object of TObject
socket: TSocket
@ -136,6 +132,52 @@ type
SockIdle, SockConnecting, SockConnected, SockListening, SockClosed,
SockUDPBound
TRequestKind* = enum
reqNil, reqReg, reqAwait, reqRead, reqWrite, reqReadLine, reqAccept, reqConnect
PAsyncProc = iterator (x: PRequest): PRequest
PRequest* = ref object
socket*: PAsyncSocket
case hasException*: bool
of true:
exc*: ref EBase
of false: nil
case kind*: TRequestKind
of reqNil:
nil
of reqReg, reqAwait:
param*: PObject
worker*: PAsyncProc
of reqRead:
count*: int ## Request
readData*: string ## Response
of reqWrite:
toWrite*: string ## Request
written: int ## Internal data
of reqReadLine:
line*: string ## Response
of reqAccept:
client*: PAsyncSocket ## Response
of reqConnect:
address*: string ## Request
port*: TPort ## Request
PWorker* = ref object
worker*: PAsyncProc
x*: PRequest
lastReq*: PRequest
case hasParent*: bool
of true:
parent*: PWorker
else: nil
PDispatcher* = ref TDispatcher
TDispatcher = object
usesDelegates: bool
delegates: seq[PDelegate]
requests: array[TRequestKind, seq[PWorker]]
proc newDelegate*(): PDelegate =
## Creates a new delegate.
new(result)
@ -382,9 +424,15 @@ proc accept*(server: PAsyncSocket): PAsyncSocket {.deprecated.} =
var address = ""
server.acceptAddr(result, address)
proc newDispatcher*(): PDispatcher =
proc newRequests(): array[TRequestKind, seq[PWorker]] =
for req in TRequestKind:
result[req] = @[]
proc newDispatcher*(useDelegates = true): PDispatcher =
new(result)
result.delegates = @[]
result.requests = newRequests()
result.usesDelegates = useDelegates
proc register*(d: PDispatcher, deleg: PDelegate) =
## Registers delegate ``deleg`` with dispatcher ``d``.
@ -569,6 +617,193 @@ proc select(readfds, writefds, exceptfds: var seq[PDelegate],
pruneSocketSet(writefds, (wr))
pruneSocketSet(exceptfds, (ex))
proc createFdSet(fd: var TFdSet, s: seq[PWorker], m: var int) =
FD_ZERO(fd)
for i in items(s):
m = max(m, int(i.lastReq.socket.getFD))
FD_SET(i.lastReq.socket.getFD, fd)
proc getWorkerSocket(worker: PWorker): PAsyncSocket =
worker.lastReq.socket
proc pruneSocketSet(s: var seq[PWorker], fd: var TFdSet) =
var i = 0
var L = s.len
while i < L:
if FD_ISSET(getWorkerSocket(s[i]).getFD, fd) != 0'i32:
s[i] = s[L-1]
dec(L)
else:
inc(i)
setLen(s, L)
proc select(readfds: var seq[PWorker], writefds: var seq[PWorker],
timeout = 500): int =
var tv {.noInit.}: TTimeVal = timeValFromMilliseconds(timeout)
var rd, wr, ex: TFdSet
var m = 0
createFdSet(rd, readfds, m)
createFdSet(wr, writefds, m)
#createFdSet(ex, @[], m)
if timeout != -1:
result = int(select(cint(m+1), addr(rd), addr(wr), addr(ex), addr(tv)))
else:
result = int(select(cint(m+1), addr(rd), addr(wr), addr(ex), nil))
pruneSocketSet(readfds, (rd))
pruneSocketSet(writefds, (wr))
#pruneSocketSet(exceptfds, (ex))
proc register*(disp: PDispatcher, worker: iterator (x: PRequest): PRequest,
param: PObject) =
#= PWorker(socket, worker, PRequest(kind: reqNil))
assert(not disp.usesDelegates,
"You need to set ``usesDelegates`` to false in the newDispatcher proc.")
var req = PWorker(
worker: worker,
x: PRequest(kind: reqReg, param: param, worker: nil),
lastReq: PRequest(kind: reqNil)
)
disp.requests[reqNil].add(req)
proc processWorkers(d: PDispatcher) =
var newRequests: array[TRequestKind, seq[PWorker]] = d.requests
newRequests[reqNil] = @[]
proc processWorker(idle: PWorker) =
let req = idle.worker(idle.x)
if req != nil:
case req.kind
of reqReg:
# Reset the lastReq. This caused odd issues in an 'accept' loop.
# My guess is that the request object got corrupted somehow, but
# I could not determine the cause. The issue that occurred was that
# 'await accept' simply stopped working after 6 clients (when connecting
# 10 simultaneously as seen in the tasyncitermacro test).
idle.lastReq = PRequest(kind: reqNil)
newRequests[reqNil].add(idle)
let newWorker = PWorker(worker: req.worker, lastReq: PRequest(kind: reqNil),
x: req)
newRequests[reqNil].add(newWorker)
of reqAwait:
let newWorker = PWorker(worker: req.worker, lastReq: PRequest(kind: reqNil),
x: req, hasParent: true, parent: idle)
# The worker which ``await``-ed this user-defined async proc;
# will be re-added to ``d.requests`` when ``newWorker`` finishes.
# We call this proc recursively so that the execution of this
# worker begins and we can immediately satisfy it's async request.
# We do not need to add newWorker to newRequests manually as it will be
# added by ``processWorker`` if necessary automatically.
processWorker(newWorker)
else:
idle.lastReq = req
newRequests[req.kind].add(idle)
else:
assert idle.worker.finished
if idle.hasParent:
# Re-add the parent worker, which is the worker which awaited this
# user-defined async proc which just finished. Do this by calling
# processWorker recursively. Same way as above.
processWorker(idle.parent)
for idle in d.requests[reqNil]:
processWorker(idle)
d.requests = newRequests
proc populateRead(d: PDispatcher): seq[PWorker] =
result = @[]
result.add d.requests[reqAccept]
result.add d.requests[reqRead]
result.add d.requests[reqReadLine]
proc populateWrite(d: PDispatcher): seq[PWorker] =
result = @[]
result.add d.requests[reqWrite]
result.add d.requests[reqConnect]
proc processRequests(requests, readWorkers, writeWorkers: seq[PWorker],
newRequests: var array[TRequestKind, seq[PWorker]]) =
for worker in requests:
#echo(worker.lastReq.kind)
var addTo = worker.lastReq.kind
template execReq(workers: var seq[PWorker], autoadd: bool,
body: stmt) {.immediate, dirty.} =
if worker notin workers:
# Worker is ready to read. Let's read.
try:
body
except:
worker.lastReq.hasException = true
worker.lastReq.exc = getCurrentException()
finally:
if autoAdd:
addTo = reqNil
case worker.lastReq.kind
of reqReadLine:
execReq readWorkers, false:
if worker.lastReq.socket.readLine(worker.lastReq.line):
addTo = reqNil
of reqAccept:
execReq readWorkers, true:
worker.lastReq.client = newAsyncSocket()
worker.lastReq.socket.accept(worker.lastReq.client)
of reqRead:
# We guarantee that all requested data will be read.
execReq readWorkers, false:
proc doRead(count: int) =
let got = worker.lastReq.socket.recvAsync(
worker.lastReq.readData, count)
assert got != -1
if got == count:
addTo = reqNil # Everything has been read
if worker.lastReq.readData.len == 0:
doRead(worker.lastReq.count)
else:
doRead(worker.lastReq.count-worker.lastReq.readData.len)
of reqWrite:
# We guarantee that all the data that is requested to be sent, will
# be sent.
execReq writeWorkers, false:
let written = worker.lastReq.written
proc doSend(toWrite: string) =
let len = toWrite.len
let sent = worker.lastReq.socket.sendAsync(toWrite)
assert sent != 0 # /Something/ should have been written.
if sent == len:
# Sent all data, request complete.
addTo = reqNil
else:
# Didn't send all data, must send the rest later.
worker.lastReq.written.inc(sent)
if written == 0:
doSend(worker.lastReq.toWrite)
else:
let toWrite = worker.lastReq.toWrite[written .. -1]
doSend(toWrite)
of reqConnect:
if not worker.lastReq.socket.isConnecting:
worker.lastReq.socket.connect(worker.lastReq.address, worker.lastReq.port)
else:
execReq writeWorkers, true:
worker.lastReq.socket.info = SockConnected
of reqReg, reqAwait:
assert false, $worker.lastReq.kind & " should have been processed already"
of reqNil:
# Nothing to do. Most likely that a new worker has just been
# registered.
newRequests[addTo].add(worker)
proc poll*(d: PDispatcher, timeout: int = 500): bool =
## This function checks for events on all the delegates in the `PDispatcher`.
## It then proceeds to call the correct event handler.
@ -582,58 +817,419 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool =
## only be executed after one or more file descriptors becomes readable or
## writeable.
result = true
var readDg, writeDg, errorDg: seq[PDelegate] = @[]
var len = d.delegates.len
var dc = 0
if d.usesDelegates:
var readDg, writeDg, errorDg: seq[PDelegate] = @[]
var len = d.delegates.len
var dc = 0
while dc < len:
let deleg = d.delegates[dc]
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and deleg.open:
readDg.add(deleg)
if (deleg.mode != fmRead) and deleg.open:
writeDg.add(deleg)
if deleg.open:
errorDg.add(deleg)
inc dc
else:
# File/socket has been closed. Remove it from dispatcher.
d.delegates[dc] = d.delegates[len-1]
dec len
while dc < len:
let deleg = d.delegates[dc]
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and deleg.open:
readDg.add(deleg)
if (deleg.mode != fmRead) and deleg.open:
writeDg.add(deleg)
if deleg.open:
errorDg.add(deleg)
inc dc
else:
# File/socket has been closed. Remove it from dispatcher.
d.delegates[dc] = d.delegates[len-1]
dec len
d.delegates.setLen(len)
d.delegates.setLen(len)
var hasDataBufferedCount = 0
for d in d.delegates:
if d.hasDataBuffered(d.deleVal):
hasDataBufferedCount.inc()
d.handleRead(d.deleVal)
if hasDataBufferedCount > 0: return True
var hasDataBufferedCount = 0
for d in d.delegates:
if d.hasDataBuffered(d.deleVal):
hasDataBufferedCount.inc()
d.handleRead(d.deleVal)
if hasDataBufferedCount > 0: return True
if readDg.len() == 0 and writeDg.len() == 0:
## TODO: Perhaps this shouldn't return if errorDg has something?
return False
if readDg.len() == 0 and writeDg.len() == 0:
## TODO: Perhaps this shouldn't return if errorDg has something?
return False
if select(readDg, writeDg, errorDg, timeout) != 0:
for i in 0..len(d.delegates)-1:
if i > len(d.delegates)-1: break # One delegate might've been removed.
let deleg = d.delegates[i]
if not deleg.open: continue # This delegate might've been closed.
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and
deleg notin readDg:
deleg.handleRead(deleg.deleVal)
if (deleg.mode != fmRead) and deleg notin writeDg:
deleg.handleWrite(deleg.deleVal)
if deleg notin errorDg:
deleg.handleError(deleg.deleVal)
if select(readDg, writeDg, errorDg, timeout) != 0:
for i in 0..len(d.delegates)-1:
if i > len(d.delegates)-1: break # One delegate might've been removed.
let deleg = d.delegates[i]
if not deleg.open: continue # This delegate might've been closed.
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and
deleg notin readDg:
deleg.handleRead(deleg.deleVal)
if (deleg.mode != fmRead) and deleg notin writeDg:
deleg.handleWrite(deleg.deleVal)
if deleg notin errorDg:
deleg.handleError(deleg.deleVal)
# Execute tasks
for i in items(d.delegates):
i.task(i.deleVal)
# Execute tasks
for i in items(d.delegates):
i.task(i.deleVal)
else:
# Async worker iterators
processWorkers(d)
var readWorkers = populateRead(d)
var writeWorkers = populateWrite(d)
#echo("ReadWorkers: ", readWorkers.len, " | WriteWorkers: ", writeWorkers.len,
# " | Idle Workers: ", d.requests[reqNil].len,
# " | Workers waiting for readLine: ", d.requests[reqReadLine].len,
# " | Workers waiting for accept: ", d.requests[reqAccept].len)
# Check buffer state.
var isDataBuffered = false
var newReadWorkers: seq[PWorker] = @[]
for i in readWorkers:
if not i.lastReq.socket.hasDataBuffered(): newReadWorkers.add(i)
else: isDataBuffered = true
readWorkers = newReadWorkers
#echo(if isDataBuffered: "Buffered" else: "Not Buffered")
var doProcess = isDataBuffered
if not doProcess: doProcess = select(readWorkers, writeWorkers, timeout) != 0
if doProcess:
var newRequests: array[TRequestKind, seq[PWorker]] = newRequests()
for req in TRequestKind:
processRequests(d.requests[req], readWorkers, writeWorkers, newRequests)
d.requests = newRequests
proc len*(disp: PDispatcher): int =
## Retrieves the amount of delegates in ``disp``.
return disp.delegates.len
# ---- Async macro
proc createRequestNode(varName,
reqArgs: string, sym: var PNimrodNode): PNimrodNode {.compiletime.} =
## Creates a constructor for a PRequest object, which will be then yielded.
result = newNimNode(nnkStmtList)
# TODO: Using gensym here causes segfaults because hasException is not
# initialised.
sym = newIdentNode(varName) #genSym(nskVar, varName)
var reqObj = newVarStmt(sym,
parseExpr("PRequest($#)" % [reqArgs]))
result.add reqObj
result.add newNimNode(nnkYieldStmt).add(sym)
# Check for exception
# TODO: Create a custom exception type, create a field which will store
# the original stack trace as given by getStackTrace. Maybe there is a way
# to override the stack trace? That'd be nice.
result.add newIfStmt(
(newDotExpr(sym, newIdentNode("hasException")),
newNimNode(nnkRaiseStmt).add(
newDotExpr(sym, newIdentNode("exc")))))
proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} =
## Transforms a var/let section
## E.g:
## let client = await(accept(server))
result = @[]
let nameIdent = n[0][0] # Var name
expectLen(n[0], 3) # IdentDefs
let insideAwait = n[0][2][1]
let reqCall = $insideAwait[0].ident
expectLen(insideAwait, 2)
let sockName = $insideAwait[1].ident
case reqCall.normalize
of "accept":
let acceptReqVar = "acceptReq"
var sym: PNimrodNode
result.add createRequestNode(acceptReqVar,
"socket: $#, kind: reqAccept, client: nil, hasException: false" % sockName, sym)
case n.kind
of nnkLetSection:
result.add newLetStmt(nameIdent, newDotExpr(sym, newIdentNode("client")))
of nnkVarSection:
result.add newVarStmt(nameIdent, newDotExpr(sym, newIdentNode("client")))
else: error "Bad node kind in toYieldVar"
of "readline":
let readReqVar = "readLineReq"
var sym: PNimrodNode
result.add createRequestNode(readReqVar,
"socket: $#, kind: reqReadLine, line: \"\"" % sockName, sym)
case n.kind
of nnkLetSection:
result.add newLetStmt(nameIdent, newDotExpr(sym, newIdentNode("line")))
of nnkVarSection:
result.add newVarStmt(nameIdent, newDotExpr(sym, newIdentNode("line")))
else: error "Bad node kind in toYieldVar"
else:
error(reqCall & " is not a valid async call")
const typeDef =
"""
type
P$#ArgObject = ref object of TObject
"""
proc transformCallWithArg(call: PNimrodNode,
sym: var PNimrodNode): PNimrodNode {.compiletime.} =
## Transforms an async await call of a user-defined proc into a
## ``reqReg`` yield.
result = newNimNode(nnkStmtList)
sym = gensym(nskVar, "argsToPass")
# Add in a call to a pre-generated proc stub so that the compiler verifies
# the params for us :)
# TODO: Inline this maybe?
var args: seq[PNimrodNode] = @[]
for i in 1 .. call[1].len-1:
args.add(call[1][i])
result.add newVarStmt(sym, newCall(call[1][0], args))
proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} =
## Transforms a call/command
if $n[0].ident != "await": error "'await' expected"
result = @[]
let callIdent = $n[1][0].ident
case callIdent.normalize
of "send":
let socketName = $n[1][1].ident
let toWrite = n[1][2]
var sym: PNimrodNode
result.add createRequestNode("sendReq",
"socket: $#, kind: reqWrite, toWrite: $#" %
[socketName, $(toWrite.toStrLit)], sym)
of "connect":
let socketName = $n[1][1].ident
let address = n[1][2]
let port = n[1][3]
var sym: PNimrodNode
result.add createRequestNode("connectReq",
"socket: $#, kind: reqConnect, address: $#, port: $#" %
[socketName, $(address.toStrLit), $(port.toStrLit)], sym)
else:
var sym: PNimrodNode
result.add(transformCallWithArg(n, sym))
# reqCustom
var yie = parseExpr("yield PRequest(socket: nil, kind: reqAwait, worker: $#)" %
[callIdent])
yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"),
sym))
result.add yie
proc toYieldReg(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} =
## Transforms the 'reg' command to a reqReg yield request.
if $n[0].ident != "reg": error "'reg' expected"
result = @[]
let callIdent = $n[1][0].ident
var sym: PNimrodNode
result.add(transformCallWithArg(n, sym))
# reqRegister
var yie = parseExpr("yield PRequest(socket: nil, kind: reqReg, worker: $#)" %
[callIdent])
yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"),
sym))
result.add yie
proc transform(n: PNimrodNode): PNimrodNode {.compiletime.} =
## Transforms body.
## Specifically it does the following:
##
## * Looks for 'await' and transforms it into a yield.
## * Handles arguments correctly.
result = newNimNode(nnkStmtList)
expectKind(n, nnkStmtList)
for i in 0 .. n.len-1:
var son = n[i]
case son.kind
of nnkVarSection, nnkLetSection:
for defs in 0 .. son.len-1:
var doAdd = true
let identDefs = son[defs]
expectKind(identDefs, nnkIdentDefs)
if identDefs[2].kind == nnkCall:
let callIdent = identDefs[2][0]
expectKind(callIdent, nnkIdent)
if $callIdent.ident == "await":
# Transform into yield.
result.add(toYieldVar(son))
doAdd = false
if doAdd:
var letOrVarSection = newNimNode(son.kind)
letOrVarSection.add(identDefs)
result.add(letOrVarSection)
of nnkWhileStmt:
son[1] = transform(son[1])
result.add(son)
of nnkForStmt:
son[2] = transform(son[2])
result.add(son)
of nnkCall, nnkCommand:
if son[0].kind == nnkIdent:
case $son[0].ident
of "await":
result.add toYieldCall(son)
of "reg":
result.add toYieldReg(son)
else:
result.add son
else:
result.add son
else:
result.add(son)
proc transformArgs(procName: string,
formalParams: PNimrodNode): PNimrodNode {.compiletime.} =
## Transforms formal params into a typedef with a dummy type
## ``ref object of TObject``. This is inserted above the proc definition.
expectKind(formalParams, nnkFormalParams)
result = parseStmt(typeDef % procName)
var RecList = newNimNode(nnkRecList)
for i in 1 .. formalParams.len-1:
expectKind(formalParams[i], nnkIdentDefs)
RecList.add(newIdentDefs(newIdentNode("dummy" & $i), formalParams[i][1]))
# TODO: Add comment with the original param name?
result[0][0][2][0][2] = RecList
proc declareArgsInBody(procName: string,
formalParams: PNimrodNode): PNimrodNode {.compiletime.} =
## Creates a local immutable var by casting the PRequest.param.
## Immutable vars are then defined as specified in the proc's params.
result = newNimNode(nnkStmtList)
var sym = genSym(ident = "passedInParams")
result.add newLetStmt(sym, parseExpr("P$#ArgObject(x.param)" % procName))
# Fields take the form ``dummy<i>``.
for i in 1 .. formalParams.len-1:
expectKind(formalParams[i], nnkIdentDefs)
result.add newLetStmt(formalParams[i][0],
newDotExpr(sym, newIdentNode("dummy" & $i)))
proc isDocumentation(n: PNimrodNode): bool {.compiletime.} =
## Determines whether this proc def is a docs stub.
result = true
for i in 0 .. n[6].len-1:
if n[6][i].kind != nnkCommentStmt:
return false
proc createVerificationProc(procName: PNimrodNode,
formalParams: PNimrodNode): PNimrodNode {.compiletime.} =
# TODO: Export this stub if our async proc is exported?
# Generate list of parameters for the proc. First param is the return type.
var params: seq[PNimrodNode] = @[newIdentNode("P$#ArgObject" % $procName.ident)]
for i in 1 .. formalParams.len-1:
params.add(formalParams[i])
# Generate body. We construct the ArgObject here, this is done so that
# default variables of the async proc can be captured.
var body = newNimNode(nnkStmtList)
if formalParams.len > 1:
body.add newCall("new", newIdentNode("result"))
else:
body.add parseExpr("nil")
for i in 1 .. formalParams.len-1:
let dotExpr = newDotExpr(newIdentNode("result"),
newIdentNode("dummy" & $i))
body.add newAssignment(dotExpr, formalParams[i][0])
result = newProc(procName, params, body)
macro async*(n: stmt): stmt {.immediate.} =
expectKind(n, nnkProcDef)
#echo(treeRepr(n))
if n.isDocumentation():
# Documentation stub?
# TODO: Give it an async tag?
# TODO: Doc strings are not generated in doc2
result = n
result[6] = parseStmt("nil")
return
result = newNimNode(nnkIteratorDef)
for i in 0 .. n.len-1:
result.add(copyNimTree(n[i]))
# Populate ``FormalParams``
assert result[3].kind == nnkFormalParams
let formalParams = newNimNode(nnkFormalParams)
formalParams.add(newIdentNode(!"PRequest")) # Return type
var params = newNimNode(nnkIdentDefs)
params.add(newIdentNode(!"x")) # First param name
params.add(newIdentNode(!"PRequest")) # First param type
params.add(newNimNode(nnkEmpty))
formalParams.add(params)
result[3] = formalParams
# Closure pragma
result[4].add(newIdentNode(!"closure"))
# Body
result[6] = newNimNode(nnkStmtList)
# Declare variables based on the params that the async proc takes.
# i.e. extract them from the PRequest object.
if n[3].len > 1:
let args = declareArgsInBody($n[0].ident, n[3])
result[6].add(args)
# Transform body
var body = transform(n[6])
result[6].add(body)
# Add typedef above the proc def for parameters.
let procDef = copyNimTree(result)
result = newNimNode(nnkStmtList)
result.add(transformArgs($n[0].ident, n[3]))
# Generate a proc to verify that the user passes the correct params.
# The proc also constructs the ArgObject, this is so that default params
# can be captured into the ArgObject.
result.add createVerificationProc(n[0], n[3])
result.add procDef
#echo treeRepr(result)
echo result.toStrLit().strVal
# ---- Async macro end
## Asynchronous IO without callbacks
## =====================
## The implementation of this is similar to what is known as ``futures``
## in other languages. The ``await`` keyword is used to call a procedure
## which is marked with an ``{.async.}`` pragma. The ``async`` procedures
## get transformed into an iterator. When you execute a procedure which may
## block using ``await`` then the execution of your procedure will stop
## until the procedure you executed has completed. While your procedure's
## execution is stopped, other async procedures can continue executing.
##
## The built-in functions which can be used using ``await`` are listed below:
##
## .. code-block:: nimrod
##
## proc send*(socket: PAsyncSocket, text: string) {.async.} =
##
## Sends ``text`` to ``socket`` asynchronously.
##
## .. code-block:: nimrod
##
## proc accept*(socket: PAsyncSocket): PAsyncSocket {.async.} =
##
## Accepts a client connecting to a server socket asynchronously.
## Returns that client.
##
## You may also define your own async procedures by annotating them with
## the ``{.async.}`` pragma, please note however that currently you cannot
## overload these procedures.
discard """ when defined(nimdoc) and isMainModule:
proc send*(socket: PAsyncSocket, text: string) {.async.} =
## Sends ``text`` to ``socket`` asynchronously.
proc accept*(socket: PAsyncSocket): PAsyncSocket {.async.} =
## Accepts a client connecting to a server socket asynchronously.
## Returns that client. """
when isMainModule:
proc testConnect(s: PAsyncSocket, no: int) =

View file

@ -0,0 +1,32 @@
discard """
msg: "type mismatch: got (int literal(234)) but expected 'string'"
"""
import sockets, asyncio, strutils
proc processRequest(client: PAsyncSocket, test: string) {.async.} =
assert test == "ahha"
assert client != nil
echo("Test = ", test)
let line = await(readLine(client))
echo("Read: ", line)
await send(client, "Goodbye.\c\L")
client.close()
proc processServer() {.async.} =
var sock = AsyncSocket()
# blocks:
sock.bindAddr(TPort(6667))
sock.listen()
# Accept loop
while true:
let client: PAsyncSocket = await(accept(sock))
assert client != nil
await processRequest(client, 234)
var disp = newDispatcher(false)
disp.register(processServer, nil)
while true:
discard disp.poll()

View file

@ -0,0 +1,78 @@
discard """
file: "tasyncitermacro.nim"
cmd: "nimrod cc --hints:on $# $#"
output: "10000"
"""
import sockets, asyncio, strutils
var globalCount = 0
proc sendCount3(client: PAsyncSocket, count: int) {.async.} =
await send(client, $count & "\c\L")
proc sendCount2(client: PAsyncSocket, count: int) {.async.} =
await sendCount3(client, count)
proc sendCount(client: PAsyncSocket, count: int) {.async.} =
# Testing multiple levels of custom await calls. Shouldn't have much of a
# performance impact.
await sendCount2(client, count)
proc processRequest(client: PAsyncSocket, paramTest: string, closeSock: bool = true) {.async.} =
doAssert paramTest == "foobarbaz"
assert client != nil
var count = 0
while true:
let line = await(readLine(client))
if line == "end": break
doAssert line.startswith("Message")
doAssert line == "Message" & $count
count.inc
await sendCount(client, count)
doAssert closeSock # Testing default params here. We won't be changing it.
if closeSock:
client.close()
proc startServer() {.async.} =
var sock = AsyncSocket()
#sock.setReuseAddr()
# The following may block, but I don't think it does.
sock.bindAddr(TPort(10235))
sock.listen()
# Accept loop
while true:
let client: PAsyncSocket = await(accept(sock))
assert client != nil
reg processRequest(client, "foobarbaz")
proc spawnClient() {.async.} =
var client = AsyncSocket()
await connect(client, "localhost", TPort(10235))
for i in 0 .. 99:
await send(client, "Message" & $i & "\c\L")
await send(client, "end\c\L")
let line = await(readLine(client))
doAssert line == "100"
globalCount.inc(100)
client.close()
proc spawnClients() {.async.} =
for i in 0 .. 99:
reg spawnClient()
var disp = newDispatcher(false)
disp.register(startServer, nil)
disp.register(spawnClients, nil)
while true:
discard disp.poll()
if globalCount == 100*100:
echo(globalCount)
quit(QuitSuccess)

185
tests/run/tasynciterraw.nim Normal file
View file

@ -0,0 +1,185 @@
discard """
file: "tasynciterraw.nim"
cmd: "nimrod cc --hints:on $# $#"
output: "10000"
"""
# This is a control test in case the macro expansion starts failing. It's the
# raw code after macro expansion in tasyncitermacro.nim.
import asyncio, sockets, strutils
var globalCount = 0
type
PsendCount3ArgObject = ref object of TObject
dummy1: PAsyncSocket
dummy2: int
proc sendCount3(client: PAsyncSocket; count: int): PsendCount3ArgObject =
new(result)
result.dummy1 = client
result.dummy2 = count
iterator sendCount3(x: PRequest): PRequest {.closure.} =
let passedInParams = PsendCount3ArgObject(x.param)
let client = passedInParams.dummy1
let count = passedInParams.dummy2
var sendReq = PRequest(socket: client, kind: reqWrite,
toWrite: $ count & "\x0D\x0A")
yield sendReq
if sendReq.hasException: raise sendReq.exc
type
PsendCount2ArgObject = ref object of TObject
dummy1: PAsyncSocket
dummy2: int
proc sendCount2(client: PAsyncSocket; count: int): PsendCount2ArgObject =
new(result)
result.dummy1 = client
result.dummy2 = count
iterator sendCount2(x: PRequest): PRequest {.closure.} =
let passedInParams = PsendCount2ArgObject(x.param)
let client = passedInParams.dummy1
let count = passedInParams.dummy2
var argsToPass = sendCount3(client, count)
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount3,
param: argsToPass)
type
PsendCountArgObject = ref object of TObject
dummy1: PAsyncSocket
dummy2: int
proc sendCount(client: PAsyncSocket; count: int): PsendCountArgObject =
new(result)
result.dummy1 = client
result.dummy2 = count
iterator sendCount(x: PRequest): PRequest {.closure.} =
let passedInParams = PsendCountArgObject(x.param)
let client = passedInParams.dummy1
let count = passedInParams.dummy2
var argsToPass = sendCount2(client, count)
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount2,
param: argsToPass)
type
PprocessRequestArgObject = ref object of TObject
dummy1: PAsyncSocket
dummy2: string
dummy3: bool
proc processRequest(client: PAsyncSocket; paramTest: string;
closeSock: bool = true): PprocessRequestArgObject =
new(result)
result.dummy1 = client
result.dummy2 = paramTest
result.dummy3 = closeSock
iterator processRequest(x: PRequest): PRequest {.closure.} =
let passedInParams = PprocessRequestArgObject(x.param)
let client = passedInParams.dummy1
let paramTest = passedInParams.dummy2
let closeSock = passedInParams.dummy3
doAssert paramTest == "foobarbaz"
assert client != nil
var count = 0
while true:
var readLineReq = PRequest(socket: client, kind: reqReadLine, line: "")
yield readLineReq
if readLineReq.hasException: raise readLineReq.exc
let line = readLineReq.line
if line == "end":
break
doAssert line.startswith("Message")
doAssert line == "Message" & $ count
count.inc
var argsToPass = sendCount(client, count)
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount,
param: argsToPass)
doAssert closeSock
if closeSock:
client.close()
type
PstartServerArgObject = ref object of TObject
proc startServer(): PstartServerArgObject =
nil
iterator startServer(x: PRequest): PRequest {.closure.} =
var sock = AsyncSocket()
#sock.setReuseAddr()
sock.bindAddr(TPort(10235))
sock.listen()
while true:
var acceptReq = PRequest(socket: sock, kind: reqAccept, client: nil,
hasException: false)
yield acceptReq
if acceptReq.hasException: raise acceptReq.exc
let client = acceptReq.client
assert client != nil
var argsToPass = processRequest(client, "foobarbaz")
yield PRequest(socket: nil, kind: reqReg, worker: processRequest,
param: argsToPass)
type
PspawnClientArgObject = ref object of TObject
proc spawnClient(): PspawnClientArgObject =
nil
iterator spawnClient(x: PRequest): PRequest {.closure.} =
var client = AsyncSocket()
var connectReq = PRequest(socket: client, kind: reqConnect,
address: "localhost", port: TPort(10235))
yield connectReq
if connectReq.hasException: raise connectReq.exc
for i in 0 .. 99:
var sendReq = PRequest(socket: client, kind: reqWrite,
toWrite: "Message" & $ i & "\x0D\x0A")
yield sendReq
if sendReq.hasException: raise sendReq.exc
var sendReq1 = PRequest(socket: client, kind: reqWrite,
toWrite: "end\x0D\x0A")
yield sendReq1
if sendReq1.hasException: raise sendReq1.exc
var readLineReq = PRequest(socket: client, kind: reqReadLine, line: "")
yield readLineReq
if readLineReq.hasException: raise readLineReq.exc
let line = readLineReq.line
doAssert line == "100"
globalCount.inc(100)
client.close()
type
PspawnClientsArgObject = ref object of TObject
proc spawnClients(): PspawnClientsArgObject =
nil
iterator spawnClients(x: PRequest): PRequest {.closure.} =
for i in 0 .. 99:
var argsToPass = spawnClient()
yield PRequest(socket: nil, kind: reqReg, worker: spawnClient,
param: argsToPass)
var disp = newDispatcher(false)
disp.register(startServer, nil)
disp.register(spawnClients, nil)
while true:
discard disp.poll()
if globalCount == 100*100:
echo(globalCount)
quit(QuitSuccess)