From 75893d1b3629e9dbea1f6cebbf35d6ffc1b9dbf6 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sat, 20 Jul 2013 23:42:10 +0100 Subject: [PATCH 01/12] First implementation of the async/await-like system with macros. --- lib/pure/asyncio.nim | 484 ++++++++++++++++++++++++++---- tests/compile/tasyncitermacro.nim | 24 ++ tests/compile/tasynciterraw.nim | 69 +++++ 3 files changed, 520 insertions(+), 57 deletions(-) create mode 100644 tests/compile/tasyncitermacro.nim create mode 100644 tests/compile/tasynciterraw.nim diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 4ff6e0ced..26df1cdfa 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,43 @@ type SockIdle, SockConnecting, SockConnected, SockListening, SockClosed, SockUDPBound + TRequestKind* = enum + reqNil, reqReg, reqRead, reqWrite, reqReadLine, reqAccept + + PRequest* = ref object + socket*: PAsyncSocket + case hasException*: bool + of true: + exc*: ref EBase + of false: nil + case kind*: TRequestKind + of reqNil: + nil + of reqReg: + param*: PObject + worker*: iterator (x: PRequest): PRequest + 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 + + PWorker* = ref object + worker: iterator (x: PRequest): PRequest {.closure.} + x: PRequest + lastReq: PRequest + + 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 +415,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 +608,87 @@ 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 pruneSocketSet(s: var seq[PWorker], fd: var TFdSet) = + var i = 0 + var L = s.len + while i < L: + if FD_ISSET(s[i].lastReq.socket.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] = @[] + for idle in d.requests[reqNil]: + let req = idle.worker(idle.x) + if req != nil: + if req.kind == reqReg: + newRequests[reqNil].add(idle) + let newWorker = PWorker(worker: req.worker, lastReq: PRequest(kind: reqNil), + x: req) + newRequests[reqNil].add(newWorker) + else: + idle.lastReq = req + newRequests[req.kind].add(idle) + else: + assert idle.worker.finished + d.requests = newRequests + +template popu(req) {.immediate, dirty.} = + for i in d.requests[req]: + result.add(i) + +proc populateRead(d: PDispatcher): seq[PWorker] = + result = @[] + + popu(reqRead) + popu(reqReadLine) + popu(reqAccept) + +proc populateWrite(d: PDispatcher): seq[PWorker] = + result = @[] + popu(reqWrite) + 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 +702,308 @@ 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.len, " ", d.requests[reqNil].len, d.requests[reqReadLine].len) + if select(readWorkers, writeWorkers, timeout) != 0: + var newRequests: array[TRequestKind, seq[PWorker]] = newRequests() + for req in TRequestKind: + for worker in d.requests[req]: + echo(req) + var addTo = req + 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 req + 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 reqReg: + assert false, "reqReg should have been processed already" + of reqNil: + assert false, "reqNil should have been processed already" + + newRequests[addTo].add(worker) + d.requests = newRequests proc len*(disp: PDispatcher): int = ## Retrieves the amount of delegates in ``disp``. return disp.delegates.len +# ---- Async macro + +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].ident # Var name + expectLen(n[0], 3) # IdentDefs + let insideAwait = n[0][2][1] + let reqCall = $insideAwait[0].ident + case reqCall + of "accept": + expectLen(insideAwait, 2) + let sockName = $insideAwait[1].ident + let acceptReqVar = "acceptReq" + # TODO: Random var names which do not conflict. or wait for gensym? + var reqObj = parseExpr( + """var $# = PRequest(socket: $#, kind: reqAccept, + client: nil)""" % + [acceptReqVar, sockName]) + result.add reqObj + result.add parseExpr("yield $#" % [acceptReqVar]) + # Check for exception + result.add parseStmt("if $1.hasException: raise $1.exc" % [acceptReqVar]) + case n.kind + of nnkLetSection: + result.add parseExpr("let $# = $#.client" % [$nameIdent, acceptReqVar]) + of nnkVarSection: + result.add parseExpr("var $# = $#.client" % [$nameIdent, acceptReqVar]) + else: error "Bad node kind in toYieldVar" + else: + error(reqCall & " is not a valid async call") + +const typeDef = + """ + type + PArgObject = ref object of TObject + """ + +proc transformCallWithArg(call: PNimrodNode): PNimrodNode {.compiletime.} = + result = parseStmt(typeDef) + var RecList = newNimNode(nnkRecList) + for i in 1 .. call[1].len-1: + var typeOf = newNimNode(nnkTypeOfExpr).add(newNimNode(nnkPar)) + typeOf[0].add(call[1][i]) + case call[1][i].kind + of nnkIdent: + RecList.add newIdentDefs(call[1][i], typeOf) + of nnkLiterals: + RecList.add newIdentDefs(newIdentNode("dummy" & $i), typeOf) + else: assert false + + result[0][0][2][0][2] = RecList + result.add parseExpr("var argsToPass: PArgObject") + result.add parseExpr("new argsToPass") + + for i in 1 .. call[1].len-1: + case call[1][i].kind + of nnkIdent: + result.add parseExpr("argsToPass.$1 = $1" % [$call[1][i].ident]) + of nnkLiterals: + let dotExpr = newDotExpr(newIdentNode("argsToPass"), + newIdentNode("dummy" & $i)) + result.add newAssignment(dotExpr, call[1][i]) + else: assert false + +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 + of "send": error("TODO") + else: + result.add(transformCallWithArg(n)) + # reqRegister + result.add parseExpr("yield PRequest(socket: $#, kind: reqReg, worker: $#, param: argsToPass)" % + [$n[1][1].ident, callIdent]) + +proc transform(n: PNimrodNode): PNimrodNode {.compiletime.} = + 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 nnkCall, nnkCommand: + if son[0].kind == nnkIdent and $son[0].ident == "await": + result.add toYieldCall(son) + else: + result.add son + else: + result.add(son) + +proc transformArgs(formalParams: PNimrodNode): PNimrodNode {.compiletime.} = + ## Transforms formal params into a typedef with a dummy type + ## ``ref object of TObject``. The PRequest.param is then casted to it, + ## and immutable vars are defined as specified the proc's params. + expectKind(formalParams, nnkFormalParams) + result = parseStmt(typeDef) + + var RecList = newNimNode(nnkRecList) + + for i in 1 .. formalParams.len-1: + expectKind(formalParams[i], nnkIdentDefs) + RecList.add(formalParams[i]) + + result[0][0][2][0][2] = RecList + + result.add(parseExpr("let passedInParams = cast[PArgObject](x.param)")) + for i in 1 .. formalParams.len-1: + expectKind(formalParams[i], nnkIdentDefs) + result.add(parseExpr("let $1: $2 = passedInParams.$1" % + [$formalParams[i][0].ident, $formalParams[i][1].ident])) + +macro async*(n: stmt): stmt {.immediate.} = + expectKind(n, nnkProcDef) + #echo("-------------") + 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 + + # Pragma + result[4].add(newIdentNode(!"closure")) + + result[6] = newNimNode(nnkStmtList) + + # Transform any params passed to the proc into a typedef. + if n[3].len > 1: + let transfArgs = transformArgs(n[3]) + + result[6].add(transfArgs) + + # Body + var body = transform(n[6]) + result[6].add(body) + + #echo treeRepr(result) + echo result.toStrLit().strVal + +# ---- Async macro end + when isMainModule: proc testConnect(s: PAsyncSocket, no: int) = diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim new file mode 100644 index 000000000..5b1c2a4dc --- /dev/null +++ b/tests/compile/tasyncitermacro.nim @@ -0,0 +1,24 @@ +import sockets, asyncio, strutils +proc processRequest(client: PAsyncSocket, test: string) {.async.} = + assert test == "ahha" + assert client != nil + echo(test) + 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, "ahha") + +var disp = newDispatcher(false) +disp.register(processServer, nil) +while true: + discard disp.poll() diff --git a/tests/compile/tasynciterraw.nim b/tests/compile/tasynciterraw.nim new file mode 100644 index 000000000..ba038f3fc --- /dev/null +++ b/tests/compile/tasynciterraw.nim @@ -0,0 +1,69 @@ +import asyncio, sockets + +const qwerty = "qwertyuiopasdfghjklzxcvbnm\c\L" + +iterator processRequest(x: PRequest): PRequest {.closure.} = + # Read first message + let client = x.socket + var lineReq = PRequest(socket: client, kind: reqReadLine, line: "") + yield lineReq + + if lineReq.line == "HELLO": + echo "Got first message." + + var longString = qwerty + for i in 0..109000: + longString.add($i & qwerty) + + var writeReq = PRequest(socket: client, kind: reqWrite, toWrite: longString) + yield writeReq + + # Read second message: No need to create another request object: + yield lineReq + + if lineReq.line == "DOM": + echo("Got second message.") + + client.close() + + +iterator processServer(x: PRequest): PRequest {.closure.} = + type + PTest = ref object of TObject + param: string + + let passedInParams = cast[PTest](x.param) + assert passedInParams.param == "blah" + + var sock = AsyncSocket() + + # blocks: + sock.bindAddr(TPort(6667)) + sock.listen() + + # Accept loop + while true: + var acceptReq = PRequest(socket: sock, kind: reqAccept, client: nil) + echo("About to yield accept") + yield acceptReq + echo("after yield") + if acceptReq.hasException: + echo("Got exception[EAssertionFailed]: ", acceptReq.exc of EAssertionFailed) + echo(acceptReq.exc.msg) + continue + + let client = acceptReq.client + assert client != nil + yield PRequest(socket: client, kind: reqReg, param: nil, worker: processRequest) + +when isMainModule: + let param = "blah" + type + PTest = ref object of TObject + param: type(param) + var test = PTest(param: param) + var disp = newDispatcher(false) + disp.register(processServer, test) + while true: + discard disp.poll() + From d15e67e0209b803cccc7aea24cc45c0d374686a7 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 21 Jul 2013 13:58:54 +0100 Subject: [PATCH 02/12] Better param passing. Implemented send, readLine and support for documentation stubs. --- lib/pure/asyncio.nim | 139 ++++++++++++++++++++---------- tests/compile/tasyncitermacro.nim | 6 +- 2 files changed, 100 insertions(+), 45 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 26df1cdfa..e7e931bf5 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -133,7 +133,7 @@ type SockUDPBound TRequestKind* = enum - reqNil, reqReg, reqRead, reqWrite, reqReadLine, reqAccept + reqNil, reqReg, reqRead, reqWrite, reqReadLine, reqAccept, reqConnect PRequest* = ref object socket*: PAsyncSocket @@ -154,10 +154,12 @@ type toWrite*: string ## Request written: int ## Internal data of reqReadLine: - line*: string ## Response + line*: string ## Response of reqAccept: client*: PAsyncSocket ## Response - + of reqConnect: + address*: string ## Request + PWorker* = ref object worker: iterator (x: PRequest): PRequest {.closure.} x: PRequest @@ -820,6 +822,8 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool = else: let toWrite = worker.lastReq.toWrite[written .. -1] doSend(toWrite) + of reqConnect: + #execReq of reqReg: assert false, "reqReg should have been processed already" of reqNil: @@ -834,6 +838,17 @@ proc len*(disp: PDispatcher): int = # ---- Async macro +proc createRequestNode(varName, + reqArgs: string): PNimrodNode {.compiletime.} = + result = newNimNode(nnkStmtList) + var reqObj = parseExpr( + """var $# = PRequest($#)""" % + [varName, reqArgs]) + result.add reqObj + result.add parseExpr("yield $#" % [varName]) + # Check for exception + result.add parseStmt("if $1.hasException: raise $1.exc" % [varName]) + proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} = ## Transforms a var/let section ## E.g: @@ -843,57 +858,50 @@ proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} = expectLen(n[0], 3) # IdentDefs let insideAwait = n[0][2][1] let reqCall = $insideAwait[0].ident - case reqCall + + expectLen(insideAwait, 2) + let sockName = $insideAwait[1].ident + case reqCall.normalize of "accept": - expectLen(insideAwait, 2) - let sockName = $insideAwait[1].ident let acceptReqVar = "acceptReq" # TODO: Random var names which do not conflict. or wait for gensym? - var reqObj = parseExpr( - """var $# = PRequest(socket: $#, kind: reqAccept, - client: nil)""" % - [acceptReqVar, sockName]) - result.add reqObj - result.add parseExpr("yield $#" % [acceptReqVar]) - # Check for exception - result.add parseStmt("if $1.hasException: raise $1.exc" % [acceptReqVar]) + result.add createRequestNode(acceptReqVar, + "socket: $#, kind: reqAccept, client: nil" % sockName) case n.kind of nnkLetSection: result.add parseExpr("let $# = $#.client" % [$nameIdent, acceptReqVar]) of nnkVarSection: result.add parseExpr("var $# = $#.client" % [$nameIdent, acceptReqVar]) else: error "Bad node kind in toYieldVar" + of "readline": + let readReqVar = "readLineReq" + # TODO: Random var names which do not conflict. or wait for gensym? + result.add createRequestNode(readReqVar, + "socket: $#, kind: reqReadLine, line: \"\"" % sockName) + case n.kind + of nnkLetSection: + result.add parseExpr("let $# = $#.line" % [$nameIdent, readReqVar]) + of nnkVarSection: + result.add parseExpr("var $# = $#.line" % [$nameIdent, readReqVar]) + else: error "Bad node kind in toYieldVar" else: error(reqCall & " is not a valid async call") const typeDef = """ type - PArgObject = ref object of TObject + P$#ArgObject = ref object of TObject """ proc transformCallWithArg(call: PNimrodNode): PNimrodNode {.compiletime.} = - result = parseStmt(typeDef) - var RecList = newNimNode(nnkRecList) - for i in 1 .. call[1].len-1: - var typeOf = newNimNode(nnkTypeOfExpr).add(newNimNode(nnkPar)) - typeOf[0].add(call[1][i]) - case call[1][i].kind - of nnkIdent: - RecList.add newIdentDefs(call[1][i], typeOf) - of nnkLiterals: - RecList.add newIdentDefs(newIdentNode("dummy" & $i), typeOf) - else: assert false + result = newNimNode(nnkStmtList) - result[0][0][2][0][2] = RecList - result.add parseExpr("var argsToPass: PArgObject") + result.add parseExpr("var argsToPass: P$#ArgObject" % [$call[1][0].ident]) result.add parseExpr("new argsToPass") for i in 1 .. call[1].len-1: case call[1][i].kind - of nnkIdent: - result.add parseExpr("argsToPass.$1 = $1" % [$call[1][i].ident]) - of nnkLiterals: + of nnkLiterals, nnkIdent: let dotExpr = newDotExpr(newIdentNode("argsToPass"), newIdentNode("dummy" & $i)) result.add newAssignment(dotExpr, call[1][i]) @@ -904,8 +912,13 @@ proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = if $n[0].ident != "await": error "'await' expected" result = @[] let callIdent = $n[1][0].ident - case callIdent - of "send": error("TODO") + case callIdent.normalize + of "send": + let socketName = $n[1][1].ident + let toWrite = n[1][2] + result.add createRequestNode("sendReq", + "socket: $#, kind: reqWrite, toWrite: $#" % + [socketName, $(toWrite.toStrLit)]) else: result.add(transformCallWithArg(n)) # reqRegister @@ -945,29 +958,53 @@ proc transform(n: PNimrodNode): PNimrodNode {.compiletime.} = else: result.add(son) -proc transformArgs(formalParams: PNimrodNode): PNimrodNode {.compiletime.} = +proc transformArgs(procName: string, + formalParams: PNimrodNode): PNimrodNode {.compiletime.} = ## Transforms formal params into a typedef with a dummy type - ## ``ref object of TObject``. The PRequest.param is then casted to it, - ## and immutable vars are defined as specified the proc's params. + ## ``ref object of TObject``. This is inserted above the proc definition. expectKind(formalParams, nnkFormalParams) - result = parseStmt(typeDef) + result = parseStmt(typeDef % procName) var RecList = newNimNode(nnkRecList) for i in 1 .. formalParams.len-1: expectKind(formalParams[i], nnkIdentDefs) - RecList.add(formalParams[i]) + RecList.add(newIdentDefs(newIdentNode("dummy" & $i), formalParams[i][1])) + # TODO: Add comment with the original param name? result[0][0][2][0][2] = RecList - result.add(parseExpr("let passedInParams = cast[PArgObject](x.param)")) +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) + result.add(parseExpr("let passedInParams = P$#ArgObject(x.param)" % procName)) + # Fields take the form ``dummy``. for i in 1 .. formalParams.len-1: expectKind(formalParams[i], nnkIdentDefs) - result.add(parseExpr("let $1: $2 = passedInParams.$1" % - [$formalParams[i][0].ident, $formalParams[i][1].ident])) + result.add(parseExpr("let $1: $2 = passedInParams.$3" % + [$formalParams[i][0].ident, $formalParams[i][1].ident, + "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 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 + #echo("-------------") result = newNimNode(nnkIteratorDef) for i in 0 .. n.len-1: @@ -989,21 +1026,35 @@ macro async*(n: stmt): stmt {.immediate.} = result[6] = newNimNode(nnkStmtList) - # Transform any params passed to the proc into a typedef. + # Declare variables based on the params that the async proc takes. if n[3].len > 1: - let transfArgs = transformArgs(n[3]) + let args = declareArgsInBody($n[0].ident, n[3]) - result[6].add(transfArgs) + result[6].add(args) # Body var body = transform(n[6]) result[6].add(body) + # Add typedef above the proc def for parameters. + if n[3].len > 1: + let procDef = copyNimTree(result) + result = newNimNode(nnkStmtList) + result.add(transformArgs($n[0].ident, n[3])) + result.add procDef + #echo treeRepr(result) echo result.toStrLit().strVal # ---- Async macro end +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/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index 5b1c2a4dc..2c6d22b5a 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -2,7 +2,11 @@ import sockets, asyncio, strutils proc processRequest(client: PAsyncSocket, test: string) {.async.} = assert test == "ahha" assert client != nil - echo(test) + echo("Test = ", test) + let line = await(readLine(client)) + echo("Read: ", line) + + await send(client, "Goodbye.\c\L") client.close() proc processServer() {.async.} = From a1a8925692174139c2c51c55a571d3f2b101a5ef Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 21 Jul 2013 14:28:02 +0100 Subject: [PATCH 03/12] Calling an async proc with invalid arguments will now fail at compile-time. --- lib/pure/asyncio.nim | 19 +++++++++++++++ tests/compile/tasyncitermacro.nim | 1 + tests/reject/tasyncmacroinvalidargs.nim | 32 +++++++++++++++++++++++++ 3 files changed, 52 insertions(+) create mode 100644 tests/reject/tasyncmacroinvalidargs.nim diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index e7e931bf5..473d46a76 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -907,6 +907,14 @@ proc transformCallWithArg(call: PNimrodNode): PNimrodNode {.compiletime.} = result.add newAssignment(dotExpr, call[1][i]) else: assert false + # Add in a call to a pre-generated proc stub so that the compiler verifies + # the params for us :) + # TODO: Make sure this doesn't actually get called, maybe inline it? + var args: seq[PNimrodNode] = @[] + for i in 1 .. call[1].len-1: + args.add(call[1][i]) + result.add newCall(call[1][0], args) + proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = ## Transforms a call/command if $n[0].ident != "await": error "'await' expected" @@ -994,6 +1002,14 @@ proc isDocumentation(n: PNimrodNode): bool {.compiletime.} = 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? + var kids: seq[PNimrodNode] = @[] + for i in 0 .. formalParams.len-1: + kids.add(formalParams[i]) + result = newProc(procName, kids) + macro async*(n: stmt): stmt {.immediate.} = expectKind(n, nnkProcDef) #echo(treeRepr(n)) @@ -1041,6 +1057,9 @@ macro async*(n: stmt): stmt {.immediate.} = let procDef = copyNimTree(result) result = newNimNode(nnkStmtList) result.add(transformArgs($n[0].ident, n[3])) + # Generate a proc stub to verify that the user passes the correct params. + result.add createVerificationProc(n[0], n[3]) + result.add procDef #echo treeRepr(result) diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index 2c6d22b5a..b20bb5ea4 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -1,4 +1,5 @@ import sockets, asyncio, strutils + proc processRequest(client: PAsyncSocket, test: string) {.async.} = assert test == "ahha" assert client != nil 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() From 1e1a583903486b0b801fd978bebbbd51f64314f2 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 21 Jul 2013 14:51:25 +0100 Subject: [PATCH 04/12] Fixes problems with async docs. --- lib/pure/asyncio.nim | 11 ++++++----- 1 file changed, 6 insertions(+), 5 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 473d46a76..fb3b55eff 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -1067,12 +1067,13 @@ macro async*(n: stmt): stmt {.immediate.} = # ---- Async macro end -proc send*(socket: PAsyncSocket, text: string) {.async.} = - ## Sends ``text`` to ``socket`` asynchronously. +when defined(nimdoc): + 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. + proc accept*(socket: PAsyncSocket): PAsyncSocket {.async.} = + ## Accepts a client connecting to a server socket asynchronously. + ## Returns that client. when isMainModule: From d8659eaaa2e445bc537388da643107f19f10ebd5 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 21 Jul 2013 15:18:44 +0100 Subject: [PATCH 05/12] Another docs fix for asyncio. --- lib/pure/asyncio.nim | 33 +++++++++++++++++++++++++++++++-- 1 file changed, 31 insertions(+), 2 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index fb3b55eff..3c312b56c 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -1067,13 +1067,42 @@ macro async*(n: stmt): stmt {.immediate.} = # ---- Async macro end -when defined(nimdoc): +## 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. + ## Returns that client. """ when isMainModule: From c5ff9650f7ec5f169b45ba02172057ce5f620c55 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Thu, 25 Jul 2013 22:13:37 +0100 Subject: [PATCH 06/12] Changed some asyncio macro code to use the new genSym. --- lib/pure/asyncio.nim | 49 ++++++++++++++++++------------- tests/compile/tasyncitermacro.nim | 2 ++ 2 files changed, 30 insertions(+), 21 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 3c312b56c..d590863f7 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -839,22 +839,27 @@ proc len*(disp: PDispatcher): int = # ---- Async macro proc createRequestNode(varName, - reqArgs: string): PNimrodNode {.compiletime.} = + reqArgs: string, sym: var PNimrodNode): PNimrodNode {.compiletime.} = result = newNimNode(nnkStmtList) - var reqObj = parseExpr( - """var $# = PRequest($#)""" % - [varName, reqArgs]) + sym = genSym(nskVar, varName) + var reqObj = newVarStmt(sym, + parseExpr("PRequest($#)" % [reqArgs])) + result.add reqObj - result.add parseExpr("yield $#" % [varName]) + result.add newNimNode(nnkYieldStmt).add(sym) # Check for exception - result.add parseStmt("if $1.hasException: raise $1.exc" % [varName]) + result.add newIfStmt( + (newDotExpr(sym, newIdentNode("hasException")), + newNimNode(nnkRaiseStmt).add( + newDotExpr(sym, newIdentNode("exc"))))) + echo treeRepr(result) 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].ident # Var name + let nameIdent = n[0][0] # Var name expectLen(n[0], 3) # IdentDefs let insideAwait = n[0][2][1] let reqCall = $insideAwait[0].ident @@ -864,25 +869,26 @@ proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} = case reqCall.normalize of "accept": let acceptReqVar = "acceptReq" - # TODO: Random var names which do not conflict. or wait for gensym? + var sym: PNimrodNode result.add createRequestNode(acceptReqVar, - "socket: $#, kind: reqAccept, client: nil" % sockName) + "socket: $#, kind: reqAccept, client: nil" % sockName, sym) + case n.kind of nnkLetSection: - result.add parseExpr("let $# = $#.client" % [$nameIdent, acceptReqVar]) + result.add newLetStmt(nameIdent, newDotExpr(sym, newIdentNode("client"))) of nnkVarSection: - result.add parseExpr("var $# = $#.client" % [$nameIdent, acceptReqVar]) + result.add newVarStmt(nameIdent, newDotExpr(sym, newIdentNode("client"))) else: error "Bad node kind in toYieldVar" of "readline": let readReqVar = "readLineReq" - # TODO: Random var names which do not conflict. or wait for gensym? + var sym: PNimrodNode result.add createRequestNode(readReqVar, - "socket: $#, kind: reqReadLine, line: \"\"" % sockName) + "socket: $#, kind: reqReadLine, line: \"\"" % sockName, sym) case n.kind of nnkLetSection: - result.add parseExpr("let $# = $#.line" % [$nameIdent, readReqVar]) + result.add newLetStmt(nameIdent, newDotExpr(sym, newIdentNode("line"))) of nnkVarSection: - result.add parseExpr("var $# = $#.line" % [$nameIdent, readReqVar]) + result.add newVarStmt(nameIdent, newDotExpr(sym, newIdentNode("line"))) else: error "Bad node kind in toYieldVar" else: error(reqCall & " is not a valid async call") @@ -924,9 +930,10 @@ proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = 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)]) + "socket: $#, kind: reqWrite, toWrite: $#" % + [socketName, $(toWrite.toStrLit)], sym) else: result.add(transformCallWithArg(n)) # reqRegister @@ -987,13 +994,13 @@ proc declareArgsInBody(procName: string, ## 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) - result.add(parseExpr("let passedInParams = P$#ArgObject(x.param)" % procName)) + 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(parseExpr("let $1: $2 = passedInParams.$3" % - [$formalParams[i][0].ident, $formalParams[i][1].ident, - "dummy" & $i])) + 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. diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index b20bb5ea4..acf2ba412 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -7,6 +7,8 @@ proc processRequest(client: PAsyncSocket, test: string) {.async.} = let line = await(readLine(client)) echo("Read: ", line) + await send(client, "Goodbye.\c\L") + await send(client, "Goodbye.\c\L") await send(client, "Goodbye.\c\L") client.close() From 2330c95b227d8e780e0cc09acee69588809334e0 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 28 Jul 2013 13:50:32 +0100 Subject: [PATCH 07/12] Implemented genSym for async function calls. --- lib/pure/asyncio.nim | 41 ++++++++++++++++++++++--------- tests/compile/tasyncitermacro.nim | 12 +++++---- 2 files changed, 36 insertions(+), 17 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index d590863f7..53054a738 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -841,18 +841,22 @@ proc len*(disp: PDispatcher): int = proc createRequestNode(varName, reqArgs: string, sym: var PNimrodNode): PNimrodNode {.compiletime.} = result = newNimNode(nnkStmtList) - sym = genSym(nskVar, varName) + # 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"))))) - echo treeRepr(result) + newDotExpr(sym, newIdentNode("exc"))))) proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} = ## Transforms a var/let section @@ -871,7 +875,7 @@ proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} = let acceptReqVar = "acceptReq" var sym: PNimrodNode result.add createRequestNode(acceptReqVar, - "socket: $#, kind: reqAccept, client: nil" % sockName, sym) + "socket: $#, kind: reqAccept, client: nil, hasException: false" % sockName, sym) case n.kind of nnkLetSection: @@ -899,16 +903,24 @@ const typeDef = P$#ArgObject = ref object of TObject """ -proc transformCallWithArg(call: PNimrodNode): PNimrodNode {.compiletime.} = +proc transformCallWithArg(call: PNimrodNode, + sym: var PNimrodNode): PNimrodNode {.compiletime.} = result = newNimNode(nnkStmtList) - result.add parseExpr("var argsToPass: P$#ArgObject" % [$call[1][0].ident]) - result.add parseExpr("new argsToPass") + sym = gensym(nskVar, "argsToPass") + + # var argsToPass: P$#ArgObject + result.add newNimNode(nnkVarSection).add( + newNimNode(nnkIdentDefs).add(sym, + newIdentNode("P$#ArgObject" % [$call[1][0].ident]), + newNimNode(nnkEmpty))) + + result.add newCall("new", sym) for i in 1 .. call[1].len-1: case call[1][i].kind of nnkLiterals, nnkIdent: - let dotExpr = newDotExpr(newIdentNode("argsToPass"), + let dotExpr = newDotExpr(sym, newIdentNode("dummy" & $i)) result.add newAssignment(dotExpr, call[1][i]) else: assert false @@ -919,7 +931,7 @@ proc transformCallWithArg(call: PNimrodNode): PNimrodNode {.compiletime.} = var args: seq[PNimrodNode] = @[] for i in 1 .. call[1].len-1: args.add(call[1][i]) - result.add newCall(call[1][0], args) + result.add newCall(call[1][0], args) proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = ## Transforms a call/command @@ -935,10 +947,15 @@ proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = "socket: $#, kind: reqWrite, toWrite: $#" % [socketName, $(toWrite.toStrLit)], sym) else: - result.add(transformCallWithArg(n)) + var sym: PNimrodNode + result.add(transformCallWithArg(n, sym)) # reqRegister - result.add parseExpr("yield PRequest(socket: $#, kind: reqReg, worker: $#, param: argsToPass)" % - [$n[1][1].ident, callIdent]) + var yie = parseExpr("yield PRequest(socket: $#, kind: reqReg, worker: $#)" % + [$n[1][1].ident, callIdent]) + yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"), + sym)) + + result.add yie proc transform(n: PNimrodNode): PNimrodNode {.compiletime.} = result = newNimNode(nnkStmtList) diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index acf2ba412..e1ec3f849 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -1,6 +1,6 @@ import sockets, asyncio, strutils -proc processRequest(client: PAsyncSocket, test: string) {.async.} = +proc processRequest(client: PAsyncSocket, test: string, closeSock: bool) {.async.} = assert test == "ahha" assert client != nil echo("Test = ", test) @@ -8,9 +8,10 @@ proc processRequest(client: PAsyncSocket, test: string) {.async.} = echo("Read: ", line) await send(client, "Goodbye.\c\L") - await send(client, "Goodbye.\c\L") - await send(client, "Goodbye.\c\L") - client.close() + #await send(client, "Goodbye.\c\L") + #await send(client, "Goodbye.\c\L") + if closeSock: + client.close() proc processServer() {.async.} = var sock = AsyncSocket() @@ -23,7 +24,8 @@ proc processServer() {.async.} = while true: let client: PAsyncSocket = await(accept(sock)) assert client != nil - await processRequest(client, "ahha") + await processRequest(client, "ahha", false) + await processRequest(client, "ahha", true) var disp = newDispatcher(false) disp.register(processServer, nil) From 426c2d6c0194c8a0e13b3145ef5870a7be591cf2 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 28 Jul 2013 14:41:09 +0100 Subject: [PATCH 08/12] Implemented default arguments for async procs. --- lib/pure/asyncio.nim | 59 +++++++++++++++++-------------- tests/compile/tasyncitermacro.nim | 4 +-- 2 files changed, 35 insertions(+), 28 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 53054a738..5bd3f578f 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -840,6 +840,7 @@ proc len*(disp: PDispatcher): int = 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. @@ -905,33 +906,19 @@ const typeDef = 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") - - # var argsToPass: P$#ArgObject - result.add newNimNode(nnkVarSection).add( - newNimNode(nnkIdentDefs).add(sym, - newIdentNode("P$#ArgObject" % [$call[1][0].ident]), - newNimNode(nnkEmpty))) - - result.add newCall("new", sym) - - for i in 1 .. call[1].len-1: - case call[1][i].kind - of nnkLiterals, nnkIdent: - let dotExpr = newDotExpr(sym, - newIdentNode("dummy" & $i)) - result.add newAssignment(dotExpr, call[1][i]) - else: assert false # Add in a call to a pre-generated proc stub so that the compiler verifies # the params for us :) - # TODO: Make sure this doesn't actually get called, maybe inline it? + # TODO: Inline this maybe? var args: seq[PNimrodNode] = @[] for i in 1 .. call[1].len-1: args.add(call[1][i]) - result.add newCall(call[1][0], args) + result.add newVarStmt(sym, newCall(call[1][0], args)) proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = ## Transforms a call/command @@ -958,6 +945,11 @@ proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = 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: @@ -1029,10 +1021,22 @@ proc isDocumentation(n: PNimrodNode): bool {.compiletime.} = proc createVerificationProc(procName: PNimrodNode, formalParams: PNimrodNode): PNimrodNode {.compiletime.} = # TODO: Export this stub if our async proc is exported? - var kids: seq[PNimrodNode] = @[] - for i in 0 .. formalParams.len-1: - kids.add(formalParams[i]) - result = newProc(procName, kids) + # 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) + body.add newCall("new", newIdentNode("result")) + + 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) @@ -1045,7 +1049,6 @@ macro async*(n: stmt): stmt {.immediate.} = result[6] = parseStmt("nil") return - #echo("-------------") result = newNimNode(nnkIteratorDef) for i in 0 .. n.len-1: result.add(copyNimTree(n[i])) @@ -1061,18 +1064,20 @@ macro async*(n: stmt): stmt {.immediate.} = formalParams.add(params) result[3] = formalParams - # Pragma + # 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) - # Body + # Transform body var body = transform(n[6]) result[6].add(body) @@ -1081,7 +1086,9 @@ macro async*(n: stmt): stmt {.immediate.} = let procDef = copyNimTree(result) result = newNimNode(nnkStmtList) result.add(transformArgs($n[0].ident, n[3])) - # Generate a proc stub to verify that the user passes the correct params. + # 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 diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index e1ec3f849..361e34c22 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -1,6 +1,6 @@ import sockets, asyncio, strutils -proc processRequest(client: PAsyncSocket, test: string, closeSock: bool) {.async.} = +proc processRequest(client: PAsyncSocket, test: string, closeSock: bool = true) {.async.} = assert test == "ahha" assert client != nil echo("Test = ", test) @@ -25,7 +25,7 @@ proc processServer() {.async.} = let client: PAsyncSocket = await(accept(sock)) assert client != nil await processRequest(client, "ahha", false) - await processRequest(client, "ahha", true) + await processRequest(client, "ahha") var disp = newDispatcher(false) disp.register(processServer, nil) From ae64089e5378eca6b4d5f3b751ce413365cfef04 Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sun, 28 Jul 2013 18:28:50 +0100 Subject: [PATCH 09/12] Fixed crash when worker has just been registered and didn't get executed yet. --- lib/pure/asyncio.nim | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 5bd3f578f..ef2a46be0 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -757,7 +757,7 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool = processWorkers(d) var readWorkers = populateRead(d) var writeWorkers = populateWrite(d) - echo(readWorkers.len, " ", d.requests[reqNil].len, d.requests[reqReadLine].len) + #echo(readWorkers.len, " ", d.requests[reqNil].len, d.requests[reqReadLine].len) if select(readWorkers, writeWorkers, timeout) != 0: var newRequests: array[TRequestKind, seq[PWorker]] = newRequests() for req in TRequestKind: @@ -827,7 +827,8 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool = of reqReg: assert false, "reqReg should have been processed already" of reqNil: - assert false, "reqNil should have been processed already" + # Nothing to do. Most likely that a new worker has just been + # registered. newRequests[addTo].add(worker) d.requests = newRequests From 70d67daa14e7359190a5a0a68863a17deb94583f Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Wed, 31 Jul 2013 19:28:37 +0100 Subject: [PATCH 10/12] Introduced 'reg' keyword, and changed the await behaviour for custom procs. --- lib/pure/asyncio.nim | 223 +++++++++++++++++++----------- tests/compile/tasyncitermacro.nim | 14 +- 2 files changed, 149 insertions(+), 88 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index ef2a46be0..ef3d1df71 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -133,7 +133,9 @@ type SockUDPBound TRequestKind* = enum - reqNil, reqReg, reqRead, reqWrite, reqReadLine, reqAccept, reqConnect + reqNil, reqReg, reqAwait, reqRead, reqWrite, reqReadLine, reqAccept, reqConnect + + PAsyncProc = iterator (x: PRequest): PRequest PRequest* = ref object socket*: PAsyncSocket @@ -144,9 +146,9 @@ type case kind*: TRequestKind of reqNil: nil - of reqReg: + of reqReg, reqAwait: param*: PObject - worker*: iterator (x: PRequest): PRequest + worker*: PAsyncProc of reqRead: count*: int ## Request readData*: string ## Response @@ -161,9 +163,13 @@ type address*: string ## Request PWorker* = ref object - worker: iterator (x: PRequest): PRequest {.closure.} - x: PRequest - lastReq: PRequest + worker*: PAsyncProc + x*: PRequest + lastReq*: PRequest + case hasParent*: bool + of true: + parent*: PWorker + else: nil PDispatcher* = ref TDispatcher TDispatcher = object @@ -615,12 +621,15 @@ proc createFdSet(fd: var TFdSet, s: seq[PWorker], m: var int) = 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(s[i].lastReq.socket.getFD, fd) != 0'i32: + if FD_ISSET(getWorkerSocket(s[i]).getFD, fd) != 0'i32: s[i] = s[L-1] dec(L) else: @@ -661,19 +670,38 @@ proc register*(disp: PDispatcher, worker: iterator (x: PRequest): PRequest, proc processWorkers(d: PDispatcher) = var newRequests: array[TRequestKind, seq[PWorker]] = d.requests newRequests[reqNil] = @[] + + + for idle in d.requests[reqNil]: let req = idle.worker(idle.x) if req != nil: - if req.kind == reqReg: + echo("Process workers, after exec: ", req.kind) + case req.kind + of reqReg: newRequests[reqNil].add(idle) let newWorker = PWorker(worker: req.worker, lastReq: PRequest(kind: reqNil), x: req) newRequests[reqNil].add(newWorker) + of reqAwait: + # For efficiency lets execute this async proc now. + let awaitReq = req.worker(req) + let newWorker = PWorker(worker: req.worker, lastReq: awaitReq, + 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. + newRequests[awaitReq.kind].add(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. + newRequests[reqNil].add(idle.parent) + echo("Await finish: ", idle.parent.lastReq.kind) + d.requests = newRequests template popu(req) {.immediate, dirty.} = @@ -683,6 +711,7 @@ template popu(req) {.immediate, dirty.} = proc populateRead(d: PDispatcher): seq[PWorker] = result = @[] + # TODO: Just add the seq[] popu(reqRead) popu(reqReadLine) popu(reqAccept) @@ -691,6 +720,79 @@ proc populateWrite(d: PDispatcher): seq[PWorker] = result = @[] popu(reqWrite) +proc processRequests(requests, readWorkers, writeWorkers: seq[PWorker], + newRequests: var array[TRequestKind, seq[PWorker]]) = + for worker in requests: + echo("Process requests: ", 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: + #execReq + 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. @@ -761,76 +863,7 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool = if select(readWorkers, writeWorkers, timeout) != 0: var newRequests: array[TRequestKind, seq[PWorker]] = newRequests() for req in TRequestKind: - for worker in d.requests[req]: - echo(req) - var addTo = req - 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 req - 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: - #execReq - of reqReg: - assert false, "reqReg should have been processed already" - of reqNil: - # Nothing to do. Most likely that a new worker has just been - # registered. - - newRequests[addTo].add(worker) + processRequests(d.requests[req], readWorkers, writeWorkers, newRequests) d.requests = newRequests proc len*(disp: PDispatcher): int = @@ -937,14 +970,29 @@ proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = else: var sym: PNimrodNode result.add(transformCallWithArg(n, sym)) - # reqRegister - var yie = parseExpr("yield PRequest(socket: $#, kind: reqReg, worker: $#)" % + # reqCustom + var yie = parseExpr("yield PRequest(socket: $#, kind: reqAwait, worker: $#)" % [$n[1][1].ident, 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: $#, kind: reqReg, worker: $#)" % + [$n[1][1].ident, 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: @@ -975,9 +1023,18 @@ proc transform(n: PNimrodNode): PNimrodNode {.compiletime.} = 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 and $son[0].ident == "await": - result.add toYieldCall(son) + 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: diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index 361e34c22..faa383f1c 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -1,15 +1,19 @@ import sockets, asyncio, strutils +proc auth(client: PAsyncSocket) {.async.} = + await send(client, "Auth\c\L") + proc processRequest(client: PAsyncSocket, test: string, closeSock: bool = true) {.async.} = assert test == "ahha" assert client != nil - echo("Test = ", test) let line = await(readLine(client)) echo("Read: ", line) + for i in 0 .. 10: + await auth(client) + #await send(client, "Auth\c\L") + await send(client, "Goodbye.\c\L") - #await send(client, "Goodbye.\c\L") - #await send(client, "Goodbye.\c\L") if closeSock: client.close() @@ -24,8 +28,8 @@ proc processServer() {.async.} = while true: let client: PAsyncSocket = await(accept(sock)) assert client != nil - await processRequest(client, "ahha", false) - await processRequest(client, "ahha") + + reg processRequest(client, "ahha") var disp = newDispatcher(false) disp.register(processServer, nil) From e42f877664bf19408a7b7438480ec0feba4d909f Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Wed, 31 Jul 2013 23:27:07 +0100 Subject: [PATCH 11/12] 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) --- lib/pure/asyncio.nim | 23 ++++++++++++++--------- tests/compile/tasyncitermacro.nim | 9 +++++++-- 2 files changed, 21 insertions(+), 11 deletions(-) diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index ef3d1df71..3010e48ff 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -671,9 +671,7 @@ proc processWorkers(d: PDispatcher) = var newRequests: array[TRequestKind, seq[PWorker]] = d.requests newRequests[reqNil] = @[] - - - for idle in d.requests[reqNil]: + proc processWorker(idle: PWorker) = let req = idle.worker(idle.x) if req != nil: echo("Process workers, after exec: ", req.kind) @@ -684,13 +682,16 @@ proc processWorkers(d: PDispatcher) = x: req) newRequests[reqNil].add(newWorker) of reqAwait: - # For efficiency lets execute this async proc now. - let awaitReq = req.worker(req) - let newWorker = PWorker(worker: req.worker, lastReq: awaitReq, + 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. - newRequests[awaitReq.kind].add(newWorker) + + # 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) @@ -698,10 +699,14 @@ proc processWorkers(d: PDispatcher) = 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. - newRequests[reqNil].add(idle.parent) + # user-defined async proc which just finished. Do this by calling + # processWorker recursively. Same way as above. + processWorker(idle.parent) echo("Await finish: ", idle.parent.lastReq.kind) + for idle in d.requests[reqNil]: + processWorker(idle) + d.requests = newRequests template popu(req) {.immediate, dirty.} = diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim index faa383f1c..4b9b7b5b2 100644 --- a/tests/compile/tasyncitermacro.nim +++ b/tests/compile/tasyncitermacro.nim @@ -1,8 +1,14 @@ import sockets, asyncio, strutils -proc auth(client: PAsyncSocket) {.async.} = +proc auth3(client: PAsyncSocket) {.async.} = await send(client, "Auth\c\L") +proc auth2(client: PAsyncSocket) {.async.} = + await auth3(client) + +proc auth(client: PAsyncSocket) {.async.} = + await auth2(client) + proc processRequest(client: PAsyncSocket, test: string, closeSock: bool = true) {.async.} = assert test == "ahha" assert client != nil @@ -11,7 +17,6 @@ proc processRequest(client: PAsyncSocket, test: string, closeSock: bool = true) for i in 0 .. 10: await auth(client) - #await send(client, "Auth\c\L") await send(client, "Goodbye.\c\L") if closeSock: From b3469c43a69c299779ba79c3f74a5f4c1ed9ebfc Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sat, 3 Aug 2013 14:19:17 +0100 Subject: [PATCH 12/12] 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. --- lib/pure/asyncio.nim | 92 ++++++++++----- tests/compile/tasyncitermacro.nim | 42 ------- tests/compile/tasynciterraw.nim | 69 ----------- tests/run/tasyncitermacro.nim | 78 +++++++++++++ tests/run/tasynciterraw.nim | 185 ++++++++++++++++++++++++++++++ 5 files changed, 325 insertions(+), 141 deletions(-) delete mode 100644 tests/compile/tasyncitermacro.nim delete mode 100644 tests/compile/tasynciterraw.nim create mode 100644 tests/run/tasyncitermacro.nim create mode 100644 tests/run/tasynciterraw.nim diff --git a/lib/pure/asyncio.nim b/lib/pure/asyncio.nim index 3010e48ff..99b0eecb0 100644 --- a/lib/pure/asyncio.nim +++ b/lib/pure/asyncio.nim @@ -161,6 +161,7 @@ type client*: PAsyncSocket ## Response of reqConnect: address*: string ## Request + port*: TPort ## Request PWorker* = ref object worker*: PAsyncProc @@ -674,9 +675,15 @@ proc processWorkers(d: PDispatcher) = proc processWorker(idle: PWorker) = let req = idle.worker(idle.x) if req != nil: - echo("Process workers, after exec: ", req.kind) 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) @@ -702,33 +709,28 @@ proc processWorkers(d: PDispatcher) = # user-defined async proc which just finished. Do this by calling # processWorker recursively. Same way as above. processWorker(idle.parent) - echo("Await finish: ", idle.parent.lastReq.kind) for idle in d.requests[reqNil]: processWorker(idle) d.requests = newRequests -template popu(req) {.immediate, dirty.} = - for i in d.requests[req]: - result.add(i) - proc populateRead(d: PDispatcher): seq[PWorker] = result = @[] - # TODO: Just add the seq[] - popu(reqRead) - popu(reqReadLine) - popu(reqAccept) + result.add d.requests[reqAccept] + result.add d.requests[reqRead] + result.add d.requests[reqReadLine] proc populateWrite(d: PDispatcher): seq[PWorker] = result = @[] - popu(reqWrite) + 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("Process requests: ", worker.lastReq.kind) + #echo(worker.lastReq.kind) var addTo = worker.lastReq.kind template execReq(workers: var seq[PWorker], autoadd: bool, body: stmt) {.immediate, dirty.} = @@ -789,7 +791,11 @@ proc processRequests(requests, readWorkers, writeWorkers: seq[PWorker], let toWrite = worker.lastReq.toWrite[written .. -1] doSend(toWrite) of reqConnect: - #execReq + 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: @@ -864,8 +870,24 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool = processWorkers(d) var readWorkers = populateRead(d) var writeWorkers = populateWrite(d) - #echo(readWorkers.len, " ", d.requests[reqNil].len, d.requests[reqReadLine].len) - if select(readWorkers, writeWorkers, timeout) != 0: + #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) @@ -972,12 +994,20 @@ proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} = 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: $#, kind: reqAwait, worker: $#)" % - [$n[1][1].ident, callIdent]) + var yie = parseExpr("yield PRequest(socket: nil, kind: reqAwait, worker: $#)" % + [callIdent]) yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"), sym)) @@ -991,8 +1021,8 @@ proc toYieldReg(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} = var sym: PNimrodNode result.add(transformCallWithArg(n, sym)) # reqRegister - var yie = parseExpr("yield PRequest(socket: $#, kind: reqReg, worker: $#)" % - [$n[1][1].ident, callIdent]) + var yie = parseExpr("yield PRequest(socket: nil, kind: reqReg, worker: $#)" % + [callIdent]) yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"), sym)) @@ -1092,7 +1122,10 @@ proc createVerificationProc(procName: PNimrodNode, # 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) - body.add newCall("new", newIdentNode("result")) + 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"), @@ -1145,16 +1178,15 @@ macro async*(n: stmt): stmt {.immediate.} = result[6].add(body) # Add typedef above the proc def for parameters. - if n[3].len > 1: - 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 + 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 diff --git a/tests/compile/tasyncitermacro.nim b/tests/compile/tasyncitermacro.nim deleted file mode 100644 index 4b9b7b5b2..000000000 --- a/tests/compile/tasyncitermacro.nim +++ /dev/null @@ -1,42 +0,0 @@ -import sockets, asyncio, strutils - -proc auth3(client: PAsyncSocket) {.async.} = - await send(client, "Auth\c\L") - -proc auth2(client: PAsyncSocket) {.async.} = - await auth3(client) - -proc auth(client: PAsyncSocket) {.async.} = - await auth2(client) - -proc processRequest(client: PAsyncSocket, test: string, closeSock: bool = true) {.async.} = - assert test == "ahha" - assert client != nil - let line = await(readLine(client)) - echo("Read: ", line) - - for i in 0 .. 10: - await auth(client) - - await send(client, "Goodbye.\c\L") - if closeSock: - 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 - - reg processRequest(client, "ahha") - -var disp = newDispatcher(false) -disp.register(processServer, nil) -while true: - discard disp.poll() diff --git a/tests/compile/tasynciterraw.nim b/tests/compile/tasynciterraw.nim deleted file mode 100644 index ba038f3fc..000000000 --- a/tests/compile/tasynciterraw.nim +++ /dev/null @@ -1,69 +0,0 @@ -import asyncio, sockets - -const qwerty = "qwertyuiopasdfghjklzxcvbnm\c\L" - -iterator processRequest(x: PRequest): PRequest {.closure.} = - # Read first message - let client = x.socket - var lineReq = PRequest(socket: client, kind: reqReadLine, line: "") - yield lineReq - - if lineReq.line == "HELLO": - echo "Got first message." - - var longString = qwerty - for i in 0..109000: - longString.add($i & qwerty) - - var writeReq = PRequest(socket: client, kind: reqWrite, toWrite: longString) - yield writeReq - - # Read second message: No need to create another request object: - yield lineReq - - if lineReq.line == "DOM": - echo("Got second message.") - - client.close() - - -iterator processServer(x: PRequest): PRequest {.closure.} = - type - PTest = ref object of TObject - param: string - - let passedInParams = cast[PTest](x.param) - assert passedInParams.param == "blah" - - var sock = AsyncSocket() - - # blocks: - sock.bindAddr(TPort(6667)) - sock.listen() - - # Accept loop - while true: - var acceptReq = PRequest(socket: sock, kind: reqAccept, client: nil) - echo("About to yield accept") - yield acceptReq - echo("after yield") - if acceptReq.hasException: - echo("Got exception[EAssertionFailed]: ", acceptReq.exc of EAssertionFailed) - echo(acceptReq.exc.msg) - continue - - let client = acceptReq.client - assert client != nil - yield PRequest(socket: client, kind: reqReg, param: nil, worker: processRequest) - -when isMainModule: - let param = "blah" - type - PTest = ref object of TObject - param: type(param) - var test = PTest(param: param) - var disp = newDispatcher(false) - disp.register(processServer, test) - 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)