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": {