From 7080a629da867f37b0e2e5f6012c5cffc150671d Mon Sep 17 00:00:00 2001 From: Ivan FB Date: Thu, 16 Jul 2026 00:14:42 +0200 Subject: [PATCH] feat: migrate the FFI layer to nim-ffi 0.2.0 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 --- library/README.md | 28 +- library/channels_api/channel_api.nim | 64 +- library/declare_lib.nim | 212 ++- library/kernel_api/debug_node_api.nim | 36 +- library/kernel_api/discovery_api.nim | 47 +- library/kernel_api/peer_manager_api.nim | 94 +- library/kernel_api/ping_api.nim | 12 +- library/kernel_api/protocols/filter_api.nim | 50 +- .../kernel_api/protocols/lightpush_api.nim | 14 +- library/kernel_api/protocols/relay_api.nim | 149 +- library/kernel_api/protocols/store_api.nim | 15 +- library/liblogosdelivery.h | 17 +- library/liblogosdelivery.nim | 5 + library/liblogosdelivery_kernel.h | 2 +- library/logos_delivery_api/debug_api.nim | 38 +- library/logos_delivery_api/messaging_api.nim | 54 +- library/logos_delivery_api/node_api.nim | 183 +-- library/rust_bindings/Cargo.toml | 13 + library/rust_bindings/build.rs | 47 + library/rust_bindings/src/api.rs | 1425 +++++++++++++++++ library/rust_bindings/src/ffi.rs | 64 + library/rust_bindings/src/lib.rs | 5 + library/rust_bindings/src/types.rs | 384 +++++ logos_delivery.nimble | 10 +- nimble.lock | 9 +- 25 files changed, 2436 insertions(+), 541 deletions(-) create mode 100644 library/rust_bindings/Cargo.toml create mode 100644 library/rust_bindings/build.rs create mode 100644 library/rust_bindings/src/api.rs create mode 100644 library/rust_bindings/src/ffi.rs create mode 100644 library/rust_bindings/src/lib.rs create mode 100644 library/rust_bindings/src/types.rs diff --git a/library/README.md b/library/README.md index 53f174759..95734e732 100644 --- a/library/README.md +++ b/library/README.md @@ -154,17 +154,39 @@ Note: The `payload` field should be base64-encoded. ### Events -#### `logosdelivery_set_event_callback` -Sets a callback that will be invoked whenever an event occurs (e.g., message received). +#### `logosdelivery_add_event_listener` +Registers a callback invoked whenever the named event fires (e.g., message received). +Listeners are registered per event name, so register once for each event you want +to observe. ```c -void logosdelivery_set_event_callback( +uint64_t logosdelivery_add_event_listener( void *ctx, + const char *eventName, FFICallBack callback, void *userData ); ``` +**Returns:** The listener id (> 0), or 0 if `callback` is NULL. + +Event names: `onMessageSent`, `onMessageError`, `onMessagePropagated`, +`onMessageReceived`, `onConnectionStatusChange`, `onTopicHealthChange`, +`onConnectionChange`, `onChannelMessageReceived`, `onChannelMessageSent`, +`onChannelMessageError`, `onReceivedMessage`. + +#### `logosdelivery_remove_event_listener` +Unregisters a previously added listener. + +```c +int logosdelivery_remove_event_listener( + void *ctx, + uint64_t listenerId +); +``` + +**Returns:** `RET_OK` when a listener was removed, `RET_ERR` otherwise. + **Important:** The callback should be fast, non-blocking, and thread-safe. ## Building diff --git a/library/channels_api/channel_api.nim b/library/channels_api/channel_api.nim index c9b96312f..d8daac02b 100644 --- a/library/channels_api/channel_api.nim +++ b/library/channels_api/channel_api.nim @@ -7,46 +7,34 @@ import logos_delivery/api/types, ../declare_lib -proc logosdelivery_channel_create( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - channelIdStr: cstring, - contentTopicStr: cstring, - senderIdStr: cstring, -) {.ffi.} = - requireInitializedNode(ctx, "ChannelCreate"): +proc logosdeliveryChannelCreate*( + lib: LogosDelivery, + channelIdStr: string, + contentTopicStr: string, + senderIdStr: string, +): Future[Result[string, string]] {.ffi.} = + requireChannels(lib, "ChannelCreate"): return err(errMsg) - requireChannels(ctx, "ChannelCreate"): - return err(errMsg) - - let id = ctx.myLib[].reliableChannelManager.createReliableChannel( - ChannelId($channelIdStr), - ContentTopic($contentTopicStr), - SdsParticipantID($senderIdStr), + let id = lib.reliableChannelManager.createReliableChannel( + ChannelId(channelIdStr), + ContentTopic(contentTopicStr), + SdsParticipantID(senderIdStr), ).valueOr: return err("ChannelCreate failed: " & $error) return ok(string(id)) -proc logosdelivery_channel_send( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - channelIdStr: cstring, - messageJson: cstring, -) {.ffi.} = +proc logosdeliveryChannelSend*( + lib: LogosDelivery, channelIdStr: string, messageJson: string +): Future[Result[string, string]] {.ffi.} = ## `messageJson` carries `{ "payload": , "ephemeral": }`. - requireInitializedNode(ctx, "ChannelSend"): - return err(errMsg) - - requireChannels(ctx, "ChannelSend"): + requireChannels(lib, "ChannelSend"): return err(errMsg) var jsonNode: JsonNode try: - jsonNode = parseJson($messageJson) + jsonNode = parseJson(messageJson) except Exception as e: return err("Failed to parse channel message JSON: " & e.msg) @@ -59,27 +47,19 @@ proc logosdelivery_channel_send( let ephemeral = jsonNode.getOrDefault("ephemeral").getBool(false) let requestId = ( - await ctx.myLib[].reliableChannelManager.send( - ChannelId($channelIdStr), payload, ephemeral - ) + await lib.reliableChannelManager.send(ChannelId(channelIdStr), payload, ephemeral) ).valueOr: return err("ChannelSend failed: " & $error) return ok($requestId) -proc logosdelivery_channel_close( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - channelIdStr: cstring, -) {.ffi.} = - requireInitializedNode(ctx, "ChannelClose"): +proc logosdeliveryChannelClose*( + lib: LogosDelivery, channelIdStr: string +): Future[Result[string, string]] {.ffi.} = + requireChannels(lib, "ChannelClose"): return err(errMsg) - requireChannels(ctx, "ChannelClose"): - return err(errMsg) - - (await ctx.myLib[].reliableChannelManager.closeChannel(ChannelId($channelIdStr))).isOkOr: + (await lib.reliableChannelManager.closeChannel(ChannelId(channelIdStr))).isOkOr: return err("ChannelClose failed: " & $error) return ok("") diff --git a/library/declare_lib.nim b/library/declare_lib.nim index 869c6f71e..b8f01cfc9 100644 --- a/library/declare_lib.nim +++ b/library/declare_lib.nim @@ -1,52 +1,186 @@ import ffi -import std/locks import results import logos_delivery +import + logos_delivery/waku/common/base64, + logos_delivery/waku/waku_core/message, + logos_delivery/waku/waku_core/message/digest -declareLibrary("logosdelivery") +declareLibrary("logosdelivery", LogosDelivery) -var eventCallbackLock: Lock -initLock(eventCallbackLock) - -template requireInitializedNode*( - ctx: ptr FFIContext[LogosDelivery], opName: string, onError: untyped -) = - if isNil(ctx): - let errMsg {.inject.} = opName & " failed: invalid context" - onError - elif isNil(ctx.myLib) or isNil(ctx.myLib[]): - let errMsg {.inject.} = opName & " failed: node is not initialized" - onError - -template requireMessaging*( - ctx: ptr FFIContext[LogosDelivery], opName: string, onError: untyped -) = - ## Use after `requireInitializedNode`. Fails if the node has no messaging client - ## (a kernel-only / fleet node). - ctx.myLib[].ensureMessaging().isOkOr: +template requireMessaging*(lib: LogosDelivery, opName: string, onError: untyped) = + ## Fails if the node has no messaging client (a kernel-only / fleet node). + lib.ensureMessaging().isOkOr: let errMsg {.inject.} = opName & " failed: " & error onError -template requireChannels*( - ctx: ptr FFIContext[LogosDelivery], opName: string, onError: untyped -) = - ## Use after `requireInitializedNode`. Fails if the node has no reliable channel - ## manager (a kernel-only / fleet node). - ctx.myLib[].ensureChannels().isOkOr: +template requireChannels*(lib: LogosDelivery, opName: string, onError: untyped) = + ## Fails if the node has no reliable channel manager (a kernel-only / fleet node). + lib.ensureChannels().isOkOr: let errMsg {.inject.} = opName & " failed: " & error onError -proc logosdelivery_set_event_callback( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.dynlib, exportc, cdecl.} = - if isNil(ctx): - echo "error: invalid context in logosdelivery_set_event_callback" - return +# Outgoing event payloads. These mirror the domain event types but hold only +# wire-friendly scalars, so they serialise through the library's ABI format and +# surface as typed listeners in the generated bindings. Byte fields are +# base64-encoded strings. +# +# The fields are deliberately unexported: genBindings copies field names +# verbatim, so an export marker would leak into the generated bindings as +# `pub payload*: String`. The `emit*` procs below are the construction path. - # prevent race conditions that might happen due incorrect usage. - eventCallbackLock.acquire() - defer: - eventCallbackLock.release() +type WakuMessagePayload* {.ffi.} = object + payload: string + contentTopic: string + version: uint32 + timestamp: int64 + ephemeral: bool + meta: string + proof: string - ctx[].eventCallback = cast[pointer](callback) - ctx[].eventUserData = userData +type MessageSentPayload* {.ffi.} = object + requestId: string + messageHash: string + +type MessageErrorPayload* {.ffi.} = object + requestId: string + messageHash: string + error: string + +type MessagePropagatedPayload* {.ffi.} = object + requestId: string + messageHash: string + +type MessageReceivedPayload* {.ffi.} = object + messageHash: string + message: WakuMessagePayload + +type ConnectionStatusChangePayload* {.ffi.} = object + connectionStatus: string + +type TopicHealthChangePayload* {.ffi.} = object + pubsubTopic: string + topicHealth: string + +type ConnectionChangePayload* {.ffi.} = object + peerId: string + peerEvent: string + +type ChannelMessageReceivedPayload* {.ffi.} = object + channelId: string + senderId: string + payload: string + +type ChannelMessageSentPayload* {.ffi.} = object + channelId: string + requestId: string + +type ChannelMessageErrorPayload* {.ffi.} = object + channelId: string + requestId: string + error: string + +type ReceivedMessagePayload* {.ffi.} = object + pubsubTopic: string + messageHash: string + wakuMessage: WakuMessagePayload + +proc onMessageSent(e: MessageSentPayload) {.ffiEvent: "onMessageSent".} +proc onMessageError(e: MessageErrorPayload) {.ffiEvent: "onMessageError".} +proc onMessagePropagated(e: MessagePropagatedPayload) {.ffiEvent: "onMessagePropagated".} +proc onMessageReceived(e: MessageReceivedPayload) {.ffiEvent: "onMessageReceived".} +proc onConnectionStatusChange( + e: ConnectionStatusChangePayload +) {.ffiEvent: "onConnectionStatusChange".} + +proc onTopicHealthChange(e: TopicHealthChangePayload) {.ffiEvent: "onTopicHealthChange".} +proc onConnectionChange(e: ConnectionChangePayload) {.ffiEvent: "onConnectionChange".} +proc onChannelMessageReceived( + e: ChannelMessageReceivedPayload +) {.ffiEvent: "onChannelMessageReceived".} + +proc onChannelMessageSent( + e: ChannelMessageSentPayload +) {.ffiEvent: "onChannelMessageSent".} + +proc onChannelMessageError( + e: ChannelMessageErrorPayload +) {.ffiEvent: "onChannelMessageError".} + +proc onReceivedMessage(e: ReceivedMessagePayload) {.ffiEvent: "onReceivedMessage".} + +proc toWakuMessagePayload(msg: WakuMessage): WakuMessagePayload = + return WakuMessagePayload( + payload: string(base64.encode(msg.payload)), + contentTopic: msg.contentTopic, + version: uint32(msg.version), + timestamp: int64(msg.timestamp), + ephemeral: msg.ephemeral, + meta: string(base64.encode(msg.meta)), + proof: string(base64.encode(msg.proof)), + ) + +proc emitMessageSent*(requestId, messageHash: string) = + onMessageSent(MessageSentPayload(requestId: requestId, messageHash: messageHash)) + +proc emitMessageError*(requestId, messageHash, error: string) = + onMessageError( + MessageErrorPayload( + requestId: requestId, messageHash: messageHash, error: error + ) + ) + +proc emitMessagePropagated*(requestId, messageHash: string) = + onMessagePropagated( + MessagePropagatedPayload(requestId: requestId, messageHash: messageHash) + ) + +proc emitMessageReceived*(messageHash: string, msg: WakuMessage) = + onMessageReceived( + MessageReceivedPayload( + messageHash: messageHash, message: toWakuMessagePayload(msg) + ) + ) + +proc emitConnectionStatusChange*(connectionStatus: string) = + onConnectionStatusChange( + ConnectionStatusChangePayload(connectionStatus: connectionStatus) + ) + +proc emitTopicHealthChange*(pubsubTopic, topicHealth: string) = + onTopicHealthChange( + TopicHealthChangePayload(pubsubTopic: pubsubTopic, topicHealth: topicHealth) + ) + +proc emitConnectionChange*(peerId, peerEvent: string) = + onConnectionChange(ConnectionChangePayload(peerId: peerId, peerEvent: peerEvent)) + +proc emitChannelMessageReceived*(channelId, senderId: string, payload: seq[byte]) = + onChannelMessageReceived( + ChannelMessageReceivedPayload( + channelId: channelId, + senderId: senderId, + payload: string(base64.encode(payload)), + ) + ) + +proc emitChannelMessageSent*(channelId, requestId: string) = + onChannelMessageSent( + ChannelMessageSentPayload(channelId: channelId, requestId: requestId) + ) + +proc emitChannelMessageError*(channelId, requestId, error: string) = + onChannelMessageError( + ChannelMessageErrorPayload( + channelId: channelId, requestId: requestId, error: error + ) + ) + +proc emitReceivedMessage*(pubSubTopic: string, msg: WakuMessage) = + onReceivedMessage( + ReceivedMessagePayload( + pubsubTopic: pubSubTopic, + messageHash: computeMessageHash(pubSubTopic, msg).to0xHex(), + wakuMessage: toWakuMessagePayload(msg), + ) + ) diff --git a/library/kernel_api/debug_node_api.nim b/library/kernel_api/debug_node_api.nim index 7d39935c6..39893355a 100644 --- a/library/kernel_api/debug_node_api.nim +++ b/library/kernel_api/debug_node_api.nim @@ -2,45 +2,33 @@ import std/strutils import chronos, results, ffi import logos_delivery, library/declare_lib -proc waku_version( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - let v = (await ctx.myLib[].waku.version()).valueOr: +proc wakuVersion*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + let v = (await lib.waku.version()).valueOr: return err(error) return ok(v) -proc waku_listen_addresses( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = +proc wakuListenAddresses*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = ## returns a comma-separated string of the listen addresses - let addrs = (await ctx.myLib[].waku.listenAddresses()).valueOr: + let addrs = (await lib.waku.listenAddresses()).valueOr: return err(error) return ok(addrs.join(",")) -proc waku_get_my_enr( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - let enrUri = (await ctx.myLib[].waku.myEnr()).valueOr: +proc wakuGetMyEnr*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + let enrUri = (await lib.waku.myEnr()).valueOr: return err(error) return ok(enrUri) -proc waku_get_my_peerid( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - let peerId = (await ctx.myLib[].waku.myPeerId()).valueOr: +proc wakuGetMyPeerid*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + let peerId = (await lib.waku.myPeerId()).valueOr: return err(error) return ok(peerId) -proc waku_get_metrics( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - let m = (await ctx.myLib[].waku.metrics()).valueOr: +proc wakuGetMetrics*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + let m = (await lib.waku.metrics()).valueOr: return err(error) return ok(m) -proc waku_is_online( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - let online = (await ctx.myLib[].waku.isOnline()).valueOr: +proc wakuIsOnline*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + let online = (await lib.waku.isOnline()).valueOr: return err(error) return ok($online) diff --git a/library/kernel_api/discovery_api.nim b/library/kernel_api/discovery_api.nim index 98c83e42d..8e0b01d53 100644 --- a/library/kernel_api/discovery_api.nim +++ b/library/kernel_api/discovery_api.nim @@ -2,58 +2,43 @@ import std/strutils import chronos, chronicles, results, ffi import logos_delivery, library/declare_lib -proc waku_discv5_update_bootnodes( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - bootnodes: cstring, -) {.ffi.} = +proc wakuDiscv5UpdateBootnodes*( + lib: LogosDelivery, bootnodes: string +): Future[Result[string, string]] {.ffi.} = ## Updates the bootnode list used for discovering new peers via DiscoveryV5 ## bootnodes - JSON array containing the bootnode ENRs i.e. `["enr:...", "enr:..."]` - (await ctx.myLib[].waku.discv5UpdateBootnodes($bootnodes)).isOkOr: + (await lib.waku.discv5UpdateBootnodes(bootnodes)).isOkOr: error "UPDATE_DISCV5_BOOTSTRAP_NODES failed", error = error return err(error) return ok("discovery request processed correctly") -proc waku_dns_discovery( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - enrTreeUrl: cstring, - nameDnsServer: cstring, - timeoutMs: cint, -) {.ffi.} = +proc wakuDnsDiscovery*( + lib: LogosDelivery, enrTreeUrl: string, nameDnsServer: string, timeoutMs: int32 +): Future[Result[string, string]] {.ffi.} = let nodes = ( - await ctx.myLib[].waku.dnsDiscovery($enrTreeUrl, $nameDnsServer, int(timeoutMs)) + await lib.waku.dnsDiscovery(enrTreeUrl, nameDnsServer, int(timeoutMs)) ).valueOr: error "GET_BOOTSTRAP_NODES failed", error = error return err(error) ## returns a comma-separated string of bootstrap nodes' multiaddresses return ok(nodes.join(",")) -proc waku_start_discv5( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - (await ctx.myLib[].waku.startDiscv5()).isOkOr: +proc wakuStartDiscv5*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + (await lib.waku.startDiscv5()).isOkOr: error "START_DISCV5 failed", error = error return err(error) return ok("discv5 started correctly") -proc waku_stop_discv5( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - (await ctx.myLib[].waku.stopDiscv5()).isOkOr: +proc wakuStopDiscv5*(lib: LogosDelivery): Future[Result[string, string]] {.ffi.} = + (await lib.waku.stopDiscv5()).isOkOr: error "STOP_DISCV5 failed", error = error return err(error) return ok("discv5 stopped correctly") -proc waku_peer_exchange_request( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - numPeers: uint64, -) {.ffi.} = - let numValidPeers = (await ctx.myLib[].waku.peerExchangeRequest(numPeers)).valueOr: +proc wakuPeerExchangeRequest*( + lib: LogosDelivery, numPeers: uint64 +): Future[Result[string, string]] {.ffi.} = + let numValidPeers = (await lib.waku.peerExchangeRequest(numPeers)).valueOr: error "waku_peer_exchange_request failed", error = error return err(error) return ok($numValidPeers) diff --git a/library/kernel_api/peer_manager_api.nim b/library/kernel_api/peer_manager_api.nim index e14b8b2c9..f13d7c760 100644 --- a/library/kernel_api/peer_manager_api.nim +++ b/library/kernel_api/peer_manager_api.nim @@ -6,75 +6,58 @@ type PeerInfo = object protocols: seq[string] addresses: seq[string] -proc waku_get_peerids_from_peerstore( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = +proc wakuGetPeeridsFromPeerstore*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = ## returns a comma-separated string of peerIDs - let peerIds = (await ctx.myLib[].waku.peerIdsFromPeerstore()).valueOr: + let peerIds = (await lib.waku.peerIdsFromPeerstore()).valueOr: return err(error) return ok(peerIds.join(",")) -proc waku_connect( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - peerMultiAddr: cstring, - timeoutMs: cuint, -) {.ffi.} = - let peers = ($peerMultiAddr).split(",") - (await ctx.myLib[].waku.connect(peers, uint32(timeoutMs))).isOkOr: +proc wakuConnect*( + lib: LogosDelivery, peerMultiAddr: string, timeoutMs: uint32 +): Future[Result[string, string]] {.ffi.} = + let peers = peerMultiAddr.split(",") + (await lib.waku.connect(peers, uint32(timeoutMs))).isOkOr: return err(error) return ok("") -proc waku_disconnect_peer_by_id( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - peerId: cstring, -) {.ffi.} = - (await ctx.myLib[].waku.disconnectPeerById($peerId)).isOkOr: +proc wakuDisconnectPeerById*( + lib: LogosDelivery, peerId: string +): Future[Result[string, string]] {.ffi.} = + (await lib.waku.disconnectPeerById(peerId)).isOkOr: error "DISCONNECT_PEER_BY_ID failed", error = error return err(error) return ok("") -proc waku_disconnect_all_peers( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - (await ctx.myLib[].waku.disconnectAllPeers()).isOkOr: +proc wakuDisconnectAllPeers*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = + (await lib.waku.disconnectAllPeers()).isOkOr: return err(error) return ok("") -proc waku_dial_peer( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - peerMultiAddr: cstring, - protocol: cstring, - timeoutMs: cuint, -) {.ffi.} = - (await ctx.myLib[].waku.dialPeer($peerMultiAddr, $protocol, int(timeoutMs))).isOkOr: +proc wakuDialPeer*( + lib: LogosDelivery, peerMultiAddr: string, protocol: string, timeoutMs: uint32 +): Future[Result[string, string]] {.ffi.} = + (await lib.waku.dialPeer(peerMultiAddr, protocol, int(timeoutMs))).isOkOr: error "DIAL_PEER failed", error = error return err(error) return ok("") -proc waku_dial_peer_by_id( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - peerId: cstring, - protocol: cstring, - timeoutMs: cuint, -) {.ffi.} = - (await ctx.myLib[].waku.dialPeerById($peerId, $protocol, int(timeoutMs))).isOkOr: +proc wakuDialPeerById*( + lib: LogosDelivery, peerId: string, protocol: string, timeoutMs: uint32 +): Future[Result[string, string]] {.ffi.} = + (await lib.waku.dialPeerById(peerId, protocol, int(timeoutMs))).isOkOr: error "DIAL_PEER_BY_ID failed", error = error return err(error) return ok("") -proc waku_get_connected_peers_info( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = +proc wakuGetConnectedPeersInfo*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = ## returns a JSON string mapping peerIDs to objects with protocols and addresses - let peers = (await ctx.myLib[].waku.connectedPeersInfo()).valueOr: + let peers = (await lib.waku.connectedPeersInfo()).valueOr: return err(error) var peersMap = initTable[string, PeerInfo]() @@ -84,21 +67,18 @@ proc waku_get_connected_peers_info( return ok($(%*peersMap)) -proc waku_get_connected_peers( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = +proc wakuGetConnectedPeers*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = ## returns a comma-separated string of peerIDs - let peerIds = (await ctx.myLib[].waku.connectedPeers()).valueOr: + let peerIds = (await lib.waku.connectedPeers()).valueOr: return err(error) return ok(peerIds.join(",")) -proc waku_get_peerids_by_protocol( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - protocol: cstring, -) {.ffi.} = +proc wakuGetPeeridsByProtocol*( + lib: LogosDelivery, protocol: string +): Future[Result[string, string]] {.ffi.} = ## returns a comma-separated string of peerIDs that mount the given protocol - let peerIds = (await ctx.myLib[].waku.peerIdsByProtocol($protocol)).valueOr: + let peerIds = (await lib.waku.peerIdsByProtocol(protocol)).valueOr: return err(error) return ok(peerIds.join(",")) diff --git a/library/kernel_api/ping_api.nim b/library/kernel_api/ping_api.nim index 6570fffd5..c7f251523 100644 --- a/library/kernel_api/ping_api.nim +++ b/library/kernel_api/ping_api.nim @@ -1,13 +1,9 @@ import chronos, results, ffi import logos_delivery, library/declare_lib -proc waku_ping_peer( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - peerAddr: cstring, - timeoutMs: cuint, -) {.ffi.} = - let rttNanos = (await ctx.myLib[].waku.pingPeer($peerAddr, int(timeoutMs))).valueOr: +proc wakuPingPeer*( + lib: LogosDelivery, peerAddr: string, timeoutMs: uint32 +): Future[Result[string, string]] {.ffi.} = + let rttNanos = (await lib.waku.pingPeer(peerAddr, int(timeoutMs))).valueOr: return err(error) return ok($rttNanos) diff --git a/library/kernel_api/protocols/filter_api.nim b/library/kernel_api/protocols/filter_api.nim index cd613c1e0..88f4cb1bb 100644 --- a/library/kernel_api/protocols/filter_api.nim +++ b/library/kernel_api/protocols/filter_api.nim @@ -9,49 +9,45 @@ import library/events/json_message_event, library/declare_lib -proc waku_filter_subscribe( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, - contentTopics: cstring, -) {.ffi.} = - proc onReceivedMessage(ctx: ptr FFIContext[LogosDelivery]): FilterPushHandler = +proc wakuFilterSubscribe*( + lib: LogosDelivery, pubSubTopic: string, contentTopics: string +): Future[Result[string, string]] {.ffi.} = + proc receivedMessageHandler(): FilterPushHandler = return proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} = - callEventCallback(ctx, "onReceivedMessage"): - $JsonMessageEvent.new(pubsubTopic, msg) + # This handler is `raises: [Defect]`, so the payload build has to be guarded. + try: + emitReceivedMessage(pubsubTopic, msg) + except Exception, CatchableError: + error "onReceivedMessage failed to emit event", + error = getCurrentExceptionMsg() ( - await ctx.myLib[].waku.filterSubscribe( - PubsubTopic($pubSubTopic), - ($contentTopics).split(",").mapIt(ContentTopic(it)), - FilterPushHandler(onReceivedMessage(ctx)), + await lib.waku.filterSubscribe( + PubsubTopic(pubSubTopic), + contentTopics.split(",").mapIt(ContentTopic(it)), + FilterPushHandler(receivedMessageHandler()), ) ).isOkOr: error "fail filter subscribe", error = error return err(error) return ok("") -proc waku_filter_unsubscribe( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, - contentTopics: cstring, -) {.ffi.} = +proc wakuFilterUnsubscribe*( + lib: LogosDelivery, pubSubTopic: string, contentTopics: string +): Future[Result[string, string]] {.ffi.} = ( - await ctx.myLib[].waku.filterUnsubscribe( - PubsubTopic($pubSubTopic), ($contentTopics).split(",").mapIt(ContentTopic(it)) + await lib.waku.filterUnsubscribe( + PubsubTopic(pubSubTopic), contentTopics.split(",").mapIt(ContentTopic(it)) ) ).isOkOr: error "fail filter unsubscribe", error = error return err(error) return ok("") -proc waku_filter_unsubscribe_all( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - (await ctx.myLib[].waku.filterUnsubscribeAll()).isOkOr: +proc wakuFilterUnsubscribeAll*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = + (await lib.waku.filterUnsubscribeAll()).isOkOr: error "fail filter unsubscribe all", error = error return err(error) return ok("") diff --git a/library/kernel_api/protocols/lightpush_api.nim b/library/kernel_api/protocols/lightpush_api.nim index c7bc8169d..a2f867466 100644 --- a/library/kernel_api/protocols/lightpush_api.nim +++ b/library/kernel_api/protocols/lightpush_api.nim @@ -7,16 +7,12 @@ import library/events/json_message_event, library/declare_lib -proc waku_lightpush_publish( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, - jsonWakuMessage: cstring, -) {.ffi.} = +proc wakuLightpushPublish*( + lib: LogosDelivery, pubSubTopic: string, jsonWakuMessage: string +): Future[Result[string, string]] {.ffi.} = var jsonMessage: JsonMessage try: - let jsonContent = parseJson($jsonWakuMessage) + let jsonContent = parseJson(jsonWakuMessage) jsonMessage = JsonMessage.fromJsonNode(jsonContent).valueOr: raise newException(JsonParsingError, $error) except JsonParsingError as exc: @@ -26,7 +22,7 @@ proc waku_lightpush_publish( return err("Problem building the WakuMessage: " & $error) let msgHashHex = ( - await ctx.myLib[].waku.lightpushPublish(PubsubTopic($pubSubTopic), msg) + await lib.waku.lightpushPublish(PubsubTopic(pubSubTopic), msg) ).valueOr: error "PUBLISH failed", error = error return err(error) diff --git a/library/kernel_api/protocols/relay_api.nim b/library/kernel_api/protocols/relay_api.nim index 7a1fe446f..be67e790a 100644 --- a/library/kernel_api/protocols/relay_api.nim +++ b/library/kernel_api/protocols/relay_api.nim @@ -8,111 +8,87 @@ import library/events/json_message_event, library/declare_lib -proc waku_relay_get_peers_in_mesh( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, -) {.ffi.} = - let peers = (await ctx.myLib[].waku.relayPeersInMesh(PubsubTopic($pubSubTopic))).valueOr: +proc wakuRelayGetPeersInMesh*( + lib: LogosDelivery, pubSubTopic: string +): Future[Result[string, string]] {.ffi.} = + let peers = (await lib.waku.relayPeersInMesh(PubsubTopic(pubSubTopic))).valueOr: error "LIST_MESH_PEERS failed", error = error return err(error) ## returns a comma-separated string of peerIDs return ok(peers.join(",")) -proc waku_relay_get_num_peers_in_mesh( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, -) {.ffi.} = - let n = (await ctx.myLib[].waku.relayNumPeersInMesh(PubsubTopic($pubSubTopic))).valueOr: +proc wakuRelayGetNumPeersInMesh*( + lib: LogosDelivery, pubSubTopic: string +): Future[Result[string, string]] {.ffi.} = + let n = (await lib.waku.relayNumPeersInMesh(PubsubTopic(pubSubTopic))).valueOr: error "NUM_MESH_PEERS failed", error = error return err(error) return ok($n) -proc waku_relay_get_connected_peers( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, -) {.ffi.} = +proc wakuRelayGetConnectedPeers*( + lib: LogosDelivery, pubSubTopic: string +): Future[Result[string, string]] {.ffi.} = ## Returns the list of all connected peers to an specific pubsub topic - let peers = (await ctx.myLib[].waku.relayConnectedPeers(PubsubTopic($pubSubTopic))).valueOr: + let peers = (await lib.waku.relayConnectedPeers(PubsubTopic(pubSubTopic))).valueOr: error "LIST_CONNECTED_PEERS failed", error = error return err(error) return ok(peers.join(",")) -proc waku_relay_get_num_connected_peers( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, -) {.ffi.} = - let n = (await ctx.myLib[].waku.relayNumConnectedPeers(PubsubTopic($pubSubTopic))).valueOr: +proc wakuRelayGetNumConnectedPeers*( + lib: LogosDelivery, pubSubTopic: string +): Future[Result[string, string]] {.ffi.} = + let n = (await lib.waku.relayNumConnectedPeers(PubsubTopic(pubSubTopic))).valueOr: error "NUM_CONNECTED_PEERS failed", error = error return err(error) return ok($n) -proc waku_relay_add_protected_shard( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - clusterId: cint, - shardId: cint, - publicKey: cstring, -) {.ffi.} = +proc wakuRelayAddProtectedShard*( + lib: LogosDelivery, clusterId: int32, shardId: int32, publicKey: string +): Future[Result[string, string]] {.ffi.} = ## Protects a shard with a public key ( - await ctx.myLib[].waku.relayAddProtectedShard( - uint16(clusterId), uint16(shardId), $publicKey + await lib.waku.relayAddProtectedShard( + uint16(clusterId), uint16(shardId), publicKey ) ).isOkOr: return err(error) return ok("") -proc waku_relay_subscribe( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, -) {.ffi.} = - proc onReceivedMessage(ctx: ptr FFIContext[LogosDelivery]): WakuRelayHandler = +proc wakuRelaySubscribe*( + lib: LogosDelivery, pubSubTopic: string +): Future[Result[string, string]] {.ffi.} = + proc receivedMessageHandler(): WakuRelayHandler = return proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} = - callEventCallback(ctx, "onReceivedMessage"): - $JsonMessageEvent.new(pubsubTopic, msg) + # This handler is `raises: [Defect]`, so the payload build has to be guarded. + try: + emitReceivedMessage(pubsubTopic, msg) + except Exception, CatchableError: + error "onReceivedMessage failed to emit event", + error = getCurrentExceptionMsg() ( - await ctx.myLib[].waku.relaySubscribe( - PubsubTopic($pubSubTopic), WakuRelayHandler(onReceivedMessage(ctx)) + await lib.waku.relaySubscribe( + PubsubTopic(pubSubTopic), WakuRelayHandler(receivedMessageHandler()) ) ).isOkOr: error "SUBSCRIBE failed", error = error return err(error) return ok("") -proc waku_relay_unsubscribe( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, -) {.ffi.} = - (await ctx.myLib[].waku.relayUnsubscribe(PubsubTopic($pubSubTopic))).isOkOr: +proc wakuRelayUnsubscribe*( + lib: LogosDelivery, pubSubTopic: string +): Future[Result[string, string]] {.ffi.} = + (await lib.waku.relayUnsubscribe(PubsubTopic(pubSubTopic))).isOkOr: error "UNSUBSCRIBE failed", error = error return err(error) return ok("") -proc waku_relay_publish( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - pubSubTopic: cstring, - jsonWakuMessage: cstring, - timeoutMs: cuint, -) {.ffi.} = +proc wakuRelayPublish*( + lib: LogosDelivery, pubSubTopic: string, jsonWakuMessage: string, timeoutMs: uint32 +): Future[Result[string, string]] {.ffi.} = var jsonMessage: JsonMessage try: - let jsonContent = parseJson($jsonWakuMessage) + let jsonContent = parseJson(jsonWakuMessage) jsonMessage = JsonMessage.fromJsonNode(jsonContent).valueOr: raise newException(JsonParsingError, $error) except JsonParsingError as exc: @@ -122,44 +98,37 @@ proc waku_relay_publish( return err("Problem building the WakuMessage: " & $error) let msgHash = ( - await ctx.myLib[].waku.relayPublish( - PubsubTopic($pubSubTopic), msg, uint32(timeoutMs) - ) + await lib.waku.relayPublish(PubsubTopic(pubSubTopic), msg, uint32(timeoutMs)) ).valueOr: error "PUBLISH failed", error = error return err(error) return ok(msgHash) -proc waku_default_pubsub_topic( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - let topic = (await ctx.myLib[].waku.defaultPubsubTopic()).valueOr: +proc wakuDefaultPubsubTopic*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = + let topic = (await lib.waku.defaultPubsubTopic()).valueOr: return err(error) return ok(string(topic)) -proc waku_content_topic( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - appName: cstring, - appVersion: cuint, - contentTopicName: cstring, - encoding: cstring, -) {.ffi.} = +proc wakuContentTopic*( + lib: LogosDelivery, + appName: string, + appVersion: uint32, + contentTopicName: string, + encoding: string, +): Future[Result[string, string]] {.ffi.} = let topic = ( - await ctx.myLib[].waku.buildContentTopic( - $appName, uint32(appVersion), $contentTopicName, $encoding + await lib.waku.buildContentTopic( + appName, uint32(appVersion), contentTopicName, encoding ) ).valueOr: return err(error) return ok(string(topic)) -proc waku_pubsub_topic( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - topicName: cstring, -) {.ffi.} = - let topic = (await ctx.myLib[].waku.buildPubsubTopic($topicName)).valueOr: +proc wakuPubsubTopic*( + lib: LogosDelivery, topicName: string +): Future[Result[string, string]] {.ffi.} = + let topic = (await lib.waku.buildPubsubTopic(topicName)).valueOr: return err(error) return ok(string(topic)) diff --git a/library/kernel_api/protocols/store_api.nim b/library/kernel_api/protocols/store_api.nim index 75d43fee1..faab152aa 100644 --- a/library/kernel_api/protocols/store_api.nim +++ b/library/kernel_api/protocols/store_api.nim @@ -64,16 +64,11 @@ func fromJsonNode(jsonContent: JsonNode): Result[StoreQueryRequest, string] = ) ) -proc waku_store_query( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - jsonQuery: cstring, - peerAddr: cstring, - timeoutMs: cint, -) {.ffi.} = +proc wakuStoreQuery*( + lib: LogosDelivery, jsonQuery: string, peerAddr: string, timeoutMs: int32 +): Future[Result[string, string]] {.ffi.} = let jsonContentRes = catch: - parseJson($jsonQuery) + parseJson(jsonQuery) if jsonContentRes.isErr(): return err("StoreRequest failed parsing store request: " & jsonContentRes.error.msg) @@ -81,7 +76,7 @@ proc waku_store_query( let storeQueryRequest = ?fromJsonNode(jsonContentRes.get()) let queryResponse = ( - await ctx.myLib[].waku.storeQuery(storeQueryRequest, $peerAddr, int(timeoutMs)) + await lib.waku.storeQuery(storeQueryRequest, peerAddr, int(timeoutMs)) ).valueOr: return err("StoreRequest failed store query: " & error) diff --git a/library/liblogosdelivery.h b/library/liblogosdelivery.h index 21d8890d2..f24d9662e 100644 --- a/library/liblogosdelivery.h +++ b/library/liblogosdelivery.h @@ -102,16 +102,25 @@ extern "C" void *userData, const char *channelId); - // Channel lifecycle events are delivered through the event callback set via - // logosdelivery_set_event_callback: "onChannelMessageReceived" (payload + // Channel lifecycle events are delivered to listeners registered via + // logosdelivery_add_event_listener: "onChannelMessageReceived" (payload // base64-encoded), "onChannelMessageSent", "onChannelMessageError". - // Sets a callback that will be invoked whenever an event occurs. + // Registers a callback invoked whenever the named event fires. Listeners are + // per event name; register once per event you care about. Returns the listener + // id (> 0) to pass to logosdelivery_remove_event_listener, or 0 if callback is + // NULL. // It is crucial that the passed callback is fast, non-blocking and potentially thread-safe. - void logosdelivery_set_event_callback(void *ctx, + uint64_t logosdelivery_add_event_listener(void *ctx, + const char *eventName, FFICallBack callback, void *userData); + // Unregisters the listener with the given id. Returns RET_OK when a listener + // was removed, RET_ERR otherwise. + int logosdelivery_remove_event_listener(void *ctx, + uint64_t listenerId); + // Retrieves the list of available node info IDs. int logosdelivery_get_available_node_info_ids(void *ctx, FFICallBack callback, diff --git a/library/liblogosdelivery.nim b/library/liblogosdelivery.nim index d3753da63..836e505b5 100644 --- a/library/liblogosdelivery.nim +++ b/library/liblogosdelivery.nim @@ -33,3 +33,8 @@ include # logosdelivery_* surface in ./logos_delivery_api/node_api. The former # waku_new / waku_start / waku_stop / waku_destroy entry points were removed to # avoid maintaining two parallel node-lifecycle APIs. + +# Must stay the last FFI call in the compilation root: it emits the C/C++/Rust +# bindings from the registries the annotations above populate. No-op unless +# built with -d:ffiGenBindings. +genBindings() diff --git a/library/liblogosdelivery_kernel.h b/library/liblogosdelivery_kernel.h index 58c67b712..3e3403253 100644 --- a/library/liblogosdelivery_kernel.h +++ b/library/liblogosdelivery_kernel.h @@ -36,7 +36,7 @@ extern "C" FFICallBack callback, void *userData); - // NOTE: event callbacks are registered via logosdelivery_set_event_callback + // NOTE: event listeners are registered via logosdelivery_add_event_listener // (declared above) which the waku_* API shares. int waku_content_topic(void *ctx, diff --git a/library/logos_delivery_api/debug_api.nim b/library/logos_delivery_api/debug_api.nim index 98d48c97c..8c0480bbe 100644 --- a/library/logos_delivery_api/debug_api.nim +++ b/library/logos_delivery_api/debug_api.nim @@ -2,41 +2,29 @@ import std/[json, strutils] import logos_delivery/waku/factory/waku_state_info import tools/confutils/[cli_args, config_option_meta] -proc logosdelivery_get_available_node_info_ids( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = +proc logosdeliveryGetAvailableNodeInfoIds*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = ## Returns the list of all available node info item ids that ## can be queried with `get_node_info_item`. - requireInitializedNode(ctx, "GetNodeInfoIds"): - return err(errMsg) + return ok($lib.waku.stateInfo.getAllPossibleInfoItemIds()) - return ok($ctx.myLib[].waku.stateInfo.getAllPossibleInfoItemIds()) - -proc logosdelivery_get_node_info( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - nodeInfoId: cstring, -) {.ffi.} = +proc logosdeliveryGetNodeInfo*( + lib: LogosDelivery, nodeInfoId: string +): Future[Result[string, string]] {.ffi.} = ## Returns the content of the node info item with the given id if it exists. - requireInitializedNode(ctx, "GetNodeInfoItem"): - return err(errMsg) - let infoItemIdEnum = try: - parseEnum[NodeInfoId]($nodeInfoId) + parseEnum[NodeInfoId](nodeInfoId) except ValueError: - return err("Invalid node info id: " & $nodeInfoId) + return err("Invalid node info id: " & nodeInfoId) - return ok(ctx.myLib[].waku.stateInfo.getNodeInfoItem(infoItemIdEnum)) + return ok(lib.waku.stateInfo.getNodeInfoItem(infoItemIdEnum)) -proc logosdelivery_get_available_configs( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = +proc logosdeliveryGetAvailableConfigs*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = ## Returns information about the accepted config items. - requireInitializedNode(ctx, "GetAvailableConfigs"): - return err(errMsg) - let optionMetas: seq[ConfigOptionMeta] = extractConfigOptionMeta(WakuNodeConf) var configOptionDetails = newJArray() diff --git a/library/logos_delivery_api/messaging_api.nim b/library/logos_delivery_api/messaging_api.nim index 4cab9bd93..c1466addf 100644 --- a/library/logos_delivery_api/messaging_api.nim +++ b/library/logos_delivery_api/messaging_api.nim @@ -8,64 +8,46 @@ import logos_delivery/api/types, ../declare_lib -proc logosdelivery_subscribe( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - contentTopicStr: cstring, -) {.ffi.} = - requireInitializedNode(ctx, "Subscribe"): - return err(errMsg) - - requireMessaging(ctx, "Subscribe"): +proc logosdeliverySubscribe*( + lib: LogosDelivery, contentTopicStr: string +): Future[Result[string, string]] {.ffi.} = + requireMessaging(lib, "Subscribe"): return err(errMsg) # ContentTopic is just a string type alias - let contentTopic = ContentTopic($contentTopicStr) + let contentTopic = ContentTopic(contentTopicStr) - (await ctx.myLib[].messagingClient.subscribe(contentTopic)).isOkOr: + (await lib.messagingClient.subscribe(contentTopic)).isOkOr: let errMsg = $error return err("Subscribe failed: " & errMsg) return ok("") -proc logosdelivery_unsubscribe( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - contentTopicStr: cstring, -) {.ffi.} = - requireInitializedNode(ctx, "Unsubscribe"): - return err(errMsg) - - requireMessaging(ctx, "Unsubscribe"): +proc logosdeliveryUnsubscribe*( + lib: LogosDelivery, contentTopicStr: string +): Future[Result[string, string]] {.ffi.} = + requireMessaging(lib, "Unsubscribe"): return err(errMsg) # ContentTopic is just a string type alias - let contentTopic = ContentTopic($contentTopicStr) + let contentTopic = ContentTopic(contentTopicStr) - ctx.myLib[].messagingClient.unsubscribe(contentTopic).isOkOr: + lib.messagingClient.unsubscribe(contentTopic).isOkOr: let errMsg = $error return err("Unsubscribe failed: " & errMsg) return ok("") -proc logosdelivery_send( - ctx: ptr FFIContext[LogosDelivery], - callback: FFICallBack, - userData: pointer, - messageJson: cstring, -) {.ffi.} = - requireInitializedNode(ctx, "Send"): - return err(errMsg) - - requireMessaging(ctx, "Send"): +proc logosdeliverySend*( + lib: LogosDelivery, messageJson: string +): Future[Result[string, string]] {.ffi.} = + requireMessaging(lib, "Send"): return err(errMsg) ## Parse the message JSON and send the message var jsonNode: JsonNode try: - jsonNode = parseJson($messageJson) + jsonNode = parseJson(messageJson) except Exception as e: return err("Failed to parse message JSON: " & e.msg) @@ -93,7 +75,7 @@ proc logosdelivery_send( ) # Send the message via the messaging layer's own API. - let requestId = (await ctx.myLib[].messagingClient.send(envelope)).valueOr: + let requestId = (await lib.messagingClient.send(envelope)).valueOr: let errMsg = $error return err("Send failed: " & errMsg) diff --git a/library/logos_delivery_api/node_api.nim b/library/logos_delivery_api/node_api.nim index 9547257b1..2f23e48e1 100644 --- a/library/logos_delivery_api/node_api.nim +++ b/library/logos_delivery_api/node_api.nim @@ -16,206 +16,129 @@ import proc `%`*(id: RequestId): JsonNode = %($id) -registerReqFFI(CreateNodeRequest, ctx: ptr FFIContext[LogosDelivery]): - proc(configJson: cstring): Future[Result[string, string]] {.async.} = - let conf = parseLogosDeliveryConf($configJson).valueOr: - error "Failed to parse Logos Delivery configuration JSON", - error = error, configJson = $configJson - return err("failed parseLogosDeliveryConf " & error) +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) - ctx.myLib[] = (await LogosDelivery.new(conf)).valueOr: - let errMsg = $error - chronicles.error "CreateNodeRequest failed", err = errMsg - return err(errMsg) - - return ok("") - -proc logosdelivery_destroy( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -): cint {.dynlib, exportc, cdecl.} = - initializeLibrary() - checkParams(ctx, callback, userData) - - ffi.destroyFFIContext(ctx).isOkOr: - let msg = "liblogosdelivery error: " & $error - callback(RET_ERR, unsafeAddr msg[0], cast[csize_t](len(msg)), userData) - return RET_ERR - - ## always need to invoke the callback although we don't retrieve value to the caller - callback(RET_OK, nil, 0, userData) - - return RET_OK - -proc logosdelivery_create_node( - configJson: cstring, callback: FFICallback, userData: pointer -): pointer {.dynlib, exportc, cdecl.} = - initializeLibrary() - - if callback.isNil(): - echo "error: missing callback in logosdelivery_create_node" - return nil - - var ctx = ffi.createFFIContext[LogosDelivery]().valueOr: - let msg = "Error in createFFIContext: " & $error - callback(RET_ERR, unsafeAddr msg[0], cast[csize_t](len(msg)), userData) - return nil - - ctx.userData = userData - - ffi.sendRequestToFFIThread( - ctx, CreateNodeRequest.ffiNewReq(callback, userData, configJson) - ).isOkOr: - let msg = "error in sendRequestToFFIThread: " & $error - callback(RET_ERR, unsafeAddr msg[0], cast[csize_t](len(msg)), userData) - # free allocated resources as they won't be available - ffi.destroyFFIContext(ctx).isOkOr: - chronicles.error "Error in destroyFFIContext after sendRequestToFFIThread during creation", - err = $error - return nil - - return ctx - -proc logosdelivery_start_node( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - requireInitializedNode(ctx, "START_NODE"): - return err(errMsg) + return await LogosDelivery.new(conf) +proc logosdeliveryStartNode*( + lib: LogosDelivery +): Future[Result[string, string]] {.ffi.} = # setting up outgoing event listeners let sentListener = MessageSentEvent.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: MessageSentEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onMessageSent"): - $newJsonEvent("message_sent", event), + emitMessageSent($event.requestId, event.messageHash), ).valueOr: chronicles.error "MessageSentEvent.listen failed", err = $error return err("MessageSentEvent.listen failed: " & $error) let errorListener = MessageErrorEvent.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: MessageErrorEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onMessageError"): - $newJsonEvent("message_error", event), + 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( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: MessagePropagatedEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onMessagePropagated"): - $newJsonEvent("message_propagated", event), + emitMessagePropagated($event.requestId, event.messageHash), ).valueOr: chronicles.error "MessagePropagatedEvent.listen failed", err = $error return err("MessagePropagatedEvent.listen failed: " & $error) let receivedListener = MessageReceivedEvent.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: MessageReceivedEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onMessageReceived"): - $newJsonEvent("message_received", event), + emitMessageReceived(event.messageHash, event.message), ).valueOr: chronicles.error "MessageReceivedEvent.listen failed", err = $error return err("MessageReceivedEvent.listen failed: " & $error) let ConnectionStatusChangeListener = EventConnectionStatusChange.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: EventConnectionStatusChange) {.async: (raises: []).} = - callEventCallback(ctx, "onConnectionStatusChange"): - $newJsonEvent("connection_status_change", event), + emitConnectionStatusChange($event.connectionStatus), ).valueOr: chronicles.error "ConnectionStatusChange.listen failed", err = $error return err("ConnectionStatusChange.listen failed: " & $error) let shardTopicHealthListener = EventShardTopicHealthChange.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: EventShardTopicHealthChange) {.async: (raises: []).} = - callEventCallback(ctx, "onTopicHealthChange"): - $( - %*{ - "eventType": "relay_topic_health_change", - "pubsubTopic": $event.topic, - "topicHealth": $event.health, - } - ), + emitTopicHealthChange($event.topic, $event.health), ).valueOr: chronicles.error "EventShardTopicHealthChange.listen failed", err = $error return err("EventShardTopicHealthChange.listen failed: " & $error) let peerEventListener = WakuPeerEvent.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: WakuPeerEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onConnectionChange"): - $( - %*{ - "eventType": "connection_change", - "peerId": $event.peerId, - "peerEvent": $event.kind, - } - ), + emitConnectionChange($event.peerId, $event.kind), ).valueOr: chronicles.error "WakuPeerEvent.listen failed", err = $error return err("WakuPeerEvent.listen failed: " & $error) let channelReceivedListener = ChannelMessageReceivedEvent.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: ChannelMessageReceivedEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onChannelMessageReceived"): - $( - %*{ - "eventType": "channel_message_received", - "channelId": string(event.channelId), - "senderId": $event.senderId, - "payload": string(base64.encode(event.payload)), - } - ), + 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( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: ChannelMessageSentEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onChannelMessageSent"): - $newJsonEvent("channel_message_sent", event), + emitChannelMessageSent(string(event.channelId), $event.requestId), ).valueOr: chronicles.error "ChannelMessageSentEvent.listen failed", err = $error return err("ChannelMessageSentEvent.listen failed: " & $error) let channelErrorListener = ChannelMessageErrorEvent.listen( - ctx.myLib[].waku.brokerCtx, + lib.waku.brokerCtx, proc(event: ChannelMessageErrorEvent) {.async: (raises: []).} = - callEventCallback(ctx, "onChannelMessageError"): - $newJsonEvent("channel_message_error", event), + emitChannelMessageError( + string(event.channelId), $event.requestId, event.error + ), ).valueOr: chronicles.error "ChannelMessageErrorEvent.listen failed", err = $error return err("ChannelMessageErrorEvent.listen failed: " & $error) - (await ctx.myLib[].start()).isOkOr: + (await lib.start()).isOkOr: let errMsg = $error chronicles.error "START_NODE failed", err = errMsg return err("failed to start: " & errMsg) return ok("") -proc logosdelivery_stop_node( - ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer -) {.ffi.} = - requireInitializedNode(ctx, "STOP_NODE"): - return err(errMsg) +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 MessageErrorEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await MessageSentEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await MessagePropagatedEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await MessageReceivedEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await EventConnectionStatusChange.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await EventShardTopicHealthChange.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await WakuPeerEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await ChannelMessageReceivedEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await ChannelMessageSentEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - await ChannelMessageErrorEvent.dropAllListeners(ctx.myLib[].waku.brokerCtx) - - (await ctx.myLib[].stop()).isOkOr: + (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 diff --git a/library/rust_bindings/Cargo.toml b/library/rust_bindings/Cargo.toml new file mode 100644 index 000000000..f47fc958b --- /dev/null +++ b/library/rust_bindings/Cargo.toml @@ -0,0 +1,13 @@ +[package] +name = "logosdelivery" +version = "0.1.0" +edition = "2021" + +[dependencies] +serde = { version = "1", features = ["derive"] } +ciborium = "0.2" +flume = { version = "0.11", default-features = false, features = ["async"] } +tokio = { version = "1", features = ["sync", "time"] } + +[dev-dependencies] +tokio = { version = "1", features = ["rt-multi-thread", "macros", "sync", "time"] } diff --git a/library/rust_bindings/build.rs b/library/rust_bindings/build.rs new file mode 100644 index 000000000..86b4359dd --- /dev/null +++ b/library/rust_bindings/build.rs @@ -0,0 +1,47 @@ +use std::path::PathBuf; +use std::process::Command; + +fn main() { + let manifest = PathBuf::from(std::env::var("CARGO_MANIFEST_DIR").unwrap()); + let nim_src = manifest.join("library/liblogosdelivery.nim"); + let nim_src = nim_src.canonicalize().unwrap_or(manifest.join("library/liblogosdelivery.nim")); + + // Walk up to find the nim-ffi repo root (directory containing nim_src's library) + // The repo root is where nim c should be run from (contains config.nims). + // We assume nim_src lives somewhere under repo_root. + // Derive repo_root as the ancestor that contains the .nimble file or config.nims. + let mut repo_root = nim_src.clone(); + loop { + repo_root = match repo_root.parent() { + Some(p) => p.to_path_buf(), + None => break, + }; + if repo_root.join("config.nims").exists() || repo_root.join("ffi.nimble").exists() { + break; + } + } + + #[cfg(target_os = "macos")] + let lib_ext = "dylib"; + #[cfg(target_os = "linux")] + let lib_ext = "so"; + + let out_lib = repo_root.join(format!("liblogosdelivery.{lib_ext}")); + + let mut cmd = Command::new("nim"); + cmd.arg("c") + .arg("--mm:orc") + .arg("-d:chronicles_log_level=WARN") + .arg("--app:lib") + .arg("--noMain") + .arg(format!("--nimMainPrefix:liblogosdelivery")) + .arg(format!("-o:{}", out_lib.display())); + cmd.arg(&nim_src).current_dir(&repo_root); + + let status = cmd.status().expect("failed to run nim compiler"); + assert!(status.success(), "Nim compilation failed"); + + println!("cargo:rustc-link-search={}", repo_root.display()); + println!("cargo:rustc-link-lib=logosdelivery"); + println!("cargo:rerun-if-changed={}", nim_src.display()); +} diff --git a/library/rust_bindings/src/api.rs b/library/rust_bindings/src/api.rs new file mode 100644 index 000000000..865bec979 --- /dev/null +++ b/library/rust_bindings/src/api.rs @@ -0,0 +1,1425 @@ +use std::os::raw::{c_char, c_int, c_void}; +use std::slice; +use std::time::Duration; +use serde::de::DeserializeOwned; +use serde::Serialize; +use super::ffi; +use super::types::*; + +fn encode_cbor(value: &T) -> Result, String> { + let mut buf = Vec::new(); + ciborium::ser::into_writer(value, &mut buf).map_err(|e| e.to_string())?; + Ok(buf) +} + +fn decode_cbor(bytes: &[u8]) -> Result { + ciborium::de::from_reader(bytes).map_err(|e| e.to_string()) +} + +type FFIResult = Result, String>; +type FFISender = flume::Sender; + +// Reconstruct the (ret, msg, len) tuple delivered by the C callback +// into a Result, String>: payload on success, UTF-8 message on error. +// `from_utf8_lossy` accepts non-UTF-8 error bytes by inserting U+FFFD; the +// alternative would be to dispatch a separate Err for invalid UTF-8, but the +// codegen contract is that Nim handlers emit `string` error payloads, so +// invalid UTF-8 here would be a Nim-side bug. +unsafe fn ffi_payload(ret: c_int, msg: *const c_char, len: usize) -> FFIResult { + let bytes = if msg.is_null() || len == 0 { + Vec::new() + } else { + slice::from_raw_parts(msg as *const u8, len).to_vec() + }; + if ret == NIMFFI_RET_OK { Ok(bytes) } + else { Err(String::from_utf8_lossy(&bytes).into_owned()) } +} + +// nim-ffi result-callback status codes (mirror ffi/ffi_types.nim). +const NIMFFI_RET_OK: c_int = 0; +const NIMFFI_RET_MISSING_CALLBACK: c_int = 2; +const NIMFFI_RET_STALE_WARN: c_int = 3; + +unsafe extern "C" fn on_result( + ret: c_int, + msg: *const c_char, + len: usize, + user_data: *mut c_void, +) { + // NIMFFI_RET_STALE_WARN (3) is a non-terminal progress ping: the request + // is still running. This wrapper only delivers the final result, so ignore + // it WITHOUT reclaiming the box — a terminal callback still owns the Sender. + if ret == NIMFFI_RET_STALE_WARN { return; } + + // Take ownership of the boxed Sender — dropping it at end of scope + // releases the only outstanding handle. + let tx = Box::from_raw(user_data as *mut FFISender); + + // `tx.send` returns Err only if the awaiting future was dropped (and with it + // the Receiver): e.g. tokio::time::timeout elapsed, a tokio::select! branch + // lost the race, or the future was dropped before being awaited. This cannot + // happen with the current rust_client demo but may occur in arbitrary + // downstream consumers, so we discard the Err safely. + // Given that this is invoked from a Nim thread, we can't propagate the error by panicking or + // returning a Result. Furthermore, an API dev may intentionally set a timeout in the await, + // in which case is also fine to discard the send error in this case because the API user will + // handle the timeout expiry in their own code. + // The important part is to ensure that the callback doesn't panic or block indefinitely if the + // receiver is gone. + let _ = tx.send(ffi_payload(ret, msg, len)); +} + +fn ffi_call_sync(timeout: Duration, f: F) -> FFIResult +where + F: FnOnce(ffi::FFICallback, *mut c_void) -> c_int, +{ + let (tx, rx) = flume::bounded::(1); + let raw = Box::into_raw(Box::new(tx)) as *mut c_void; + let ret = f(on_result, raw); + if ret == NIMFFI_RET_MISSING_CALLBACK { + // Callback will never fire; reclaim the box to avoid a leak. + drop(unsafe { Box::from_raw(raw as *mut FFISender) }); + return Err("RET_MISSING_CALLBACK (internal error)".into()); + } + match rx.recv_timeout(timeout) { + Ok(payload) => payload, + Err(flume::RecvTimeoutError::Timeout) => + Err(format!("timed out after {:?}", timeout)), + Err(flume::RecvTimeoutError::Disconnected) => + Err("callback channel disconnected before delivery".into()), + } +} + +async fn ffi_call_async(timeout: Duration, f: F) -> FFIResult +where + F: FnOnce(ffi::FFICallback, *mut c_void) -> c_int, +{ + let (tx, rx) = flume::bounded::(1); + let raw = Box::into_raw(Box::new(tx)) as *mut c_void; + let ret = f(on_result, raw); + if ret == NIMFFI_RET_MISSING_CALLBACK { + drop(unsafe { Box::from_raw(raw as *mut FFISender) }); + return Err("RET_MISSING_CALLBACK (internal error)".into()); + } + match tokio::time::timeout(timeout, rx.recv_async()).await { + Ok(Ok(payload)) => payload, + Ok(Err(_)) => Err("callback channel disconnected before delivery".into()), + Err(_) => Err(format!("timed out after {:?}", timeout)), + } +} + +struct OnMessageSentHandler { + f: Box, +} + +unsafe extern "C" fn on_message_sent_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnMessageSentHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: MessageSentPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnMessageErrorHandler { + f: Box, +} + +unsafe extern "C" fn on_message_error_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnMessageErrorHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: MessageErrorPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnMessagePropagatedHandler { + f: Box, +} + +unsafe extern "C" fn on_message_propagated_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnMessagePropagatedHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: MessagePropagatedPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnMessageReceivedHandler { + f: Box, +} + +unsafe extern "C" fn on_message_received_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnMessageReceivedHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: MessageReceivedPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnConnectionStatusChangeHandler { + f: Box, +} + +unsafe extern "C" fn on_connection_status_change_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnConnectionStatusChangeHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: ConnectionStatusChangePayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnTopicHealthChangeHandler { + f: Box, +} + +unsafe extern "C" fn on_topic_health_change_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnTopicHealthChangeHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: TopicHealthChangePayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnConnectionChangeHandler { + f: Box, +} + +unsafe extern "C" fn on_connection_change_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnConnectionChangeHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: ConnectionChangePayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnChannelMessageReceivedHandler { + f: Box, +} + +unsafe extern "C" fn on_channel_message_received_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnChannelMessageReceivedHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: ChannelMessageReceivedPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnChannelMessageSentHandler { + f: Box, +} + +unsafe extern "C" fn on_channel_message_sent_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnChannelMessageSentHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: ChannelMessageSentPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnChannelMessageErrorHandler { + f: Box, +} + +unsafe extern "C" fn on_channel_message_error_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnChannelMessageErrorHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: ChannelMessageErrorPayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +struct OnReceivedMessageHandler { + f: Box, +} + +unsafe extern "C" fn on_received_message_trampoline( + ret: c_int, msg: *const c_char, len: usize, ud: *mut c_void, +) { + if ud.is_null() || ret != 0 || msg.is_null() || len == 0 { + return; + } + let h = &*(ud as *const OnReceivedMessageHandler); + let bytes = slice::from_raw_parts(msg as *const u8, len); + #[derive(serde::Deserialize)] + struct Envelope { payload: ReceivedMessagePayload } + if let Ok(env) = ciborium::de::from_reader::(bytes) { + (h.f)(&env.payload); + } +} + +#[derive(Debug, Clone, Copy)] +pub struct ListenerHandle { pub id: u64 } + +/// High-level context for `LogosDelivery`. +pub struct LogosDeliveryCtx { + ptr: *mut c_void, + timeout: Duration, + listeners: std::sync::Mutex>>, +} + +// SAFETY: The `ptr` field points to an FFIContext owned by the Nim runtime. +// Every call through the generated FFI proc goes through +// `sendRequestToFFIThread` on the Nim side, which only enqueues the request +// onto a mutex-guarded MPSC queue (sound from any number of threads) and +// wakes the single FFI thread that dispatches every handler. The context is +// thus never mutated non-atomically from the caller's thread. The Nim-side +// reentrancy guard (`onFFIThread` threadvar) prevents handlers from +// re-entering the dispatcher. These invariants make it sound to mark the +// wrapper as Send + Sync. +unsafe impl Send for LogosDeliveryCtx {} +unsafe impl Sync for LogosDeliveryCtx {} + +impl Drop for LogosDeliveryCtx { + fn drop(&mut self) { + if !self.ptr.is_null() { + unsafe { ffi::logosdelivery_destroy(self.ptr); } + self.ptr = std::ptr::null_mut(); + } + } +} + +impl LogosDeliveryCtx { + pub fn create(config_json: String, timeout: Duration) -> Result { + let req = LogosdeliveryCreateNodeCtorReq { config_json }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(timeout, |cb, ud| unsafe { + let _ = ffi::logosdelivery_create_node(req_bytes.as_ptr(), req_bytes.len(), cb, ud); + 0 + })?; + let addr_str: String = decode_cbor(&raw_bytes)?; + let addr: usize = addr_str.parse().map_err(|e: std::num::ParseIntError| e.to_string())?; + Ok(Self { ptr: addr as *mut c_void, timeout, listeners: std::sync::Mutex::new(std::collections::HashMap::new()) }) + } + + pub async fn new_async(config_json: String, timeout: Duration) -> Result { + let req = LogosdeliveryCreateNodeCtorReq { config_json }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_async(timeout, move |cb, ud| unsafe { + let _ = ffi::logosdelivery_create_node(req_bytes.as_ptr(), req_bytes.len(), cb, ud); + 0 + }).await?; + let addr_str: String = decode_cbor(&raw_bytes)?; + let addr: usize = addr_str.parse().map_err(|e: std::num::ParseIntError| e.to_string())?; + Ok(Self { ptr: addr as *mut c_void, timeout, listeners: std::sync::Mutex::new(std::collections::HashMap::new()) }) + } + + fn add_listener_inner( + &self, + event_name: *const c_char, + callback: ffi::FFICallback, + raw: *mut c_void, + owned: Box, + ) -> ListenerHandle { + let id = unsafe { + ffi::logosdelivery_add_event_listener(self.ptr, event_name, callback, raw) + }; + if id != 0 { + self.listeners.lock().unwrap().insert(id, owned); + } + ListenerHandle { id } + } + + /// Register a typed listener for `onMessageSent`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_message_sent_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&MessageSentPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnMessageSentHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnMessageSentHandler as *mut c_void; + self.add_listener_inner(b"onMessageSent\0".as_ptr() as *const c_char, on_message_sent_trampoline, raw, owned) + } + + /// Register a typed listener for `onMessageError`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_message_error_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&MessageErrorPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnMessageErrorHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnMessageErrorHandler as *mut c_void; + self.add_listener_inner(b"onMessageError\0".as_ptr() as *const c_char, on_message_error_trampoline, raw, owned) + } + + /// Register a typed listener for `onMessagePropagated`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_message_propagated_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&MessagePropagatedPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnMessagePropagatedHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnMessagePropagatedHandler as *mut c_void; + self.add_listener_inner(b"onMessagePropagated\0".as_ptr() as *const c_char, on_message_propagated_trampoline, raw, owned) + } + + /// Register a typed listener for `onMessageReceived`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_message_received_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&MessageReceivedPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnMessageReceivedHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnMessageReceivedHandler as *mut c_void; + self.add_listener_inner(b"onMessageReceived\0".as_ptr() as *const c_char, on_message_received_trampoline, raw, owned) + } + + /// Register a typed listener for `onConnectionStatusChange`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_connection_status_change_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&ConnectionStatusChangePayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnConnectionStatusChangeHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnConnectionStatusChangeHandler as *mut c_void; + self.add_listener_inner(b"onConnectionStatusChange\0".as_ptr() as *const c_char, on_connection_status_change_trampoline, raw, owned) + } + + /// Register a typed listener for `onTopicHealthChange`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_topic_health_change_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&TopicHealthChangePayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnTopicHealthChangeHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnTopicHealthChangeHandler as *mut c_void; + self.add_listener_inner(b"onTopicHealthChange\0".as_ptr() as *const c_char, on_topic_health_change_trampoline, raw, owned) + } + + /// Register a typed listener for `onConnectionChange`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_connection_change_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&ConnectionChangePayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnConnectionChangeHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnConnectionChangeHandler as *mut c_void; + self.add_listener_inner(b"onConnectionChange\0".as_ptr() as *const c_char, on_connection_change_trampoline, raw, owned) + } + + /// Register a typed listener for `onChannelMessageReceived`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_channel_message_received_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&ChannelMessageReceivedPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnChannelMessageReceivedHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnChannelMessageReceivedHandler as *mut c_void; + self.add_listener_inner(b"onChannelMessageReceived\0".as_ptr() as *const c_char, on_channel_message_received_trampoline, raw, owned) + } + + /// Register a typed listener for `onChannelMessageSent`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_channel_message_sent_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&ChannelMessageSentPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnChannelMessageSentHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnChannelMessageSentHandler as *mut c_void; + self.add_listener_inner(b"onChannelMessageSent\0".as_ptr() as *const c_char, on_channel_message_sent_trampoline, raw, owned) + } + + /// Register a typed listener for `onChannelMessageError`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_channel_message_error_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&ChannelMessageErrorPayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnChannelMessageErrorHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnChannelMessageErrorHandler as *mut c_void; + self.add_listener_inner(b"onChannelMessageError\0".as_ptr() as *const c_char, on_channel_message_error_trampoline, raw, owned) + } + + /// Register a typed listener for `onReceivedMessage`. The returned handle can be + /// passed to `remove_event_listener` to unregister. + pub fn add_on_received_message_listener(&self, handler: F) -> ListenerHandle + where F: Fn(&ReceivedMessagePayload) + Send + Sync + 'static, + { + let owned: Box = Box::new(OnReceivedMessageHandler { f: Box::new(handler) }); + let raw = &*owned as *const OnReceivedMessageHandler as *mut c_void; + self.add_listener_inner(b"onReceivedMessage\0".as_ptr() as *const c_char, on_received_message_trampoline, raw, owned) + } + + /// Remove a previously-registered listener by handle. Returns true + /// if the listener existed and was removed; false otherwise. + pub fn remove_event_listener(&self, handle: ListenerHandle) -> bool { + if handle.id == 0 { return false; } + let rc = unsafe { + ffi::logosdelivery_remove_event_listener(self.ptr, handle.id) + }; + self.listeners.lock().unwrap().remove(&handle.id); + rc == 0 + } + + pub fn start_node(&self) -> Result { + let req = LogosdeliveryStartNodeReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_start_node(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn start_node_async(&self) -> Result { + let req = LogosdeliveryStartNodeReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_start_node(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn stop_node(&self) -> Result { + let req = LogosdeliveryStopNodeReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_stop_node(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn stop_node_async(&self) -> Result { + let req = LogosdeliveryStopNodeReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_stop_node(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn subscribe(&self, content_topic_str: String) -> Result { + let req = LogosdeliverySubscribeReq { content_topic_str }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_subscribe(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn subscribe_async(&self, content_topic_str: String) -> Result { + let req = LogosdeliverySubscribeReq { content_topic_str }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_subscribe(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn unsubscribe(&self, content_topic_str: String) -> Result { + let req = LogosdeliveryUnsubscribeReq { content_topic_str }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_unsubscribe(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn unsubscribe_async(&self, content_topic_str: String) -> Result { + let req = LogosdeliveryUnsubscribeReq { content_topic_str }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_unsubscribe(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn send(&self, message_json: String) -> Result { + let req = LogosdeliverySendReq { message_json }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_send(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn send_async(&self, message_json: String) -> Result { + let req = LogosdeliverySendReq { message_json }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_send(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn get_available_node_info_ids(&self) -> Result { + let req = LogosdeliveryGetAvailableNodeInfoIdsReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_get_available_node_info_ids(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn get_available_node_info_ids_async(&self) -> Result { + let req = LogosdeliveryGetAvailableNodeInfoIdsReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_get_available_node_info_ids(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn get_node_info(&self, node_info_id: String) -> Result { + let req = LogosdeliveryGetNodeInfoReq { node_info_id }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_get_node_info(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn get_node_info_async(&self, node_info_id: String) -> Result { + let req = LogosdeliveryGetNodeInfoReq { node_info_id }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_get_node_info(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn get_available_configs(&self) -> Result { + let req = LogosdeliveryGetAvailableConfigsReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_get_available_configs(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn get_available_configs_async(&self) -> Result { + let req = LogosdeliveryGetAvailableConfigsReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_get_available_configs(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_peerids_from_peerstore(&self) -> Result { + let req = WakuGetPeeridsFromPeerstoreReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_peerids_from_peerstore(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_peerids_from_peerstore_async(&self) -> Result { + let req = WakuGetPeeridsFromPeerstoreReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_peerids_from_peerstore(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_connect(&self, peer_multi_addr: String, timeout_ms: u32) -> Result { + let req = WakuConnectReq { peer_multi_addr, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_connect(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_connect_async(&self, peer_multi_addr: String, timeout_ms: u32) -> Result { + let req = WakuConnectReq { peer_multi_addr, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_connect(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_disconnect_peer_by_id(&self, peer_id: String) -> Result { + let req = WakuDisconnectPeerByIdReq { peer_id }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_disconnect_peer_by_id(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_disconnect_peer_by_id_async(&self, peer_id: String) -> Result { + let req = WakuDisconnectPeerByIdReq { peer_id }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_disconnect_peer_by_id(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_disconnect_all_peers(&self) -> Result { + let req = WakuDisconnectAllPeersReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_disconnect_all_peers(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_disconnect_all_peers_async(&self) -> Result { + let req = WakuDisconnectAllPeersReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_disconnect_all_peers(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_dial_peer(&self, peer_multi_addr: String, protocol: String, timeout_ms: u32) -> Result { + let req = WakuDialPeerReq { peer_multi_addr, protocol, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_dial_peer(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_dial_peer_async(&self, peer_multi_addr: String, protocol: String, timeout_ms: u32) -> Result { + let req = WakuDialPeerReq { peer_multi_addr, protocol, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_dial_peer(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_dial_peer_by_id(&self, peer_id: String, protocol: String, timeout_ms: u32) -> Result { + let req = WakuDialPeerByIdReq { peer_id, protocol, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_dial_peer_by_id(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_dial_peer_by_id_async(&self, peer_id: String, protocol: String, timeout_ms: u32) -> Result { + let req = WakuDialPeerByIdReq { peer_id, protocol, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_dial_peer_by_id(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_connected_peers_info(&self) -> Result { + let req = WakuGetConnectedPeersInfoReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_connected_peers_info(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_connected_peers_info_async(&self) -> Result { + let req = WakuGetConnectedPeersInfoReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_connected_peers_info(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_connected_peers(&self) -> Result { + let req = WakuGetConnectedPeersReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_connected_peers(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_connected_peers_async(&self) -> Result { + let req = WakuGetConnectedPeersReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_connected_peers(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_peerids_by_protocol(&self, protocol: String) -> Result { + let req = WakuGetPeeridsByProtocolReq { protocol }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_peerids_by_protocol(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_peerids_by_protocol_async(&self, protocol: String) -> Result { + let req = WakuGetPeeridsByProtocolReq { protocol }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_peerids_by_protocol(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_discv5_update_bootnodes(&self, bootnodes: String) -> Result { + let req = WakuDiscv5UpdateBootnodesReq { bootnodes }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_discv5_update_bootnodes(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_discv5_update_bootnodes_async(&self, bootnodes: String) -> Result { + let req = WakuDiscv5UpdateBootnodesReq { bootnodes }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_discv5_update_bootnodes(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_dns_discovery(&self, enr_tree_url: String, name_dns_server: String, timeout_ms: i32) -> Result { + let req = WakuDnsDiscoveryReq { enr_tree_url, name_dns_server, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_dns_discovery(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_dns_discovery_async(&self, enr_tree_url: String, name_dns_server: String, timeout_ms: i32) -> Result { + let req = WakuDnsDiscoveryReq { enr_tree_url, name_dns_server, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_dns_discovery(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_start_discv5(&self) -> Result { + let req = WakuStartDiscv5Req {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_start_discv5(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_start_discv5_async(&self) -> Result { + let req = WakuStartDiscv5Req {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_start_discv5(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_stop_discv5(&self) -> Result { + let req = WakuStopDiscv5Req {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_stop_discv5(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_stop_discv5_async(&self) -> Result { + let req = WakuStopDiscv5Req {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_stop_discv5(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_peer_exchange_request(&self, num_peers: u64) -> Result { + let req = WakuPeerExchangeRequestReq { num_peers }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_peer_exchange_request(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_peer_exchange_request_async(&self, num_peers: u64) -> Result { + let req = WakuPeerExchangeRequestReq { num_peers }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_peer_exchange_request(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_version(&self) -> Result { + let req = WakuVersionReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_version(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_version_async(&self) -> Result { + let req = WakuVersionReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_version(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_listen_addresses(&self) -> Result { + let req = WakuListenAddressesReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_listen_addresses(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_listen_addresses_async(&self) -> Result { + let req = WakuListenAddressesReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_listen_addresses(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_my_enr(&self) -> Result { + let req = WakuGetMyEnrReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_my_enr(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_my_enr_async(&self) -> Result { + let req = WakuGetMyEnrReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_my_enr(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_my_peerid(&self) -> Result { + let req = WakuGetMyPeeridReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_my_peerid(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_my_peerid_async(&self) -> Result { + let req = WakuGetMyPeeridReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_my_peerid(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_get_metrics(&self) -> Result { + let req = WakuGetMetricsReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_get_metrics(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_get_metrics_async(&self) -> Result { + let req = WakuGetMetricsReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_get_metrics(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_is_online(&self) -> Result { + let req = WakuIsOnlineReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_is_online(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_is_online_async(&self) -> Result { + let req = WakuIsOnlineReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_is_online(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_ping_peer(&self, peer_addr: String, timeout_ms: u32) -> Result { + let req = WakuPingPeerReq { peer_addr, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_ping_peer(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_ping_peer_async(&self, peer_addr: String, timeout_ms: u32) -> Result { + let req = WakuPingPeerReq { peer_addr, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_ping_peer(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_get_peers_in_mesh(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetPeersInMeshReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_get_peers_in_mesh(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_get_peers_in_mesh_async(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetPeersInMeshReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_get_peers_in_mesh(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_get_num_peers_in_mesh(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetNumPeersInMeshReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_get_num_peers_in_mesh(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_get_num_peers_in_mesh_async(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetNumPeersInMeshReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_get_num_peers_in_mesh(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_get_connected_peers(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetConnectedPeersReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_get_connected_peers(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_get_connected_peers_async(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetConnectedPeersReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_get_connected_peers(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_get_num_connected_peers(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetNumConnectedPeersReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_get_num_connected_peers(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_get_num_connected_peers_async(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayGetNumConnectedPeersReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_get_num_connected_peers(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_add_protected_shard(&self, cluster_id: i32, shard_id: i32, public_key: String) -> Result { + let req = WakuRelayAddProtectedShardReq { cluster_id, shard_id, public_key }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_add_protected_shard(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_add_protected_shard_async(&self, cluster_id: i32, shard_id: i32, public_key: String) -> Result { + let req = WakuRelayAddProtectedShardReq { cluster_id, shard_id, public_key }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_add_protected_shard(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_subscribe(&self, pub_sub_topic: String) -> Result { + let req = WakuRelaySubscribeReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_subscribe(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_subscribe_async(&self, pub_sub_topic: String) -> Result { + let req = WakuRelaySubscribeReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_subscribe(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_unsubscribe(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayUnsubscribeReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_unsubscribe(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_unsubscribe_async(&self, pub_sub_topic: String) -> Result { + let req = WakuRelayUnsubscribeReq { pub_sub_topic }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_unsubscribe(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_relay_publish(&self, pub_sub_topic: String, json_waku_message: String, timeout_ms: u32) -> Result { + let req = WakuRelayPublishReq { pub_sub_topic, json_waku_message, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_relay_publish(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_relay_publish_async(&self, pub_sub_topic: String, json_waku_message: String, timeout_ms: u32) -> Result { + let req = WakuRelayPublishReq { pub_sub_topic, json_waku_message, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_relay_publish(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_default_pubsub_topic(&self) -> Result { + let req = WakuDefaultPubsubTopicReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_default_pubsub_topic(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_default_pubsub_topic_async(&self) -> Result { + let req = WakuDefaultPubsubTopicReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_default_pubsub_topic(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_content_topic(&self, app_name: String, app_version: u32, content_topic_name: String, encoding: String) -> Result { + let req = WakuContentTopicReq { app_name, app_version, content_topic_name, encoding }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_content_topic(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_content_topic_async(&self, app_name: String, app_version: u32, content_topic_name: String, encoding: String) -> Result { + let req = WakuContentTopicReq { app_name, app_version, content_topic_name, encoding }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_content_topic(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_pubsub_topic(&self, topic_name: String) -> Result { + let req = WakuPubsubTopicReq { topic_name }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_pubsub_topic(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_pubsub_topic_async(&self, topic_name: String) -> Result { + let req = WakuPubsubTopicReq { topic_name }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_pubsub_topic(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_store_query(&self, json_query: String, peer_addr: String, timeout_ms: i32) -> Result { + let req = WakuStoreQueryReq { json_query, peer_addr, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_store_query(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_store_query_async(&self, json_query: String, peer_addr: String, timeout_ms: i32) -> Result { + let req = WakuStoreQueryReq { json_query, peer_addr, timeout_ms }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_store_query(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_lightpush_publish(&self, pub_sub_topic: String, json_waku_message: String) -> Result { + let req = WakuLightpushPublishReq { pub_sub_topic, json_waku_message }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_lightpush_publish(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_lightpush_publish_async(&self, pub_sub_topic: String, json_waku_message: String) -> Result { + let req = WakuLightpushPublishReq { pub_sub_topic, json_waku_message }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_lightpush_publish(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_filter_subscribe(&self, pub_sub_topic: String, content_topics: String) -> Result { + let req = WakuFilterSubscribeReq { pub_sub_topic, content_topics }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_filter_subscribe(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_filter_subscribe_async(&self, pub_sub_topic: String, content_topics: String) -> Result { + let req = WakuFilterSubscribeReq { pub_sub_topic, content_topics }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_filter_subscribe(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_filter_unsubscribe(&self, pub_sub_topic: String, content_topics: String) -> Result { + let req = WakuFilterUnsubscribeReq { pub_sub_topic, content_topics }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_filter_unsubscribe(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_filter_unsubscribe_async(&self, pub_sub_topic: String, content_topics: String) -> Result { + let req = WakuFilterUnsubscribeReq { pub_sub_topic, content_topics }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_filter_unsubscribe(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn waku_filter_unsubscribe_all(&self) -> Result { + let req = WakuFilterUnsubscribeAllReq {}; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::waku_filter_unsubscribe_all(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn waku_filter_unsubscribe_all_async(&self) -> Result { + let req = WakuFilterUnsubscribeAllReq {}; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::waku_filter_unsubscribe_all(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn channel_create(&self, channel_id_str: String, content_topic_str: String, sender_id_str: String) -> Result { + let req = LogosdeliveryChannelCreateReq { channel_id_str, content_topic_str, sender_id_str }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_channel_create(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn channel_create_async(&self, channel_id_str: String, content_topic_str: String, sender_id_str: String) -> Result { + let req = LogosdeliveryChannelCreateReq { channel_id_str, content_topic_str, sender_id_str }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_channel_create(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn channel_send(&self, channel_id_str: String, message_json: String) -> Result { + let req = LogosdeliveryChannelSendReq { channel_id_str, message_json }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_channel_send(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn channel_send_async(&self, channel_id_str: String, message_json: String) -> Result { + let req = LogosdeliveryChannelSendReq { channel_id_str, message_json }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_channel_send(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + + pub fn channel_close(&self, channel_id_str: String) -> Result { + let req = LogosdeliveryChannelCloseReq { channel_id_str }; + let req_bytes = encode_cbor(&req)?; + let raw_bytes = ffi_call_sync(self.timeout, |cb, ud| unsafe { + ffi::logosdelivery_channel_close(self.ptr, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + })?; + decode_cbor::(&raw_bytes) + } + + pub async fn channel_close_async(&self, channel_id_str: String) -> Result { + let req = LogosdeliveryChannelCloseReq { channel_id_str }; + let req_bytes = encode_cbor(&req)?; + let ptr = self.ptr as usize; + let raw_bytes = ffi_call_async(self.timeout, move |cb, ud| unsafe { + ffi::logosdelivery_channel_close(ptr as *mut c_void, cb, ud, req_bytes.as_ptr(), req_bytes.len()) + }).await?; + decode_cbor::(&raw_bytes) + } + +} diff --git a/library/rust_bindings/src/ffi.rs b/library/rust_bindings/src/ffi.rs new file mode 100644 index 000000000..68f024d76 --- /dev/null +++ b/library/rust_bindings/src/ffi.rs @@ -0,0 +1,64 @@ +use std::os::raw::{c_char, c_int, c_void}; + +pub type FFICallback = unsafe extern "C" fn( + ret: c_int, + msg: *const c_char, + len: usize, + user_data: *mut c_void, +); + +#[link(name = "logosdelivery")] +extern "C" { + pub fn logosdelivery_create_node(req_cbor: *const u8, req_cbor_len: usize, callback: FFICallback, user_data: *mut c_void) -> *mut c_void; + pub fn logosdelivery_start_node(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_stop_node(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_destroy(ctx: *mut c_void) -> c_int; + pub fn logosdelivery_subscribe(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_unsubscribe(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_send(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_get_available_node_info_ids(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_get_node_info(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_get_available_configs(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_peerids_from_peerstore(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_connect(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_disconnect_peer_by_id(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_disconnect_all_peers(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_dial_peer(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_dial_peer_by_id(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_connected_peers_info(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_connected_peers(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_peerids_by_protocol(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_discv5_update_bootnodes(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_dns_discovery(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_start_discv5(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_stop_discv5(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_peer_exchange_request(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_version(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_listen_addresses(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_my_enr(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_my_peerid(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_get_metrics(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_is_online(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_ping_peer(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_get_peers_in_mesh(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_get_num_peers_in_mesh(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_get_connected_peers(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_get_num_connected_peers(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_add_protected_shard(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_subscribe(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_unsubscribe(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_relay_publish(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_default_pubsub_topic(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_content_topic(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_pubsub_topic(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_store_query(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_lightpush_publish(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_filter_subscribe(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_filter_unsubscribe(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn waku_filter_unsubscribe_all(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_channel_create(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_channel_send(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_channel_close(ctx: *mut c_void, callback: FFICallback, user_data: *mut c_void, req_cbor: *const u8, req_cbor_len: usize) -> c_int; + pub fn logosdelivery_add_event_listener(ctx: *mut c_void, event_name: *const c_char, callback: FFICallback, user_data: *mut c_void) -> u64; + pub fn logosdelivery_remove_event_listener(ctx: *mut c_void, listener_id: u64) -> c_int; +} diff --git a/library/rust_bindings/src/lib.rs b/library/rust_bindings/src/lib.rs new file mode 100644 index 000000000..29c439a46 --- /dev/null +++ b/library/rust_bindings/src/lib.rs @@ -0,0 +1,5 @@ +mod ffi; +mod types; +mod api; +pub use types::*; +pub use api::*; diff --git a/library/rust_bindings/src/types.rs b/library/rust_bindings/src/types.rs new file mode 100644 index 000000000..1b48e059d --- /dev/null +++ b/library/rust_bindings/src/types.rs @@ -0,0 +1,384 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuMessagePayload { + pub payload: String, + #[serde(rename = "contentTopic")] + pub content_topic: String, + pub version: u32, + pub timestamp: i64, + pub ephemeral: bool, + pub meta: String, + pub proof: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessageSentPayload { + #[serde(rename = "requestId")] + pub request_id: String, + #[serde(rename = "messageHash")] + pub message_hash: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessageErrorPayload { + #[serde(rename = "requestId")] + pub request_id: String, + #[serde(rename = "messageHash")] + pub message_hash: String, + pub error: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessagePropagatedPayload { + #[serde(rename = "requestId")] + pub request_id: String, + #[serde(rename = "messageHash")] + pub message_hash: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessageReceivedPayload { + #[serde(rename = "messageHash")] + pub message_hash: String, + pub message: WakuMessagePayload, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ConnectionStatusChangePayload { + #[serde(rename = "connectionStatus")] + pub connection_status: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct TopicHealthChangePayload { + #[serde(rename = "pubsubTopic")] + pub pubsub_topic: String, + #[serde(rename = "topicHealth")] + pub topic_health: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ConnectionChangePayload { + #[serde(rename = "peerId")] + pub peer_id: String, + #[serde(rename = "peerEvent")] + pub peer_event: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ChannelMessageReceivedPayload { + #[serde(rename = "channelId")] + pub channel_id: String, + #[serde(rename = "senderId")] + pub sender_id: String, + pub payload: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ChannelMessageSentPayload { + #[serde(rename = "channelId")] + pub channel_id: String, + #[serde(rename = "requestId")] + pub request_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ChannelMessageErrorPayload { + #[serde(rename = "channelId")] + pub channel_id: String, + #[serde(rename = "requestId")] + pub request_id: String, + pub error: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ReceivedMessagePayload { + #[serde(rename = "pubsubTopic")] + pub pubsub_topic: String, + #[serde(rename = "messageHash")] + pub message_hash: String, + #[serde(rename = "wakuMessage")] + pub waku_message: WakuMessagePayload, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryCreateNodeCtorReq { + #[serde(rename = "configJson")] + pub config_json: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryStartNodeReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryStopNodeReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliverySubscribeReq { + #[serde(rename = "contentTopicStr")] + pub content_topic_str: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryUnsubscribeReq { + #[serde(rename = "contentTopicStr")] + pub content_topic_str: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliverySendReq { + #[serde(rename = "messageJson")] + pub message_json: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryGetAvailableNodeInfoIdsReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryGetNodeInfoReq { + #[serde(rename = "nodeInfoId")] + pub node_info_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryGetAvailableConfigsReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetPeeridsFromPeerstoreReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuConnectReq { + #[serde(rename = "peerMultiAddr")] + pub peer_multi_addr: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: u32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDisconnectPeerByIdReq { + #[serde(rename = "peerId")] + pub peer_id: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDisconnectAllPeersReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDialPeerReq { + #[serde(rename = "peerMultiAddr")] + pub peer_multi_addr: String, + pub protocol: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: u32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDialPeerByIdReq { + #[serde(rename = "peerId")] + pub peer_id: String, + pub protocol: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: u32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetConnectedPeersInfoReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetConnectedPeersReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetPeeridsByProtocolReq { + pub protocol: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDiscv5UpdateBootnodesReq { + pub bootnodes: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDnsDiscoveryReq { + #[serde(rename = "enrTreeUrl")] + pub enr_tree_url: String, + #[serde(rename = "nameDnsServer")] + pub name_dns_server: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: i32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuStartDiscv5Req {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuStopDiscv5Req {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuPeerExchangeRequestReq { + #[serde(rename = "numPeers")] + pub num_peers: u64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuVersionReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuListenAddressesReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetMyEnrReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetMyPeeridReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuGetMetricsReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuIsOnlineReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuPingPeerReq { + #[serde(rename = "peerAddr")] + pub peer_addr: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: u32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayGetPeersInMeshReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayGetNumPeersInMeshReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayGetConnectedPeersReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayGetNumConnectedPeersReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayAddProtectedShardReq { + #[serde(rename = "clusterId")] + pub cluster_id: i32, + #[serde(rename = "shardId")] + pub shard_id: i32, + #[serde(rename = "publicKey")] + pub public_key: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelaySubscribeReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayUnsubscribeReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuRelayPublishReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, + #[serde(rename = "jsonWakuMessage")] + pub json_waku_message: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: u32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuDefaultPubsubTopicReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuContentTopicReq { + #[serde(rename = "appName")] + pub app_name: String, + #[serde(rename = "appVersion")] + pub app_version: u32, + #[serde(rename = "contentTopicName")] + pub content_topic_name: String, + pub encoding: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuPubsubTopicReq { + #[serde(rename = "topicName")] + pub topic_name: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuStoreQueryReq { + #[serde(rename = "jsonQuery")] + pub json_query: String, + #[serde(rename = "peerAddr")] + pub peer_addr: String, + #[serde(rename = "timeoutMs")] + pub timeout_ms: i32, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuLightpushPublishReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, + #[serde(rename = "jsonWakuMessage")] + pub json_waku_message: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuFilterSubscribeReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, + #[serde(rename = "contentTopics")] + pub content_topics: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuFilterUnsubscribeReq { + #[serde(rename = "pubSubTopic")] + pub pub_sub_topic: String, + #[serde(rename = "contentTopics")] + pub content_topics: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct WakuFilterUnsubscribeAllReq {} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryChannelCreateReq { + #[serde(rename = "channelIdStr")] + pub channel_id_str: String, + #[serde(rename = "contentTopicStr")] + pub content_topic_str: String, + #[serde(rename = "senderIdStr")] + pub sender_id_str: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryChannelSendReq { + #[serde(rename = "channelIdStr")] + pub channel_id_str: String, + #[serde(rename = "messageJson")] + pub message_json: String, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct LogosdeliveryChannelCloseReq { + #[serde(rename = "channelIdStr")] + pub channel_id_str: String, +} diff --git a/logos_delivery.nimble b/logos_delivery.nimble index c79e3eced..9903ebb63 100644 --- a/logos_delivery.nimble +++ b/logos_delivery.nimble @@ -61,7 +61,7 @@ requires "nim >= 2.2.4", # Packages not on nimble (use git URLs) -requires "https://github.com/logos-messaging/nim-ffi#v0.1.3" +requires "https://github.com/logos-messaging/nim-ffi#master" requires "https://github.com/logos-messaging/nim-sds.git#b12f5ee07c5b764303b51fb948b32a4ade1de3b5" @@ -504,6 +504,14 @@ task liblogosdeliveryStaticLinux, "Generate bindings": task liblogosdeliveryStaticMac, "Generate bindings": buildLibStaticMac("liblogosdelivery", "library") +task liblogosdeliveryGenBindingsRust, "Emit the Rust bindings for liblogosdelivery": + # --compileOnly is enough: genBindings() writes the files during macro + # expansion, so nothing needs linking. + exec "nim c --threads:on --mm:refc --skipParentCfg:off -d:discv5_protocol_id=d5waku " & + "-d:ffiGenBindings -d:targetLang=rust -d:ffiOutputDir=library/rust_bindings " & + "-d:ffiSrcPath=library/liblogosdelivery.nim " & getMyCPU() & getNimParams() & + " --compileOnly library/liblogosdelivery.nim" + ### Formatting tasks task nphchanges, "Run nph on .nim/.nims/.nimble files changed on this branch/PR": diff --git a/nimble.lock b/nimble.lock index 03bb7a434..cb10f6047 100644 --- a/nimble.lock +++ b/nimble.lock @@ -643,18 +643,19 @@ } }, "ffi": { - "version": "0.1.3", - "vcsRevision": "06111de155253b34e47ed2aaed1d61d08d62cc1b", + "version": "#master", + "vcsRevision": "9ed1fedf96e675ff4fd60c3054b6818dbbd1d54b", "url": "https://github.com/logos-messaging/nim-ffi", "downloadMethod": "git", "dependencies": [ "nim", "chronos", "chronicles", - "taskpools" + "taskpools", + "cbor_serialization" ], "checksums": { - "sha1": "6f9d49375ea1dc71add55c72ac80a808f238e5b0" + "sha1": "4819b421c88c25785e6ab3d9f3c2b3f105b60b17" } }, "boringssl": {