Compare commits
18 commits
devel
...
asyncmacro
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b3469c43a6 | ||
|
|
ab534804c6 | ||
|
|
e42f877664 | ||
|
|
70d67daa14 | ||
|
|
ae64089e53 | ||
|
|
426c2d6c01 | ||
|
|
2330c95b22 | ||
|
|
adc74f8ab0 | ||
|
|
3ead23cb26 | ||
|
|
c5ff9650f7 | ||
|
|
2cedc1d1af | ||
|
|
3e7ea859ec | ||
|
|
d8659eaaa2 | ||
|
|
1e1a583903 | ||
|
|
a1a8925692 | ||
|
|
d15e67e020 | ||
|
|
84ea209fe1 | ||
|
|
75893d1b36 |
4 changed files with 948 additions and 57 deletions
|
|
@ -6,7 +6,7 @@
|
||||||
# distribution, for details about the copyright.
|
# 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
|
## This module implements an asynchronous event loop together with asynchronous sockets
|
||||||
## which use this event loop.
|
## 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
|
## on with the events. The type that you set userArg to must be inheriting from
|
||||||
## TObject!
|
## TObject!
|
||||||
##
|
##
|
||||||
## **Note:** If you want to provide async ability to your module please do not
|
## **Note:** If you are after using this module to provide async functionality
|
||||||
## use the ``TDelegate`` object, instead use ``PAsyncSocket``. It is possible
|
## for one of your modules then it is best to use PAsyncSocket if your module
|
||||||
## that in the future this type's fields will not be exported therefore breaking
|
## only requires sockets. For non-socket objects with a select-like interface
|
||||||
## your code.
|
## a TDelegate implementation should be created, if one doesn't already exist.
|
||||||
##
|
##
|
||||||
## **Warning:** The API of this module is unstable, and therefore is subject
|
## **Warning:** The API of this module is unstable, and therefore is subject
|
||||||
## to change.
|
## to change.
|
||||||
|
|
@ -109,10 +109,6 @@ type
|
||||||
|
|
||||||
PDelegate* = ref TDelegate
|
PDelegate* = ref TDelegate
|
||||||
|
|
||||||
PDispatcher* = ref TDispatcher
|
|
||||||
TDispatcher = object
|
|
||||||
delegates: seq[PDelegate]
|
|
||||||
|
|
||||||
PAsyncSocket* = ref TAsyncSocket
|
PAsyncSocket* = ref TAsyncSocket
|
||||||
TAsyncSocket* = object of TObject
|
TAsyncSocket* = object of TObject
|
||||||
socket: TSocket
|
socket: TSocket
|
||||||
|
|
@ -136,6 +132,52 @@ type
|
||||||
SockIdle, SockConnecting, SockConnected, SockListening, SockClosed,
|
SockIdle, SockConnecting, SockConnected, SockListening, SockClosed,
|
||||||
SockUDPBound
|
SockUDPBound
|
||||||
|
|
||||||
|
TRequestKind* = enum
|
||||||
|
reqNil, reqReg, reqAwait, reqRead, reqWrite, reqReadLine, reqAccept, reqConnect
|
||||||
|
|
||||||
|
PAsyncProc = iterator (x: PRequest): PRequest
|
||||||
|
|
||||||
|
PRequest* = ref object
|
||||||
|
socket*: PAsyncSocket
|
||||||
|
case hasException*: bool
|
||||||
|
of true:
|
||||||
|
exc*: ref EBase
|
||||||
|
of false: nil
|
||||||
|
case kind*: TRequestKind
|
||||||
|
of reqNil:
|
||||||
|
nil
|
||||||
|
of reqReg, reqAwait:
|
||||||
|
param*: PObject
|
||||||
|
worker*: PAsyncProc
|
||||||
|
of reqRead:
|
||||||
|
count*: int ## Request
|
||||||
|
readData*: string ## Response
|
||||||
|
of reqWrite:
|
||||||
|
toWrite*: string ## Request
|
||||||
|
written: int ## Internal data
|
||||||
|
of reqReadLine:
|
||||||
|
line*: string ## Response
|
||||||
|
of reqAccept:
|
||||||
|
client*: PAsyncSocket ## Response
|
||||||
|
of reqConnect:
|
||||||
|
address*: string ## Request
|
||||||
|
port*: TPort ## Request
|
||||||
|
|
||||||
|
PWorker* = ref object
|
||||||
|
worker*: PAsyncProc
|
||||||
|
x*: PRequest
|
||||||
|
lastReq*: PRequest
|
||||||
|
case hasParent*: bool
|
||||||
|
of true:
|
||||||
|
parent*: PWorker
|
||||||
|
else: nil
|
||||||
|
|
||||||
|
PDispatcher* = ref TDispatcher
|
||||||
|
TDispatcher = object
|
||||||
|
usesDelegates: bool
|
||||||
|
delegates: seq[PDelegate]
|
||||||
|
requests: array[TRequestKind, seq[PWorker]]
|
||||||
|
|
||||||
proc newDelegate*(): PDelegate =
|
proc newDelegate*(): PDelegate =
|
||||||
## Creates a new delegate.
|
## Creates a new delegate.
|
||||||
new(result)
|
new(result)
|
||||||
|
|
@ -382,9 +424,15 @@ proc accept*(server: PAsyncSocket): PAsyncSocket {.deprecated.} =
|
||||||
var address = ""
|
var address = ""
|
||||||
server.acceptAddr(result, 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)
|
new(result)
|
||||||
result.delegates = @[]
|
result.delegates = @[]
|
||||||
|
result.requests = newRequests()
|
||||||
|
result.usesDelegates = useDelegates
|
||||||
|
|
||||||
proc register*(d: PDispatcher, deleg: PDelegate) =
|
proc register*(d: PDispatcher, deleg: PDelegate) =
|
||||||
## Registers delegate ``deleg`` with dispatcher ``d``.
|
## Registers delegate ``deleg`` with dispatcher ``d``.
|
||||||
|
|
@ -569,6 +617,193 @@ proc select(readfds, writefds, exceptfds: var seq[PDelegate],
|
||||||
pruneSocketSet(writefds, (wr))
|
pruneSocketSet(writefds, (wr))
|
||||||
pruneSocketSet(exceptfds, (ex))
|
pruneSocketSet(exceptfds, (ex))
|
||||||
|
|
||||||
|
proc createFdSet(fd: var TFdSet, s: seq[PWorker], m: var int) =
|
||||||
|
FD_ZERO(fd)
|
||||||
|
for i in items(s):
|
||||||
|
m = max(m, int(i.lastReq.socket.getFD))
|
||||||
|
FD_SET(i.lastReq.socket.getFD, fd)
|
||||||
|
|
||||||
|
proc getWorkerSocket(worker: PWorker): PAsyncSocket =
|
||||||
|
worker.lastReq.socket
|
||||||
|
|
||||||
|
proc pruneSocketSet(s: var seq[PWorker], fd: var TFdSet) =
|
||||||
|
var i = 0
|
||||||
|
var L = s.len
|
||||||
|
while i < L:
|
||||||
|
if FD_ISSET(getWorkerSocket(s[i]).getFD, fd) != 0'i32:
|
||||||
|
s[i] = s[L-1]
|
||||||
|
dec(L)
|
||||||
|
else:
|
||||||
|
inc(i)
|
||||||
|
setLen(s, L)
|
||||||
|
|
||||||
|
proc select(readfds: var seq[PWorker], writefds: var seq[PWorker],
|
||||||
|
timeout = 500): int =
|
||||||
|
var tv {.noInit.}: TTimeVal = timeValFromMilliseconds(timeout)
|
||||||
|
|
||||||
|
var rd, wr, ex: TFdSet
|
||||||
|
var m = 0
|
||||||
|
createFdSet(rd, readfds, m)
|
||||||
|
createFdSet(wr, writefds, m)
|
||||||
|
#createFdSet(ex, @[], m)
|
||||||
|
|
||||||
|
if timeout != -1:
|
||||||
|
result = int(select(cint(m+1), addr(rd), addr(wr), addr(ex), addr(tv)))
|
||||||
|
else:
|
||||||
|
result = int(select(cint(m+1), addr(rd), addr(wr), addr(ex), nil))
|
||||||
|
|
||||||
|
pruneSocketSet(readfds, (rd))
|
||||||
|
pruneSocketSet(writefds, (wr))
|
||||||
|
#pruneSocketSet(exceptfds, (ex))
|
||||||
|
|
||||||
|
proc register*(disp: PDispatcher, worker: iterator (x: PRequest): PRequest,
|
||||||
|
param: PObject) =
|
||||||
|
#= PWorker(socket, worker, PRequest(kind: reqNil))
|
||||||
|
assert(not disp.usesDelegates,
|
||||||
|
"You need to set ``usesDelegates`` to false in the newDispatcher proc.")
|
||||||
|
var req = PWorker(
|
||||||
|
worker: worker,
|
||||||
|
x: PRequest(kind: reqReg, param: param, worker: nil),
|
||||||
|
lastReq: PRequest(kind: reqNil)
|
||||||
|
)
|
||||||
|
disp.requests[reqNil].add(req)
|
||||||
|
|
||||||
|
proc processWorkers(d: PDispatcher) =
|
||||||
|
var newRequests: array[TRequestKind, seq[PWorker]] = d.requests
|
||||||
|
newRequests[reqNil] = @[]
|
||||||
|
|
||||||
|
proc processWorker(idle: PWorker) =
|
||||||
|
let req = idle.worker(idle.x)
|
||||||
|
if req != nil:
|
||||||
|
case req.kind
|
||||||
|
of reqReg:
|
||||||
|
# Reset the lastReq. This caused odd issues in an 'accept' loop.
|
||||||
|
# My guess is that the request object got corrupted somehow, but
|
||||||
|
# I could not determine the cause. The issue that occurred was that
|
||||||
|
# 'await accept' simply stopped working after 6 clients (when connecting
|
||||||
|
# 10 simultaneously as seen in the tasyncitermacro test).
|
||||||
|
idle.lastReq = PRequest(kind: reqNil)
|
||||||
|
|
||||||
|
newRequests[reqNil].add(idle)
|
||||||
|
let newWorker = PWorker(worker: req.worker, lastReq: PRequest(kind: reqNil),
|
||||||
|
x: req)
|
||||||
|
newRequests[reqNil].add(newWorker)
|
||||||
|
of reqAwait:
|
||||||
|
let newWorker = PWorker(worker: req.worker, lastReq: PRequest(kind: reqNil),
|
||||||
|
x: req, hasParent: true, parent: idle)
|
||||||
|
# The worker which ``await``-ed this user-defined async proc;
|
||||||
|
# will be re-added to ``d.requests`` when ``newWorker`` finishes.
|
||||||
|
|
||||||
|
# We call this proc recursively so that the execution of this
|
||||||
|
# worker begins and we can immediately satisfy it's async request.
|
||||||
|
# We do not need to add newWorker to newRequests manually as it will be
|
||||||
|
# added by ``processWorker`` if necessary automatically.
|
||||||
|
processWorker(newWorker)
|
||||||
|
else:
|
||||||
|
idle.lastReq = req
|
||||||
|
newRequests[req.kind].add(idle)
|
||||||
|
else:
|
||||||
|
assert idle.worker.finished
|
||||||
|
if idle.hasParent:
|
||||||
|
# Re-add the parent worker, which is the worker which awaited this
|
||||||
|
# user-defined async proc which just finished. Do this by calling
|
||||||
|
# processWorker recursively. Same way as above.
|
||||||
|
processWorker(idle.parent)
|
||||||
|
|
||||||
|
for idle in d.requests[reqNil]:
|
||||||
|
processWorker(idle)
|
||||||
|
|
||||||
|
d.requests = newRequests
|
||||||
|
|
||||||
|
proc populateRead(d: PDispatcher): seq[PWorker] =
|
||||||
|
result = @[]
|
||||||
|
|
||||||
|
result.add d.requests[reqAccept]
|
||||||
|
result.add d.requests[reqRead]
|
||||||
|
result.add d.requests[reqReadLine]
|
||||||
|
|
||||||
|
proc populateWrite(d: PDispatcher): seq[PWorker] =
|
||||||
|
result = @[]
|
||||||
|
result.add d.requests[reqWrite]
|
||||||
|
result.add d.requests[reqConnect]
|
||||||
|
|
||||||
|
proc processRequests(requests, readWorkers, writeWorkers: seq[PWorker],
|
||||||
|
newRequests: var array[TRequestKind, seq[PWorker]]) =
|
||||||
|
for worker in requests:
|
||||||
|
#echo(worker.lastReq.kind)
|
||||||
|
var addTo = worker.lastReq.kind
|
||||||
|
template execReq(workers: var seq[PWorker], autoadd: bool,
|
||||||
|
body: stmt) {.immediate, dirty.} =
|
||||||
|
if worker notin workers:
|
||||||
|
# Worker is ready to read. Let's read.
|
||||||
|
try:
|
||||||
|
body
|
||||||
|
except:
|
||||||
|
worker.lastReq.hasException = true
|
||||||
|
worker.lastReq.exc = getCurrentException()
|
||||||
|
finally:
|
||||||
|
if autoAdd:
|
||||||
|
addTo = reqNil
|
||||||
|
|
||||||
|
case worker.lastReq.kind
|
||||||
|
of reqReadLine:
|
||||||
|
execReq readWorkers, false:
|
||||||
|
if worker.lastReq.socket.readLine(worker.lastReq.line):
|
||||||
|
addTo = reqNil
|
||||||
|
of reqAccept:
|
||||||
|
execReq readWorkers, true:
|
||||||
|
worker.lastReq.client = newAsyncSocket()
|
||||||
|
worker.lastReq.socket.accept(worker.lastReq.client)
|
||||||
|
of reqRead:
|
||||||
|
# We guarantee that all requested data will be read.
|
||||||
|
execReq readWorkers, false:
|
||||||
|
proc doRead(count: int) =
|
||||||
|
let got = worker.lastReq.socket.recvAsync(
|
||||||
|
worker.lastReq.readData, count)
|
||||||
|
assert got != -1
|
||||||
|
if got == count:
|
||||||
|
addTo = reqNil # Everything has been read
|
||||||
|
if worker.lastReq.readData.len == 0:
|
||||||
|
doRead(worker.lastReq.count)
|
||||||
|
else:
|
||||||
|
doRead(worker.lastReq.count-worker.lastReq.readData.len)
|
||||||
|
|
||||||
|
of reqWrite:
|
||||||
|
# We guarantee that all the data that is requested to be sent, will
|
||||||
|
# be sent.
|
||||||
|
|
||||||
|
execReq writeWorkers, false:
|
||||||
|
let written = worker.lastReq.written
|
||||||
|
proc doSend(toWrite: string) =
|
||||||
|
let len = toWrite.len
|
||||||
|
let sent = worker.lastReq.socket.sendAsync(toWrite)
|
||||||
|
assert sent != 0 # /Something/ should have been written.
|
||||||
|
if sent == len:
|
||||||
|
# Sent all data, request complete.
|
||||||
|
addTo = reqNil
|
||||||
|
else:
|
||||||
|
# Didn't send all data, must send the rest later.
|
||||||
|
worker.lastReq.written.inc(sent)
|
||||||
|
|
||||||
|
if written == 0:
|
||||||
|
doSend(worker.lastReq.toWrite)
|
||||||
|
else:
|
||||||
|
let toWrite = worker.lastReq.toWrite[written .. -1]
|
||||||
|
doSend(toWrite)
|
||||||
|
of reqConnect:
|
||||||
|
if not worker.lastReq.socket.isConnecting:
|
||||||
|
worker.lastReq.socket.connect(worker.lastReq.address, worker.lastReq.port)
|
||||||
|
else:
|
||||||
|
execReq writeWorkers, true:
|
||||||
|
worker.lastReq.socket.info = SockConnected
|
||||||
|
of reqReg, reqAwait:
|
||||||
|
assert false, $worker.lastReq.kind & " should have been processed already"
|
||||||
|
of reqNil:
|
||||||
|
# Nothing to do. Most likely that a new worker has just been
|
||||||
|
# registered.
|
||||||
|
|
||||||
|
newRequests[addTo].add(worker)
|
||||||
|
|
||||||
proc poll*(d: PDispatcher, timeout: int = 500): bool =
|
proc poll*(d: PDispatcher, timeout: int = 500): bool =
|
||||||
## This function checks for events on all the delegates in the `PDispatcher`.
|
## This function checks for events on all the delegates in the `PDispatcher`.
|
||||||
## It then proceeds to call the correct event handler.
|
## It then proceeds to call the correct event handler.
|
||||||
|
|
@ -582,58 +817,419 @@ proc poll*(d: PDispatcher, timeout: int = 500): bool =
|
||||||
## only be executed after one or more file descriptors becomes readable or
|
## only be executed after one or more file descriptors becomes readable or
|
||||||
## writeable.
|
## writeable.
|
||||||
result = true
|
result = true
|
||||||
var readDg, writeDg, errorDg: seq[PDelegate] = @[]
|
if d.usesDelegates:
|
||||||
var len = d.delegates.len
|
var readDg, writeDg, errorDg: seq[PDelegate] = @[]
|
||||||
var dc = 0
|
var len = d.delegates.len
|
||||||
|
var dc = 0
|
||||||
while dc < len:
|
|
||||||
let deleg = d.delegates[dc]
|
while dc < len:
|
||||||
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and deleg.open:
|
let deleg = d.delegates[dc]
|
||||||
readDg.add(deleg)
|
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and deleg.open:
|
||||||
if (deleg.mode != fmRead) and deleg.open:
|
readDg.add(deleg)
|
||||||
writeDg.add(deleg)
|
if (deleg.mode != fmRead) and deleg.open:
|
||||||
if deleg.open:
|
writeDg.add(deleg)
|
||||||
errorDg.add(deleg)
|
if deleg.open:
|
||||||
inc dc
|
errorDg.add(deleg)
|
||||||
else:
|
inc dc
|
||||||
# File/socket has been closed. Remove it from dispatcher.
|
else:
|
||||||
d.delegates[dc] = d.delegates[len-1]
|
# File/socket has been closed. Remove it from dispatcher.
|
||||||
dec len
|
d.delegates[dc] = d.delegates[len-1]
|
||||||
|
dec len
|
||||||
d.delegates.setLen(len)
|
|
||||||
|
d.delegates.setLen(len)
|
||||||
var hasDataBufferedCount = 0
|
|
||||||
for d in d.delegates:
|
var hasDataBufferedCount = 0
|
||||||
if d.hasDataBuffered(d.deleVal):
|
for d in d.delegates:
|
||||||
hasDataBufferedCount.inc()
|
if d.hasDataBuffered(d.deleVal):
|
||||||
d.handleRead(d.deleVal)
|
hasDataBufferedCount.inc()
|
||||||
if hasDataBufferedCount > 0: return True
|
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?
|
if readDg.len() == 0 and writeDg.len() == 0:
|
||||||
return False
|
## 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 select(readDg, writeDg, errorDg, timeout) != 0:
|
||||||
if i > len(d.delegates)-1: break # One delegate might've been removed.
|
for i in 0..len(d.delegates)-1:
|
||||||
let deleg = d.delegates[i]
|
if i > len(d.delegates)-1: break # One delegate might've been removed.
|
||||||
if not deleg.open: continue # This delegate might've been closed.
|
let deleg = d.delegates[i]
|
||||||
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and
|
if not deleg.open: continue # This delegate might've been closed.
|
||||||
deleg notin readDg:
|
if (deleg.mode != fmWrite or deleg.mode != fmAppend) and
|
||||||
deleg.handleRead(deleg.deleVal)
|
deleg notin readDg:
|
||||||
if (deleg.mode != fmRead) and deleg notin writeDg:
|
deleg.handleRead(deleg.deleVal)
|
||||||
deleg.handleWrite(deleg.deleVal)
|
if (deleg.mode != fmRead) and deleg notin writeDg:
|
||||||
if deleg notin errorDg:
|
deleg.handleWrite(deleg.deleVal)
|
||||||
deleg.handleError(deleg.deleVal)
|
if deleg notin errorDg:
|
||||||
|
deleg.handleError(deleg.deleVal)
|
||||||
# Execute tasks
|
|
||||||
for i in items(d.delegates):
|
# Execute tasks
|
||||||
i.task(i.deleVal)
|
for i in items(d.delegates):
|
||||||
|
i.task(i.deleVal)
|
||||||
|
else:
|
||||||
|
# Async worker iterators
|
||||||
|
processWorkers(d)
|
||||||
|
var readWorkers = populateRead(d)
|
||||||
|
var writeWorkers = populateWrite(d)
|
||||||
|
#echo("ReadWorkers: ", readWorkers.len, " | WriteWorkers: ", writeWorkers.len,
|
||||||
|
# " | Idle Workers: ", d.requests[reqNil].len,
|
||||||
|
# " | Workers waiting for readLine: ", d.requests[reqReadLine].len,
|
||||||
|
# " | Workers waiting for accept: ", d.requests[reqAccept].len)
|
||||||
|
|
||||||
|
# Check buffer state.
|
||||||
|
var isDataBuffered = false
|
||||||
|
var newReadWorkers: seq[PWorker] = @[]
|
||||||
|
for i in readWorkers:
|
||||||
|
if not i.lastReq.socket.hasDataBuffered(): newReadWorkers.add(i)
|
||||||
|
else: isDataBuffered = true
|
||||||
|
readWorkers = newReadWorkers
|
||||||
|
|
||||||
|
#echo(if isDataBuffered: "Buffered" else: "Not Buffered")
|
||||||
|
|
||||||
|
var doProcess = isDataBuffered
|
||||||
|
if not doProcess: doProcess = select(readWorkers, writeWorkers, timeout) != 0
|
||||||
|
if doProcess:
|
||||||
|
var newRequests: array[TRequestKind, seq[PWorker]] = newRequests()
|
||||||
|
for req in TRequestKind:
|
||||||
|
processRequests(d.requests[req], readWorkers, writeWorkers, newRequests)
|
||||||
|
d.requests = newRequests
|
||||||
|
|
||||||
proc len*(disp: PDispatcher): int =
|
proc len*(disp: PDispatcher): int =
|
||||||
## Retrieves the amount of delegates in ``disp``.
|
## Retrieves the amount of delegates in ``disp``.
|
||||||
return disp.delegates.len
|
return disp.delegates.len
|
||||||
|
|
||||||
|
# ---- Async macro
|
||||||
|
|
||||||
|
proc createRequestNode(varName,
|
||||||
|
reqArgs: string, sym: var PNimrodNode): PNimrodNode {.compiletime.} =
|
||||||
|
## Creates a constructor for a PRequest object, which will be then yielded.
|
||||||
|
result = newNimNode(nnkStmtList)
|
||||||
|
# TODO: Using gensym here causes segfaults because hasException is not
|
||||||
|
# initialised.
|
||||||
|
sym = newIdentNode(varName) #genSym(nskVar, varName)
|
||||||
|
var reqObj = newVarStmt(sym,
|
||||||
|
parseExpr("PRequest($#)" % [reqArgs]))
|
||||||
|
|
||||||
|
result.add reqObj
|
||||||
|
result.add newNimNode(nnkYieldStmt).add(sym)
|
||||||
|
# Check for exception
|
||||||
|
# TODO: Create a custom exception type, create a field which will store
|
||||||
|
# the original stack trace as given by getStackTrace. Maybe there is a way
|
||||||
|
# to override the stack trace? That'd be nice.
|
||||||
|
result.add newIfStmt(
|
||||||
|
(newDotExpr(sym, newIdentNode("hasException")),
|
||||||
|
newNimNode(nnkRaiseStmt).add(
|
||||||
|
newDotExpr(sym, newIdentNode("exc")))))
|
||||||
|
|
||||||
|
proc toYieldVar(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} =
|
||||||
|
## Transforms a var/let section
|
||||||
|
## E.g:
|
||||||
|
## let client = await(accept(server))
|
||||||
|
result = @[]
|
||||||
|
let nameIdent = n[0][0] # Var name
|
||||||
|
expectLen(n[0], 3) # IdentDefs
|
||||||
|
let insideAwait = n[0][2][1]
|
||||||
|
let reqCall = $insideAwait[0].ident
|
||||||
|
|
||||||
|
expectLen(insideAwait, 2)
|
||||||
|
let sockName = $insideAwait[1].ident
|
||||||
|
case reqCall.normalize
|
||||||
|
of "accept":
|
||||||
|
let acceptReqVar = "acceptReq"
|
||||||
|
var sym: PNimrodNode
|
||||||
|
result.add createRequestNode(acceptReqVar,
|
||||||
|
"socket: $#, kind: reqAccept, client: nil, hasException: false" % sockName, sym)
|
||||||
|
|
||||||
|
case n.kind
|
||||||
|
of nnkLetSection:
|
||||||
|
result.add newLetStmt(nameIdent, newDotExpr(sym, newIdentNode("client")))
|
||||||
|
of nnkVarSection:
|
||||||
|
result.add newVarStmt(nameIdent, newDotExpr(sym, newIdentNode("client")))
|
||||||
|
else: error "Bad node kind in toYieldVar"
|
||||||
|
of "readline":
|
||||||
|
let readReqVar = "readLineReq"
|
||||||
|
var sym: PNimrodNode
|
||||||
|
result.add createRequestNode(readReqVar,
|
||||||
|
"socket: $#, kind: reqReadLine, line: \"\"" % sockName, sym)
|
||||||
|
case n.kind
|
||||||
|
of nnkLetSection:
|
||||||
|
result.add newLetStmt(nameIdent, newDotExpr(sym, newIdentNode("line")))
|
||||||
|
of nnkVarSection:
|
||||||
|
result.add newVarStmt(nameIdent, newDotExpr(sym, newIdentNode("line")))
|
||||||
|
else: error "Bad node kind in toYieldVar"
|
||||||
|
else:
|
||||||
|
error(reqCall & " is not a valid async call")
|
||||||
|
|
||||||
|
const typeDef =
|
||||||
|
"""
|
||||||
|
type
|
||||||
|
P$#ArgObject = ref object of TObject
|
||||||
|
"""
|
||||||
|
|
||||||
|
proc transformCallWithArg(call: PNimrodNode,
|
||||||
|
sym: var PNimrodNode): PNimrodNode {.compiletime.} =
|
||||||
|
## Transforms an async await call of a user-defined proc into a
|
||||||
|
## ``reqReg`` yield.
|
||||||
|
result = newNimNode(nnkStmtList)
|
||||||
|
|
||||||
|
sym = gensym(nskVar, "argsToPass")
|
||||||
|
|
||||||
|
# Add in a call to a pre-generated proc stub so that the compiler verifies
|
||||||
|
# the params for us :)
|
||||||
|
# TODO: Inline this maybe?
|
||||||
|
var args: seq[PNimrodNode] = @[]
|
||||||
|
for i in 1 .. call[1].len-1:
|
||||||
|
args.add(call[1][i])
|
||||||
|
result.add newVarStmt(sym, newCall(call[1][0], args))
|
||||||
|
|
||||||
|
proc toYieldCall(n: PNimrodNode): seq[PNimrodNode] {.compileTime.} =
|
||||||
|
## Transforms a call/command
|
||||||
|
if $n[0].ident != "await": error "'await' expected"
|
||||||
|
result = @[]
|
||||||
|
let callIdent = $n[1][0].ident
|
||||||
|
case callIdent.normalize
|
||||||
|
of "send":
|
||||||
|
let socketName = $n[1][1].ident
|
||||||
|
let toWrite = n[1][2]
|
||||||
|
var sym: PNimrodNode
|
||||||
|
result.add createRequestNode("sendReq",
|
||||||
|
"socket: $#, kind: reqWrite, toWrite: $#" %
|
||||||
|
[socketName, $(toWrite.toStrLit)], sym)
|
||||||
|
of "connect":
|
||||||
|
let socketName = $n[1][1].ident
|
||||||
|
let address = n[1][2]
|
||||||
|
let port = n[1][3]
|
||||||
|
var sym: PNimrodNode
|
||||||
|
result.add createRequestNode("connectReq",
|
||||||
|
"socket: $#, kind: reqConnect, address: $#, port: $#" %
|
||||||
|
[socketName, $(address.toStrLit), $(port.toStrLit)], sym)
|
||||||
|
else:
|
||||||
|
var sym: PNimrodNode
|
||||||
|
result.add(transformCallWithArg(n, sym))
|
||||||
|
# reqCustom
|
||||||
|
var yie = parseExpr("yield PRequest(socket: nil, kind: reqAwait, worker: $#)" %
|
||||||
|
[callIdent])
|
||||||
|
yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"),
|
||||||
|
sym))
|
||||||
|
|
||||||
|
result.add yie
|
||||||
|
|
||||||
|
proc toYieldReg(n: PNimrodNode): seq[PNimrodNode] {.compiletime.} =
|
||||||
|
## Transforms the 'reg' command to a reqReg yield request.
|
||||||
|
if $n[0].ident != "reg": error "'reg' expected"
|
||||||
|
result = @[]
|
||||||
|
let callIdent = $n[1][0].ident
|
||||||
|
var sym: PNimrodNode
|
||||||
|
result.add(transformCallWithArg(n, sym))
|
||||||
|
# reqRegister
|
||||||
|
var yie = parseExpr("yield PRequest(socket: nil, kind: reqReg, worker: $#)" %
|
||||||
|
[callIdent])
|
||||||
|
yie[0].add(newNimNode(nnkExprColonExpr).add(newIdentNode("param"),
|
||||||
|
sym))
|
||||||
|
|
||||||
|
result.add yie
|
||||||
|
|
||||||
|
proc transform(n: PNimrodNode): PNimrodNode {.compiletime.} =
|
||||||
|
## Transforms body.
|
||||||
|
## Specifically it does the following:
|
||||||
|
##
|
||||||
|
## * Looks for 'await' and transforms it into a yield.
|
||||||
|
## * Handles arguments correctly.
|
||||||
|
result = newNimNode(nnkStmtList)
|
||||||
|
expectKind(n, nnkStmtList)
|
||||||
|
for i in 0 .. n.len-1:
|
||||||
|
var son = n[i]
|
||||||
|
case son.kind
|
||||||
|
of nnkVarSection, nnkLetSection:
|
||||||
|
for defs in 0 .. son.len-1:
|
||||||
|
var doAdd = true
|
||||||
|
let identDefs = son[defs]
|
||||||
|
expectKind(identDefs, nnkIdentDefs)
|
||||||
|
if identDefs[2].kind == nnkCall:
|
||||||
|
let callIdent = identDefs[2][0]
|
||||||
|
expectKind(callIdent, nnkIdent)
|
||||||
|
if $callIdent.ident == "await":
|
||||||
|
# Transform into yield.
|
||||||
|
result.add(toYieldVar(son))
|
||||||
|
doAdd = false
|
||||||
|
if doAdd:
|
||||||
|
var letOrVarSection = newNimNode(son.kind)
|
||||||
|
letOrVarSection.add(identDefs)
|
||||||
|
result.add(letOrVarSection)
|
||||||
|
of nnkWhileStmt:
|
||||||
|
son[1] = transform(son[1])
|
||||||
|
result.add(son)
|
||||||
|
of nnkForStmt:
|
||||||
|
son[2] = transform(son[2])
|
||||||
|
result.add(son)
|
||||||
|
of nnkCall, nnkCommand:
|
||||||
|
if son[0].kind == nnkIdent:
|
||||||
|
case $son[0].ident
|
||||||
|
of "await":
|
||||||
|
result.add toYieldCall(son)
|
||||||
|
of "reg":
|
||||||
|
result.add toYieldReg(son)
|
||||||
|
else:
|
||||||
|
result.add son
|
||||||
|
else:
|
||||||
|
result.add son
|
||||||
|
else:
|
||||||
|
result.add(son)
|
||||||
|
|
||||||
|
proc transformArgs(procName: string,
|
||||||
|
formalParams: PNimrodNode): PNimrodNode {.compiletime.} =
|
||||||
|
## Transforms formal params into a typedef with a dummy type
|
||||||
|
## ``ref object of TObject``. This is inserted above the proc definition.
|
||||||
|
expectKind(formalParams, nnkFormalParams)
|
||||||
|
result = parseStmt(typeDef % procName)
|
||||||
|
|
||||||
|
var RecList = newNimNode(nnkRecList)
|
||||||
|
|
||||||
|
for i in 1 .. formalParams.len-1:
|
||||||
|
expectKind(formalParams[i], nnkIdentDefs)
|
||||||
|
RecList.add(newIdentDefs(newIdentNode("dummy" & $i), formalParams[i][1]))
|
||||||
|
# TODO: Add comment with the original param name?
|
||||||
|
|
||||||
|
result[0][0][2][0][2] = RecList
|
||||||
|
|
||||||
|
proc declareArgsInBody(procName: string,
|
||||||
|
formalParams: PNimrodNode): PNimrodNode {.compiletime.} =
|
||||||
|
## Creates a local immutable var by casting the PRequest.param.
|
||||||
|
## Immutable vars are then defined as specified in the proc's params.
|
||||||
|
result = newNimNode(nnkStmtList)
|
||||||
|
var sym = genSym(ident = "passedInParams")
|
||||||
|
result.add newLetStmt(sym, parseExpr("P$#ArgObject(x.param)" % procName))
|
||||||
|
# Fields take the form ``dummy<i>``.
|
||||||
|
for i in 1 .. formalParams.len-1:
|
||||||
|
expectKind(formalParams[i], nnkIdentDefs)
|
||||||
|
result.add newLetStmt(formalParams[i][0],
|
||||||
|
newDotExpr(sym, newIdentNode("dummy" & $i)))
|
||||||
|
|
||||||
|
proc isDocumentation(n: PNimrodNode): bool {.compiletime.} =
|
||||||
|
## Determines whether this proc def is a docs stub.
|
||||||
|
result = true
|
||||||
|
for i in 0 .. n[6].len-1:
|
||||||
|
if n[6][i].kind != nnkCommentStmt:
|
||||||
|
return false
|
||||||
|
|
||||||
|
proc createVerificationProc(procName: PNimrodNode,
|
||||||
|
formalParams: PNimrodNode): PNimrodNode {.compiletime.} =
|
||||||
|
# TODO: Export this stub if our async proc is exported?
|
||||||
|
# Generate list of parameters for the proc. First param is the return type.
|
||||||
|
var params: seq[PNimrodNode] = @[newIdentNode("P$#ArgObject" % $procName.ident)]
|
||||||
|
for i in 1 .. formalParams.len-1:
|
||||||
|
params.add(formalParams[i])
|
||||||
|
|
||||||
|
# Generate body. We construct the ArgObject here, this is done so that
|
||||||
|
# default variables of the async proc can be captured.
|
||||||
|
var body = newNimNode(nnkStmtList)
|
||||||
|
if formalParams.len > 1:
|
||||||
|
body.add newCall("new", newIdentNode("result"))
|
||||||
|
else:
|
||||||
|
body.add parseExpr("nil")
|
||||||
|
|
||||||
|
for i in 1 .. formalParams.len-1:
|
||||||
|
let dotExpr = newDotExpr(newIdentNode("result"),
|
||||||
|
newIdentNode("dummy" & $i))
|
||||||
|
body.add newAssignment(dotExpr, formalParams[i][0])
|
||||||
|
|
||||||
|
result = newProc(procName, params, body)
|
||||||
|
|
||||||
|
macro async*(n: stmt): stmt {.immediate.} =
|
||||||
|
expectKind(n, nnkProcDef)
|
||||||
|
#echo(treeRepr(n))
|
||||||
|
if n.isDocumentation():
|
||||||
|
# Documentation stub?
|
||||||
|
# TODO: Give it an async tag?
|
||||||
|
# TODO: Doc strings are not generated in doc2
|
||||||
|
result = n
|
||||||
|
result[6] = parseStmt("nil")
|
||||||
|
return
|
||||||
|
|
||||||
|
result = newNimNode(nnkIteratorDef)
|
||||||
|
for i in 0 .. n.len-1:
|
||||||
|
result.add(copyNimTree(n[i]))
|
||||||
|
|
||||||
|
# Populate ``FormalParams``
|
||||||
|
assert result[3].kind == nnkFormalParams
|
||||||
|
let formalParams = newNimNode(nnkFormalParams)
|
||||||
|
formalParams.add(newIdentNode(!"PRequest")) # Return type
|
||||||
|
var params = newNimNode(nnkIdentDefs)
|
||||||
|
params.add(newIdentNode(!"x")) # First param name
|
||||||
|
params.add(newIdentNode(!"PRequest")) # First param type
|
||||||
|
params.add(newNimNode(nnkEmpty))
|
||||||
|
formalParams.add(params)
|
||||||
|
result[3] = formalParams
|
||||||
|
|
||||||
|
# Closure pragma
|
||||||
|
result[4].add(newIdentNode(!"closure"))
|
||||||
|
|
||||||
|
# Body
|
||||||
|
result[6] = newNimNode(nnkStmtList)
|
||||||
|
|
||||||
|
# Declare variables based on the params that the async proc takes.
|
||||||
|
# i.e. extract them from the PRequest object.
|
||||||
|
if n[3].len > 1:
|
||||||
|
let args = declareArgsInBody($n[0].ident, n[3])
|
||||||
|
|
||||||
|
result[6].add(args)
|
||||||
|
|
||||||
|
# Transform body
|
||||||
|
var body = transform(n[6])
|
||||||
|
result[6].add(body)
|
||||||
|
|
||||||
|
# Add typedef above the proc def for parameters.
|
||||||
|
let procDef = copyNimTree(result)
|
||||||
|
result = newNimNode(nnkStmtList)
|
||||||
|
result.add(transformArgs($n[0].ident, n[3]))
|
||||||
|
# Generate a proc to verify that the user passes the correct params.
|
||||||
|
# The proc also constructs the ArgObject, this is so that default params
|
||||||
|
# can be captured into the ArgObject.
|
||||||
|
result.add createVerificationProc(n[0], n[3])
|
||||||
|
|
||||||
|
result.add procDef
|
||||||
|
|
||||||
|
#echo treeRepr(result)
|
||||||
|
echo result.toStrLit().strVal
|
||||||
|
|
||||||
|
# ---- Async macro end
|
||||||
|
|
||||||
|
## Asynchronous IO without callbacks
|
||||||
|
## =====================
|
||||||
|
## The implementation of this is similar to what is known as ``futures``
|
||||||
|
## in other languages. The ``await`` keyword is used to call a procedure
|
||||||
|
## which is marked with an ``{.async.}`` pragma. The ``async`` procedures
|
||||||
|
## get transformed into an iterator. When you execute a procedure which may
|
||||||
|
## block using ``await`` then the execution of your procedure will stop
|
||||||
|
## until the procedure you executed has completed. While your procedure's
|
||||||
|
## execution is stopped, other async procedures can continue executing.
|
||||||
|
##
|
||||||
|
## The built-in functions which can be used using ``await`` are listed below:
|
||||||
|
##
|
||||||
|
## .. code-block:: nimrod
|
||||||
|
##
|
||||||
|
## proc send*(socket: PAsyncSocket, text: string) {.async.} =
|
||||||
|
##
|
||||||
|
## Sends ``text`` to ``socket`` asynchronously.
|
||||||
|
##
|
||||||
|
## .. code-block:: nimrod
|
||||||
|
##
|
||||||
|
## proc accept*(socket: PAsyncSocket): PAsyncSocket {.async.} =
|
||||||
|
##
|
||||||
|
## Accepts a client connecting to a server socket asynchronously.
|
||||||
|
## Returns that client.
|
||||||
|
##
|
||||||
|
## You may also define your own async procedures by annotating them with
|
||||||
|
## the ``{.async.}`` pragma, please note however that currently you cannot
|
||||||
|
## overload these procedures.
|
||||||
|
|
||||||
|
discard """ when defined(nimdoc) and isMainModule:
|
||||||
|
proc send*(socket: PAsyncSocket, text: string) {.async.} =
|
||||||
|
## Sends ``text`` to ``socket`` asynchronously.
|
||||||
|
|
||||||
|
proc accept*(socket: PAsyncSocket): PAsyncSocket {.async.} =
|
||||||
|
## Accepts a client connecting to a server socket asynchronously.
|
||||||
|
## Returns that client. """
|
||||||
|
|
||||||
when isMainModule:
|
when isMainModule:
|
||||||
|
|
||||||
proc testConnect(s: PAsyncSocket, no: int) =
|
proc testConnect(s: PAsyncSocket, no: int) =
|
||||||
|
|
|
||||||
32
tests/reject/tasyncmacroinvalidargs.nim
Normal file
32
tests/reject/tasyncmacroinvalidargs.nim
Normal file
|
|
@ -0,0 +1,32 @@
|
||||||
|
discard """
|
||||||
|
msg: "type mismatch: got (int literal(234)) but expected 'string'"
|
||||||
|
"""
|
||||||
|
import sockets, asyncio, strutils
|
||||||
|
|
||||||
|
proc processRequest(client: PAsyncSocket, test: string) {.async.} =
|
||||||
|
assert test == "ahha"
|
||||||
|
assert client != nil
|
||||||
|
echo("Test = ", test)
|
||||||
|
let line = await(readLine(client))
|
||||||
|
echo("Read: ", line)
|
||||||
|
|
||||||
|
await send(client, "Goodbye.\c\L")
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
proc processServer() {.async.} =
|
||||||
|
var sock = AsyncSocket()
|
||||||
|
|
||||||
|
# blocks:
|
||||||
|
sock.bindAddr(TPort(6667))
|
||||||
|
sock.listen()
|
||||||
|
|
||||||
|
# Accept loop
|
||||||
|
while true:
|
||||||
|
let client: PAsyncSocket = await(accept(sock))
|
||||||
|
assert client != nil
|
||||||
|
await processRequest(client, 234)
|
||||||
|
|
||||||
|
var disp = newDispatcher(false)
|
||||||
|
disp.register(processServer, nil)
|
||||||
|
while true:
|
||||||
|
discard disp.poll()
|
||||||
78
tests/run/tasyncitermacro.nim
Normal file
78
tests/run/tasyncitermacro.nim
Normal file
|
|
@ -0,0 +1,78 @@
|
||||||
|
discard """
|
||||||
|
file: "tasyncitermacro.nim"
|
||||||
|
cmd: "nimrod cc --hints:on $# $#"
|
||||||
|
output: "10000"
|
||||||
|
"""
|
||||||
|
import sockets, asyncio, strutils
|
||||||
|
|
||||||
|
var globalCount = 0
|
||||||
|
|
||||||
|
proc sendCount3(client: PAsyncSocket, count: int) {.async.} =
|
||||||
|
await send(client, $count & "\c\L")
|
||||||
|
|
||||||
|
proc sendCount2(client: PAsyncSocket, count: int) {.async.} =
|
||||||
|
await sendCount3(client, count)
|
||||||
|
|
||||||
|
proc sendCount(client: PAsyncSocket, count: int) {.async.} =
|
||||||
|
# Testing multiple levels of custom await calls. Shouldn't have much of a
|
||||||
|
# performance impact.
|
||||||
|
await sendCount2(client, count)
|
||||||
|
|
||||||
|
proc processRequest(client: PAsyncSocket, paramTest: string, closeSock: bool = true) {.async.} =
|
||||||
|
doAssert paramTest == "foobarbaz"
|
||||||
|
assert client != nil
|
||||||
|
var count = 0
|
||||||
|
while true:
|
||||||
|
let line = await(readLine(client))
|
||||||
|
if line == "end": break
|
||||||
|
doAssert line.startswith("Message")
|
||||||
|
doAssert line == "Message" & $count
|
||||||
|
count.inc
|
||||||
|
|
||||||
|
await sendCount(client, count)
|
||||||
|
|
||||||
|
doAssert closeSock # Testing default params here. We won't be changing it.
|
||||||
|
if closeSock:
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
proc startServer() {.async.} =
|
||||||
|
var sock = AsyncSocket()
|
||||||
|
#sock.setReuseAddr()
|
||||||
|
|
||||||
|
# The following may block, but I don't think it does.
|
||||||
|
sock.bindAddr(TPort(10235))
|
||||||
|
sock.listen()
|
||||||
|
|
||||||
|
# Accept loop
|
||||||
|
while true:
|
||||||
|
let client: PAsyncSocket = await(accept(sock))
|
||||||
|
assert client != nil
|
||||||
|
|
||||||
|
reg processRequest(client, "foobarbaz")
|
||||||
|
|
||||||
|
proc spawnClient() {.async.} =
|
||||||
|
var client = AsyncSocket()
|
||||||
|
await connect(client, "localhost", TPort(10235))
|
||||||
|
|
||||||
|
for i in 0 .. 99:
|
||||||
|
await send(client, "Message" & $i & "\c\L")
|
||||||
|
await send(client, "end\c\L")
|
||||||
|
|
||||||
|
let line = await(readLine(client))
|
||||||
|
doAssert line == "100"
|
||||||
|
globalCount.inc(100)
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
proc spawnClients() {.async.} =
|
||||||
|
for i in 0 .. 99:
|
||||||
|
reg spawnClient()
|
||||||
|
|
||||||
|
var disp = newDispatcher(false)
|
||||||
|
disp.register(startServer, nil)
|
||||||
|
disp.register(spawnClients, nil)
|
||||||
|
while true:
|
||||||
|
discard disp.poll()
|
||||||
|
|
||||||
|
if globalCount == 100*100:
|
||||||
|
echo(globalCount)
|
||||||
|
quit(QuitSuccess)
|
||||||
185
tests/run/tasynciterraw.nim
Normal file
185
tests/run/tasynciterraw.nim
Normal file
|
|
@ -0,0 +1,185 @@
|
||||||
|
discard """
|
||||||
|
file: "tasynciterraw.nim"
|
||||||
|
cmd: "nimrod cc --hints:on $# $#"
|
||||||
|
output: "10000"
|
||||||
|
"""
|
||||||
|
# This is a control test in case the macro expansion starts failing. It's the
|
||||||
|
# raw code after macro expansion in tasyncitermacro.nim.
|
||||||
|
|
||||||
|
import asyncio, sockets, strutils
|
||||||
|
|
||||||
|
var globalCount = 0
|
||||||
|
|
||||||
|
type
|
||||||
|
PsendCount3ArgObject = ref object of TObject
|
||||||
|
dummy1: PAsyncSocket
|
||||||
|
dummy2: int
|
||||||
|
|
||||||
|
proc sendCount3(client: PAsyncSocket; count: int): PsendCount3ArgObject =
|
||||||
|
new(result)
|
||||||
|
result.dummy1 = client
|
||||||
|
result.dummy2 = count
|
||||||
|
|
||||||
|
iterator sendCount3(x: PRequest): PRequest {.closure.} =
|
||||||
|
let passedInParams = PsendCount3ArgObject(x.param)
|
||||||
|
let client = passedInParams.dummy1
|
||||||
|
let count = passedInParams.dummy2
|
||||||
|
var sendReq = PRequest(socket: client, kind: reqWrite,
|
||||||
|
toWrite: $ count & "\x0D\x0A")
|
||||||
|
yield sendReq
|
||||||
|
if sendReq.hasException: raise sendReq.exc
|
||||||
|
|
||||||
|
|
||||||
|
type
|
||||||
|
PsendCount2ArgObject = ref object of TObject
|
||||||
|
dummy1: PAsyncSocket
|
||||||
|
dummy2: int
|
||||||
|
|
||||||
|
proc sendCount2(client: PAsyncSocket; count: int): PsendCount2ArgObject =
|
||||||
|
new(result)
|
||||||
|
result.dummy1 = client
|
||||||
|
result.dummy2 = count
|
||||||
|
|
||||||
|
iterator sendCount2(x: PRequest): PRequest {.closure.} =
|
||||||
|
let passedInParams = PsendCount2ArgObject(x.param)
|
||||||
|
let client = passedInParams.dummy1
|
||||||
|
let count = passedInParams.dummy2
|
||||||
|
var argsToPass = sendCount3(client, count)
|
||||||
|
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount3,
|
||||||
|
param: argsToPass)
|
||||||
|
|
||||||
|
|
||||||
|
type
|
||||||
|
PsendCountArgObject = ref object of TObject
|
||||||
|
dummy1: PAsyncSocket
|
||||||
|
dummy2: int
|
||||||
|
|
||||||
|
proc sendCount(client: PAsyncSocket; count: int): PsendCountArgObject =
|
||||||
|
new(result)
|
||||||
|
result.dummy1 = client
|
||||||
|
result.dummy2 = count
|
||||||
|
|
||||||
|
iterator sendCount(x: PRequest): PRequest {.closure.} =
|
||||||
|
let passedInParams = PsendCountArgObject(x.param)
|
||||||
|
let client = passedInParams.dummy1
|
||||||
|
let count = passedInParams.dummy2
|
||||||
|
var argsToPass = sendCount2(client, count)
|
||||||
|
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount2,
|
||||||
|
param: argsToPass)
|
||||||
|
|
||||||
|
|
||||||
|
type
|
||||||
|
PprocessRequestArgObject = ref object of TObject
|
||||||
|
dummy1: PAsyncSocket
|
||||||
|
dummy2: string
|
||||||
|
dummy3: bool
|
||||||
|
|
||||||
|
proc processRequest(client: PAsyncSocket; paramTest: string;
|
||||||
|
closeSock: bool = true): PprocessRequestArgObject =
|
||||||
|
new(result)
|
||||||
|
result.dummy1 = client
|
||||||
|
result.dummy2 = paramTest
|
||||||
|
result.dummy3 = closeSock
|
||||||
|
|
||||||
|
iterator processRequest(x: PRequest): PRequest {.closure.} =
|
||||||
|
let passedInParams = PprocessRequestArgObject(x.param)
|
||||||
|
let client = passedInParams.dummy1
|
||||||
|
let paramTest = passedInParams.dummy2
|
||||||
|
let closeSock = passedInParams.dummy3
|
||||||
|
doAssert paramTest == "foobarbaz"
|
||||||
|
assert client != nil
|
||||||
|
var count = 0
|
||||||
|
while true:
|
||||||
|
var readLineReq = PRequest(socket: client, kind: reqReadLine, line: "")
|
||||||
|
yield readLineReq
|
||||||
|
if readLineReq.hasException: raise readLineReq.exc
|
||||||
|
let line = readLineReq.line
|
||||||
|
if line == "end":
|
||||||
|
break
|
||||||
|
doAssert line.startswith("Message")
|
||||||
|
doAssert line == "Message" & $ count
|
||||||
|
count.inc
|
||||||
|
var argsToPass = sendCount(client, count)
|
||||||
|
yield PRequest(socket: nil, kind: reqAwait, worker: sendCount,
|
||||||
|
param: argsToPass)
|
||||||
|
doAssert closeSock
|
||||||
|
if closeSock:
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
type
|
||||||
|
PstartServerArgObject = ref object of TObject
|
||||||
|
|
||||||
|
proc startServer(): PstartServerArgObject =
|
||||||
|
nil
|
||||||
|
|
||||||
|
iterator startServer(x: PRequest): PRequest {.closure.} =
|
||||||
|
var sock = AsyncSocket()
|
||||||
|
#sock.setReuseAddr()
|
||||||
|
sock.bindAddr(TPort(10235))
|
||||||
|
sock.listen()
|
||||||
|
while true:
|
||||||
|
var acceptReq = PRequest(socket: sock, kind: reqAccept, client: nil,
|
||||||
|
hasException: false)
|
||||||
|
yield acceptReq
|
||||||
|
if acceptReq.hasException: raise acceptReq.exc
|
||||||
|
let client = acceptReq.client
|
||||||
|
assert client != nil
|
||||||
|
var argsToPass = processRequest(client, "foobarbaz")
|
||||||
|
yield PRequest(socket: nil, kind: reqReg, worker: processRequest,
|
||||||
|
param: argsToPass)
|
||||||
|
|
||||||
|
|
||||||
|
type
|
||||||
|
PspawnClientArgObject = ref object of TObject
|
||||||
|
|
||||||
|
proc spawnClient(): PspawnClientArgObject =
|
||||||
|
nil
|
||||||
|
|
||||||
|
iterator spawnClient(x: PRequest): PRequest {.closure.} =
|
||||||
|
var client = AsyncSocket()
|
||||||
|
var connectReq = PRequest(socket: client, kind: reqConnect,
|
||||||
|
address: "localhost", port: TPort(10235))
|
||||||
|
yield connectReq
|
||||||
|
if connectReq.hasException: raise connectReq.exc
|
||||||
|
for i in 0 .. 99:
|
||||||
|
var sendReq = PRequest(socket: client, kind: reqWrite,
|
||||||
|
toWrite: "Message" & $ i & "\x0D\x0A")
|
||||||
|
yield sendReq
|
||||||
|
if sendReq.hasException: raise sendReq.exc
|
||||||
|
|
||||||
|
|
||||||
|
var sendReq1 = PRequest(socket: client, kind: reqWrite,
|
||||||
|
toWrite: "end\x0D\x0A")
|
||||||
|
yield sendReq1
|
||||||
|
if sendReq1.hasException: raise sendReq1.exc
|
||||||
|
|
||||||
|
var readLineReq = PRequest(socket: client, kind: reqReadLine, line: "")
|
||||||
|
yield readLineReq
|
||||||
|
if readLineReq.hasException: raise readLineReq.exc
|
||||||
|
|
||||||
|
let line = readLineReq.line
|
||||||
|
doAssert line == "100"
|
||||||
|
globalCount.inc(100)
|
||||||
|
client.close()
|
||||||
|
|
||||||
|
type
|
||||||
|
PspawnClientsArgObject = ref object of TObject
|
||||||
|
|
||||||
|
proc spawnClients(): PspawnClientsArgObject =
|
||||||
|
nil
|
||||||
|
|
||||||
|
iterator spawnClients(x: PRequest): PRequest {.closure.} =
|
||||||
|
for i in 0 .. 99:
|
||||||
|
var argsToPass = spawnClient()
|
||||||
|
yield PRequest(socket: nil, kind: reqReg, worker: spawnClient,
|
||||||
|
param: argsToPass)
|
||||||
|
|
||||||
|
var disp = newDispatcher(false)
|
||||||
|
disp.register(startServer, nil)
|
||||||
|
disp.register(spawnClients, nil)
|
||||||
|
while true:
|
||||||
|
discard disp.poll()
|
||||||
|
|
||||||
|
if globalCount == 100*100:
|
||||||
|
echo(globalCount)
|
||||||
|
quit(QuitSuccess)
|
||||||
Loading…
Add table
Add a link
Reference in a new issue