From ff38a8e69bea1525cb084aedbcfb3bd19ccd10a3 Mon Sep 17 00:00:00 2001 From: Gabriel Cruz Date: Tue, 2 Jun 2026 15:32:14 -0300 Subject: [PATCH] chore: context lifecycle --- ffi/ffi_context.nim | 417 +++++++++++++++++++++++++++++---------- ffi/ffi_context_pool.nim | 14 +- ffi/ffi_events.nim | 151 +++++++++++--- 3 files changed, 452 insertions(+), 130 deletions(-) diff --git a/ffi/ffi_context.nim b/ffi/ffi_context.nim index 7dcb01b..2858e1c 100644 --- a/ffi/ffi_context.nim +++ b/ffi/ffi_context.nim @@ -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] = diff --git a/ffi/ffi_context_pool.nim b/ffi/ffi_context_pool.nim index 547a1ff..f5b8c56 100644 --- a/ffi/ffi_context_pool.nim +++ b/ffi/ffi_context_pool.nim @@ -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 = diff --git a/ffi/ffi_events.nim b/ffi/ffi_events.nim index 09cfa0f..1e8390e 100644 --- a/ffi/ffi_events.nim +++ b/ffi/ffi_events.nim @@ -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)