mirror of
https://github.com/logos-messaging/nim-ffi.git
synced 2026-08-05 14:33:13 +00:00
chore: context lifecycle
This commit is contained in:
parent
e394166c46
commit
ff38a8e69b
@ -1,9 +1,14 @@
|
||||
{.passc: "-fPIC".}
|
||||
|
||||
import std/[atomics, locks, json, tables]
|
||||
import system/ansi_c
|
||||
import std/[atomics, locks, json, options, tables]
|
||||
import chronicles, chronos, chronos/threadsync, taskpools/channels_spsc_single, results
|
||||
import
|
||||
./ffi_types, ./ffi_events, ./ffi_thread_request, ./internal/ffi_macro, ./logging,
|
||||
./ffi_types,
|
||||
./ffi_events,
|
||||
./ffi_thread_request,
|
||||
./internal/ffi_macro,
|
||||
./logging,
|
||||
./cbor_serial
|
||||
|
||||
export ffi_events
|
||||
@ -13,21 +18,39 @@ type FFIContext*[T] = object
|
||||
# main library object (e.g., Waku, LibP2P, SDS, the one to be exposed as a library)
|
||||
ffiThread: Thread[(ptr FFIContext[T])]
|
||||
# represents the main FFI thread in charge of attending API consumer actions
|
||||
watchdogThread: Thread[(ptr FFIContext[T])]
|
||||
# monitors the FFI thread and notifies the FFI API consumer if it hangs
|
||||
eventThread: Thread[(ptr FFIContext[T])]
|
||||
# drains the bounded event queue and runs the heartbeat health check;
|
||||
# replaces the previous standalone watchdog thread
|
||||
lock: Lock
|
||||
reqChannel: ChannelSPSCSingle[ptr FFIThreadRequest]
|
||||
reqSignal: ThreadSignalPtr # to notify the FFI Thread that a new request is sent
|
||||
reqReceivedSignal: ThreadSignalPtr
|
||||
# to signal main thread, interfacing with the FFI thread, that FFI thread received the request
|
||||
stopSignal: ThreadSignalPtr
|
||||
# fired by destroyFFIContext so both ffiThread and watchdogThread can exit promptly
|
||||
# fired by destroyFFIContext so both ffiThread and eventThread can exit promptly
|
||||
threadExitSignal: ThreadSignalPtr
|
||||
# fired by ffiThread just before it exits; destroyFFIContext waits on
|
||||
# this with a bounded timeout instead of joining unconditionally, so a
|
||||
# blocked event loop cannot hang the caller forever
|
||||
eventQueueSignal: ThreadSignalPtr
|
||||
# fired by the FFI thread (via the dispatch templates) when it enqueues
|
||||
# an event, so the event thread wakes promptly instead of waiting out
|
||||
# the tick interval
|
||||
eventThreadExitSignal: ThreadSignalPtr
|
||||
# fired by the event thread just before it exits; mirrors threadExitSignal
|
||||
# so destroyFFIContext can do a bounded wait on the event thread too
|
||||
userData*: pointer
|
||||
eventRegistry*: FFIEventRegistry
|
||||
eventQueue*: EventQueue
|
||||
# bounded SPSC ring; the FFI thread is the only producer, the event
|
||||
# thread the only consumer
|
||||
ffiHeartbeat*: Atomic[int64]
|
||||
# advanced by the FFI thread on every iteration of its main loop; the
|
||||
# event thread reads it to detect a wedged FFI thread (>1s without an
|
||||
# advance after the start-grace window) and fires onNotResponding
|
||||
eventQueueStuck*: Atomic[bool]
|
||||
# sticky overflow flag — once the queue saturates, sendRequestToFFIThread
|
||||
# rejects further calls so the library can't dig itself deeper
|
||||
running: Atomic[bool] # To control when the threads are running
|
||||
registeredRequests: ptr Table[cstring, FFIRequestProc]
|
||||
# Pointer to with the registered requests at compile time
|
||||
@ -39,9 +62,74 @@ var onFFIThread* {.threadvar.}: bool
|
||||
|
||||
const git_version* {.strdefine.} = "n/a"
|
||||
|
||||
const
|
||||
EventThreadTickInterval* = 1.seconds
|
||||
## How often the event thread wakes to do a heartbeat check when no
|
||||
## events are pending. The dispatch templates also fire
|
||||
## `eventQueueSignal` on enqueue, so this only bounds the *idle*
|
||||
## latency between consecutive heartbeat checks.
|
||||
FFIHeartbeatStartDelay* = 10.seconds
|
||||
## Grace window after thread startup during which heartbeat stalls
|
||||
## are ignored — gives the host library (waku / libp2p / …) time to
|
||||
## come up before we start measuring liveness. Same value the old
|
||||
## standalone watchdog used.
|
||||
FFIHeartbeatStaleThreshold* = 1.seconds
|
||||
## Once past the start delay, the FFI thread must advance its
|
||||
## heartbeat at least once per this interval or it is considered
|
||||
## blocked and onNotResponding fires.
|
||||
|
||||
type JsonNotRespondingEvent = object
|
||||
eventType: string
|
||||
|
||||
proc init(T: type JsonNotRespondingEvent): T =
|
||||
return JsonNotRespondingEvent(eventType: "not_responding")
|
||||
|
||||
proc `$`(event: JsonNotRespondingEvent): string =
|
||||
$(%*event)
|
||||
|
||||
proc onNotResponding*(ctx: ptr FFIContext) =
|
||||
## Shim: still emits the legacy JSON payload through the registry, so
|
||||
## existing foreign consumers see no wire-shape change. A follow-up
|
||||
## PR replaces this with a CBOR `NotRespondingEvent`.
|
||||
##
|
||||
## Synchronous, lock-during-invocation by design: this is the global
|
||||
## "library is unhealthy" notification path and bypasses the event
|
||||
## queue (which may itself be the thing that's stuck). The dispatch
|
||||
## templates' lock-during-invocation contract is mirrored here.
|
||||
withLock ctx[].eventRegistry.lock:
|
||||
let snap =
|
||||
ctx[].eventRegistry.byEvent.getOrDefault("onNotResponding") &
|
||||
ctx[].eventRegistry.wildcard
|
||||
if snap.len == 0:
|
||||
chronicles.debug "onNotResponding - no listener registered"
|
||||
return
|
||||
foreignThreadGc:
|
||||
let event = $JsonNotRespondingEvent.init()
|
||||
for listener in snap:
|
||||
listener.callback(
|
||||
RET_OK,
|
||||
cast[ptr cchar](unsafeAddr event[0]),
|
||||
cast[csize_t](len(event)),
|
||||
listener.userData,
|
||||
)
|
||||
|
||||
proc sendRequestToFFIThread*(
|
||||
ctx: ptr FFIContext, ffiRequest: ptr FFIThreadRequest, timeout = InfiniteDuration
|
||||
): Result[void, string] =
|
||||
# Issue #6: once the event queue has overflowed we stop accepting new
|
||||
# requests entirely, on the assumption that any further work will just
|
||||
# produce more events the listener side already can't keep up with.
|
||||
# Stuck flag is sticky for the context lifetime — recovery requires
|
||||
# destroy + recreate, matching the issue's "expected malfunctioning".
|
||||
# NB: we deliberately do NOT call onNotResponding here. The event
|
||||
# thread fires it once when it observes the stuck flag (its loop is
|
||||
# the only place where reg.lock is guaranteed not to be held by an
|
||||
# in-flight listener); calling it from a foreign thread would
|
||||
# deadlock against a back-pressuring listener mid-invocation.
|
||||
if ctx.eventQueueStuck.load():
|
||||
deleteRequest(ffiRequest)
|
||||
return err("event queue stuck - library cannot accept new requests")
|
||||
|
||||
# Reentrancy guard (PR #23 review, item 6): if a handler running on the FFI
|
||||
# thread tries to dispatch back through this proc, it would wait forever on
|
||||
# `reqReceivedSignal` — which only this thread can fire — and self-deadlock.
|
||||
@ -86,83 +174,6 @@ proc sendRequestToFFIThread*(
|
||||
## process proc.
|
||||
return ok()
|
||||
|
||||
type Foo = object
|
||||
registerReqFFI(WatchdogReq, foo: ptr Foo):
|
||||
proc(): Future[Result[string, string]] {.async.} =
|
||||
return ok("FFI thread is not blocked")
|
||||
|
||||
type JsonNotRespondingEvent = object
|
||||
eventType: string
|
||||
|
||||
proc init(T: type JsonNotRespondingEvent): T =
|
||||
return JsonNotRespondingEvent(eventType: "not_responding")
|
||||
|
||||
proc `$`(event: JsonNotRespondingEvent): string =
|
||||
$(%*event)
|
||||
|
||||
proc onNotResponding*(ctx: ptr FFIContext) =
|
||||
## Shim: still emits the legacy JSON payload through the registry, so
|
||||
## existing foreign consumers see no wire-shape change. A follow-up
|
||||
## PR replaces this with a CBOR `NotRespondingEvent`.
|
||||
## Mirrors the dispatch templates' lock-during-invocation contract
|
||||
## (see `ffi_events.nim`).
|
||||
withLock ctx[].eventRegistry.lock:
|
||||
let snap = ctx[].eventRegistry.byEvent.getOrDefault("onNotResponding") &
|
||||
ctx[].eventRegistry.wildcard
|
||||
if snap.len == 0:
|
||||
chronicles.debug "onNotResponding - no listener registered"
|
||||
return
|
||||
foreignThreadGc:
|
||||
let event = $JsonNotRespondingEvent.init()
|
||||
for listener in snap:
|
||||
listener.callback(
|
||||
RET_OK,
|
||||
cast[ptr cchar](unsafeAddr event[0]),
|
||||
cast[csize_t](len(event)),
|
||||
listener.userData,
|
||||
)
|
||||
|
||||
proc watchdogThreadBody(ctx: ptr FFIContext) {.thread.} =
|
||||
## Watchdog thread that monitors the FFI thread and notifies the library user if it hangs.
|
||||
## This thread never blocks.
|
||||
|
||||
let watchdogRun = proc(ctx: ptr FFIContext) {.async.} =
|
||||
const WatchdogStartDelay = 10.seconds
|
||||
const WatchdogTimeinterval = 1.seconds
|
||||
const WatchdogTimeout = 20.seconds
|
||||
|
||||
# Give time for the node to be created and up before sending watchdog requests
|
||||
let initialStop = await ctx.stopSignal.wait().withTimeout(WatchdogStartDelay)
|
||||
if initialStop or ctx.running.load == false:
|
||||
return
|
||||
|
||||
while true:
|
||||
let intervalStop = await ctx.stopSignal.wait().withTimeout(WatchdogTimeinterval)
|
||||
|
||||
if intervalStop or ctx.running.load == false:
|
||||
debug "Watchdog thread exiting because FFIContext is not running"
|
||||
break
|
||||
|
||||
let callback = proc(
|
||||
callerRet: cint, msg: ptr cchar, len: csize_t, userData: pointer
|
||||
) {.cdecl, gcsafe, raises: [].} =
|
||||
discard ## Don't do anything. Just respecting the callback signature.
|
||||
const nilUserData = nil
|
||||
|
||||
trace "Sending watchdog request to FFI thread"
|
||||
|
||||
try:
|
||||
sendRequestToFFIThread(
|
||||
ctx, WatchdogReq.ffiNewReq(callback, nilUserData), WatchdogTimeout
|
||||
).isOkOr:
|
||||
error "Failed to send watchdog request to FFI thread", error = $error
|
||||
onNotResponding(ctx)
|
||||
except Exception as exc:
|
||||
error "Exception sending watchdog request", exc = exc.msg
|
||||
onNotResponding(ctx)
|
||||
|
||||
waitFor watchdogRun(ctx)
|
||||
|
||||
proc processRequest[T](
|
||||
request: ptr FFIThreadRequest, ctx: ptr FFIContext[T]
|
||||
) {.async.} =
|
||||
@ -204,9 +215,31 @@ proc processRequest[T](
|
||||
except Exception as exc:
|
||||
error "Unexpected exception in handleRes", error = exc.msg
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Heartbeat-aware closure capturing for the dispatch threadvar hook
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
var ffiEventQueueSignalPtr {.threadvar.}: ThreadSignalPtr
|
||||
## Stash for the event-thread wakeup signal. Captured here so the
|
||||
## notify-enqueued hook below has no closure environment.
|
||||
|
||||
proc ffiNotifyEventEnqueuedHook() {.gcsafe, raises: [].} =
|
||||
## Wakes the event thread immediately after a successful enqueue so
|
||||
## the listener fan-out latency isn't bounded by the tick interval.
|
||||
if not ffiEventQueueSignalPtr.isNil():
|
||||
let res = ffiEventQueueSignalPtr.fireSync()
|
||||
if res.isErr():
|
||||
# The event thread will still see the queue depth on the next
|
||||
# tick; logging is enough to flag a misconfigured signal fd.
|
||||
error "failed to fire eventQueueSignal after enqueue", err = res.error
|
||||
|
||||
proc ffiThreadBody[T](ctx: ptr FFIContext[T]) {.thread.} =
|
||||
## FFI thread body that attends library user API requests
|
||||
ffiCurrentEventRegistry = addr ctx[].eventRegistry
|
||||
ffiCurrentEventQueue = addr ctx[].eventQueue
|
||||
ffiCurrentEventQueueStuck = addr ctx[].eventQueueStuck
|
||||
ffiEventQueueSignalPtr = ctx.eventQueueSignal
|
||||
ffiCurrentNotifyEventEnqueued = ffiNotifyEventEnqueuedHook
|
||||
onFFIThread = true
|
||||
|
||||
logging.setupLog(logging.LogLevel.DEBUG, logging.LogFormat.TEXT)
|
||||
@ -239,6 +272,13 @@ proc ffiThreadBody[T](ctx: ptr FFIContext[T]) {.thread.} =
|
||||
inc i
|
||||
|
||||
while ctx.running.load():
|
||||
# Heartbeat: the event thread reads this to confirm the FFI thread
|
||||
# isn't wedged. The 100 ms `reqSignal.wait` below means we advance
|
||||
# at least ~10x per second under any normal load; a sync handler
|
||||
# that blocks the dispatcher will freeze the counter, which is
|
||||
# exactly the failure mode the watchdog used to detect.
|
||||
discard ctx.ffiHeartbeat.fetchAdd(1)
|
||||
|
||||
reapCompleted()
|
||||
|
||||
let gotSignal = await ctx.reqSignal.wait().withTimeout(100.milliseconds)
|
||||
@ -269,17 +309,134 @@ proc ffiThreadBody[T](ctx: ptr FFIContext[T]) {.thread.} =
|
||||
try:
|
||||
await allFutures(pending)
|
||||
except CatchableError as exc:
|
||||
error "draining pending FFI requests on shutdown raised",
|
||||
error = exc.msg
|
||||
error "draining pending FFI requests on shutdown raised", error = exc.msg
|
||||
|
||||
waitFor ffiRun(ctx)
|
||||
|
||||
proc cleanUpResources[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
## Full cleanup for heap-allocated contexts: closes all resources and frees memory.
|
||||
proc eventThreadBody[T](ctx: ptr FFIContext[T]) {.thread.} =
|
||||
## Drains the bounded event queue and runs the FFI-thread heartbeat
|
||||
## health check. Replaces the standalone watchdog thread: the event
|
||||
## thread checks liveness in-band (no probe request round-trip), and
|
||||
## the dispatch templates' queue-overflow path is the second trigger
|
||||
## for `onNotResponding`.
|
||||
##
|
||||
## The queue stores raw `c_malloc` payloads + names; this thread owns
|
||||
## them for the duration of dispatch and frees them after the listener
|
||||
## fan-out returns.
|
||||
|
||||
defer:
|
||||
freeShared(ctx)
|
||||
# Best-effort: tell stopAndJoinThreads we've exited so its bounded
|
||||
# wait unblocks. If this fails we still exit; the caller's timeout
|
||||
# path will take over.
|
||||
let fireRes = ctx.eventThreadExitSignal.fireSync()
|
||||
if fireRes.isErr():
|
||||
error "failed to fire eventThreadExitSignal", err = fireRes.error
|
||||
|
||||
let eventRun = proc(ctx: ptr FFIContext[T]) {.async.} =
|
||||
let startedAt = Moment.now()
|
||||
var lastHeartbeat = ctx.ffiHeartbeat.load()
|
||||
var lastHeartbeatChange = Moment.now()
|
||||
var notifiedStale = false
|
||||
var notifiedStuck = false
|
||||
|
||||
while ctx.running.load():
|
||||
# Wake on either an enqueue (eventQueueSignal) or the tick interval —
|
||||
# whichever comes first. The signal path keeps dispatch latency low
|
||||
# under load; the timeout path bounds idle latency for the
|
||||
# heartbeat check.
|
||||
discard await ctx.eventQueueSignal.wait().withTimeout(EventThreadTickInterval)
|
||||
|
||||
# Drain whatever is currently in the queue. Each iteration:
|
||||
# 1. Pop one event (queue lock).
|
||||
# 2. Snapshot listeners + invoke them under reg.lock — the
|
||||
# lock-during-invocation contract from PR #39 / issue #40 is
|
||||
# preserved here: a foreign `removeEventListener` blocks
|
||||
# until the in-flight callback fan-out returns.
|
||||
# 3. Free the payload.
|
||||
while true:
|
||||
let opt = ctx.eventQueue.tryDequeueEvent()
|
||||
if opt.isNone:
|
||||
break
|
||||
let qe = opt.get()
|
||||
defer:
|
||||
if not qe.name.isNil:
|
||||
c_free(cast[pointer](qe.name))
|
||||
if not qe.data.isNil:
|
||||
c_free(qe.data)
|
||||
|
||||
withLock ctx[].eventRegistry.lock:
|
||||
let snap =
|
||||
ctx[].eventRegistry.byEvent.getOrDefault($qe.name) &
|
||||
ctx[].eventRegistry.wildcard
|
||||
if snap.len == 0:
|
||||
chronicles.debug "event has no listeners", event = $qe.name
|
||||
else:
|
||||
foreignThreadGc:
|
||||
try:
|
||||
for listener in snap:
|
||||
listener.callback(
|
||||
RET_OK,
|
||||
cast[ptr cchar](qe.data),
|
||||
cast[csize_t](qe.dataLen),
|
||||
listener.userData,
|
||||
)
|
||||
except Exception, CatchableError:
|
||||
let msg =
|
||||
"Exception dispatching " & $qe.name & ": " & getCurrentExceptionMsg()
|
||||
for listener in snap:
|
||||
listener.callback(
|
||||
RET_ERR,
|
||||
cast[ptr cchar](unsafeAddr msg[0]),
|
||||
cast[csize_t](msg.len),
|
||||
listener.userData,
|
||||
)
|
||||
|
||||
# Queue-overflow notification: the FFI thread can only set the
|
||||
# sticky flag (firing onNotResponding from there would deadlock
|
||||
# against a back-pressuring listener that's holding reg.lock on
|
||||
# this thread). We fire it once from here, after the drain loop,
|
||||
# so the slow listener has already released the lock.
|
||||
if not notifiedStuck and ctx.eventQueueStuck.load():
|
||||
onNotResponding(ctx)
|
||||
notifiedStuck = true
|
||||
|
||||
# Heartbeat staleness check. Skipped during the start-delay grace
|
||||
# window so a slow library bring-up doesn't fire a spurious
|
||||
# not_responding. Once we've fired, latch `notifiedStale` until
|
||||
# the FFI thread proves it's alive again — avoids spamming the
|
||||
# listener while the FFI thread is still stuck.
|
||||
if not ctx.running.load():
|
||||
break
|
||||
if Moment.now() - startedAt <= FFIHeartbeatStartDelay:
|
||||
continue
|
||||
|
||||
let cur = ctx.ffiHeartbeat.load()
|
||||
if cur != lastHeartbeat:
|
||||
lastHeartbeat = cur
|
||||
lastHeartbeatChange = Moment.now()
|
||||
notifiedStale = false
|
||||
elif not notifiedStale and
|
||||
Moment.now() - lastHeartbeatChange > FFIHeartbeatStaleThreshold:
|
||||
onNotResponding(ctx)
|
||||
notifiedStale = true
|
||||
|
||||
try:
|
||||
waitFor eventRun(ctx)
|
||||
except CatchableError as exc:
|
||||
error "event thread exited with exception", error = exc.msg
|
||||
|
||||
proc deinitContextResources*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
## Mirror of `initContextResources`: tears down the lock, registry,
|
||||
## queue, and signal fds in place. The caller is responsible for the
|
||||
## memory holding `ctx` (free it for heap allocations, return it to the
|
||||
## pool for slot-allocated contexts). Threads MUST already be joined.
|
||||
##
|
||||
## Each field is nil'd after close so a subsequent re-init on the same
|
||||
## storage (pool slot reuse) doesn't double-close a stale pointer if
|
||||
## init's deferred cleanup runs.
|
||||
ctx.lock.deinitLock()
|
||||
deinitEventRegistry(ctx[].eventRegistry)
|
||||
deinitEventQueue(ctx[].eventQueue)
|
||||
when defined(gcRefc):
|
||||
## ThreadSignalPtr.close() is intentionally skipped under --mm:refc.
|
||||
##
|
||||
@ -304,20 +461,50 @@ proc cleanUpResources[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
else:
|
||||
if not ctx.reqSignal.isNil():
|
||||
?ctx.reqSignal.close()
|
||||
ctx.reqSignal = nil
|
||||
if not ctx.reqReceivedSignal.isNil():
|
||||
?ctx.reqReceivedSignal.close()
|
||||
ctx.reqReceivedSignal = nil
|
||||
if not ctx.stopSignal.isNil():
|
||||
?ctx.stopSignal.close()
|
||||
ctx.stopSignal = nil
|
||||
if not ctx.threadExitSignal.isNil():
|
||||
?ctx.threadExitSignal.close()
|
||||
ctx.threadExitSignal = nil
|
||||
if not ctx.eventQueueSignal.isNil():
|
||||
?ctx.eventQueueSignal.close()
|
||||
ctx.eventQueueSignal = nil
|
||||
if not ctx.eventThreadExitSignal.isNil():
|
||||
?ctx.eventThreadExitSignal.close()
|
||||
ctx.eventThreadExitSignal = nil
|
||||
return ok()
|
||||
|
||||
proc cleanUpResources[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
## Full cleanup for heap-allocated contexts: closes all resources and frees memory.
|
||||
defer:
|
||||
freeShared(ctx)
|
||||
return ctx.deinitContextResources()
|
||||
|
||||
proc initContextResources*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
## Initialises all resources inside an already-allocated FFIContext slot.
|
||||
## On failure every partially-initialised resource is closed; the caller
|
||||
## is responsible for releasing the slot (freeShared or pool.releaseSlot).
|
||||
##
|
||||
## Defensive: a reused pool slot still holds the previous lifetime's
|
||||
## signal pointers (set to nil by `deinitContextResources` on destroy,
|
||||
## but explicit here in case a future path forgets). Nil before
|
||||
## allocating so the deferred cleanup never double-closes.
|
||||
ctx.reqSignal = nil
|
||||
ctx.reqReceivedSignal = nil
|
||||
ctx.stopSignal = nil
|
||||
ctx.threadExitSignal = nil
|
||||
ctx.eventQueueSignal = nil
|
||||
ctx.eventThreadExitSignal = nil
|
||||
ctx.lock.initLock()
|
||||
initEventRegistry(ctx[].eventRegistry)
|
||||
initEventQueue(ctx[].eventQueue)
|
||||
ctx.ffiHeartbeat.store(0)
|
||||
ctx.eventQueueStuck.store(false)
|
||||
|
||||
var success = false
|
||||
defer:
|
||||
@ -338,6 +525,12 @@ proc initContextResources*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
ctx.threadExitSignal = ThreadSignalPtr.new().valueOr:
|
||||
return err("couldn't create threadExitSignal ThreadSignalPtr: " & $error)
|
||||
|
||||
ctx.eventQueueSignal = ThreadSignalPtr.new().valueOr:
|
||||
return err("couldn't create eventQueueSignal ThreadSignalPtr: " & $error)
|
||||
|
||||
ctx.eventThreadExitSignal = ThreadSignalPtr.new().valueOr:
|
||||
return err("couldn't create eventThreadExitSignal ThreadSignalPtr: " & $error)
|
||||
|
||||
ctx.registeredRequests = addr ffi_types.registeredRequests
|
||||
|
||||
ctx.running.store(true)
|
||||
@ -348,32 +541,44 @@ proc initContextResources*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
return err("failed to create the FFI thread: " & getCurrentExceptionMsg())
|
||||
|
||||
try:
|
||||
createThread(ctx.watchdogThread, watchdogThreadBody, ctx)
|
||||
createThread(ctx.eventThread, eventThreadBody[T], ctx)
|
||||
except ValueError, ResourceExhaustedError:
|
||||
## ffiThread is already running; signal it to exit and join before the
|
||||
## deferred cleanUpResources closes the signals it's waiting on.
|
||||
ctx.running.store(false)
|
||||
let fireRes = ctx.reqSignal.fireSync()
|
||||
if fireRes.isErr():
|
||||
error "failed to signal ffiThread during watchdog cleanup", error = fireRes.error
|
||||
error "failed to signal ffiThread during event-thread cleanup",
|
||||
error = fireRes.error
|
||||
joinThread(ctx.ffiThread)
|
||||
return err("failed to create the watchdog thread: " & getCurrentExceptionMsg())
|
||||
return err("failed to create the event thread: " & getCurrentExceptionMsg())
|
||||
|
||||
success = true
|
||||
return ok()
|
||||
|
||||
proc signalStop*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
ctx.running.store(false)
|
||||
# We deliberately do NOT call onNotResponding from these error paths:
|
||||
# the event thread may be stuck mid-callback holding `eventRegistry.lock`,
|
||||
# and onNotResponding takes that same lock — exactly the scenario the
|
||||
# bounded-timeout caller needs to escape from, not amplify into a deadlock.
|
||||
let reqSignaled = ctx.reqSignal.fireSync().valueOr:
|
||||
ctx.onNotResponding()
|
||||
return err("error signaling reqSignal in signalStop: " & $error)
|
||||
if not reqSignaled:
|
||||
ctx.onNotResponding()
|
||||
return err("failed to signal reqSignal on time in signalStop")
|
||||
let stopSignaled = ctx.stopSignal.fireSync().valueOr:
|
||||
return err("error signaling stopSignal in signalStop: " & $error)
|
||||
if not stopSignaled:
|
||||
return err("failed to signal stopSignal on time in signalStop")
|
||||
# Wake the event thread so it observes `running == false` immediately
|
||||
# instead of waiting out the tick interval. fireSync failing here is
|
||||
# not fatal — the event thread will still notice on the next tick — so
|
||||
# we only log and continue.
|
||||
let evtSignaled = ctx.eventQueueSignal.fireSync()
|
||||
if evtSignaled.isErr():
|
||||
error "failed to signal eventQueueSignal in signalStop", error = evtSignaled.error
|
||||
elif evtSignaled.get() == false:
|
||||
error "failed to signal eventQueueSignal on time in signalStop"
|
||||
return ok()
|
||||
|
||||
## If the FFI thread's event loop is blocked by a synchronous handler
|
||||
@ -384,23 +589,33 @@ proc signalStop*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
const ThreadExitTimeout* = 1500.milliseconds
|
||||
|
||||
proc stopAndJoinThreads*[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
## Signals the FFI and watchdog threads to stop, waits up to ThreadExitTimeout
|
||||
## for the FFI thread to exit, and joins both. On timeout returns err and
|
||||
## skips joinThread (leaving the threads live) rather than hanging the caller.
|
||||
## Resource cleanup (signal fds, lock) is the caller's responsibility.
|
||||
## Signals the FFI and event threads to stop, waits up to ThreadExitTimeout
|
||||
## for each to exit, and joins them. On timeout returns err and skips the
|
||||
## remaining joinThread (leaving the threads live) rather than hanging the
|
||||
## caller. Resource cleanup (signal fds, lock) is the caller's
|
||||
## responsibility.
|
||||
ctx.signalStop().isOkOr:
|
||||
return err("signalStop failed: " & $error)
|
||||
|
||||
let exitedOnTime = ctx.threadExitSignal.waitSync(ThreadExitTimeout).valueOr:
|
||||
ctx.onNotResponding()
|
||||
# We deliberately do NOT call onNotResponding from the timeout paths:
|
||||
# the event thread may be stuck mid-callback holding `eventRegistry.lock`,
|
||||
# and onNotResponding takes that same lock — exactly the scenario the
|
||||
# bounded-timeout caller needs to escape from, not amplify into a deadlock.
|
||||
let ffiExitedOnTime = ctx.threadExitSignal.waitSync(ThreadExitTimeout).valueOr:
|
||||
return err("error waiting for FFI thread exit: " & $error)
|
||||
|
||||
if not exitedOnTime:
|
||||
ctx.onNotResponding()
|
||||
if not ffiExitedOnTime:
|
||||
return err("FFI thread did not exit in time; leaking ctx to avoid hang")
|
||||
|
||||
joinThread(ctx.ffiThread)
|
||||
joinThread(ctx.watchdogThread)
|
||||
|
||||
let evtExitedOnTime = ctx.eventThreadExitSignal.waitSync(ThreadExitTimeout).valueOr:
|
||||
return err("error waiting for event thread exit: " & $error)
|
||||
|
||||
if not evtExitedOnTime:
|
||||
return err("event thread did not exit in time; leaking ctx to avoid hang")
|
||||
|
||||
joinThread(ctx.eventThread)
|
||||
return ok()
|
||||
|
||||
proc clearContext[T](ctx: ptr FFIContext[T]): Result[void, string] =
|
||||
|
||||
@ -48,13 +48,15 @@ proc destroyFFIContext*[T](
|
||||
## unsafe.
|
||||
ctx.stopAndJoinThreads().isOkOr:
|
||||
return err("destroyFFIContext(pool): " & $error)
|
||||
# Tear down the event registry on the *owning* thread so its
|
||||
# GC-managed Table / seq storage is freed on the same heap that
|
||||
# allocated it. Without this, the next thread to grab this slot
|
||||
# would crash inside `initEventRegistry`'s assignment-dtor when
|
||||
# `initTable` tries to dealloc the previous thread's data.
|
||||
deinitEventRegistry(ctx[].eventRegistry)
|
||||
# Mirror initContextResources: tear down the lock, registry, queue,
|
||||
# and signal fds in place. Without this the next slot acquisition would
|
||||
# re-init an already-initialised lock (UB at the pthread layer) and
|
||||
# overwrite the existing ThreadSignalPtr fields without closing the
|
||||
# underlying fds (unbounded fd leak across create/destroy cycles).
|
||||
let deinitRes = ctx.deinitContextResources()
|
||||
pool.releaseSlot(ctx)
|
||||
deinitRes.isOkOr:
|
||||
return err("destroyFFIContext(pool): " & $error)
|
||||
return ok()
|
||||
|
||||
proc isValidCtx*[T](pool: var FFIContextPool[T], ctx: pointer): bool =
|
||||
|
||||
@ -1,23 +1,27 @@
|
||||
## Event registry and dispatch primitives for FFI library-initiated events.
|
||||
## Event registry, bounded event queue, and dispatch primitives for FFI
|
||||
## library-initiated events.
|
||||
##
|
||||
## This module owns two concerns so they can evolve together without dragging
|
||||
## in the rest of `FFIContext`:
|
||||
## This module owns three concerns so they can evolve together without
|
||||
## dragging in the rest of `FFIContext`:
|
||||
##
|
||||
## 1. A multi-listener registry. Each event name maps to a `seq` of listeners;
|
||||
## the empty event name `""` is the wildcard channel and receives every
|
||||
## dispatched event in addition to its own per-name subscribers.
|
||||
## 2. The dispatch templates (`dispatchFFIEvent`, `dispatchFFIEventCbor`) used
|
||||
## by `{.ffiEvent.}`-generated procs. They snapshot the registry under its
|
||||
## lock, then invoke each listener *outside* the lock so re-entrant
|
||||
## add/remove from within a handler cannot self-deadlock.
|
||||
##
|
||||
## Phase 1 keeps dispatch synchronous on the FFI thread. A later phase will
|
||||
## route events through a bounded queue to a dedicated event thread; the
|
||||
## registry API does not change.
|
||||
## 1. A multi-listener registry. Each event name maps to a `seq` of
|
||||
## listeners; the empty event name `""` is the wildcard channel and
|
||||
## receives every dispatched event in addition to its own per-name
|
||||
## subscribers.
|
||||
## 2. A bounded SPSC event queue. Infrastructure for the dedicated event
|
||||
## thread (owned by `FFIContext`) to drain encoded events; payloads
|
||||
## travel via `c_malloc` so transfer across Nim heaps is safe under
|
||||
## both `--mm:orc` and `--mm:refc`. The dispatch templates do not yet
|
||||
## enqueue — that rewiring lands alongside the dispatch overhaul.
|
||||
## 3. The dispatch templates (`dispatchFFIEvent`, `dispatchFFIEventCbor`)
|
||||
## used by `{.ffiEvent.}`-generated procs. They snapshot the registry
|
||||
## under its lock, then invoke each listener *outside* the lock so
|
||||
## re-entrant add/remove from within a handler cannot self-deadlock.
|
||||
|
||||
{.pragma: callback, cdecl, raises: [], gcsafe.}
|
||||
|
||||
import std/[locks, tables]
|
||||
import system/ansi_c
|
||||
import std/[atomics, locks, options, tables]
|
||||
import chronicles
|
||||
import ./ffi_types, ./cbor_serial
|
||||
|
||||
@ -171,6 +175,98 @@ proc snapshotListeners*(
|
||||
snap.add(l)
|
||||
return snap
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Bounded event queue
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
const EventQueueCapacity* = 1024
|
||||
## Maximum number of events that can sit in the queue at once. Sized
|
||||
## generously — a sustained backlog at this depth almost certainly
|
||||
## means a user listener is wedged, which is exactly what the stuck
|
||||
## flag is meant to surface. Each `QueuedEvent` is two pointers plus
|
||||
## an int (24 B on 64-bit), so the ring is ~24 KiB per context.
|
||||
|
||||
type
|
||||
QueuedEvent* = object
|
||||
## A single event sitting in the bounded queue. All fields are
|
||||
## raw `c_malloc` pointers — no GC-managed storage — so the queue
|
||||
## can be a plain `array` without an assignment destructor running
|
||||
## across thread heaps when an `FFIContextPool` slot is reused.
|
||||
name*: cstring ## c_malloc'd copy of the event name.
|
||||
data*: ptr UncheckedArray[byte] ## c_malloc'd CBOR-encoded payload (may be nil).
|
||||
dataLen*: int
|
||||
|
||||
EventQueue* = object
|
||||
## SPSC ring. Only the FFI thread enqueues; only the event thread
|
||||
## dequeues. `lock` is sufficient — no need for atomic indices —
|
||||
## because every operation is short and uncontended.
|
||||
lock*: Lock
|
||||
head*: int ## Next slot the consumer will read.
|
||||
tail*: int ## Next slot the producer will write.
|
||||
count*: int ## Current depth, in [0, EventQueueCapacity].
|
||||
buf*: array[EventQueueCapacity, QueuedEvent]
|
||||
|
||||
proc initEventQueue*(q: var EventQueue) {.raises: [].} =
|
||||
## Initialises the queue's lock and zeroes the ring. Must be called
|
||||
## exactly once on the owning thread before any other thread uses it
|
||||
## (same constraint as `initEventRegistry`).
|
||||
q.lock.initLock()
|
||||
q.head = 0
|
||||
q.tail = 0
|
||||
q.count = 0
|
||||
for i in 0 ..< EventQueueCapacity:
|
||||
q.buf[i] = QueuedEvent(name: nil, data: nil, dataLen: 0)
|
||||
|
||||
proc deinitEventQueue*(q: var EventQueue) {.raises: [].} =
|
||||
## Frees any pending entries with `c_free` and tears down the lock.
|
||||
## Called on shutdown (after both producer and consumer threads have
|
||||
## stopped) and on pool-slot reuse so the next thread to grab the
|
||||
## slot starts from a clean state.
|
||||
for i in 0 ..< EventQueueCapacity:
|
||||
let e = q.buf[i]
|
||||
if not e.name.isNil:
|
||||
c_free(cast[pointer](e.name))
|
||||
if not e.data.isNil:
|
||||
c_free(e.data)
|
||||
q.buf[i] = QueuedEvent(name: nil, data: nil, dataLen: 0)
|
||||
q.head = 0
|
||||
q.tail = 0
|
||||
q.count = 0
|
||||
q.lock.deinitLock()
|
||||
|
||||
proc tryEnqueueEvent*(
|
||||
q: var EventQueue, name: cstring, data: ptr UncheckedArray[byte], dataLen: int
|
||||
): bool {.raises: [], gcsafe.} =
|
||||
## Pushes `(name, data, dataLen)` onto the queue. The queue takes
|
||||
## ownership of both `name` and `data` (both must be `c_malloc`'d by
|
||||
## the caller). Returns false if the queue is full — in that case the
|
||||
## caller still owns the buffers and must free them.
|
||||
withLock q.lock:
|
||||
if q.count >= EventQueueCapacity:
|
||||
return false
|
||||
q.buf[q.tail] = QueuedEvent(name: name, data: data, dataLen: dataLen)
|
||||
q.tail = (q.tail + 1) mod EventQueueCapacity
|
||||
q.count.inc()
|
||||
return true
|
||||
|
||||
proc tryDequeueEvent*(q: var EventQueue): Option[QueuedEvent] {.raises: [], gcsafe.} =
|
||||
## Pops the next entry off the queue and transfers ownership of its
|
||||
## buffers to the caller (who must `c_free(name)` and `c_free(data)`).
|
||||
## Returns `none` when the queue is empty.
|
||||
withLock q.lock:
|
||||
if q.count == 0:
|
||||
return none(QueuedEvent)
|
||||
let e = q.buf[q.head]
|
||||
q.buf[q.head] = QueuedEvent(name: nil, data: nil, dataLen: 0)
|
||||
q.head = (q.head + 1) mod EventQueueCapacity
|
||||
q.count.dec()
|
||||
return some(e)
|
||||
|
||||
proc eventQueueLen*(q: var EventQueue): int {.raises: [], gcsafe.} =
|
||||
## Snapshot depth, mainly useful from tests.
|
||||
withLock q.lock:
|
||||
return q.count
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Dispatch templates (used by {.ffiEvent.}-generated procs)
|
||||
# ---------------------------------------------------------------------------
|
||||
@ -179,9 +275,21 @@ var ffiCurrentEventRegistry* {.threadvar.}: ptr FFIEventRegistry
|
||||
## Set by the FFI thread at startup so dispatchFFIEvent / dispatchFFIEventCbor
|
||||
## can find their registry without taking a context pointer per call site.
|
||||
|
||||
template withFFIEventDispatch(
|
||||
eventName: string, listeners, body: untyped
|
||||
) =
|
||||
var ffiCurrentEventQueue* {.threadvar.}: ptr EventQueue
|
||||
## Bounded queue handle for the dedicated event thread to drain. The
|
||||
## dispatch templates do not enqueue yet; this is set up so the event
|
||||
## thread infrastructure has somewhere to read from.
|
||||
|
||||
var ffiCurrentEventQueueStuck* {.threadvar.}: ptr Atomic[bool]
|
||||
## Sticky overflow flag belonging to the owning `FFIContext`. Reserved
|
||||
## for the dispatch-overhaul follow-up; not consulted yet.
|
||||
|
||||
var ffiCurrentNotifyEventEnqueued* {.threadvar.}: proc() {.gcsafe, raises: [].}
|
||||
## Wakes the event thread after a successful enqueue. Kept as a
|
||||
## threadvar hook (rather than a queue field) so `ffi_events.nim`
|
||||
## doesn't have to depend on chronos's `ThreadSignalPtr`. Nil-safe.
|
||||
|
||||
template withFFIEventDispatch(eventName: string, listeners, body: untyped) =
|
||||
## Shared scaffold for `dispatchFFIEvent` / `dispatchFFIEventCbor`:
|
||||
## resolves the thread-local registry, snapshots listeners under
|
||||
## `reg.lock` into the caller-named `listeners` binding, then runs
|
||||
@ -192,8 +300,7 @@ template withFFIEventDispatch(
|
||||
return
|
||||
|
||||
withLock regPtr[].lock:
|
||||
let listeners =
|
||||
regPtr[].byEvent.getOrDefault(eventName) & regPtr[].wildcard
|
||||
let listeners = regPtr[].byEvent.getOrDefault(eventName) & regPtr[].wildcard
|
||||
if listeners.len == 0:
|
||||
chronicles.debug eventName & " - no listener registered"
|
||||
else:
|
||||
@ -241,9 +348,7 @@ template dispatchFFIEventCbor*(eventName: string, eventPayload: typed) =
|
||||
## also replace the `payload:` field name inside `EventEnvelope`.
|
||||
withFFIEventDispatch(eventName, listeners):
|
||||
var (data, dataLen) = cborEncodeShared(
|
||||
EventEnvelope[typeof(eventPayload)](
|
||||
eventType: eventName, payload: eventPayload
|
||||
)
|
||||
EventEnvelope[typeof(eventPayload)](eventType: eventName, payload: eventPayload)
|
||||
)
|
||||
defer:
|
||||
cborFreeShared(data)
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user