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.
This commit is contained in:
parent
ab534804c6
commit
b3469c43a6
5 changed files with 325 additions and 141 deletions
|
|
@ -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)
|
||||
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,7 +1178,6 @@ 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]))
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
@ -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()
|
||||
|
||||
78
tests/run/tasyncitermacro.nim
Normal file
78
tests/run/tasyncitermacro.nim
Normal file
|
|
@ -0,0 +1,78 @@
|
|||
discard """
|
||||
file: "tasyncitermacro.nim"
|
||||
cmd: "nimrod cc --hints:on $# $#"
|
||||
output: "10000"
|
||||
"""
|
||||
import sockets, asyncio, strutils
|
||||
|
||||
var globalCount = 0
|
||||
|
||||
proc sendCount3(client: PAsyncSocket, count: int) {.async.} =
|
||||
await send(client, $count & "\c\L")
|
||||
|
||||
proc sendCount2(client: PAsyncSocket, count: int) {.async.} =
|
||||
await sendCount3(client, count)
|
||||
|
||||
proc sendCount(client: PAsyncSocket, count: int) {.async.} =
|
||||
# Testing multiple levels of custom await calls. Shouldn't have much of a
|
||||
# performance impact.
|
||||
await sendCount2(client, count)
|
||||
|
||||
proc processRequest(client: PAsyncSocket, paramTest: string, closeSock: bool = true) {.async.} =
|
||||
doAssert paramTest == "foobarbaz"
|
||||
assert client != nil
|
||||
var count = 0
|
||||
while true:
|
||||
let line = await(readLine(client))
|
||||
if line == "end": break
|
||||
doAssert line.startswith("Message")
|
||||
doAssert line == "Message" & $count
|
||||
count.inc
|
||||
|
||||
await sendCount(client, count)
|
||||
|
||||
doAssert closeSock # Testing default params here. We won't be changing it.
|
||||
if closeSock:
|
||||
client.close()
|
||||
|
||||
proc startServer() {.async.} =
|
||||
var sock = AsyncSocket()
|
||||
#sock.setReuseAddr()
|
||||
|
||||
# The following may block, but I don't think it does.
|
||||
sock.bindAddr(TPort(10235))
|
||||
sock.listen()
|
||||
|
||||
# Accept loop
|
||||
while true:
|
||||
let client: PAsyncSocket = await(accept(sock))
|
||||
assert client != nil
|
||||
|
||||
reg processRequest(client, "foobarbaz")
|
||||
|
||||
proc spawnClient() {.async.} =
|
||||
var client = AsyncSocket()
|
||||
await connect(client, "localhost", TPort(10235))
|
||||
|
||||
for i in 0 .. 99:
|
||||
await send(client, "Message" & $i & "\c\L")
|
||||
await send(client, "end\c\L")
|
||||
|
||||
let line = await(readLine(client))
|
||||
doAssert line == "100"
|
||||
globalCount.inc(100)
|
||||
client.close()
|
||||
|
||||
proc spawnClients() {.async.} =
|
||||
for i in 0 .. 99:
|
||||
reg spawnClient()
|
||||
|
||||
var disp = newDispatcher(false)
|
||||
disp.register(startServer, nil)
|
||||
disp.register(spawnClients, nil)
|
||||
while true:
|
||||
discard disp.poll()
|
||||
|
||||
if globalCount == 100*100:
|
||||
echo(globalCount)
|
||||
quit(QuitSuccess)
|
||||
185
tests/run/tasynciterraw.nim
Normal file
185
tests/run/tasynciterraw.nim
Normal file
|
|
@ -0,0 +1,185 @@
|
|||
discard """
|
||||
file: "tasynciterraw.nim"
|
||||
cmd: "nimrod cc --hints:on $# $#"
|
||||
output: "10000"
|
||||
"""
|
||||
# This is a control test in case the macro expansion starts failing. It's the
|
||||
# raw code after macro expansion in tasyncitermacro.nim.
|
||||
|
||||
import asyncio, sockets, strutils
|
||||
|
||||
var globalCount = 0
|
||||
|
||||
type
|
||||
PsendCount3ArgObject = ref object of TObject
|
||||
dummy1: PAsyncSocket
|
||||
dummy2: int
|
||||
|
||||
proc sendCount3(client: PAsyncSocket; count: int): PsendCount3ArgObject =
|
||||
new(result)
|
||||
result.dummy1 = client
|
||||
result.dummy2 = count
|
||||
|
||||
iterator sendCount3(x: PRequest): PRequest {.closure.} =
|
||||
let passedInParams = PsendCount3ArgObject(x.param)
|
||||
let client = passedInParams.dummy1
|
||||
let count = passedInParams.dummy2
|
||||
var sendReq = PRequest(socket: client, kind: reqWrite,
|
||||
toWrite: $ count & "\x0D\x0A")
|
||||
yield sendReq
|
||||
if sendReq.hasException: raise sendReq.exc
|
||||
|
||||
|
||||
type
|
||||
PsendCount2ArgObject = ref object of TObject
|
||||
dummy1: PAsyncSocket
|
||||
dummy2: int
|
||||
|
||||
proc sendCount2(client: PAsyncSocket; count: int): PsendCount2ArgObject =
|
||||
new(result)
|
||||
result.dummy1 = client
|
||||
result.dummy2 = count
|
||||
|
||||
iterator sendCount2(x: PRequest): PRequest {.closure.} =
|
||||
let passedInParams = PsendCount2ArgObject(x.param)
|
||||
let client = passedInParams.dummy1
|
||||
let count = passedInParams.dummy2
|
||||
var argsToPass = sendCount3(client, count)
|
||||
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount3,
|
||||
param: argsToPass)
|
||||
|
||||
|
||||
type
|
||||
PsendCountArgObject = ref object of TObject
|
||||
dummy1: PAsyncSocket
|
||||
dummy2: int
|
||||
|
||||
proc sendCount(client: PAsyncSocket; count: int): PsendCountArgObject =
|
||||
new(result)
|
||||
result.dummy1 = client
|
||||
result.dummy2 = count
|
||||
|
||||
iterator sendCount(x: PRequest): PRequest {.closure.} =
|
||||
let passedInParams = PsendCountArgObject(x.param)
|
||||
let client = passedInParams.dummy1
|
||||
let count = passedInParams.dummy2
|
||||
var argsToPass = sendCount2(client, count)
|
||||
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount2,
|
||||
param: argsToPass)
|
||||
|
||||
|
||||
type
|
||||
PprocessRequestArgObject = ref object of TObject
|
||||
dummy1: PAsyncSocket
|
||||
dummy2: string
|
||||
dummy3: bool
|
||||
|
||||
proc processRequest(client: PAsyncSocket; paramTest: string;
|
||||
closeSock: bool = true): PprocessRequestArgObject =
|
||||
new(result)
|
||||
result.dummy1 = client
|
||||
result.dummy2 = paramTest
|
||||
result.dummy3 = closeSock
|
||||
|
||||
iterator processRequest(x: PRequest): PRequest {.closure.} =
|
||||
let passedInParams = PprocessRequestArgObject(x.param)
|
||||
let client = passedInParams.dummy1
|
||||
let paramTest = passedInParams.dummy2
|
||||
let closeSock = passedInParams.dummy3
|
||||
doAssert paramTest == "foobarbaz"
|
||||
assert client != nil
|
||||
var count = 0
|
||||
while true:
|
||||
var readLineReq = PRequest(socket: client, kind: reqReadLine, line: "")
|
||||
yield readLineReq
|
||||
if readLineReq.hasException: raise readLineReq.exc
|
||||
let line = readLineReq.line
|
||||
if line == "end":
|
||||
break
|
||||
doAssert line.startswith("Message")
|
||||
doAssert line == "Message" & $ count
|
||||
count.inc
|
||||
var argsToPass = sendCount(client, count)
|
||||
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount,
|
||||
param: argsToPass)
|
||||
doAssert closeSock
|
||||
if closeSock:
|
||||
client.close()
|
||||
|
||||
type
|
||||
PstartServerArgObject = ref object of TObject
|
||||
|
||||
proc startServer(): PstartServerArgObject =
|
||||
nil
|
||||
|
||||
iterator startServer(x: PRequest): PRequest {.closure.} =
|
||||
var sock = AsyncSocket()
|
||||
#sock.setReuseAddr()
|
||||
sock.bindAddr(TPort(10235))
|
||||
sock.listen()
|
||||
while true:
|
||||
var acceptReq = PRequest(socket: sock, kind: reqAccept, client: nil,
|
||||
hasException: false)
|
||||
yield acceptReq
|
||||
if acceptReq.hasException: raise acceptReq.exc
|
||||
let client = acceptReq.client
|
||||
assert client != nil
|
||||
var argsToPass = processRequest(client, "foobarbaz")
|
||||
yield PRequest(socket: nil, kind: reqReg, worker: processRequest,
|
||||
param: argsToPass)
|
||||
|
||||
|
||||
type
|
||||
PspawnClientArgObject = ref object of TObject
|
||||
|
||||
proc spawnClient(): PspawnClientArgObject =
|
||||
nil
|
||||
|
||||
iterator spawnClient(x: PRequest): PRequest {.closure.} =
|
||||
var client = AsyncSocket()
|
||||
var connectReq = PRequest(socket: client, kind: reqConnect,
|
||||
address: "localhost", port: TPort(10235))
|
||||
yield connectReq
|
||||
if connectReq.hasException: raise connectReq.exc
|
||||
for i in 0 .. 99:
|
||||
var sendReq = PRequest(socket: client, kind: reqWrite,
|
||||
toWrite: "Message" & $ i & "\x0D\x0A")
|
||||
yield sendReq
|
||||
if sendReq.hasException: raise sendReq.exc
|
||||
|
||||
|
||||
var sendReq1 = PRequest(socket: client, kind: reqWrite,
|
||||
toWrite: "end\x0D\x0A")
|
||||
yield sendReq1
|
||||
if sendReq1.hasException: raise sendReq1.exc
|
||||
|
||||
var readLineReq = PRequest(socket: client, kind: reqReadLine, line: "")
|
||||
yield readLineReq
|
||||
if readLineReq.hasException: raise readLineReq.exc
|
||||
|
||||
let line = readLineReq.line
|
||||
doAssert line == "100"
|
||||
globalCount.inc(100)
|
||||
client.close()
|
||||
|
||||
type
|
||||
PspawnClientsArgObject = ref object of TObject
|
||||
|
||||
proc spawnClients(): PspawnClientsArgObject =
|
||||
nil
|
||||
|
||||
iterator spawnClients(x: PRequest): PRequest {.closure.} =
|
||||
for i in 0 .. 99:
|
||||
var argsToPass = spawnClient()
|
||||
yield PRequest(socket: nil, kind: reqReg, worker: spawnClient,
|
||||
param: argsToPass)
|
||||
|
||||
var disp = newDispatcher(false)
|
||||
disp.register(startServer, nil)
|
||||
disp.register(spawnClients, nil)
|
||||
while true:
|
||||
discard disp.poll()
|
||||
|
||||
if globalCount == 100*100:
|
||||
echo(globalCount)
|
||||
quit(QuitSuccess)
|
||||
Loading…
Add table
Add a link
Reference in a new issue