diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 4ff6e0ced..99b0eecb0 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -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 - - 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) - - 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 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) + 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 + + 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 + + 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) + + # 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``. + 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) = diff --git a/tests/reject/tasyncmacroinvalidargs.nim b/tests/reject/tasyncmacroinvalidargs.nim new file mode 100644 index 000000000..9a4d02627 --- /dev/null +++ b/tests/reject/tasyncmacroinvalidargs.nim @@ -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() diff --git a/tests/run/tasyncitermacro.nim b/tests/run/tasyncitermacro.nim new file mode 100644 index 000000000..f638b0e5a --- /dev/null +++ b/tests/run/tasyncitermacro.nim @@ -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) diff --git a/tests/run/tasynciterraw.nim b/tests/run/tasynciterraw.nim new file mode 100644 index 000000000..6cedc63d9 --- /dev/null +++ b/tests/run/tasynciterraw.nim @@ -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)