From b8cab5d8b9f823ff5739ad8998019416353fa2c4 Mon Sep 17 00:00:00 2001 From: Gabriel Cruz Date: Fri, 5 Jun 2026 14:06:12 -0300 Subject: [PATCH] chore: address comments --- ffi/codegen/cpp.nim | 14 ++++++----- ffi/codegen/rust.nim | 12 +++++----- ffi/ffi_context.nim | 55 +++++++++++--------------------------------- ffi/ffi_events.nim | 51 ++++++++++++++++++++-------------------- 4 files changed, 52 insertions(+), 80 deletions(-) diff --git a/ffi/codegen/cpp.nim b/ffi/codegen/cpp.nim index f335c13..3401a83 100644 --- a/ffi/codegen/cpp.nim +++ b/ffi/codegen/cpp.nim @@ -23,8 +23,10 @@ const proc genericInnerType(typeName, prefix: string): string = if typeName.startsWith(prefix) and typeName.endsWith("]"): - return typeName[prefix.len .. typeName.len - 2] - "" + let start = prefix.len + let lastIndex = typeName.len - 2 + return typeName[start .. lastIndex] + return "" proc nimTypeToCpp*(typeName: string): string = let trimmed = typeName.strip() @@ -54,7 +56,7 @@ proc stripLibPrefixCpp(procName, libName: string): string = let prefix = libName & "_" if procName.startsWith(prefix): return procName[prefix.len .. ^1] - procName + return procName proc reqStructName(p: FFIProcMeta): string = let camel = snakeToPascalCase(p.procName) @@ -132,7 +134,7 @@ proc cppBracedInit(structName: string, fieldNames: seq[string]): string = ## ## Empty `fieldNames` collapses cleanly because `join` on an empty seq ## returns "", so the result is the well-formed empty-init `Name{}`. - structName & "{" & fieldNames.join(", ") & "}" + return structName & "{" & fieldNames.join(", ") & "}" proc emitEventDispatcher( lines: var seq[string], ctxTypeName, libName: string, events: seq[FFIEventMeta] @@ -573,11 +575,11 @@ proc generateCppHeader*( lines.add("};") lines.add("") - lines.join("\n") + return lines.join("\n") proc generateCppCMakeLists*(libName: string, nimSrcRelPath: string): string = let src = nimSrcRelPath.replace("\\", "/") - CMakeListsTpl.multiReplace(("{{LIB}}", libName), ("{{SRC}}", src)) + return CMakeListsTpl.multiReplace(("{{LIB}}", libName), ("{{SRC}}", src)) proc generateCppBindings*( procs: seq[FFIProcMeta], diff --git a/ffi/codegen/rust.nim b/ffi/codegen/rust.nim index 39972b8..c8d278e 100644 --- a/ffi/codegen/rust.nim +++ b/ffi/codegen/rust.nim @@ -45,7 +45,7 @@ proc deriveLibName*(procs: seq[FFIProcMeta]): string = let parts = first.split('_') if parts.len > 0: return parts[0] - "unknown" + return "unknown" proc stripLibPrefix*(procName: string, libName: string): string = ## Strips the library prefix from a proc name. @@ -53,7 +53,7 @@ proc stripLibPrefix*(procName: string, libName: string): string = let prefix = libName & "_" if procName.startsWith(prefix): return procName[prefix.len .. ^1] - procName + return procName proc reqStructName(p: FFIProcMeta): string = ## Mirrors the Nim macro: Req or CtorReq for ctors. @@ -146,7 +146,7 @@ fn main() { [escapedSrc, libName] proc generateLibRs*(): string = - """mod ffi; + return """mod ffi; mod types; mod api; pub use types::*; @@ -221,7 +221,7 @@ proc generateFFIRs*(procs: seq[FFIProcMeta]): string = ) lines.add("}") - lines.join("\n") & "\n" + return lines.join("\n") & "\n" proc generateTypesRs*(types: seq[FFITypeMeta], procs: seq[FFIProcMeta]): string = ## Generates types.rs with Rust structs for all user-declared FFI types and @@ -267,7 +267,7 @@ proc generateTypesRs*(types: seq[FFITypeMeta], procs: seq[FFIProcMeta]): string lines.add("}") lines.add("") - lines.join("\n") + return lines.join("\n") proc generateApiRs*( procs: seq[FFIProcMeta], libName: string, events: seq[FFIEventMeta] = @[] @@ -780,7 +780,7 @@ proc generateApiRs*( lines.add("") lines.add("}") - lines.join("\n") & "\n" + return lines.join("\n") & "\n" proc generateRustCrate*( procs: seq[FFIProcMeta], diff --git a/ffi/ffi_context.nim b/ffi/ffi_context.nim index f9284d7..c630be3 100644 --- a/ffi/ffi_context.nim +++ b/ffi/ffi_context.nim @@ -27,13 +27,12 @@ type FFIContext*[T] = object # to signal main thread, interfacing with the FFI thread, that FFI thread received the request stopSignal: ThreadSignalPtr threadExitSignal: ThreadSignalPtr # bounds destroyFFIContext's wait so a blocked loop cannot hang the caller - eventQueueSignal: ThreadSignalPtr # wakes the event thread on enqueue + eventQueueSignal: ThreadSignalPtr # wakes the event thread on enqueue (used once dispatch is rewired in PR #69) eventThreadExitSignal: ThreadSignalPtr # mirrors threadExitSignal for the event thread userData*: pointer eventRegistry*: FFIEventRegistry eventQueue*: EventQueue ffiHeartbeat*: Atomic[int64] # advanced each FFI-thread loop; event thread reads for liveness - eventQueueStuck*: Atomic[bool] # sticky overflow flag; recovery is destroy+recreate running: Atomic[bool] # To control when the threads are running registeredRequests: ptr Table[cstring, FFIRequestProc] # Pointer to with the registered requests at compile time @@ -71,7 +70,7 @@ proc dispatchToListeners[T]( return foreignThreadGc: try: - notifyListenersOk(listeners, data, dataLen) + notifyListeners(listeners, RET_OK, data, dataLen) except Exception, CatchableError: notifyListenersErr( listeners, @@ -88,18 +87,12 @@ proc onNotResponding*(ctx: ptr FFIContext) = chronicles.error "onNotResponding - encode failed", err = exc.msg return let dataPtr: pointer = - if event.len > 0: unsafeAddr event[0] - else: nil + if event.len > 0: unsafeAddr event[0] else: nil ctx.dispatchToListeners(NotRespondingEventName, dataPtr, event.len) proc sendRequestToFFIThread*( ctx: ptr FFIContext, ffiRequest: ptr FFIThreadRequest, timeout = InfiniteDuration ): Result[void, string] = - # Event-queue overflow refuses further requests; the event thread fires onNotResponding to avoid deadlocking on reg.lock here. - if ctx.eventQueueStuck.load(): - deleteRequest(ffiRequest) - return err("event queue stuck - library cannot accept new requests") - # Reentrancy guard: only this thread can fire `reqReceivedSignal`, so a handler dispatching back would self-deadlock. if onFFIThread: deleteRequest(ffiRequest) @@ -182,22 +175,9 @@ proc processRequest[T]( except Exception as exc: error "Unexpected exception in handleRes", error = exc.msg -var ffiEventQueueSignalPtr {.threadvar.}: ThreadSignalPtr - ## Stashed so the hook below has no closure environment. - -proc ffiNotifyEventEnqueuedHook() {.gcsafe, raises: [].} = - if not ffiEventQueueSignalPtr.isNil(): - let res = ffiEventQueueSignalPtr.fireSync() - if res.isErr(): - 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) @@ -266,22 +246,19 @@ proc ffiThreadBody[T](ctx: ptr FFIContext[T]) {.thread.} = waitFor ffiRun(ctx) -proc freeQueuedEventPayload(qe: QueuedEvent) = - if not qe.name.isNil: - c_free(cast[pointer](qe.name)) - if not qe.data.isNil: - c_free(qe.data) - proc dispatchQueuedEvent[T](ctx: ptr FFIContext[T], qe: QueuedEvent) = ## Frees `qe`'s c_malloc buffers on exit. defer: - freeQueuedEventPayload(qe) + if not qe.name.isNil(): + c_free(cast[pointer](qe.name)) + if not qe.data.isNil(): + c_free(qe.data) ctx.dispatchToListeners($qe.name, qe.data, qe.dataLen) proc drainEventQueue[T](ctx: ptr FFIContext[T]) = while true: let opt = ctx.eventQueue.tryDequeueEvent() - if opt.isNone: + if opt.isNone(): break ctx.dispatchQueuedEvent(opt.get()) @@ -291,9 +268,9 @@ type HeartbeatMonitor = object lastValue: int64 notifiedStale: bool -proc initHeartbeatMonitor[T](ctx: ptr FFIContext[T]): HeartbeatMonitor = +proc init(T: type HeartbeatMonitor, ctx: ptr FFIContext): T = let now = Moment.now() - HeartbeatMonitor( + T( startedAt: now, lastChange: now, lastValue: ctx.ffiHeartbeat.load(), @@ -316,20 +293,15 @@ proc check[T](hb: var HeartbeatMonitor, ctx: ptr FFIContext[T]) = hb.notifiedStale = true proc eventRun[T](ctx: ptr FFIContext[T]) {.async.} = - var hb = initHeartbeatMonitor(ctx) - var notifiedStuck = false + var hb = HeartbeatMonitor.init(ctx) while ctx.running.load(): - # Wake on enqueue or tick — whichever first. + # Wake on enqueue or tick — whichever first. The enqueue path lands in PR #69; + # until then the wait always times out and we fall through to the heartbeat check. discard await ctx.eventQueueSignal.wait().withTimeout(EventThreadTickInterval) ctx.drainEventQueue() - # Fires here (after drain releases reg.lock) — from the FFI thread it'd deadlock on a back-pressuring listener. - if not notifiedStuck and ctx.eventQueueStuck.load(): - onNotResponding(ctx) - notifiedStuck = true - if not ctx.running.load(): break hb.check(ctx) @@ -418,7 +390,6 @@ proc initContextResources*[T](ctx: ptr FFIContext[T]): Result[void, string] = initEventRegistry(ctx[].eventRegistry) initEventQueue(ctx[].eventQueue) ctx.ffiHeartbeat.store(0) - ctx.eventQueueStuck.store(false) var success = false defer: diff --git a/ffi/ffi_events.nim b/ffi/ffi_events.nim index 96cc0cd..9662154 100644 --- a/ffi/ffi_events.nim +++ b/ffi/ffi_events.nim @@ -6,7 +6,7 @@ {.pragma: callback, cdecl, raises: [], gcsafe.} import system/ansi_c -import std/[atomics, locks, sequtils, options, tables] +import std/[locks, sequtils, options, tables] import chronicles import ./ffi_types, ./cbor_serial @@ -203,34 +203,34 @@ proc eventQueueLen*(q: var EventQueue): int {.raises: [], gcsafe.} = return q.count -proc notifyListenersOk*( - listeners: seq[FFIEventListener], data: pointer, dataLen: int +const emptyListenerPayload: cstring = "" + ## Non-nil zero-length buffer handed to listeners when a payload is + ## empty, so a consumer doing `std::string(data, len)` / `memcpy` never + ## receives a nil pointer (which is UB even at len 0). + +proc notifyListeners*( + listeners: seq[FFIEventListener], retCode: cint, data: pointer, dataLen: int ) = - ## Fans out a successful payload to every listener in the snapshot. + ## Fans out a payload to every listener in the snapshot. Empty payloads + ## are delivered as the non-nil `emptyListenerPayload` sentinel so a + ## consumer doing `std::string(data, len)` / `memcpy` never receives nil. + let n = max(dataLen, 0) let dataPtr = - if dataLen > 0: cast[ptr cchar](data) - else: nil + if n > 0 and not data.isNil(): cast[ptr cchar](data) + else: cast[ptr cchar](emptyListenerPayload) for listener in listeners: - listener.callback( - RET_OK, dataPtr, cast[csize_t](dataLen), listener.userData - ) + listener.callback(retCode, dataPtr, cast[csize_t](n), listener.userData) proc notifyListenersErr*(listeners: seq[FFIEventListener], msg: string) = - ## Fans out an error message to every listener in the snapshot. - let dataPtr = - if msg.len > 0: cast[ptr cchar](unsafeAddr msg[0]) - else: nil - for listener in listeners: - listener.callback( - RET_ERR, dataPtr, cast[csize_t](msg.len), listener.userData - ) + ## Error fan-out: adapts the message string to `notifyListeners`, which + ## supplies the non-nil pointer for the empty-message case. + let p = + if msg.len > 0: cast[pointer](unsafeAddr msg[0]) else: nil + notifyListeners(listeners, RET_ERR, p, msg.len) var ffiCurrentEventRegistry* {.threadvar.}: ptr FFIEventRegistry -var ffiCurrentEventQueue* {.threadvar.}: ptr EventQueue -var ffiCurrentEventQueueStuck* {.threadvar.}: ptr Atomic[bool] -var ffiCurrentNotifyEventEnqueued* {.threadvar.}: proc() {.gcsafe, raises: [].} - ## Hook (not a queue field) so this module doesn't depend on chronos's - ## ThreadSignalPtr. Nil-safe. + ## 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) = ## Resolves the thread-local registry, snapshots listeners under @@ -266,9 +266,8 @@ template dispatchFFIEvent*(eventName: string, body: untyped) = withFFIEventDispatch(eventName, listeners): let event = body let dataPtr: pointer = - if event.len > 0: unsafeAddr event[0] - else: nil - notifyListenersOk(listeners, dataPtr, event.len) + if event.len > 0: unsafeAddr event[0] else: nil + notifyListeners(listeners, RET_OK, dataPtr, event.len) template dispatchFFIEventCbor*(eventName: string, eventPayload: typed) = ## Typed CBOR variant of `dispatchFFIEvent`. Wraps `eventPayload` in an @@ -284,4 +283,4 @@ template dispatchFFIEventCbor*(eventName: string, eventPayload: typed) = ) defer: cborFreeShared(data) - notifyListenersOk(listeners, data, dataLen) + notifyListeners(listeners, RET_OK, data, dataLen)