mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-20 11:40:02 +00:00
nim-ffi 0.2.0 reshapes the authoring model: `.ffi.` procs take the library value plus typed params instead of threading (ctx, callback, userData) by hand, the macro validates the context itself, and payloads ride the wire as CBOR rather than ad-hoc JSON strings. The old idiom no longer compiles against it, so the whole surface moves at once. Proc names are camelCase chosen so the generated snake_case export matches the previous C symbol exactly (wakuRelayPublish -> waku_relay_publish), keeping the ABI names stable. Node lifecycle now uses the dedicated pragmas: `.ffiCtor.` for create_node (LogosDelivery.new already returns the Future[Result[...]] the contract wants) and `.ffiDtor.` for destroy. Contexts come from the macro-emitted FFIContextPool, which caps live contexts at 32. Events become typed `.ffiEvent.` procs over `.ffi.` payload objects. The payloads carry wire-friendly scalars rather than the domain types, which are not serialisable; byte fields stay base64. This is what makes the generated bindings emit typed listeners instead of leaving consumers to register by name and parse JSON themselves. The payload fields are deliberately unexported. genBindings copies field names verbatim, so an export marker leaks into the generated Rust as `pub payload*: String` and the file does not parse -- a nim-ffi bug (it strips the marker from type names but not fields, so its single-file examples never hit it). Construction therefore lives behind the emit* procs in declare_lib, which also keeps event emission in one place and collapses each listener body to a single call. `requireInitializedNode` is gone: the macro rejects a null/invalid ctx before the handler runs, so all 14 call sites were redundant. Relay and filter push handlers are declared `raises: [Defect]`, so the emit call is wrapped explicitly -- the dispatch path no longer guards the body for us. genBindings() emits the C/C++/Rust bindings and must stay last in the compilation root; it is a no-op without -d:ffiGenBindings. The Rust output is checked in so consumers can vendor it directly. Known gaps, tracked separately: the hand-written liblogosdelivery.h / _kernel.h still declare the pre-CBOR signatures and need generating or dropping, and nimble resolves cbor_serialization 0.4.0 while the lock and nim-ffi both pin 0.3.0. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
145 lines
5.8 KiB
Nim
145 lines
5.8 KiB
Nim
import std/json
|
|
import chronos, chronicles, results, ffi
|
|
import libp2p/peerid # pull PeerId pretty string formatting
|
|
import logos_delivery/waku/common/base64
|
|
import
|
|
logos_delivery,
|
|
logos_delivery/waku/node/waku_node,
|
|
logos_delivery/api/types,
|
|
logos_delivery/waku/api/events/health_events,
|
|
logos_delivery/waku/api/events/peer_events,
|
|
logos_delivery/api/conf/logos_delivery_conf_json,
|
|
../declare_lib,
|
|
../json_event
|
|
|
|
# Add JSON serialization for RequestId
|
|
proc `%`*(id: RequestId): JsonNode =
|
|
%($id)
|
|
|
|
proc logosdeliveryCreateNode*(
|
|
configJson: string
|
|
): Future[Result[LogosDelivery, string]] {.ffiCtor.} =
|
|
let conf = parseLogosDeliveryConf(configJson).valueOr:
|
|
error "Failed to parse Logos Delivery configuration JSON",
|
|
error = error, configJson = configJson
|
|
return err("failed parseLogosDeliveryConf " & error)
|
|
|
|
return await LogosDelivery.new(conf)
|
|
|
|
proc logosdeliveryStartNode*(
|
|
lib: LogosDelivery
|
|
): Future[Result[string, string]] {.ffi.} =
|
|
# setting up outgoing event listeners
|
|
let sentListener = MessageSentEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: MessageSentEvent) {.async: (raises: []).} =
|
|
emitMessageSent($event.requestId, event.messageHash),
|
|
).valueOr:
|
|
chronicles.error "MessageSentEvent.listen failed", err = $error
|
|
return err("MessageSentEvent.listen failed: " & $error)
|
|
|
|
let errorListener = MessageErrorEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: MessageErrorEvent) {.async: (raises: []).} =
|
|
emitMessageError($event.requestId, event.messageHash, event.error),
|
|
).valueOr:
|
|
chronicles.error "MessageErrorEvent.listen failed", err = $error
|
|
return err("MessageErrorEvent.listen failed: " & $error)
|
|
|
|
let propagatedListener = MessagePropagatedEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: MessagePropagatedEvent) {.async: (raises: []).} =
|
|
emitMessagePropagated($event.requestId, event.messageHash),
|
|
).valueOr:
|
|
chronicles.error "MessagePropagatedEvent.listen failed", err = $error
|
|
return err("MessagePropagatedEvent.listen failed: " & $error)
|
|
|
|
let receivedListener = MessageReceivedEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: MessageReceivedEvent) {.async: (raises: []).} =
|
|
emitMessageReceived(event.messageHash, event.message),
|
|
).valueOr:
|
|
chronicles.error "MessageReceivedEvent.listen failed", err = $error
|
|
return err("MessageReceivedEvent.listen failed: " & $error)
|
|
|
|
let ConnectionStatusChangeListener = EventConnectionStatusChange.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: EventConnectionStatusChange) {.async: (raises: []).} =
|
|
emitConnectionStatusChange($event.connectionStatus),
|
|
).valueOr:
|
|
chronicles.error "ConnectionStatusChange.listen failed", err = $error
|
|
return err("ConnectionStatusChange.listen failed: " & $error)
|
|
|
|
let shardTopicHealthListener = EventShardTopicHealthChange.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: EventShardTopicHealthChange) {.async: (raises: []).} =
|
|
emitTopicHealthChange($event.topic, $event.health),
|
|
).valueOr:
|
|
chronicles.error "EventShardTopicHealthChange.listen failed", err = $error
|
|
return err("EventShardTopicHealthChange.listen failed: " & $error)
|
|
|
|
let peerEventListener = WakuPeerEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: WakuPeerEvent) {.async: (raises: []).} =
|
|
emitConnectionChange($event.peerId, $event.kind),
|
|
).valueOr:
|
|
chronicles.error "WakuPeerEvent.listen failed", err = $error
|
|
return err("WakuPeerEvent.listen failed: " & $error)
|
|
|
|
let channelReceivedListener = ChannelMessageReceivedEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: ChannelMessageReceivedEvent) {.async: (raises: []).} =
|
|
emitChannelMessageReceived(
|
|
string(event.channelId), $event.senderId, event.payload
|
|
),
|
|
).valueOr:
|
|
chronicles.error "ChannelMessageReceivedEvent.listen failed", err = $error
|
|
return err("ChannelMessageReceivedEvent.listen failed: " & $error)
|
|
|
|
let channelSentListener = ChannelMessageSentEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: ChannelMessageSentEvent) {.async: (raises: []).} =
|
|
emitChannelMessageSent(string(event.channelId), $event.requestId),
|
|
).valueOr:
|
|
chronicles.error "ChannelMessageSentEvent.listen failed", err = $error
|
|
return err("ChannelMessageSentEvent.listen failed: " & $error)
|
|
|
|
let channelErrorListener = ChannelMessageErrorEvent.listen(
|
|
lib.waku.brokerCtx,
|
|
proc(event: ChannelMessageErrorEvent) {.async: (raises: []).} =
|
|
emitChannelMessageError(
|
|
string(event.channelId), $event.requestId, event.error
|
|
),
|
|
).valueOr:
|
|
chronicles.error "ChannelMessageErrorEvent.listen failed", err = $error
|
|
return err("ChannelMessageErrorEvent.listen failed: " & $error)
|
|
|
|
(await lib.start()).isOkOr:
|
|
let errMsg = $error
|
|
chronicles.error "START_NODE failed", err = errMsg
|
|
return err("failed to start: " & errMsg)
|
|
return ok("")
|
|
|
|
proc logosdeliveryStopNode*(
|
|
lib: LogosDelivery
|
|
): Future[Result[string, string]] {.ffi.} =
|
|
await MessageErrorEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await MessageSentEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await MessagePropagatedEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await MessageReceivedEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await EventConnectionStatusChange.dropAllListeners(lib.waku.brokerCtx)
|
|
await EventShardTopicHealthChange.dropAllListeners(lib.waku.brokerCtx)
|
|
await WakuPeerEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await ChannelMessageReceivedEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await ChannelMessageSentEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
await ChannelMessageErrorEvent.dropAllListeners(lib.waku.brokerCtx)
|
|
|
|
(await lib.stop()).isOkOr:
|
|
let errMsg = $error
|
|
chronicles.error "STOP_NODE failed", err = errMsg
|
|
return err("failed to stop: " & errMsg)
|
|
return ok("")
|
|
|
|
proc logosdeliveryDestroy*(lib: LogosDelivery) {.ffiDtor.} =
|
|
discard
|