Compare commits
18 commits
devel
...
asyncmacro
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b3469c43a6 | ||
|
|
ab534804c6 | ||
|
|
e42f877664 | ||
|
|
70d67daa14 | ||
|
|
ae64089e53 | ||
|
|
426c2d6c01 | ||
|
|
2330c95b22 | ||
|
|
adc74f8ab0 | ||
|
|
3ead23cb26 | ||
|
|
c5ff9650f7 | ||
|
|
2cedc1d1af | ||
|
|
3e7ea859ec | ||
|
|
d8659eaaa2 | ||
|
|
1e1a583903 | ||
|
|
a1a8925692 | ||
|
|
d15e67e020 | ||
|
|
84ea209fe1 | ||
|
|
75893d1b36 |
4 changed files with 948 additions and 57 deletions
|
|
@ -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,6 +817,7 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool =
|
|||
## only be executed after one or more file descriptors becomes readable or
|
||||
## writeable.
|
||||
result = true
|
||||
if d.usesDelegates:
|
||||
var readDg, writeDg, errorDg: seq[PDelegate] = @[]
|
||||
var len = d.delegates.len
|
||||
var dc = 0
|
||||
|
|
@ -629,11 +865,371 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool =
|
|||
# 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) =
|
||||
|
|
|
|||
32
tests/reject/tasyncmacroinvalidargs.nim
Normal file
32
tests/reject/tasyncmacroinvalidargs.nim
Normal 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()
|
||||
78
tests/run/tasyncitermacro.nim
Normal file
78
tests/run/tasyncitermacro.nim
Normal 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
185
tests/run/tasynciterraw.nim
Normal 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)
|
||||
Loading…
Add table
Add a link
Reference in a new issue