mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-26 06:23:14 +00:00
Migrate FFI to nim-ffi v0.2.0-rc.3 (typed procs, direct params)
Rewrite the C FFI over the new per-layer APIs using nim-ffi v0.2.0 typed
{.ffiCtor.}/{.ffiDtor.}/{.ffi.}/{.ffiEvent.} + CBOR, replacing the
hand-written cstring/JSON bridge and the declare/lifecycle scaffolding.
Final API shape from the start (no intermediate request wrappers):
- every {.ffi.} proc takes its parameters directly; the macro bundles them
into the per-proc CBOR request, so no `*Request` objects are needed
- WakuMessage rides the wire directly (its fields are CBOR-serializable)
- multi-parameter {.ffiEvent.} via the rc.3 envelope synthesis
- events are fed by internal nim-broker listeners (no AppCallbacks)
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
parent
aca652008a
commit
e59703fad7
@ -1,33 +0,0 @@
|
||||
import ffi
|
||||
import std/locks
|
||||
import logos_delivery
|
||||
|
||||
declareLibrary("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
|
||||
|
||||
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
|
||||
|
||||
# prevent race conditions that might happen due incorrect usage.
|
||||
eventCallbackLock.acquire()
|
||||
defer:
|
||||
eventCallbackLock.release()
|
||||
|
||||
ctx[].eventCallback = cast[pointer](callback)
|
||||
ctx[].eventUserData = userData
|
||||
12
library/events/connection_change_events.nim
Normal file
12
library/events/connection_change_events.nim
Normal file
@ -0,0 +1,12 @@
|
||||
## Per-peer connection changes (connected/disconnected/…), fed by WakuPeerEvent.
|
||||
|
||||
proc onConnectionChange*(
|
||||
peerId: string, event: string
|
||||
) {.ffiEvent: "on_connection_change".}
|
||||
|
||||
proc listenConnectionChangeEvents(self: LogosDelivery) =
|
||||
discard WakuPeerEvent.listen(
|
||||
self.waku.brokerCtx,
|
||||
proc(e: WakuPeerEvent) {.async: (raises: []).} =
|
||||
onConnectionChange($e.peerId, $e.kind),
|
||||
)
|
||||
12
library/events/connection_status_events.nim
Normal file
12
library/events/connection_status_events.nim
Normal file
@ -0,0 +1,12 @@
|
||||
## Node connectivity (online/offline) status, fed by EventConnectionStatusChange.
|
||||
|
||||
proc onConnectionStatusChange*(
|
||||
status: string
|
||||
) {.ffiEvent: "on_connection_status_change".}
|
||||
|
||||
proc listenConnectionStatusEvents(self: LogosDelivery) =
|
||||
discard EventConnectionStatusChange.listen(
|
||||
self.waku.brokerCtx,
|
||||
proc(e: EventConnectionStatusChange) {.async: (raises: []).} =
|
||||
onConnectionStatusChange($e.connectionStatus),
|
||||
)
|
||||
@ -1,6 +0,0 @@
|
||||
type JsonEvent* = ref object of RootObj # https://rfc.vac.dev/spec/36/#jsonsignal-type
|
||||
eventType* {.requiresInit.}: string
|
||||
|
||||
method `$`*(jsonEvent: JsonEvent): string {.base.} =
|
||||
discard
|
||||
# All events should implement this
|
||||
@ -1,17 +0,0 @@
|
||||
import system, std/json, libp2p/[connmanager, peerid]
|
||||
|
||||
import ../../logos_delivery/waku/common/base64, ./json_base_event
|
||||
|
||||
type JsonConnectionChangeEvent* = ref object of JsonEvent
|
||||
peerId*: string
|
||||
peerEvent*: PeerEventKind
|
||||
|
||||
proc new*(
|
||||
T: type JsonConnectionChangeEvent, peerId: string, peerEvent: PeerEventKind
|
||||
): T =
|
||||
return JsonConnectionChangeEvent(
|
||||
eventType: "connection_change", peerId: peerId, peerEvent: peerEvent
|
||||
)
|
||||
|
||||
method `$`*(jsonConnectionChangeEvent: JsonConnectionChangeEvent): string =
|
||||
$(%*jsonConnectionChangeEvent)
|
||||
@ -1,15 +0,0 @@
|
||||
{.push raises: [].}
|
||||
|
||||
import system, std/json
|
||||
import ./json_base_event
|
||||
import ../../logos_delivery/api/types
|
||||
|
||||
type JsonConnectionStatusChangeEvent* = ref object of JsonEvent
|
||||
status*: ConnectionStatus
|
||||
|
||||
proc new*(T: type JsonConnectionStatusChangeEvent, status: ConnectionStatus): T =
|
||||
return
|
||||
JsonConnectionStatusChangeEvent(eventType: "node_health_change", status: status)
|
||||
|
||||
method `$`*(event: JsonConnectionStatusChangeEvent): string =
|
||||
$(%*event)
|
||||
@ -1,106 +0,0 @@
|
||||
import system, results, std/json, std/strutils
|
||||
import stew/byteutils
|
||||
import
|
||||
../../logos_delivery/waku/common/base64,
|
||||
../../logos_delivery/waku/waku_core/message,
|
||||
../../logos_delivery/waku/waku_core/message/message,
|
||||
../utils,
|
||||
./json_base_event
|
||||
|
||||
type JsonMessage* = ref object # https://rfc.vac.dev/spec/36/#jsonmessage-type
|
||||
payload*: Base64String
|
||||
contentTopic*: string
|
||||
version*: uint
|
||||
timestamp*: int64
|
||||
ephemeral*: bool
|
||||
meta*: Base64String
|
||||
proof*: Base64String
|
||||
|
||||
func fromJsonNode*(
|
||||
T: type JsonMessage, jsonContent: JsonNode
|
||||
): Result[JsonMessage, string] =
|
||||
# Visit https://rfc.vac.dev/spec/14/ for further details
|
||||
|
||||
# Check if required fields exist
|
||||
if not jsonContent.hasKey("payload"):
|
||||
return err("Missing required field in WakuMessage: payload")
|
||||
if not jsonContent.hasKey("contentTopic"):
|
||||
return err("Missing required field in WakuMessage: contentTopic")
|
||||
|
||||
ok(
|
||||
JsonMessage(
|
||||
payload: Base64String(jsonContent["payload"].getStr()),
|
||||
contentTopic: jsonContent["contentTopic"].getStr(),
|
||||
version: uint32(jsonContent{"version"}.getInt()),
|
||||
timestamp: (?jsonContent.getProtoInt64("timestamp")).get(0),
|
||||
ephemeral: jsonContent{"ephemeral"}.getBool(),
|
||||
meta: Base64String(jsonContent{"meta"}.getStr()),
|
||||
proof: Base64String(jsonContent{"proof"}.getStr()),
|
||||
)
|
||||
)
|
||||
|
||||
proc toWakuMessage*(self: JsonMessage): Result[WakuMessage, string] =
|
||||
let payload = base64.decode(self.payload).valueOr:
|
||||
return err("invalid payload format: " & error)
|
||||
|
||||
let meta = base64.decode(self.meta).valueOr:
|
||||
return err("invalid meta format: " & error)
|
||||
|
||||
let proof = base64.decode(self.proof).valueOr:
|
||||
return err("invalid proof format: " & error)
|
||||
|
||||
ok(
|
||||
WakuMessage(
|
||||
payload: payload,
|
||||
meta: meta,
|
||||
contentTopic: self.contentTopic,
|
||||
version: uint32(self.version),
|
||||
timestamp: self.timestamp,
|
||||
ephemeral: self.ephemeral,
|
||||
proof: proof,
|
||||
)
|
||||
)
|
||||
|
||||
proc `%`*(value: Base64String): JsonNode =
|
||||
%(value.string)
|
||||
|
||||
type JsonMessageEvent* = ref object of JsonEvent
|
||||
pubsubTopic*: string
|
||||
messageHash*: string
|
||||
wakuMessage*: JsonMessage
|
||||
|
||||
proc new*(T: type JsonMessageEvent, pubSubTopic: string, msg: WakuMessage): T =
|
||||
# Returns a WakuMessage event as indicated in
|
||||
# https://github.com/vacp2p/rfc/blob/master/content/docs/rfcs/36/README.md#jsonmessageevent-type
|
||||
|
||||
var payload = newSeq[byte](len(msg.payload))
|
||||
if len(msg.payload) != 0:
|
||||
copyMem(addr payload[0], unsafeAddr msg.payload[0], len(msg.payload))
|
||||
|
||||
var meta = newSeq[byte](len(msg.meta))
|
||||
if len(msg.meta) != 0:
|
||||
copyMem(addr meta[0], unsafeAddr msg.meta[0], len(msg.meta))
|
||||
|
||||
var proof = newSeq[byte](len(msg.proof))
|
||||
if len(msg.proof) != 0:
|
||||
copyMem(addr proof[0], unsafeAddr msg.proof[0], len(msg.proof))
|
||||
|
||||
let msgHash = computeMessageHash(pubSubTopic, msg)
|
||||
|
||||
return JsonMessageEvent(
|
||||
eventType: "message",
|
||||
pubSubTopic: pubSubTopic,
|
||||
messageHash: msgHash.to0xHex(),
|
||||
wakuMessage: JsonMessage(
|
||||
payload: base64.encode(payload),
|
||||
contentTopic: msg.contentTopic,
|
||||
version: msg.version,
|
||||
timestamp: int64(msg.timestamp),
|
||||
ephemeral: msg.ephemeral,
|
||||
meta: base64.encode(meta),
|
||||
proof: base64.encode(proof),
|
||||
),
|
||||
)
|
||||
|
||||
method `$`*(jsonMessage: JsonMessageEvent): string =
|
||||
$(%*jsonMessage)
|
||||
@ -1,20 +0,0 @@
|
||||
import system, results, std/json
|
||||
import stew/byteutils
|
||||
import ../../logos_delivery/waku/common/base64, ./json_base_event
|
||||
import ../../logos_delivery/waku/waku_relay
|
||||
|
||||
type JsonTopicHealthChangeEvent* = ref object of JsonEvent
|
||||
pubsubTopic*: string
|
||||
topicHealth*: TopicHealth
|
||||
|
||||
proc new*(
|
||||
T: type JsonTopicHealthChangeEvent, pubsubTopic: string, topicHealth: TopicHealth
|
||||
): T =
|
||||
return JsonTopicHealthChangeEvent(
|
||||
eventType: "relay_topic_health_change",
|
||||
pubsubTopic: pubsubTopic,
|
||||
topicHealth: topicHealth,
|
||||
)
|
||||
|
||||
method `$`*(jsonTopicHealthChange: JsonTopicHealthChangeEvent): string =
|
||||
$(%*jsonTopicHealthChange)
|
||||
49
library/events/message_events.nim
Normal file
49
library/events/message_events.nim
Normal file
@ -0,0 +1,49 @@
|
||||
## Message events: send lifecycle (sent/error/propagated/received) plus raw
|
||||
## inbound network messages. Each FFI event is fed by an internal broker event.
|
||||
|
||||
proc onMessageSent*(
|
||||
requestId: string, messageHash: string
|
||||
) {.ffiEvent: "on_message_sent".}
|
||||
|
||||
proc onMessageError*(
|
||||
requestId: string, messageHash: string, error: string
|
||||
) {.ffiEvent: "on_message_error".}
|
||||
|
||||
proc onMessagePropagated*(
|
||||
requestId: string, messageHash: string
|
||||
) {.ffiEvent: "on_message_propagated".}
|
||||
|
||||
proc onMessageReceived*(messageHash: string) {.ffiEvent: "on_message_received".}
|
||||
|
||||
proc onNetworkMessage*(
|
||||
pubsubTopic: string, message: WakuMessage
|
||||
) {.ffiEvent: "on_network_message".}
|
||||
|
||||
proc listenMessageEvents(self: LogosDelivery) =
|
||||
let brokerCtx = self.waku.brokerCtx
|
||||
|
||||
discard MessageSentEvent.listen(
|
||||
brokerCtx,
|
||||
proc(e: MessageSentEvent) {.async: (raises: []).} =
|
||||
onMessageSent($e.requestId, e.messageHash),
|
||||
)
|
||||
discard MessageErrorEvent.listen(
|
||||
brokerCtx,
|
||||
proc(e: MessageErrorEvent) {.async: (raises: []).} =
|
||||
onMessageError($e.requestId, e.messageHash, e.error),
|
||||
)
|
||||
discard MessagePropagatedEvent.listen(
|
||||
brokerCtx,
|
||||
proc(e: MessagePropagatedEvent) {.async: (raises: []).} =
|
||||
onMessagePropagated($e.requestId, e.messageHash),
|
||||
)
|
||||
discard MessageReceivedEvent.listen(
|
||||
brokerCtx,
|
||||
proc(e: MessageReceivedEvent) {.async: (raises: []).} =
|
||||
onMessageReceived(e.messageHash),
|
||||
)
|
||||
discard MessageSeenEvent.listen(
|
||||
brokerCtx,
|
||||
proc(e: MessageSeenEvent) {.async: (raises: []).} =
|
||||
onNetworkMessage(string(e.topic), e.message),
|
||||
)
|
||||
12
library/events/topic_health_events.nim
Normal file
12
library/events/topic_health_events.nim
Normal file
@ -0,0 +1,12 @@
|
||||
## Per-shard (pubsub topic) health changes, fed by EventShardTopicHealthChange.
|
||||
|
||||
proc onTopicHealthChange*(
|
||||
pubsubTopic: string, health: string
|
||||
) {.ffiEvent: "on_topic_health_change".}
|
||||
|
||||
proc listenTopicHealthEvents(self: LogosDelivery) =
|
||||
discard EventShardTopicHealthChange.listen(
|
||||
self.waku.brokerCtx,
|
||||
proc(e: EventShardTopicHealthChange) {.async: (raises: []).} =
|
||||
onTopicHealthChange(string(e.topic), $e.health),
|
||||
)
|
||||
@ -1,46 +1,32 @@
|
||||
import std/strutils
|
||||
import chronos, results, ffi
|
||||
import logos_delivery, library/declare_lib
|
||||
## The waku api getters are synchronous and can't fail, so the bodies just wrap
|
||||
## the value; the `{.ffi.}` macro wraps it into the `Future` it must expose.
|
||||
|
||||
proc waku_version(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
let v = (await ctx.myLib[].waku.version()).valueOr:
|
||||
return err(error)
|
||||
return ok(v)
|
||||
proc version*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok(self.waku.version())
|
||||
|
||||
proc waku_listen_addresses(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
## returns a comma-separated string of the listen addresses
|
||||
let addrs = (await ctx.myLib[].waku.listenAddresses()).valueOr:
|
||||
return err(error)
|
||||
return ok(addrs.join(","))
|
||||
proc listen_addresses*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok(self.waku.listenAddresses().join(","))
|
||||
|
||||
proc waku_get_my_enr(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
let enrUri = (await ctx.myLib[].waku.myEnr()).valueOr:
|
||||
return err(error)
|
||||
return ok(enrUri)
|
||||
proc get_my_enr*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok(self.waku.myEnr())
|
||||
|
||||
proc waku_get_my_peerid(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
let peerId = (await ctx.myLib[].waku.myPeerId()).valueOr:
|
||||
return err(error)
|
||||
return ok(peerId)
|
||||
proc get_my_peerid*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok(self.waku.myPeerId())
|
||||
|
||||
proc waku_get_metrics(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
let m = (await ctx.myLib[].waku.metrics()).valueOr:
|
||||
return err(error)
|
||||
return ok(m)
|
||||
proc get_metrics*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok(self.waku.metrics())
|
||||
|
||||
proc waku_is_online(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
let online = (await ctx.myLib[].waku.isOnline()).valueOr:
|
||||
return err(error)
|
||||
return ok($online)
|
||||
proc is_online*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok($self.waku.isOnline())
|
||||
|
||||
@ -1,59 +1,31 @@
|
||||
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.} =
|
||||
## 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:
|
||||
error "UPDATE_DISCV5_BOOTSTRAP_NODES failed", error = error
|
||||
proc discv5_update_bootnodes*(
|
||||
self: LogosDelivery, bootnodes: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## `bootnodes` is a JSON array of ENRs, e.g. `["enr:...", "enr:..."]`.
|
||||
(await self.waku.discv5UpdateBootnodes(bootnodes)).isOkOr:
|
||||
return err(error)
|
||||
return ok("discovery request processed correctly")
|
||||
return ok("")
|
||||
|
||||
proc waku_dns_discovery(
|
||||
ctx: ptr FFIContext[LogosDelivery],
|
||||
callback: FFICallBack,
|
||||
userData: pointer,
|
||||
enrTreeUrl: cstring,
|
||||
nameDnsServer: cstring,
|
||||
timeoutMs: cint,
|
||||
) {.ffi.} =
|
||||
let nodes = (
|
||||
await ctx.myLib[].waku.dnsDiscovery($enrTreeUrl, $nameDnsServer, int(timeoutMs))
|
||||
).valueOr:
|
||||
error "GET_BOOTSTRAP_NODES failed", error = error
|
||||
proc dns_discovery*(
|
||||
self: LogosDelivery, enrTreeUrl: string, nameDnsServer: string, timeoutMs: int
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let nodes = (await self.waku.dnsDiscovery(enrTreeUrl, nameDnsServer, timeoutMs)).valueOr:
|
||||
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:
|
||||
error "START_DISCV5 failed", error = error
|
||||
proc start_discv5*(self: LogosDelivery): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.waku.startDiscv5()).isOkOr:
|
||||
return err(error)
|
||||
return ok("discv5 started correctly")
|
||||
return ok("")
|
||||
|
||||
proc waku_stop_discv5(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
(await ctx.myLib[].waku.stopDiscv5()).isOkOr:
|
||||
error "STOP_DISCV5 failed", error = error
|
||||
proc stop_discv5*(self: LogosDelivery): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.waku.stopDiscv5()).isOkOr:
|
||||
return err(error)
|
||||
return ok("discv5 stopped correctly")
|
||||
return ok("")
|
||||
|
||||
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:
|
||||
error "waku_peer_exchange_request failed", error = error
|
||||
proc peer_exchange_request*(
|
||||
self: LogosDelivery, numPeers: uint64
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let n = (await self.waku.peerExchangeRequest(numPeers)).valueOr:
|
||||
return err(error)
|
||||
return ok($numValidPeers)
|
||||
return ok($n)
|
||||
|
||||
34
library/kernel_api/node_info_api.nim
Normal file
34
library/kernel_api/node_info_api.nim
Normal file
@ -0,0 +1,34 @@
|
||||
import std/json
|
||||
import logos_delivery/waku/factory/waku_state_info
|
||||
import tools/confutils/[cli_args, config_option_meta]
|
||||
|
||||
proc get_available_node_info_ids*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## All node-info item ids that can be queried with `get_node_info`.
|
||||
return ok($self.waku.stateInfo.getAllPossibleInfoItemIds())
|
||||
|
||||
proc get_node_info*(
|
||||
self: LogosDelivery, nodeInfoId: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let infoItemIdEnum =
|
||||
try:
|
||||
parseEnum[NodeInfoId](nodeInfoId)
|
||||
except ValueError:
|
||||
return err("Invalid node info id: " & nodeInfoId)
|
||||
return ok(self.waku.stateInfo.getNodeInfoItem(infoItemIdEnum))
|
||||
|
||||
proc get_available_configs*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let optionMetas: seq[ConfigOptionMeta] = extractConfigOptionMeta(WakuNodeConf)
|
||||
var configOptionDetails = newJArray()
|
||||
for meta in optionMetas:
|
||||
configOptionDetails.add(
|
||||
%*{
|
||||
meta.fieldName: meta.typeName & "(" & meta.defaultValue & ")", "desc": meta.desc
|
||||
}
|
||||
)
|
||||
var jsonNode = newJObject()
|
||||
jsonNode["configOptions"] = configOptionDetails
|
||||
return ok(pretty(jsonNode))
|
||||
@ -1,82 +0,0 @@
|
||||
import logos_delivery/waku/compat/option_valueor
|
||||
import std/[options, json, strutils, net]
|
||||
import chronos, chronicles, results, confutils, confutils/std/net, ffi
|
||||
|
||||
import
|
||||
logos_delivery/waku/node/peer_manager/peer_manager,
|
||||
tools/confutils/cli_args,
|
||||
logos_delivery/waku/waku,
|
||||
logos_delivery/waku/factory/node_factory,
|
||||
logos_delivery/waku/factory/app_callbacks,
|
||||
logos_delivery/waku/rest_api/endpoint/builder,
|
||||
library/declare_lib
|
||||
|
||||
proc createWaku(
|
||||
configJson: cstring, appCallbacks: AppCallbacks = nil
|
||||
): Future[Result[LogosDelivery, string]] {.async.} =
|
||||
var conf = defaultWakuNodeConf().valueOr:
|
||||
return err("Failed creating node: " & error)
|
||||
|
||||
var errorResp: string
|
||||
|
||||
var jsonNode: JsonNode
|
||||
try:
|
||||
jsonNode = parseJson($configJson)
|
||||
except Exception:
|
||||
return err(
|
||||
"exception in createWaku when calling parseJson: " & getCurrentExceptionMsg() &
|
||||
" configJson string: " & $configJson
|
||||
)
|
||||
|
||||
for confField, confValue in fieldPairs(conf):
|
||||
if jsonNode.contains(confField):
|
||||
# Make sure string doesn't contain the leading or trailing " character
|
||||
let formattedString = ($jsonNode[confField]).strip(chars = {'\"'})
|
||||
# Override conf field with the value set in the json-string
|
||||
try:
|
||||
confValue = parseCmdArg(typeof(confValue), formattedString)
|
||||
except Exception:
|
||||
return err(
|
||||
"exception in createWaku when parsing configuration. exc: " &
|
||||
getCurrentExceptionMsg() & ". string that could not be parsed: " &
|
||||
formattedString & ". expected type: " & $typeof(confValue)
|
||||
)
|
||||
|
||||
# Don't send relay app callbacks if relay is disabled
|
||||
if not conf.relay and not appCallbacks.isNil():
|
||||
appCallbacks.relayHandler = nil
|
||||
appCallbacks.topicHealthChangeHandler = nil
|
||||
|
||||
conf.rest = false ## libwaku never runs the REST server
|
||||
|
||||
let logosRes = (await LogosDelivery.new(conf, appCallbacks)).valueOr:
|
||||
error "LogosDelivery initialization failed", error = error
|
||||
return err("Failed setting up LogosDelivery: " & $error)
|
||||
|
||||
return ok(logosRes)
|
||||
|
||||
registerReqFFI(CreateNodeWithCallbacksRequest, ctx: ptr FFIContext[LogosDelivery]):
|
||||
proc(
|
||||
configJson: cstring, appCallbacks: AppCallbacks
|
||||
): Future[Result[string, string]] {.async.} =
|
||||
ctx.myLib[] = (await createWaku(configJson, cast[AppCallbacks](appCallbacks))).valueOr:
|
||||
error "CreateNodeWithCallbacksRequest failed", error = error
|
||||
return err($error)
|
||||
|
||||
return ok("")
|
||||
|
||||
proc waku_start(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
(await ctx.myLib[].start()).isOkOr:
|
||||
error "START_NODE failed", error = error
|
||||
return err("failed to start: " & $error)
|
||||
return ok("")
|
||||
|
||||
proc waku_stop(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
(await ctx.myLib[].stop()).isOkOr:
|
||||
error "STOP_NODE failed", error = error
|
||||
return err("failed to stop: " & $error)
|
||||
return ok("")
|
||||
@ -1,104 +1,76 @@
|
||||
import std/[strutils, tables, json]
|
||||
import chronicles, chronos, results, ffi
|
||||
import logos_delivery, library/declare_lib
|
||||
import std/sequtils
|
||||
|
||||
type PeerInfo = object
|
||||
protocols: seq[string]
|
||||
addresses: seq[string]
|
||||
type ConnectedPeersInfoResponse {.ffi.} = object
|
||||
peers: seq[PeerConnInfoFFI]
|
||||
|
||||
proc waku_get_peerids_from_peerstore(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
## returns a comma-separated string of peerIDs
|
||||
let peerIds = (await ctx.myLib[].waku.peerIdsFromPeerstore()).valueOr:
|
||||
proc get_peerids_from_peerstore*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let ids = (await self.waku.peerIdsFromPeerstore()).valueOr:
|
||||
return err(error)
|
||||
return ok(peerIds.join(","))
|
||||
return ok(ids.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 connect_peers*(
|
||||
self: LogosDelivery, peers: seq[string], timeoutMs: uint32
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## `peers` are multiaddrs.
|
||||
(await self.waku.connect(peers, 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:
|
||||
error "DISCONNECT_PEER_BY_ID failed", error = error
|
||||
proc disconnect_peer_by_id*(
|
||||
self: LogosDelivery, peerId: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.waku.disconnectPeerById(peerId)).isOkOr:
|
||||
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 disconnect_all_peers*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.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:
|
||||
error "DIAL_PEER failed", error = error
|
||||
proc dial_peer*(
|
||||
self: LogosDelivery, peer: string, protocol: string, timeoutMs: int
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## `peer` is a multiaddr.
|
||||
(await self.waku.dialPeer(peer, protocol, timeoutMs)).isOkOr:
|
||||
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:
|
||||
error "DIAL_PEER_BY_ID failed", error = error
|
||||
proc dial_peer_by_id*(
|
||||
self: LogosDelivery, peer: string, protocol: string, timeoutMs: int
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## `peer` is a peerId.
|
||||
(await self.waku.dialPeerById(peer, protocol, timeoutMs)).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
|
||||
proc waku_get_connected_peers_info(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
## returns a JSON string mapping peerIDs to objects with protocols and addresses
|
||||
let peers = (await ctx.myLib[].waku.connectedPeersInfo()).valueOr:
|
||||
proc get_connected_peers_info*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[ConnectedPeersInfoResponse, string]] {.ffi.} =
|
||||
let infos = (await self.waku.connectedPeersInfo()).valueOr:
|
||||
return err(error)
|
||||
return ok(
|
||||
ConnectedPeersInfoResponse(
|
||||
peers: infos.mapIt(
|
||||
PeerConnInfoFFI(
|
||||
peerId: it.peerId, protocols: it.protocols, addresses: it.addresses
|
||||
)
|
||||
)
|
||||
)
|
||||
)
|
||||
|
||||
var peersMap = initTable[string, PeerInfo]()
|
||||
for peer in peers:
|
||||
peersMap[peer.peerId] =
|
||||
PeerInfo(protocols: peer.protocols, addresses: peer.addresses)
|
||||
|
||||
return ok($(%*peersMap))
|
||||
|
||||
proc waku_get_connected_peers(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
## returns a comma-separated string of peerIDs
|
||||
let peerIds = (await ctx.myLib[].waku.connectedPeers()).valueOr:
|
||||
proc get_connected_peers*(self: LogosDelivery): Future[Result[string, string]] {.ffi.} =
|
||||
let ids = (await self.waku.connectedPeers()).valueOr:
|
||||
return err(error)
|
||||
return ok(peerIds.join(","))
|
||||
return ok(ids.join(","))
|
||||
|
||||
proc waku_get_peerids_by_protocol(
|
||||
ctx: ptr FFIContext[LogosDelivery],
|
||||
callback: FFICallBack,
|
||||
userData: pointer,
|
||||
protocol: cstring,
|
||||
) {.ffi.} =
|
||||
## returns a comma-separated string of peerIDs that mount the given protocol
|
||||
let peerIds = (await ctx.myLib[].waku.peerIdsByProtocol($protocol)).valueOr:
|
||||
proc get_peerids_by_protocol*(
|
||||
self: LogosDelivery, protocol: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let ids = (await self.waku.peerIdsByProtocol(protocol)).valueOr:
|
||||
return err(error)
|
||||
return ok(peerIds.join(","))
|
||||
return ok(ids.join(","))
|
||||
|
||||
@ -1,13 +1,7 @@
|
||||
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 ping_peer*(
|
||||
self: LogosDelivery, peerAddr: string, timeoutMs: int
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## Returns the round-trip time in nanoseconds.
|
||||
let rtt = (await self.waku.pingPeer(peerAddr, timeoutMs)).valueOr:
|
||||
return err(error)
|
||||
return ok($rttNanos)
|
||||
return ok($rtt)
|
||||
|
||||
@ -1,57 +1,39 @@
|
||||
import std/[strutils, sequtils]
|
||||
import chronicles, chronos, results, ffi
|
||||
import
|
||||
logos_delivery,
|
||||
logos_delivery/waku/waku_core/message/message,
|
||||
logos_delivery/waku/waku_core/subscription/push_handler,
|
||||
logos_delivery/waku/waku_core/topics/pubsub_topic,
|
||||
logos_delivery/waku/waku_core/topics/content_topic,
|
||||
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 =
|
||||
return proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} =
|
||||
callEventCallback(ctx, "onReceivedMessage"):
|
||||
$JsonMessageEvent.new(pubsubTopic, msg)
|
||||
import std/sequtils
|
||||
import logos_delivery/waku/waku_core/subscription/push_handler
|
||||
|
||||
proc filter_subscribe*(
|
||||
self: LogosDelivery, pubsubTopic: string, contentTopics: seq[string]
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
# `filterSubscribe` re-registers the filter push handler, so it must keep
|
||||
# feeding MessageSeenEvent — the single source the ctor's listener delivers
|
||||
# to the foreign side (see liblogosdelivery.nim).
|
||||
let brokerCtx = self.waku.brokerCtx
|
||||
let pushHandler = proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} =
|
||||
MessageSeenEvent.emit(brokerCtx, pubsubTopic, msg)
|
||||
(
|
||||
await ctx.myLib[].waku.filterSubscribe(
|
||||
PubsubTopic($pubSubTopic),
|
||||
($contentTopics).split(",").mapIt(ContentTopic(it)),
|
||||
FilterPushHandler(onReceivedMessage(ctx)),
|
||||
await self.waku.filterSubscribe(
|
||||
PubsubTopic(pubsubTopic),
|
||||
contentTopics.mapIt(ContentTopic(it)),
|
||||
FilterPushHandler(pushHandler),
|
||||
)
|
||||
).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 filter_unsubscribe*(
|
||||
self: LogosDelivery, pubsubTopic: string, contentTopics: seq[string]
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(
|
||||
await ctx.myLib[].waku.filterUnsubscribe(
|
||||
PubsubTopic($pubSubTopic), ($contentTopics).split(",").mapIt(ContentTopic(it))
|
||||
await self.waku.filterUnsubscribe(
|
||||
PubsubTopic(pubsubTopic), contentTopics.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:
|
||||
error "fail filter unsubscribe all", error = error
|
||||
proc filter_unsubscribe_all*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.waku.filterUnsubscribeAll()).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
|
||||
@ -1,34 +1,7 @@
|
||||
import std/[json, strformat]
|
||||
import chronicles, chronos, results, ffi
|
||||
import
|
||||
logos_delivery,
|
||||
logos_delivery/waku/waku_core/message,
|
||||
logos_delivery/waku/waku_core/topics/pubsub_topic,
|
||||
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.} =
|
||||
var jsonMessage: JsonMessage
|
||||
try:
|
||||
let jsonContent = parseJson($jsonWakuMessage)
|
||||
jsonMessage = JsonMessage.fromJsonNode(jsonContent).valueOr:
|
||||
raise newException(JsonParsingError, $error)
|
||||
except JsonParsingError as exc:
|
||||
return err(fmt"Error parsing json message: {exc.msg}")
|
||||
|
||||
let msg = json_message_event.toWakuMessage(jsonMessage).valueOr:
|
||||
return err("Problem building the WakuMessage: " & $error)
|
||||
|
||||
let msgHashHex = (
|
||||
await ctx.myLib[].waku.lightpushPublish(PubsubTopic($pubSubTopic), msg)
|
||||
).valueOr:
|
||||
error "PUBLISH failed", error = error
|
||||
proc lightpush_publish*(
|
||||
self: LogosDelivery, pubsubTopic: string, message: WakuMessage
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## Returns the published message hash.
|
||||
let hash = (await self.waku.lightpushPublish(PubsubTopic(pubsubTopic), message)).valueOr:
|
||||
return err(error)
|
||||
|
||||
return ok(msgHashHex)
|
||||
return ok(hash)
|
||||
|
||||
@ -1,163 +1,85 @@
|
||||
import std/[strutils, json]
|
||||
import chronicles, chronos, results, ffi
|
||||
import
|
||||
logos_delivery,
|
||||
logos_delivery/waku/waku_core/topics/pubsub_topic,
|
||||
logos_delivery/waku/waku_core/message,
|
||||
logos_delivery/waku/waku_relay/protocol,
|
||||
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:
|
||||
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:
|
||||
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.} =
|
||||
## Returns the list of all connected peers to an specific pubsub topic
|
||||
let peers = (await ctx.myLib[].waku.relayConnectedPeers(PubsubTopic($pubSubTopic))).valueOr:
|
||||
error "LIST_CONNECTED_PEERS failed", error = error
|
||||
proc relay_get_peers_in_mesh*(
|
||||
self: LogosDelivery, pubsubTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let peers = (await self.waku.relayPeersInMesh(PubsubTopic(pubsubTopic))).valueOr:
|
||||
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:
|
||||
error "NUM_CONNECTED_PEERS failed", error = error
|
||||
proc relay_get_num_peers_in_mesh*(
|
||||
self: LogosDelivery, pubsubTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let n = (await self.waku.relayNumPeersInMesh(PubsubTopic(pubsubTopic))).valueOr:
|
||||
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.} =
|
||||
## Protects a shard with a public key
|
||||
(
|
||||
await ctx.myLib[].waku.relayAddProtectedShard(
|
||||
uint16(clusterId), uint16(shardId), $publicKey
|
||||
)
|
||||
).isOkOr:
|
||||
proc relay_get_connected_peers*(
|
||||
self: LogosDelivery, pubsubTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let peers = (await self.waku.relayConnectedPeers(PubsubTopic(pubsubTopic))).valueOr:
|
||||
return err(error)
|
||||
return ok(peers.join(","))
|
||||
|
||||
proc relay_get_num_connected_peers*(
|
||||
self: LogosDelivery, pubsubTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let n = (await self.waku.relayNumConnectedPeers(PubsubTopic(pubsubTopic))).valueOr:
|
||||
return err(error)
|
||||
return ok($n)
|
||||
|
||||
proc relay_add_protected_shard*(
|
||||
self: LogosDelivery, clusterId: uint16, shardId: uint16, publicKey: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.waku.relayAddProtectedShard(clusterId, 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 =
|
||||
return proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} =
|
||||
callEventCallback(ctx, "onReceivedMessage"):
|
||||
$JsonMessageEvent.new(pubsubTopic, msg)
|
||||
|
||||
(
|
||||
await ctx.myLib[].waku.relaySubscribe(
|
||||
PubsubTopic($pubSubTopic), WakuRelayHandler(onReceivedMessage(ctx))
|
||||
)
|
||||
).isOkOr:
|
||||
error "SUBSCRIBE failed", error = error
|
||||
proc relay_subscribe*(
|
||||
self: LogosDelivery, pubsubTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
# Just establishes the subscription; delivery flows through the global
|
||||
# MessageSeenEvent listener (see the ctor in liblogosdelivery.nim).
|
||||
(await self.waku.relaySubscribe(PubsubTopic(pubsubTopic))).isOkOr:
|
||||
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:
|
||||
error "UNSUBSCRIBE failed", error = error
|
||||
proc relay_unsubscribe*(
|
||||
self: LogosDelivery, pubsubTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.waku.relayUnsubscribe(PubsubTopic(pubsubTopic))).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
|
||||
proc waku_relay_publish(
|
||||
ctx: ptr FFIContext[LogosDelivery],
|
||||
callback: FFICallBack,
|
||||
userData: pointer,
|
||||
pubSubTopic: cstring,
|
||||
jsonWakuMessage: cstring,
|
||||
timeoutMs: cuint,
|
||||
) {.ffi.} =
|
||||
var jsonMessage: JsonMessage
|
||||
try:
|
||||
let jsonContent = parseJson($jsonWakuMessage)
|
||||
jsonMessage = JsonMessage.fromJsonNode(jsonContent).valueOr:
|
||||
raise newException(JsonParsingError, $error)
|
||||
except JsonParsingError as exc:
|
||||
return err("Error parsing json message: " & exc.msg)
|
||||
|
||||
let msg = json_message_event.toWakuMessage(jsonMessage).valueOr:
|
||||
return err("Problem building the WakuMessage: " & $error)
|
||||
|
||||
let msgHash = (
|
||||
await ctx.myLib[].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:
|
||||
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.} =
|
||||
let topic = (
|
||||
await ctx.myLib[].waku.buildContentTopic(
|
||||
$appName, uint32(appVersion), $contentTopicName, $encoding
|
||||
)
|
||||
proc relay_publish*(
|
||||
self: LogosDelivery, pubsubTopic: string, message: WakuMessage, timeoutMs: uint32
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
## Returns the published message hash (0x-hex).
|
||||
let hash = (
|
||||
await self.waku.relayPublish(PubsubTopic(pubsubTopic), message, timeoutMs)
|
||||
).valueOr:
|
||||
return err(error)
|
||||
return ok(string(topic))
|
||||
return ok(hash)
|
||||
|
||||
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 relay_default_pubsub_topic*(
|
||||
self: LogosDelivery
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
return ok(string(self.waku.defaultPubsubTopic()))
|
||||
|
||||
proc relay_content_topic*(
|
||||
self: LogosDelivery,
|
||||
appName: string,
|
||||
appVersion: uint32,
|
||||
contentTopicName: string,
|
||||
encoding: string,
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let contentTopic = self.waku.buildContentTopic(
|
||||
appName, appVersion, contentTopicName, encoding
|
||||
).valueOr:
|
||||
return err(error)
|
||||
return ok(string(topic))
|
||||
return ok(string(contentTopic))
|
||||
|
||||
proc relay_pubsub_topic*(
|
||||
self: LogosDelivery, topicName: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let pubsubTopic = self.waku.buildPubsubTopic(topicName).valueOr:
|
||||
return err(error)
|
||||
return ok(string(pubsubTopic))
|
||||
|
||||
@ -1,14 +1,12 @@
|
||||
## The query/response are complex types, so this keeps the JSON bridge: the
|
||||
## request carries the query as a JSON string, the response is returned as JSON.
|
||||
import std/[json, sugar, options]
|
||||
import chronos, chronicles, results, ffi
|
||||
import
|
||||
logos_delivery,
|
||||
library/utils,
|
||||
logos_delivery/waku/waku_core/message/digest,
|
||||
logos_delivery/waku/waku_store/common,
|
||||
logos_delivery/waku/common/paging,
|
||||
library/declare_lib
|
||||
import logos_delivery/waku/waku_core/message/digest
|
||||
import logos_delivery/waku/waku_store/common
|
||||
import logos_delivery/waku/common/paging
|
||||
import library/utils
|
||||
|
||||
func fromJsonNode(jsonContent: JsonNode): Result[StoreQueryRequest, string] =
|
||||
func storeQueryFromJson(jsonContent: JsonNode): Result[StoreQueryRequest, string] =
|
||||
var contentTopics: seq[string]
|
||||
if jsonContent.contains("contentTopics"):
|
||||
contentTopics = collect(newSeq):
|
||||
@ -64,26 +62,19 @@ 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.} =
|
||||
let jsonContentRes = catch:
|
||||
parseJson($jsonQuery)
|
||||
proc store_query*(
|
||||
self: LogosDelivery, queryJson: string, peer: string, timeoutMs: int
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let jsonContent =
|
||||
try:
|
||||
parseJson(queryJson)
|
||||
except CatchableError as e:
|
||||
return err("StoreRequest failed parsing store request: " & e.msg)
|
||||
|
||||
if jsonContentRes.isErr():
|
||||
return err("StoreRequest failed parsing store request: " & jsonContentRes.error.msg)
|
||||
let storeQueryRequest = storeQueryFromJson(jsonContent).valueOr:
|
||||
return err(error)
|
||||
|
||||
let storeQueryRequest = ?fromJsonNode(jsonContentRes.get())
|
||||
|
||||
let queryResponse = (
|
||||
await ctx.myLib[].waku.storeQuery(storeQueryRequest, $peerAddr, int(timeoutMs))
|
||||
).valueOr:
|
||||
let queryResponse = (await self.waku.storeQuery(storeQueryRequest, peer, timeoutMs)).valueOr:
|
||||
return err("StoreRequest failed store query: " & error)
|
||||
|
||||
let res = $(%*(queryResponse.toHex()))
|
||||
return ok(res) ## returning the response in json format
|
||||
return ok($(%*(queryResponse.toHex()))) ## response in json format
|
||||
|
||||
@ -1,113 +1,87 @@
|
||||
import logos_delivery/waku/compat/option_valueor
|
||||
import std/[atomics, options, macros]
|
||||
import chronicles, chronos, chronos/threadsync, ffi
|
||||
import
|
||||
logos_delivery/waku/waku_core/message/message,
|
||||
logos_delivery/waku/waku_core/topics/pubsub_topic,
|
||||
logos_delivery/waku/waku_relay,
|
||||
logos_delivery,
|
||||
logos_delivery/waku/waku,
|
||||
logos_delivery/waku/node/waku_node,
|
||||
logos_delivery/waku/node/health_monitor/health_status,
|
||||
../logos_delivery/waku/factory/app_callbacks,
|
||||
./events/json_message_event,
|
||||
./events/json_topic_health_change_event,
|
||||
./events/json_connection_change_event,
|
||||
./events/json_connection_status_change_event,
|
||||
./declare_lib
|
||||
## C FFI library root (nim-ffi v0.2.0).
|
||||
##
|
||||
## The FFI context owns one `LogosDelivery` (the per-layer concentrator). The
|
||||
## v0.2.0 framework generates the C ABI, CBOR (de)serialization and the request
|
||||
## channel from the `{.ffiCtor.}` / `{.ffiDtor.}` / `{.ffi.}` / `{.ffiEvent.}`
|
||||
## annotations below and in the included api modules; `genBindings()` (last
|
||||
## call) emits the foreign-language bindings under `-d:ffiGenBindings`.
|
||||
import ffi
|
||||
import std/strutils
|
||||
import chronos, results, chronicles
|
||||
|
||||
################################################################################
|
||||
## Include different APIs, i.e. all procs with {.ffi.} pragma
|
||||
import logos_delivery
|
||||
import logos_delivery/api/types
|
||||
import tools/confutils/conf_from_json
|
||||
import logos_delivery/waku/api/events/peer_events
|
||||
import logos_delivery/waku/waku_core
|
||||
|
||||
declareLibrary("logosdelivery", LogosDelivery, defaultABIFormat = "cbor")
|
||||
|
||||
# --- shared wire types -----------------------------------------------------
|
||||
type PeerConnInfoFFI* {.ffi.} = object
|
||||
peerId: string
|
||||
protocols: seq[string]
|
||||
addresses: seq[string]
|
||||
|
||||
# --- library-initiated events (one {.ffi.} type-set + listener per file) -----
|
||||
include
|
||||
./logos_delivery_api/node_api,
|
||||
./logos_delivery_api/messaging_api,
|
||||
./logos_delivery_api/debug_api,
|
||||
./kernel_api/peer_manager_api,
|
||||
./kernel_api/discovery_api,
|
||||
./kernel_api/node_lifecycle_api,
|
||||
./events/message_events,
|
||||
./events/connection_status_events,
|
||||
./events/topic_health_events,
|
||||
./events/connection_change_events
|
||||
|
||||
proc listenInternalEvents(self: LogosDelivery) =
|
||||
## Feed every FFI event from an internal nim-broker event.
|
||||
## Listener handles are discarded on purpose: the listeners live for the node's lifetime.
|
||||
self.listenMessageEvents()
|
||||
self.listenConnectionStatusEvents()
|
||||
self.listenTopicHealthEvents()
|
||||
self.listenConnectionChangeEvents()
|
||||
|
||||
# --- constructor / destructor ----------------------------------------------
|
||||
proc logosdelivery_create*(
|
||||
configJson: string
|
||||
): Future[Result[LogosDelivery, string]] {.ffiCtor.} =
|
||||
let conf = parseNodeConfFromJson(configJson).valueOr:
|
||||
return err("failed to parse node config: " & error)
|
||||
|
||||
let logos = (await LogosDelivery.new(conf)).valueOr:
|
||||
return err("failed to create LogosDelivery: " & error)
|
||||
|
||||
logos.listenInternalEvents()
|
||||
|
||||
return ok(logos)
|
||||
|
||||
proc logosdelivery_destroy*(self: LogosDelivery) {.ffiDtor.} =
|
||||
## The framework drains the FFI thread and frees the context; callers stop the
|
||||
## node via `logosdelivery_stop` first.
|
||||
discard
|
||||
|
||||
# --- lifecycle -------------------------------------------------------------
|
||||
proc start*(self: LogosDelivery): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.start()).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
|
||||
proc stop*(self: LogosDelivery): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.stop()).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
|
||||
# --- operations (typed {.ffi.} procs, grouped per protocol) ----------------
|
||||
include
|
||||
./messaging_api/subscriptions_api,
|
||||
./messaging_api/send_api,
|
||||
./kernel_api/node_info_api,
|
||||
./kernel_api/debug_node_api,
|
||||
./kernel_api/ping_api,
|
||||
./kernel_api/peer_manager_api,
|
||||
./kernel_api/discovery_api,
|
||||
./kernel_api/protocols/relay_api,
|
||||
./kernel_api/protocols/store_api,
|
||||
./kernel_api/protocols/lightpush_api,
|
||||
./kernel_api/protocols/store_api,
|
||||
./kernel_api/protocols/filter_api
|
||||
|
||||
################################################################################
|
||||
### Exported procs (former libwaku API)
|
||||
|
||||
proc waku_new(
|
||||
configJson: cstring, callback: FFICallback, userData: pointer
|
||||
): pointer {.dynlib, exportc, cdecl.} =
|
||||
initializeLibrary()
|
||||
|
||||
## Creates a new instance of the WakuNode.
|
||||
if isNil(callback):
|
||||
echo "error: missing callback in waku_new"
|
||||
return nil
|
||||
|
||||
## Create the Waku thread that will keep waiting for req from the main thread.
|
||||
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
|
||||
|
||||
proc onReceivedMessage(ctx: ptr FFIContext): WakuRelayHandler =
|
||||
return proc(pubsubTopic: PubsubTopic, msg: WakuMessage) {.async.} =
|
||||
callEventCallback(ctx, "onReceivedMessage"):
|
||||
$JsonMessageEvent.new(pubsubTopic, msg)
|
||||
|
||||
proc onTopicHealthChange(ctx: ptr FFIContext): TopicHealthChangeHandler =
|
||||
return proc(pubsubTopic: PubsubTopic, topicHealth: TopicHealth) {.async.} =
|
||||
callEventCallback(ctx, "onTopicHealthChange"):
|
||||
$JsonTopicHealthChangeEvent.new(pubsubTopic, topicHealth)
|
||||
|
||||
proc onConnectionChange(ctx: ptr FFIContext): ConnectionChangeHandler =
|
||||
return proc(peerId: PeerId, peerEvent: PeerEventKind) {.async.} =
|
||||
callEventCallback(ctx, "onConnectionChange"):
|
||||
$JsonConnectionChangeEvent.new($peerId, peerEvent)
|
||||
|
||||
proc onConnectionStatusChange(ctx: ptr FFIContext): ConnectionStatusChangeHandler =
|
||||
return proc(status: ConnectionStatus) {.async.} =
|
||||
callEventCallback(ctx, "onConnectionStatusChange"):
|
||||
$JsonConnectionStatusChangeEvent.new(status)
|
||||
|
||||
let appCallbacks = AppCallbacks(
|
||||
relayHandler: onReceivedMessage(ctx),
|
||||
topicHealthChangeHandler: onTopicHealthChange(ctx),
|
||||
connectionChangeHandler: onConnectionChange(ctx),
|
||||
connectionStatusChangeHandler: onConnectionStatusChange(ctx),
|
||||
)
|
||||
|
||||
ffi.sendRequestToFFIThread(
|
||||
ctx,
|
||||
CreateNodeWithCallbacksRequest.ffiNewReq(
|
||||
callback, userData, configJson, appCallbacks
|
||||
),
|
||||
).isOkOr:
|
||||
let msg = "error in sendRequestToFFIThread: " & $error
|
||||
callback(RET_ERR, unsafeAddr msg[0], cast[csize_t](len(msg)), userData)
|
||||
return nil
|
||||
|
||||
return ctx
|
||||
|
||||
proc waku_destroy(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
): cint {.dynlib, exportc, cdecl.} =
|
||||
initializeLibrary()
|
||||
checkParams(ctx, callback, userData)
|
||||
|
||||
ffi.destroyFFIContext(ctx).isOkOr:
|
||||
let msg = "libwaku 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
|
||||
|
||||
# ### End of exported procs
|
||||
# ################################################################################
|
||||
# genBindings() MUST be the last top-level call — after every {.ffi.},
|
||||
# {.ffiCtor.}, {.ffiDtor.} and {.ffiEvent.} pragma (incl. the included files).
|
||||
genBindings()
|
||||
|
||||
@ -1,56 +0,0 @@
|
||||
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.} =
|
||||
## 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($ctx.myLib[].waku.stateInfo.getAllPossibleInfoItemIds())
|
||||
|
||||
proc logosdelivery_get_node_info(
|
||||
ctx: ptr FFIContext[LogosDelivery],
|
||||
callback: FFICallBack,
|
||||
userData: pointer,
|
||||
nodeInfoId: cstring,
|
||||
) {.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)
|
||||
except ValueError:
|
||||
return err("Invalid node info id: " & $nodeInfoId)
|
||||
|
||||
return ok(ctx.myLib[].waku.stateInfo.getNodeInfoItem(infoItemIdEnum))
|
||||
|
||||
proc logosdelivery_get_available_configs(
|
||||
ctx: ptr FFIContext[LogosDelivery], callback: FFICallBack, userData: pointer
|
||||
) {.ffi.} =
|
||||
## Returns information about the accepted config items.
|
||||
requireInitializedNode(ctx, "GetAvailableConfigs"):
|
||||
return err(errMsg)
|
||||
|
||||
let optionMetas: seq[ConfigOptionMeta] = extractConfigOptionMeta(WakuNodeConf)
|
||||
var configOptionDetails = newJArray()
|
||||
|
||||
# for confField, confValue in fieldPairs(conf):
|
||||
# defaultConfig[confField] = $confValue
|
||||
|
||||
for meta in optionMetas:
|
||||
configOptionDetails.add(
|
||||
%*{
|
||||
meta.fieldName: meta.typeName & "(" & meta.defaultValue & ")", "desc": meta.desc
|
||||
}
|
||||
)
|
||||
|
||||
var jsonNode = newJObject()
|
||||
jsonNode["configOptions"] = configOptionDetails
|
||||
let asString = pretty(jsonNode)
|
||||
return ok(pretty(jsonNode))
|
||||
@ -1,91 +0,0 @@
|
||||
import std/[json]
|
||||
import chronos, results, ffi
|
||||
import stew/byteutils
|
||||
import
|
||||
logos_delivery/waku/common/base64,
|
||||
logos_delivery/waku/waku,
|
||||
logos_delivery/waku/waku_core/topics/content_topic,
|
||||
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)
|
||||
|
||||
# ContentTopic is just a string type alias
|
||||
let contentTopic = ContentTopic($contentTopicStr)
|
||||
|
||||
(await ctx.myLib[].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)
|
||||
|
||||
# ContentTopic is just a string type alias
|
||||
let contentTopic = ContentTopic($contentTopicStr)
|
||||
|
||||
ctx.myLib[].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)
|
||||
|
||||
## Parse the message JSON and send the message
|
||||
var jsonNode: JsonNode
|
||||
try:
|
||||
jsonNode = parseJson($messageJson)
|
||||
except Exception as e:
|
||||
return err("Failed to parse message JSON: " & e.msg)
|
||||
|
||||
# Extract content topic
|
||||
if not jsonNode.hasKey("contentTopic"):
|
||||
return err("Missing contentTopic field")
|
||||
|
||||
# ContentTopic is just a string type alias
|
||||
let contentTopic = ContentTopic(jsonNode["contentTopic"].getStr())
|
||||
|
||||
# Extract payload (expect base64 encoded string)
|
||||
if not jsonNode.hasKey("payload"):
|
||||
return err("Missing payload field")
|
||||
|
||||
let payloadStr = jsonNode["payload"].getStr()
|
||||
let payload = base64.decode(Base64String(payloadStr)).valueOr:
|
||||
return err("invalid payload format: " & error)
|
||||
|
||||
# Extract ephemeral flag
|
||||
let ephemeral = jsonNode.getOrDefault("ephemeral").getBool(false)
|
||||
|
||||
# Create message envelope
|
||||
let envelope = MessageEnvelope.init(
|
||||
contentTopic = contentTopic, payload = payload, ephemeral = ephemeral
|
||||
)
|
||||
|
||||
# Send the message via the messaging layer's own API.
|
||||
let requestId = (await ctx.myLib[].messagingClient.send(envelope)).valueOr:
|
||||
let errMsg = $error
|
||||
return err("Send failed: " & errMsg)
|
||||
|
||||
return ok($requestId)
|
||||
@ -1,149 +0,0 @@
|
||||
import std/json
|
||||
import chronos, chronicles, results, ffi
|
||||
import
|
||||
logos_delivery,
|
||||
logos_delivery/waku/node/waku_node,
|
||||
logos_delivery/api/types,
|
||||
logos_delivery/waku/api/events/health_events,
|
||||
tools/confutils/conf_from_json,
|
||||
../declare_lib,
|
||||
../json_event
|
||||
|
||||
# Add JSON serialization for RequestId
|
||||
proc `%`*(id: RequestId): JsonNode =
|
||||
%($id)
|
||||
|
||||
registerReqFFI(CreateNodeRequest, ctx: ptr FFIContext[LogosDelivery]):
|
||||
proc(configJson: cstring): Future[Result[string, string]] {.async.} =
|
||||
let conf = parseNodeConfFromJson($configJson).valueOr:
|
||||
error "Failed to assemble WakuNodeConf from JSON",
|
||||
error = error, configJson = $configJson
|
||||
return err("failed parseNodeConfFromJson " & 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)
|
||||
|
||||
# setting up outgoing event listeners
|
||||
let sentListener = MessageSentEvent.listen(
|
||||
ctx.myLib[].waku.brokerCtx,
|
||||
proc(event: MessageSentEvent) {.async: (raises: []).} =
|
||||
callEventCallback(ctx, "onMessageSent"):
|
||||
$newJsonEvent("message_sent", event),
|
||||
).valueOr:
|
||||
chronicles.error "MessageSentEvent.listen failed", err = $error
|
||||
return err("MessageSentEvent.listen failed: " & $error)
|
||||
|
||||
let errorListener = MessageErrorEvent.listen(
|
||||
ctx.myLib[].waku.brokerCtx,
|
||||
proc(event: MessageErrorEvent) {.async: (raises: []).} =
|
||||
callEventCallback(ctx, "onMessageError"):
|
||||
$newJsonEvent("message_error", event),
|
||||
).valueOr:
|
||||
chronicles.error "MessageErrorEvent.listen failed", err = $error
|
||||
return err("MessageErrorEvent.listen failed: " & $error)
|
||||
|
||||
let propagatedListener = MessagePropagatedEvent.listen(
|
||||
ctx.myLib[].waku.brokerCtx,
|
||||
proc(event: MessagePropagatedEvent) {.async: (raises: []).} =
|
||||
callEventCallback(ctx, "onMessagePropagated"):
|
||||
$newJsonEvent("message_propagated", event),
|
||||
).valueOr:
|
||||
chronicles.error "MessagePropagatedEvent.listen failed", err = $error
|
||||
return err("MessagePropagatedEvent.listen failed: " & $error)
|
||||
|
||||
let receivedListener = MessageReceivedEvent.listen(
|
||||
ctx.myLib[].waku.brokerCtx,
|
||||
proc(event: MessageReceivedEvent) {.async: (raises: []).} =
|
||||
callEventCallback(ctx, "onMessageReceived"):
|
||||
$newJsonEvent("message_received", event),
|
||||
).valueOr:
|
||||
chronicles.error "MessageReceivedEvent.listen failed", err = $error
|
||||
return err("MessageReceivedEvent.listen failed: " & $error)
|
||||
|
||||
let ConnectionStatusChangeListener = EventConnectionStatusChange.listen(
|
||||
ctx.myLib[].waku.brokerCtx,
|
||||
proc(event: EventConnectionStatusChange) {.async: (raises: []).} =
|
||||
callEventCallback(ctx, "onConnectionStatusChange"):
|
||||
$newJsonEvent("connection_status_change", event),
|
||||
).valueOr:
|
||||
chronicles.error "ConnectionStatusChange.listen failed", err = $error
|
||||
return err("ConnectionStatusChange.listen failed: " & $error)
|
||||
|
||||
(await ctx.myLib[].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)
|
||||
|
||||
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 ctx.myLib[].stop()).isOkOr:
|
||||
let errMsg = $error
|
||||
chronicles.error "STOP_NODE failed", err = errMsg
|
||||
return err("failed to stop: " & errMsg)
|
||||
return ok("")
|
||||
9
library/messaging_api/send_api.nim
Normal file
9
library/messaging_api/send_api.nim
Normal file
@ -0,0 +1,9 @@
|
||||
proc messaging_send*(
|
||||
self: LogosDelivery, contentTopic: string, payload: seq[byte], ephemeral: bool
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
let envelope = MessageEnvelope.init(
|
||||
contentTopic = ContentTopic(contentTopic), payload = payload, ephemeral = ephemeral
|
||||
)
|
||||
let requestId = (await self.messagingClient.send(envelope)).valueOr:
|
||||
return err(error)
|
||||
return ok($requestId)
|
||||
13
library/messaging_api/subscriptions_api.nim
Normal file
13
library/messaging_api/subscriptions_api.nim
Normal file
@ -0,0 +1,13 @@
|
||||
proc subscribe*(
|
||||
self: LogosDelivery, contentTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
(await self.messagingClient.subscribe(ContentTopic(contentTopic))).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
|
||||
proc unsubscribe*(
|
||||
self: LogosDelivery, contentTopic: string
|
||||
): Future[Result[string, string]] {.ffi.} =
|
||||
self.messagingClient.unsubscribe(ContentTopic(contentTopic)).isOkOr:
|
||||
return err(error)
|
||||
return ok("")
|
||||
@ -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#v0.2.0-rc.3"
|
||||
|
||||
requires "https://github.com/logos-messaging/nim-sds.git#b12f5ee07c5b764303b51fb948b32a4ade1de3b5"
|
||||
|
||||
|
||||
@ -1,36 +1,24 @@
|
||||
## Waku layer API — debug / info operations.
|
||||
## Waku layer API — debug / info getters (all synchronous).
|
||||
{.push raises: [].}
|
||||
|
||||
import results, chronos, chronicles, metrics
|
||||
import metrics
|
||||
import eth/p2p/discoveryv5/enr
|
||||
|
||||
import logos_delivery/waku/waku
|
||||
import logos_delivery/waku/[waku_core, node/waku_node]
|
||||
import logos_delivery/waku/node/waku_node
|
||||
|
||||
proc version*(self: Waku): Future[Result[string, string]] {.async.} =
|
||||
return ok(WakuNodeVersionString)
|
||||
proc version*(self: Waku): string =
|
||||
return WakuNodeVersionString
|
||||
|
||||
proc listenAddresses*(self: Waku): Future[Result[seq[string], string]] {.async.} =
|
||||
try:
|
||||
return ok(self.node.info().listenAddresses)
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
proc listenAddresses*(self: Waku): seq[string] =
|
||||
return self.node.info().listenAddresses
|
||||
|
||||
proc myEnr*(self: Waku): Future[Result[string, string]] {.async.} =
|
||||
try:
|
||||
return ok(self.node.enr.toURI())
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
proc myEnr*(self: Waku): string =
|
||||
return self.node.enr.toURI()
|
||||
|
||||
proc myPeerId*(self: Waku): Future[Result[string, string]] {.async.} =
|
||||
try:
|
||||
return ok($self.node.peerId())
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
proc myPeerId*(self: Waku): string =
|
||||
return $self.node.peerId()
|
||||
|
||||
proc metrics*(self: Waku): Future[Result[string, string]] {.async.} =
|
||||
proc metrics*(self: Waku): string =
|
||||
{.gcsafe.}:
|
||||
try:
|
||||
return ok(defaultRegistry.toText())
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
return defaultRegistry.toText()
|
||||
|
||||
@ -6,8 +6,5 @@ import results, chronos, chronicles
|
||||
import logos_delivery/waku/waku
|
||||
import logos_delivery/waku/[node/health_monitor, node/health_monitor/online_monitor]
|
||||
|
||||
proc isOnline*(self: Waku): Future[Result[bool, string]] {.async.} =
|
||||
try:
|
||||
return ok(self.healthMonitor.onlineMonitor.amIOnline())
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
proc isOnline*(self: Waku): bool =
|
||||
return self.healthMonitor.onlineMonitor.amIOnline()
|
||||
|
||||
@ -2,26 +2,24 @@
|
||||
{.push raises: [].}
|
||||
|
||||
import std/strformat
|
||||
import results, chronos
|
||||
import results
|
||||
|
||||
import logos_delivery/waku/waku
|
||||
import logos_delivery/waku/waku_core
|
||||
|
||||
proc buildContentTopic*(
|
||||
self: Waku, appName: string, appVersion: uint32, name: string, encoding: string
|
||||
): Future[Result[ContentTopic, string]] {.async.} =
|
||||
): Result[ContentTopic, string] =
|
||||
try:
|
||||
return ok(ContentTopic(fmt"/{appName}/{appVersion}/{name}/{encoding}"))
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
|
||||
proc buildPubsubTopic*(
|
||||
self: Waku, topicName: string
|
||||
): Future[Result[PubsubTopic, string]] {.async.} =
|
||||
proc buildPubsubTopic*(self: Waku, topicName: string): Result[PubsubTopic, string] =
|
||||
try:
|
||||
return ok(PubsubTopic(fmt"/waku/2/{topicName}"))
|
||||
except CatchableError as e:
|
||||
return err(e.msg)
|
||||
|
||||
proc defaultPubsubTopic*(self: Waku): Future[Result[PubsubTopic, string]] {.async.} =
|
||||
return ok(DefaultPubsubTopic)
|
||||
proc defaultPubsubTopic*(self: Waku): PubsubTopic =
|
||||
return DefaultPubsubTopic
|
||||
|
||||
@ -643,18 +643,19 @@
|
||||
}
|
||||
},
|
||||
"ffi": {
|
||||
"version": "0.1.3",
|
||||
"vcsRevision": "06111de155253b34e47ed2aaed1d61d08d62cc1b",
|
||||
"version": "#v0.2.0-rc.3",
|
||||
"vcsRevision": "8f15afce5c377a0e5ee53c35b228025b903604ea",
|
||||
"url": "https://github.com/logos-messaging/nim-ffi",
|
||||
"downloadMethod": "git",
|
||||
"dependencies": [
|
||||
"nim",
|
||||
"chronos",
|
||||
"chronicles",
|
||||
"taskpools"
|
||||
"taskpools",
|
||||
"cbor_serialization"
|
||||
],
|
||||
"checksums": {
|
||||
"sha1": "6f9d49375ea1dc71add55c72ac80a808f238e5b0"
|
||||
"sha1": "7f00eaaa01ce59a0c1603e6fb8757ba712f9a53e"
|
||||
}
|
||||
},
|
||||
"boringssl": {
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user