Merge branch 'cheatfate-fixsharray' into devel
This commit is contained in:
commit
79993d3a77
1 changed files with 51 additions and 13 deletions
|
|
@ -46,10 +46,12 @@ when hasThreadSupport:
|
||||||
SelectorImpl[T] = object
|
SelectorImpl[T] = object
|
||||||
kqFD : cint
|
kqFD : cint
|
||||||
maxFD : int
|
maxFD : int
|
||||||
changes: seq[KEvent]
|
changes: ptr SharedArray[KEvent]
|
||||||
fds: ptr SharedArray[SelectorKey[T]]
|
fds: ptr SharedArray[SelectorKey[T]]
|
||||||
count: int
|
count: int
|
||||||
changesLock: Lock
|
changesLock: Lock
|
||||||
|
changesSize: int
|
||||||
|
changesLength: int
|
||||||
sock: cint
|
sock: cint
|
||||||
Selector*[T] = ptr SelectorImpl[T]
|
Selector*[T] = ptr SelectorImpl[T]
|
||||||
else:
|
else:
|
||||||
|
|
@ -94,14 +96,17 @@ proc newSelector*[T](): Selector[T] =
|
||||||
when hasThreadSupport:
|
when hasThreadSupport:
|
||||||
result = cast[Selector[T]](allocShared0(sizeof(SelectorImpl[T])))
|
result = cast[Selector[T]](allocShared0(sizeof(SelectorImpl[T])))
|
||||||
result.fds = allocSharedArray[SelectorKey[T]](maxFD)
|
result.fds = allocSharedArray[SelectorKey[T]](maxFD)
|
||||||
|
result.changes = allocSharedArray[KEvent](MAX_KQUEUE_EVENTS)
|
||||||
|
result.changesSize = MAX_KQUEUE_EVENTS
|
||||||
initLock(result.changesLock)
|
initLock(result.changesLock)
|
||||||
else:
|
else:
|
||||||
result = Selector[T]()
|
result = Selector[T]()
|
||||||
result.fds = newSeq[SelectorKey[T]](maxFD)
|
result.fds = newSeq[SelectorKey[T]](maxFD)
|
||||||
|
result.changes = newSeqOfCap[KEvent](MAX_KQUEUE_EVENTS)
|
||||||
|
|
||||||
result.kqFD = kqFD
|
result.kqFD = kqFD
|
||||||
result.maxFD = maxFD.int
|
result.maxFD = maxFD.int
|
||||||
result.changes = newSeqOfCap[KEvent](MAX_KQUEUE_EVENTS)
|
|
||||||
# we allocating empty socket to duplicate it handle in future, to get unique
|
# we allocating empty socket to duplicate it handle in future, to get unique
|
||||||
# indexes for `fds` array. This is needed to properly identify
|
# indexes for `fds` array. This is needed to properly identify
|
||||||
# {Event.Timer, Event.Signal, Event.Process} events.
|
# {Event.Timer, Event.Signal, Event.Process} events.
|
||||||
|
|
@ -162,20 +167,44 @@ else:
|
||||||
template withChangeLock(s, body: untyped) =
|
template withChangeLock(s, body: untyped) =
|
||||||
body
|
body
|
||||||
|
|
||||||
template modifyKQueue[T](s: Selector[T], nident: uint, nfilter: cshort,
|
when hasThreadSupport:
|
||||||
|
template modifyKQueue[T](s: Selector[T], nident: uint, nfilter: cshort,
|
||||||
nflags: cushort, nfflags: cuint, ndata: int,
|
nflags: cushort, nfflags: cuint, ndata: int,
|
||||||
nudata: pointer) =
|
nudata: pointer) =
|
||||||
mixin withChangeLock
|
mixin withChangeLock
|
||||||
s.withChangeLock():
|
s.withChangeLock():
|
||||||
|
if s.changesLength == s.changesSize:
|
||||||
|
# if cache array is full, we allocating new with size * 2
|
||||||
|
let newSize = s.changesSize shl 1
|
||||||
|
let rdata = allocSharedArray[KEvent](newSize)
|
||||||
|
copyMem(rdata, s.changes, s.changesSize * sizeof(KEvent))
|
||||||
|
s.changesSize = newSize
|
||||||
|
s.changes[s.changesLength] = KEvent(ident: nident,
|
||||||
|
filter: nfilter, flags: nflags,
|
||||||
|
fflags: nfflags, data: ndata,
|
||||||
|
udata: nudata)
|
||||||
|
inc(s.changesLength)
|
||||||
|
|
||||||
|
when not declared(CACHE_EVENTS):
|
||||||
|
template flushKQueue[T](s: Selector[T]) =
|
||||||
|
mixin withChangeLock
|
||||||
|
s.withChangeLock():
|
||||||
|
if s.changesLength > 0:
|
||||||
|
if kevent(s.kqFD, addr(s.changes[0]), cint(s.changesLength),
|
||||||
|
nil, 0, nil) == -1:
|
||||||
|
raiseIOSelectorsError(osLastError())
|
||||||
|
s.changesLength = 0
|
||||||
|
else:
|
||||||
|
template modifyKQueue[T](s: Selector[T], nident: uint, nfilter: cshort,
|
||||||
|
nflags: cushort, nfflags: cuint, ndata: int,
|
||||||
|
nudata: pointer) =
|
||||||
s.changes.add(KEvent(ident: nident,
|
s.changes.add(KEvent(ident: nident,
|
||||||
filter: nfilter, flags: nflags,
|
filter: nfilter, flags: nflags,
|
||||||
fflags: nfflags, data: ndata,
|
fflags: nfflags, data: ndata,
|
||||||
udata: nudata))
|
udata: nudata))
|
||||||
|
|
||||||
when not declared(CACHE_EVENTS):
|
when not declared(CACHE_EVENTS):
|
||||||
template flushKQueue[T](s: Selector[T]) =
|
template flushKQueue[T](s: Selector[T]) =
|
||||||
mixin withChangeLock
|
|
||||||
s.withChangeLock():
|
|
||||||
let length = cint(len(s.changes))
|
let length = cint(len(s.changes))
|
||||||
if length > 0:
|
if length > 0:
|
||||||
if kevent(s.kqFD, addr(s.changes[0]), length,
|
if kevent(s.kqFD, addr(s.changes[0]), length,
|
||||||
|
|
@ -432,7 +461,16 @@ proc selectInto*[T](s: Selector[T], timeout: int,
|
||||||
when not declared(CACHE_EVENTS):
|
when not declared(CACHE_EVENTS):
|
||||||
count = kevent(s.kqFD, nil, cint(0), addr(resTable[0]), cint(maxres), ptv)
|
count = kevent(s.kqFD, nil, cint(0), addr(resTable[0]), cint(maxres), ptv)
|
||||||
else:
|
else:
|
||||||
|
when hasThreadSupport:
|
||||||
s.withChangeLock():
|
s.withChangeLock():
|
||||||
|
if s.changesLength > 0:
|
||||||
|
count = kevent(s.kqFD, addr(s.changes[0]), cint(s.changesLength),
|
||||||
|
addr(resTable[0]), cint(maxres), ptv)
|
||||||
|
s.changesLength = 0
|
||||||
|
else:
|
||||||
|
count = kevent(s.kqFD, nil, cint(0), addr(resTable[0]), cint(maxres),
|
||||||
|
ptv)
|
||||||
|
else:
|
||||||
let length = cint(len(s.changes))
|
let length = cint(len(s.changes))
|
||||||
if length > 0:
|
if length > 0:
|
||||||
count = kevent(s.kqFD, addr(s.changes[0]), length,
|
count = kevent(s.kqFD, addr(s.changes[0]), length,
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue