Merge branch 'bigbreak' of https://github.com/Araq/Nimrod into bigbreak
This commit is contained in:
commit
2011805829
2 changed files with 46 additions and 46 deletions
|
|
@ -326,8 +326,8 @@ when defined(windows) or defined(nimdoc):
|
||||||
TCompletionKey = Dword
|
TCompletionKey = Dword
|
||||||
|
|
||||||
TCompletionData* = object
|
TCompletionData* = object
|
||||||
sock*: TAsyncFD # TODO: Rename this.
|
fd*: TAsyncFD # TODO: Rename this.
|
||||||
cb*: proc (sock: TAsyncFD, bytesTransferred: Dword,
|
cb*: proc (fd: TAsyncFD, bytesTransferred: Dword,
|
||||||
errcode: OSErrorCode) {.closure,gcsafe.}
|
errcode: OSErrorCode) {.closure,gcsafe.}
|
||||||
|
|
||||||
PDispatcher* = ref object of PDispatcherBase
|
PDispatcher* = ref object of PDispatcherBase
|
||||||
|
|
@ -357,18 +357,18 @@ when defined(windows) or defined(nimdoc):
|
||||||
if gDisp.isNil: gDisp = newDispatcher()
|
if gDisp.isNil: gDisp = newDispatcher()
|
||||||
result = gDisp
|
result = gDisp
|
||||||
|
|
||||||
proc register*(sock: TAsyncFD) =
|
proc register*(fd: TAsyncFD) =
|
||||||
## Registers ``sock`` with the dispatcher.
|
## Registers ``fd`` with the dispatcher.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
if createIoCompletionPort(sock.THandle, p.ioPort,
|
if createIoCompletionPort(fd.THandle, p.ioPort,
|
||||||
cast[TCompletionKey](sock), 1) == 0:
|
cast[TCompletionKey](fd), 1) == 0:
|
||||||
raiseOSError(osLastError())
|
raiseOSError(osLastError())
|
||||||
p.handles.incl(sock)
|
p.handles.incl(fd)
|
||||||
|
|
||||||
proc verifyPresence(sock: TAsyncFD) =
|
proc verifyPresence(fd: TAsyncFD) =
|
||||||
## Ensures that socket has been registered with the dispatcher.
|
## Ensures that file descriptor has been registered with the dispatcher.
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
if sock notin p.handles:
|
if fd notin p.handles:
|
||||||
raise newException(ValueError,
|
raise newException(ValueError,
|
||||||
"Operation performed on a socket which has not been registered with" &
|
"Operation performed on a socket which has not been registered with" &
|
||||||
" the dispatcher yet.")
|
" the dispatcher yet.")
|
||||||
|
|
@ -394,16 +394,16 @@ when defined(windows) or defined(nimdoc):
|
||||||
# TODO: http://www.serverframework.com/handling-multiple-pending-socket-read-and-write-operations.html
|
# TODO: http://www.serverframework.com/handling-multiple-pending-socket-read-and-write-operations.html
|
||||||
if res:
|
if res:
|
||||||
# This is useful for ensuring the reliability of the overlapped struct.
|
# This is useful for ensuring the reliability of the overlapped struct.
|
||||||
assert customOverlapped.data.sock == lpCompletionKey.TAsyncFD
|
assert customOverlapped.data.fd == lpCompletionKey.TAsyncFD
|
||||||
|
|
||||||
customOverlapped.data.cb(customOverlapped.data.sock,
|
customOverlapped.data.cb(customOverlapped.data.fd,
|
||||||
lpNumberOfBytesTransferred, OSErrorCode(-1))
|
lpNumberOfBytesTransferred, OSErrorCode(-1))
|
||||||
GC_unref(customOverlapped)
|
GC_unref(customOverlapped)
|
||||||
else:
|
else:
|
||||||
let errCode = osLastError()
|
let errCode = osLastError()
|
||||||
if customOverlapped != nil:
|
if customOverlapped != nil:
|
||||||
assert customOverlapped.data.sock == lpCompletionKey.TAsyncFD
|
assert customOverlapped.data.fd == lpCompletionKey.TAsyncFD
|
||||||
customOverlapped.data.cb(customOverlapped.data.sock,
|
customOverlapped.data.cb(customOverlapped.data.fd,
|
||||||
lpNumberOfBytesTransferred, errCode)
|
lpNumberOfBytesTransferred, errCode)
|
||||||
GC_unref(customOverlapped)
|
GC_unref(customOverlapped)
|
||||||
else:
|
else:
|
||||||
|
|
@ -506,8 +506,8 @@ when defined(windows) or defined(nimdoc):
|
||||||
# http://blogs.msdn.com/b/oldnewthing/archive/2011/02/02/10123392.aspx
|
# http://blogs.msdn.com/b/oldnewthing/archive/2011/02/02/10123392.aspx
|
||||||
var ol = PCustomOverlapped()
|
var ol = PCustomOverlapped()
|
||||||
GC_ref(ol)
|
GC_ref(ol)
|
||||||
ol.data = TCompletionData(sock: socket, cb:
|
ol.data = TCompletionData(fd: socket, cb:
|
||||||
proc (sock: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
proc (fd: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
||||||
if not retFuture.finished:
|
if not retFuture.finished:
|
||||||
if errcode == OSErrorCode(-1):
|
if errcode == OSErrorCode(-1):
|
||||||
retFuture.complete()
|
retFuture.complete()
|
||||||
|
|
@ -570,8 +570,8 @@ when defined(windows) or defined(nimdoc):
|
||||||
var flagsio = flags.toOSFlags().Dword
|
var flagsio = flags.toOSFlags().Dword
|
||||||
var ol = PCustomOverlapped()
|
var ol = PCustomOverlapped()
|
||||||
GC_ref(ol)
|
GC_ref(ol)
|
||||||
ol.data = TCompletionData(sock: socket, cb:
|
ol.data = TCompletionData(fd: socket, cb:
|
||||||
proc (sock: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
proc (fd: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
||||||
if not retFuture.finished:
|
if not retFuture.finished:
|
||||||
if errcode == OSErrorCode(-1):
|
if errcode == OSErrorCode(-1):
|
||||||
if bytesCount == 0 and dataBuf.buf[0] == '\0':
|
if bytesCount == 0 and dataBuf.buf[0] == '\0':
|
||||||
|
|
@ -648,8 +648,8 @@ when defined(windows) or defined(nimdoc):
|
||||||
var bytesReceived, lowFlags: Dword
|
var bytesReceived, lowFlags: Dword
|
||||||
var ol = PCustomOverlapped()
|
var ol = PCustomOverlapped()
|
||||||
GC_ref(ol)
|
GC_ref(ol)
|
||||||
ol.data = TCompletionData(sock: socket, cb:
|
ol.data = TCompletionData(fd: socket, cb:
|
||||||
proc (sock: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
proc (fd: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
||||||
if not retFuture.finished:
|
if not retFuture.finished:
|
||||||
if errcode == OSErrorCode(-1):
|
if errcode == OSErrorCode(-1):
|
||||||
retFuture.complete()
|
retFuture.complete()
|
||||||
|
|
@ -737,8 +737,8 @@ when defined(windows) or defined(nimdoc):
|
||||||
|
|
||||||
var ol = PCustomOverlapped()
|
var ol = PCustomOverlapped()
|
||||||
GC_ref(ol)
|
GC_ref(ol)
|
||||||
ol.data = TCompletionData(sock: socket, cb:
|
ol.data = TCompletionData(fd: socket, cb:
|
||||||
proc (sock: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
proc (fd: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
||||||
if not retFuture.finished:
|
if not retFuture.finished:
|
||||||
if errcode == OSErrorCode(-1):
|
if errcode == OSErrorCode(-1):
|
||||||
completeAccept()
|
completeAccept()
|
||||||
|
|
@ -806,10 +806,10 @@ else:
|
||||||
|
|
||||||
type
|
type
|
||||||
TAsyncFD* = distinct cint
|
TAsyncFD* = distinct cint
|
||||||
TCallback = proc (sock: TAsyncFD): bool {.closure,gcsafe.}
|
TCallback = proc (fd: TAsyncFD): bool {.closure,gcsafe.}
|
||||||
|
|
||||||
PData* = ref object of PObject
|
PData* = ref object of PObject
|
||||||
sock: TAsyncFD
|
fd: TAsyncFD
|
||||||
readCBs: seq[TCallback]
|
readCBs: seq[TCallback]
|
||||||
writeCBs: seq[TCallback]
|
writeCBs: seq[TCallback]
|
||||||
|
|
||||||
|
|
@ -828,15 +828,15 @@ else:
|
||||||
if gDisp.isNil: gDisp = newDispatcher()
|
if gDisp.isNil: gDisp = newDispatcher()
|
||||||
result = gDisp
|
result = gDisp
|
||||||
|
|
||||||
proc update(sock: TAsyncFD, events: set[TEvent]) =
|
proc update(fd: TAsyncFD, events: set[TEvent]) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
assert sock.SocketHandle in p.selector
|
assert fd.SocketHandle in p.selector
|
||||||
discard p.selector.update(sock.SocketHandle, events)
|
discard p.selector.update(fd.SocketHandle, events)
|
||||||
|
|
||||||
proc register*(sock: TAsyncFD) =
|
proc register*(fd: TAsyncFD) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
var data = PData(sock: sock, readCBs: @[], writeCBs: @[])
|
var data = PData(fd: fd, readCBs: @[], writeCBs: @[])
|
||||||
p.selector.register(sock.SocketHandle, {}, data.PObject)
|
p.selector.register(fd.SocketHandle, {}, data.PObject)
|
||||||
|
|
||||||
proc newAsyncRawSocket*(domain: cint, typ: cint, protocol: cint): TAsyncFD =
|
proc newAsyncRawSocket*(domain: cint, typ: cint, protocol: cint): TAsyncFD =
|
||||||
result = newRawSocket(domain, typ, protocol).TAsyncFD
|
result = newRawSocket(domain, typ, protocol).TAsyncFD
|
||||||
|
|
@ -858,26 +858,26 @@ else:
|
||||||
proc unregister*(fd: TAsyncFD) =
|
proc unregister*(fd: TAsyncFD) =
|
||||||
getGlobalDispatcher().selector.unregister(fd.SocketHandle)
|
getGlobalDispatcher().selector.unregister(fd.SocketHandle)
|
||||||
|
|
||||||
proc addRead*(sock: TAsyncFD, cb: TCallback) =
|
proc addRead*(fd: TAsyncFD, cb: TCallback) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
if sock.SocketHandle notin p.selector:
|
if fd.SocketHandle notin p.selector:
|
||||||
raise newException(EInvalidValue, "File descriptor not registered.")
|
raise newException(EInvalidValue, "File descriptor not registered.")
|
||||||
p.selector[sock.SocketHandle].data.PData.readCBs.add(cb)
|
p.selector[fd.SocketHandle].data.PData.readCBs.add(cb)
|
||||||
update(sock, p.selector[sock.SocketHandle].events + {EvRead})
|
update(fd, p.selector[fd.SocketHandle].events + {EvRead})
|
||||||
|
|
||||||
proc addWrite*(sock: TAsyncFD, cb: TCallback) =
|
proc addWrite*(fd: TAsyncFD, cb: TCallback) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
if sock.SocketHandle notin p.selector:
|
if fd.SocketHandle notin p.selector:
|
||||||
raise newException(EInvalidValue, "File descriptor not registered.")
|
raise newException(EInvalidValue, "File descriptor not registered.")
|
||||||
p.selector[sock.SocketHandle].data.PData.writeCBs.add(cb)
|
p.selector[fd.SocketHandle].data.PData.writeCBs.add(cb)
|
||||||
update(sock, p.selector[sock.SocketHandle].events + {EvWrite})
|
update(fd, p.selector[fd.SocketHandle].events + {EvWrite})
|
||||||
|
|
||||||
proc poll*(timeout = 500) =
|
proc poll*(timeout = 500) =
|
||||||
let p = getGlobalDispatcher()
|
let p = getGlobalDispatcher()
|
||||||
for info in p.selector.select(timeout):
|
for info in p.selector.select(timeout):
|
||||||
let data = PData(info.key.data)
|
let data = PData(info.key.data)
|
||||||
assert data.sock == info.key.fd.TAsyncFD
|
assert data.fd == info.key.fd.TAsyncFD
|
||||||
#echo("In poll ", data.sock.cint)
|
#echo("In poll ", data.fd.cint)
|
||||||
if EvRead in info.events:
|
if EvRead in info.events:
|
||||||
# Callback may add items to ``data.readCBs`` which causes issues if
|
# Callback may add items to ``data.readCBs`` which causes issues if
|
||||||
# we are iterating over ``data.readCBs`` at the same time. We therefore
|
# we are iterating over ``data.readCBs`` at the same time. We therefore
|
||||||
|
|
@ -885,7 +885,7 @@ else:
|
||||||
let currentCBs = data.readCBs
|
let currentCBs = data.readCBs
|
||||||
data.readCBs = @[]
|
data.readCBs = @[]
|
||||||
for cb in currentCBs:
|
for cb in currentCBs:
|
||||||
if not cb(data.sock):
|
if not cb(data.fd):
|
||||||
# Callback wants to be called again.
|
# Callback wants to be called again.
|
||||||
data.readCBs.add(cb)
|
data.readCBs.add(cb)
|
||||||
|
|
||||||
|
|
@ -893,7 +893,7 @@ else:
|
||||||
let currentCBs = data.writeCBs
|
let currentCBs = data.writeCBs
|
||||||
data.writeCBs = @[]
|
data.writeCBs = @[]
|
||||||
for cb in currentCBs:
|
for cb in currentCBs:
|
||||||
if not cb(data.sock):
|
if not cb(data.fd):
|
||||||
# Callback wants to be called again.
|
# Callback wants to be called again.
|
||||||
data.writeCBs.add(cb)
|
data.writeCBs.add(cb)
|
||||||
|
|
||||||
|
|
@ -902,7 +902,7 @@ else:
|
||||||
if data.readCBs.len != 0: newEvents = {EvRead}
|
if data.readCBs.len != 0: newEvents = {EvRead}
|
||||||
if data.writeCBs.len != 0: newEvents = newEvents + {EvWrite}
|
if data.writeCBs.len != 0: newEvents = newEvents + {EvWrite}
|
||||||
if newEvents != info.key.events:
|
if newEvents != info.key.events:
|
||||||
update(data.sock, newEvents)
|
update(data.fd, newEvents)
|
||||||
else:
|
else:
|
||||||
# FD no longer a part of the selector. Likely been closed
|
# FD no longer a part of the selector. Likely been closed
|
||||||
# (e.g. socket disconnected).
|
# (e.g. socket disconnected).
|
||||||
|
|
@ -914,7 +914,7 @@ else:
|
||||||
af = AF_INET): Future[void] =
|
af = AF_INET): Future[void] =
|
||||||
var retFuture = newFuture[void]("connect")
|
var retFuture = newFuture[void]("connect")
|
||||||
|
|
||||||
proc cb(sock: TAsyncFD): bool =
|
proc cb(fd: TAsyncFD): bool =
|
||||||
# We have connected.
|
# We have connected.
|
||||||
retFuture.complete()
|
retFuture.complete()
|
||||||
return true
|
return true
|
||||||
|
|
|
||||||
|
|
@ -125,7 +125,7 @@ proc read*(f: AsyncFile, size: int): Future[string] =
|
||||||
|
|
||||||
var ol = PCustomOverlapped()
|
var ol = PCustomOverlapped()
|
||||||
GC_ref(ol)
|
GC_ref(ol)
|
||||||
ol.data = TCompletionData(sock: f.fd, cb:
|
ol.data = TCompletionData(fd: f.fd, cb:
|
||||||
proc (fd: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
proc (fd: TAsyncFD, bytesCount: Dword, errcode: OSErrorCode) =
|
||||||
if not retFuture.finished:
|
if not retFuture.finished:
|
||||||
if errcode == OSErrorCode(-1):
|
if errcode == OSErrorCode(-1):
|
||||||
|
|
@ -251,7 +251,7 @@ proc write*(f: AsyncFile, data: string): Future[void] =
|
||||||
|
|
||||||
var ol = PCustomOverlapped()
|
var ol = PCustomOverlapped()
|
||||||
GC_ref(ol)
|
GC_ref(ol)
|
||||||
ol.data = TCompletionData(sock: f.fd, cb:
|
ol.data = TCompletionData(fd: f.fd, cb:
|
||||||
proc (fd: TAsyncFD, bytesCount: DWord, errcode: OSErrorCode) =
|
proc (fd: TAsyncFD, bytesCount: DWord, errcode: OSErrorCode) =
|
||||||
if not retFuture.finished:
|
if not retFuture.finished:
|
||||||
if errcode == OSErrorCode(-1):
|
if errcode == OSErrorCode(-1):
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue