From b3469c43a69c299779ba79c3f74a5f4c1ed9ebfc Mon Sep 17 00:00:00 2001 From: Dominik Picheta Date: Sat, 3 Aug 2013 14:19:17 +0100 Subject: [PATCH] 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)