From 9827be59905bb86abae59705202b5b16816bf805 Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Fri, 17 Jul 2026 18:14:52 +0200 Subject: [PATCH] feat: logos_delivery_node app + messaging REST API with event observability (#4014) * WIP logosdeliverynode app initial commit * WIP - extra cli option * WIP: messaging client REST endpoints * Add event poll for messaging rest with cache mechanism * Messaging rest tests * test: assert 404 via raw string client in messaging REST test presto's typed REST client raises RestDecodingError when it cannot decode a non-2xx text error body into the response type. Add a RestResponse[string] stub (messagingGetSendEventsByIdRawV1) and point the "already-polled id -> 404" assertion at it, matching the relay REST test pattern. Co-Authored-By: Claude Opus 4.8 * remove customized cli args as confutils has no support for it * Introduce --entry-layer and re-introduce --mode flags into cli args, applied new driver into LogosDelivery + tests * Add messaging REST client test * Add docker image build of logosdeliverynode for CI builds * Fix tests * Potential fix for pull request finding Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> * Refactor Messaging REST API to better match Messaging Send and Receive APIs * chore: migrate messaging REST API to Opt[T] Follow-up to the rebase onto master's repo-wide Option[T] -> Opt[T] change (#4035). Converts the code this branch adds to the new convention: - messaging/rest_api/types.nim: MessagingJsonEnvelope fields to Opt[T], Opt.some/Opt.none, and json_serialization/pkg/results instead of json_serialization/std/options. - tests: WakuNodeConf.clusterId is now Opt[uint16]; DTO fields are Opt. `Option[ContentBody]` in the handlers is presto's own API and stays as-is. Co-Authored-By: Claude Opus 4.8 --------- Co-authored-by: Claude Opus 4.8 Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- .github/workflows/container-image.yml | 11 +- Makefile | 15 +- .../logos_delivery_node/logosdeliverynode.nim | 98 +++++++ apps/logos_delivery_node/nim.cfg | 10 + logos_delivery.nimble | 4 + .../api/conf/logos_delivery_conf.nim | 12 +- .../api/conf/logos_delivery_conf_json.nim | 45 ++- logos_delivery/api/conf/messaging_conf.nim | 13 +- logos_delivery/api/conf/modes.nim | 18 ++ logos_delivery/logos_delivery.nim | 64 +++- logos_delivery/messaging/rest_api/client.nim | 63 ++++ .../messaging/rest_api/event_cache.nim | 115 ++++++++ .../messaging/rest_api/handlers.nim | 178 +++++++++++ logos_delivery/messaging/rest_api/types.nim | 276 ++++++++++++++++++ tests/api/test_all.nim | 4 +- tests/api/test_conf.nim | 40 +-- tests/api/test_entry_layer.nim | 128 ++++++++ tests/api/test_messaging_rest.nim | 174 +++++++++++ tools/confutils/cli_args.nim | 17 +- 19 files changed, 1226 insertions(+), 59 deletions(-) create mode 100644 apps/logos_delivery_node/logosdeliverynode.nim create mode 100644 apps/logos_delivery_node/nim.cfg create mode 100644 logos_delivery/api/conf/modes.nim create mode 100644 logos_delivery/messaging/rest_api/client.nim create mode 100644 logos_delivery/messaging/rest_api/event_cache.nim create mode 100644 logos_delivery/messaging/rest_api/handlers.nim create mode 100644 logos_delivery/messaging/rest_api/types.nim create mode 100644 tests/api/test_entry_layer.nim create mode 100644 tests/api/test_messaging_rest.nim diff --git a/.github/workflows/container-image.yml b/.github/workflows/container-image.yml index 4b0e7dcd2..432858cf5 100644 --- a/.github/workflows/container-image.yml +++ b/.github/workflows/container-image.yml @@ -87,19 +87,25 @@ jobs: id: build if: ${{ steps.secrets.outcome == 'success' }} run: | - make -j${NPROC} V=1 POSTGRES=1 NIMFLAGS="-d:disableMarchNative -d:chronicles_colors:none" wakunode2 + make -j${NPROC} V=1 POSTGRES=1 NIMFLAGS="-d:disableMarchNative -d:chronicles_colors:none" wakunode2 logosdeliverynode SHORT_REF=$(git rev-parse --short HEAD) TAG=$([ "${PR_NUMBER}" == "" ] && echo "${SHORT_REF}" || echo "${PR_NUMBER}") IMAGE=quay.io/wakuorg/nwaku-pr:${TAG} + LD_IMAGE=quay.io/wakuorg/nwaku-pr:${TAG}-logosdeliverynode echo "image=${IMAGE}" >> $GITHUB_OUTPUT + echo "ld_image=${LD_IMAGE}" >> $GITHUB_OUTPUT echo "commit_hash=$(git rev-parse HEAD)" >> $GITHUB_OUTPUT docker login -u ${QUAY_USER} -p ${QUAY_PASSWORD} quay.io docker build -t ${IMAGE} -f docker/binaries/Dockerfile.bn.amd64 --label quay.expires-after=30d . docker push ${IMAGE} + + # logosdeliverynode image (same generic Dockerfile, selected via MAKE_TARGET) + docker build -t ${LD_IMAGE} --build-arg MAKE_TARGET=logosdeliverynode -f docker/binaries/Dockerfile.bn.amd64 --label quay.expires-after=30d . + docker push ${LD_IMAGE} env: QUAY_PASSWORD: ${{ secrets.QUAY_PASSWORD }} QUAY_USER: ${{ secrets.QUAY_USER }} @@ -110,10 +116,11 @@ jobs: if: ${{ github.event_name == 'pull_request' && steps.secrets.outcome == 'success' }} with: message: | - You can find the image built from this PR at + You can find the images built from this PR at ``` ${{steps.build.outputs.image}} + ${{steps.build.outputs.ld_image}} ``` Built from ${{ steps.build.outputs.commit_hash }} diff --git a/Makefile b/Makefile index daa8a8ad7..c7bda841d 100644 --- a/Makefile +++ b/Makefile @@ -54,7 +54,7 @@ endif .PHONY: all test clean examples deps nimble install-nim install-nimble # default target -all: | wakunode2 liblogosdelivery +all: | wakunode2 logosdeliverynode liblogosdelivery examples: | example2 chat2 chat2bridge @@ -219,7 +219,7 @@ testcommon: | build-deps build ########## ## Waku ## ########## -.PHONY: testwaku wakunode2 testwakunode2 example2 chat2 chat2bridge liteprotocoltester +.PHONY: testwaku wakunode2 logosdeliverynode testwakunode2 example2 chat2 chat2bridge liteprotocoltester testwaku: | build-deps build rln-deps librln echo -e $(BUILD_MSG) "build/$@" && \ @@ -236,6 +236,17 @@ else $(NIMBLE) wakunode2 endif +# Windows: build with nim directly — `nimble ` re-clones git deps every +# build and they intermittently hang on the MSYS2 runner. Flags mirror logos_delivery.nimble. +logosdeliverynode: | build-deps build deps librln +ifeq ($(detected_OS),Windows) + echo -e $(BUILD_MSG) "build/$@" && \ + nim c --out:build/logosdeliverynode --mm:refc --cpu:amd64 $(NIM_PARAMS) -d:chronicles_log_level=TRACE apps/logos_delivery_node/logosdeliverynode.nim +else + echo -e $(BUILD_MSG) "build/$@" && \ + $(NIMBLE) logosdeliverynode +endif + benchmarks: | build-deps build deps librln echo -e $(BUILD_MSG) "build/$@" && \ $(NIMBLE) benchmarks diff --git a/apps/logos_delivery_node/logosdeliverynode.nim b/apps/logos_delivery_node/logosdeliverynode.nim new file mode 100644 index 000000000..1a7f1e0df --- /dev/null +++ b/apps/logos_delivery_node/logosdeliverynode.nim @@ -0,0 +1,98 @@ +{.push raises: [].} + +import + std/[options, strutils, sequtils, net], + chronicles, + chronos, + metrics, + system/ansi_c, + libp2p/crypto/crypto +import + ../../tools/confutils/cli_args, + logos_delivery/logos_delivery, + logos_delivery/waku/common/logging + +logScope: + topics = "logosdeliverynode main" + +const git_version* {.strdefine.} = "n/a" + +{.pop.} + # @TODO confutils.nim(775, 17) Error: can raise an unlisted exception: ref IOError +when isMainModule: + ## Node setup happens in 6 phases: + ## 1. Set up storage + ## 2. Initialize node + ## 3. Mount and initialize configured protocols + ## 4. Start node and mounted protocols + ## 5. Start monitoring tools and external interfaces + ## 6. Setup graceful shutdown hooks + + const versionString = "version / git commit hash: " & git_version + + var wakuNodeConf = WakuNodeConf.load(version = versionString).valueOr: + error "failure while loading the configuration", error = error + quit(QuitFailure) + + ## Also called within LogosDelivery.new. The call to startRestServerEssentials + ## needs the following line + logging.setupLog(wakuNodeConf.logLevel, wakuNodeConf.logFormat) + + case wakuNodeConf.cmd + of generateRlnKeystore: + error "generateRlnKeystore not supported by logos_delivery_node; use wakunode2" + quit(QuitFailure) + of noCommand: + # `LogosDelivery` derives the per-layer config from `WakuNodeConf` itself + # (it runs `toWakuConf` internally), then builds the full stack bottom-up: + # Waku <- MessagingClient <- ReliableChannelManager + var node = (waitFor LogosDelivery.new(wakuNodeConf)).valueOr: + error "LogosDelivery initialization failed", error = error + quit(QuitFailure) + + (waitFor node.start()).isOkOr: + error "Starting LogosDelivery failed", error = error + quit(QuitFailure) + + info "Setting up shutdown hooks" + proc asyncStopper(node: LogosDelivery) {.async: (raises: [Exception]).} = + (await node.stop()).isOkOr: + error "LogosDelivery shutdown failed", error = error + quit(QuitSuccess) + + # Handle Ctrl-C SIGINT + proc handleCtrlC() {.noconv.} = + when defined(windows): + # workaround for https://github.com/nim-lang/Nim/issues/4057 + setupForeignThreadGc() + notice "Shutting down after receiving SIGINT" + asyncSpawn asyncStopper(node) + + setControlCHook(handleCtrlC) + + # Handle SIGTERM + when defined(posix): + proc handleSigterm(signal: cint) {.noconv.} = + notice "Shutting down after receiving SIGTERM" + asyncSpawn asyncStopper(node) + + c_signal(ansi_c.SIGTERM, handleSigterm) + + # Handle SIGSEGV + when defined(posix): + proc handleSigsegv(signal: cint) {.noconv.} = + # Require --debugger:native + fatal "Shutting down after receiving SIGSEGV" + + # Not available in -d:release mode + writeStackTrace() + + (waitFor node.stop()).isOkOr: + error "LogosDelivery shutdown failed", error = error + quit(QuitFailure) + + c_signal(ansi_c.SIGSEGV, handleSigsegv) + + info "Node setup complete" + + runForever() diff --git a/apps/logos_delivery_node/nim.cfg b/apps/logos_delivery_node/nim.cfg new file mode 100644 index 000000000..a6fab9c9b --- /dev/null +++ b/apps/logos_delivery_node/nim.cfg @@ -0,0 +1,10 @@ +-d:chronicles_line_numbers +-d:discv5_protocol_id="d5waku" +-d:chronicles_runtime_filtering=on +-d:chronicles_sinks="textlines,json" +-d:chronicles_default_output_device=dynamic +# Disabling the following topics from nim-eth and nim-dnsdisc since some types cannot be serialized +-d:chronicles_disabled_topics="eth,dnsdisc.client" +# Results in empty output for some reason +#-d:"chronicles_enabled_topics=GossipSub:TRACE,WakuRelay:TRACE" +path = "../.." diff --git a/logos_delivery.nimble b/logos_delivery.nimble index 8c4d32cb5..275bd1a89 100644 --- a/logos_delivery.nimble +++ b/logos_delivery.nimble @@ -388,6 +388,10 @@ task wakunode2, "Build Waku v2 cli node": let name = "wakunode2" buildBinary name, "apps/wakunode2/", " -d:chronicles_log_level=TRACE " +task logosdeliverynode, "Build Logos Delivery cli node": + let name = "logosdeliverynode" + buildBinary name, "apps/logos_delivery_node/", " -d:chronicles_log_level=TRACE " + task benchmarks, "Some benchmarks": let name = "benchmarks" buildBinary name, "apps/benchmarks/", "-p:../.." diff --git a/logos_delivery/api/conf/logos_delivery_conf.nim b/logos_delivery/api/conf/logos_delivery_conf.nim index bec9f9a78..eb2e5b9fa 100644 --- a/logos_delivery/api/conf/logos_delivery_conf.nim +++ b/logos_delivery/api/conf/logos_delivery_conf.nim @@ -2,10 +2,11 @@ import results +import logos_delivery/api/conf/modes import logos_delivery/api/conf/messaging_conf import logos_delivery/api/conf/channels_conf -export messaging_conf, channels_conf +export modes, messaging_conf, channels_conf type LogosDeliveryConf* = object ## Aggregates the per-layer config objects. A layer is mounted iff its config @@ -19,11 +20,14 @@ proc init*(T: type LogosDeliveryConf, kernelConf: KernelConf): LogosDeliveryConf proc init*( T: type LogosDeliveryConf, + entryLayer: EntryLayer = EntryLayer.channels, mode: LogosDeliveryMode, preset: string, messagingOverrides: MessagingClientConf, channelsOverrides: ReliableChannelManagerConf, ): ConfResult[LogosDeliveryConf] = + ## Structured (preset + overrides) entry. Only `messaging` / `channels` layers + ## reach here; the `kernel` layer uses `init(kernelConf)` (raw, mode ignored). let merged = merge(?resolvePreset(preset), messagingOverrides) var kernelConf = ?toWakuNodeConf(merged, mode) kernelConf.preset = preset @@ -31,7 +35,11 @@ proc init*( LogosDeliveryConf( kernelConf: KernelConf(kernelConf), messagingConf: Opt.some(merged), - channelsConf: Opt.some(channelsOverrides), + channelsConf: + if entryLayer == EntryLayer.channels: + Opt.some(channelsOverrides) + else: + Opt.none(ReliableChannelManagerConf), ) ) diff --git a/logos_delivery/api/conf/logos_delivery_conf_json.nim b/logos_delivery/api/conf/logos_delivery_conf_json.nim index 6440a63ac..307ca51c6 100644 --- a/logos_delivery/api/conf/logos_delivery_conf_json.nim +++ b/logos_delivery/api/conf/logos_delivery_conf_json.nim @@ -8,6 +8,7 @@ import logos_delivery/api/conf/logos_delivery_conf const # Lowercased, since `collectJsonFields` keys the object case-insensitively. + KeyEntryLayer = "entrylayer" KeyMode = "mode" KeyPreset = "preset" KeyKernelConf = "kernelconf" @@ -24,10 +25,21 @@ proc parseMode(s: string): Result[LogosDeliveryMode, string] = return ok(LogosDeliveryMode.Core) of "edge": return ok(LogosDeliveryMode.Edge) - of "fleet": - return ok(LogosDeliveryMode.Fleet) else: - return err("invalid mode: '" & s & "' (expected 'Core', 'Edge' or 'Fleet')") + return err("invalid mode: '" & s & "' (expected 'Core' or 'Edge')") + +proc parseEntryLayer(s: string): Result[EntryLayer, string] = + case s.strip().toLowerAscii() + of "kernel": + return ok(EntryLayer.kernel) + of "messaging": + return ok(EntryLayer.messaging) + of "channels": + return ok(EntryLayer.channels) + else: + return err( + "invalid entryLayer: '" & s & "' (expected 'kernel', 'messaging' or 'channels')" + ) proc parseOverrides[T](defaults: T, node: JsonNode, label: string): Result[T, string] = ## Parse the JSON object `node` as overrides on top of `defaults`. @@ -106,16 +118,25 @@ proc parseLogosDeliveryConf*(jsonStr: string): ConfResult[LogosDeliveryConf] = mode = ?parseMode(v.getStr()) top.del(KeyMode) - if mode == LogosDeliveryMode.Fleet: - # Kernel-only: a raw kernelConf and no upper layers. + var entryLayer = EntryLayer.channels + if top.hasKey(KeyEntryLayer): + let (_, v) = top.getOrDefault(KeyEntryLayer) + if v.kind != JString: + return err("entryLayer must be a string") + entryLayer = ?parseEntryLayer(v.getStr()) + top.del(KeyEntryLayer) + + if entryLayer == EntryLayer.kernel: + # Kernel-only: a raw kernelConf and no upper layers; mode is ignored. if not top.hasKey(KeyKernelConf): - return err("fleet mode requires a 'kernelConf' object") + return err("kernel entry layer requires a 'kernelConf' object") let (_, v) = top.getOrDefault(KeyKernelConf) let kernel = ?parseOverrides(?defaultWakuNodeConf(), v, "kernelConf") top.del(KeyKernelConf) if top.len > 0: - return - err(unknownKeysError(top, "fleet mode takes only 'kernelConf'; unexpected")) + return err( + unknownKeysError(top, "kernel entry layer takes only 'kernelConf'; unexpected") + ) return ok(LogosDeliveryConf.init(KernelConf(kernel))) # [Legacy flat JSON config] A wrapper key marks our structured shape. Otherwise any @@ -158,6 +179,12 @@ proc parseLogosDeliveryConf*(jsonStr: string): ConfResult[LogosDeliveryConf] = if top.len > 0: return err(unknownKeysError(top, "Unrecognized configuration option(s) found")) - return LogosDeliveryConf.init(mode, preset, messagingOverrides, channelsOverrides) + return LogosDeliveryConf.init( + entryLayer = entryLayer, + mode = mode, + preset = preset, + messagingOverrides = messagingOverrides, + channelsOverrides = channelsOverrides, + ) {.pop.} diff --git a/logos_delivery/api/conf/messaging_conf.nim b/logos_delivery/api/conf/messaging_conf.nim index a8c91acd7..ef254654a 100644 --- a/logos_delivery/api/conf/messaging_conf.nim +++ b/logos_delivery/api/conf/messaging_conf.nim @@ -7,10 +7,8 @@ import logos_delivery/waku/factory/networks_config export kernel_conf -type LogosDeliveryMode* {.pure.} = enum - Edge # client-only node - Core # full service node - Fleet # kernel-only node from a raw kernel config +# `LogosDeliveryMode` and `EntryLayer` are defined at the leaf (`cli_args`) so +# they can appear on `WakuNodeConf`; re-exported here via `kernel_conf`. type MessagingClientConf* = object clusterId* {.name: "cluster-id".}: Opt[uint16] ## Network cluster id. @@ -69,9 +67,6 @@ proc applyMode*(conf: var WakuNodeConf, mode: LogosDeliveryMode): ConfResult[voi conf.filter = false conf.lightpush = false conf.store = false - of LogosDeliveryMode.Fleet: - return - err("fleet mode takes a raw kernel config; use LogosDelivery.new(kernelConf)") return ok() proc toWakuNodeConf*( @@ -80,6 +75,10 @@ proc toWakuNodeConf*( ## Mode sets the protocol flags; set fields map to their kernel counterpart. var conf = ?defaultWakuNodeConf() ?applyMode(conf, mode) + # Keep the `mode` field consistent with the applied flags so a later + # `LogosDelivery.new(WakuNodeConf)` re-application is idempotent instead of + # clobbering these flags with the field's default (`Core`). + conf.mode = mode if self.store.isSome(): conf.store = self.store.get() diff --git a/logos_delivery/api/conf/modes.nim b/logos_delivery/api/conf/modes.nim new file mode 100644 index 000000000..56058aed9 --- /dev/null +++ b/logos_delivery/api/conf/modes.nim @@ -0,0 +1,18 @@ +## Leaf module for the app-level mode / entry-layer enums. +## +## These appear on `WakuNodeConf` (in the leaf `tools/confutils/cli_args`), so +## they must live in a module that `cli_args` can import without a cycle — i.e. +## a module that imports nothing from the config/api layers. `logos_delivery_conf` +## re-exports them so consumers still get them from there. + +type LogosDeliveryMode* {.pure.} = enum + ## Drives the kernel-internal protocol mountings. Applied only for the + ## `messaging` / `channels` entry layers; ignored when `entryLayer == kernel`. + Edge # client-only node + Core # full service node + +type EntryLayer* {.pure.} = enum + ## Selects which API layer `LogosDelivery` instantiates. + kernel # transport kernel only; ignores `mode` and uses the config as-is + messaging # kernel + messaging client + channels # kernel + messaging + reliable channels diff --git a/logos_delivery/logos_delivery.nim b/logos_delivery/logos_delivery.nim index 35e9f6a97..55867e582 100644 --- a/logos_delivery/logos_delivery.nim +++ b/logos_delivery/logos_delivery.nim @@ -39,6 +39,8 @@ import logos_delivery/messaging/[messaging_client, messaging_client_lifecycle] export messaging_client import logos_delivery/messaging/api/[subscription, send] export subscription, send +import logos_delivery/messaging/rest_api/handlers as messaging_rest_api +export messaging_rest_api import logos_delivery/api/events/messaging_client_events export messaging_client_events import logos_delivery/api/conf/messaging_conf @@ -106,21 +108,38 @@ proc new*( proc new*( T: type LogosDelivery, conf: WakuNodeConf, appCallbacks: AppCallbacks = nil ): Future[Result[LogosDelivery, string]] {.async.} = - ## Builds the full stack from a kernel `WakuNodeConf`. - return await LogosDelivery.new( - LogosDeliveryConf( - kernelConf: KernelConf(conf), - messagingConf: Opt.some(MessagingClientConf()), - channelsConf: Opt.some(ReliableChannelManagerConf()), - ), - appCallbacks, + ## Builds the stack from a kernel `WakuNodeConf`, selecting which API layers to + ## instantiate by `conf.entryLayer`: + ## kernel -> transport only; `conf.mode` is ignored and the config is used as-is + ## messaging -> kernel + messaging client + ## channels -> kernel + messaging + reliable channels + ## For `messaging`/`channels`, `conf.mode` (Edge/Core) sets the kernel protocol + ## flags first (messaging-level concern); for `kernel` it is skipped. + var kernelConf = conf + if conf.entryLayer != EntryLayer.kernel: + applyMode(kernelConf, conf.mode).isOkOr: + return err("failed to apply mode: " & error) + + let ldConf = LogosDeliveryConf( + kernelConf: KernelConf(kernelConf), + messagingConf: + if conf.entryLayer == EntryLayer.kernel: + Opt.none(MessagingClientConf) + else: + Opt.some(MessagingClientConf()), + channelsConf: + if conf.entryLayer == EntryLayer.channels: + Opt.some(ReliableChannelManagerConf()) + else: + Opt.none(ReliableChannelManagerConf), ) + return await LogosDelivery.new(ldConf, appCallbacks) proc new*( T: type LogosDelivery, kernelConf: KernelConf, appCallbacks: AppCallbacks = nil ): Future[Result[LogosDelivery, string]] {.async.} = - ## Fleet mode: mounts the kernel only from a raw `KernelConf`; no messaging client, - ## no channel manager. + ## Kernel entry layer: mounts the kernel only from a raw `KernelConf`; no + ## messaging client, no channel manager. return await LogosDelivery.new(LogosDeliveryConf.init(kernelConf), appCallbacks) proc new*( @@ -143,14 +162,23 @@ proc new*( proc new*( T: type LogosDelivery, + entryLayer: EntryLayer = EntryLayer.channels, mode: LogosDeliveryMode = LogosDeliveryMode.Core, preset: string = "", messagingOverrides: MessagingClientConf = MessagingClientConf(), channelsOverrides: ReliableChannelManagerConf = ReliableChannelManagerConf(), appCallbacks: AppCallbacks = nil, ): Future[Result[LogosDelivery, string]] {.async.} = - ## Messaging entry point (app dev). Builds the full stack from preset, mode and overrides. - let conf = LogosDeliveryConf.init(mode, preset, messagingOverrides, channelsOverrides).valueOr: + ## Messaging entry point (app dev). Builds the stack from preset, mode and + ## overrides; `entryLayer` selects messaging vs channels (use `new(kernelConf)` + ## for a kernel-only node). + let conf = LogosDeliveryConf.init( + entryLayer = entryLayer, + mode = mode, + preset = preset, + messagingOverrides = messagingOverrides, + channelsOverrides = channelsOverrides, + ).valueOr: return err("failed to synthesize configuration: " & error) return await LogosDelivery.new(conf, appCallbacks) @@ -165,6 +193,10 @@ proc start*(self: LogosDelivery): Future[Result[void, string]] {.async.} = if not self.messagingClient.isNil(): self.messagingClient.start().isOkOr: return err("failed to start MessagingClient: " & error) + # Mount the messaging REST endpoints onto the kernel's REST router (no-op if + # REST is disabled). Done here rather than in MessagingClient.start so the + # core messaging module need not depend on the REST layer above it. + self.messagingClient.mountRestApi() if not self.reliableChannelManager.isNil(): self.reliableChannelManager.start().isOkOr: @@ -190,15 +222,15 @@ proc isOnline*(self: LogosDelivery): Future[Result[bool, string]] {.async.} = return await self.waku.isOnline() proc ensureMessaging*(self: LogosDelivery): Result[void, string] = - ## Fails if the node has no messaging client (a kernel-only / fleet node). + ## Fails if the node has no messaging client (a kernel-only node). if self.isNil() or self.messagingClient.isNil(): - return err("node has no messaging client (kernel-only/fleet node)") + return err("node has no messaging client (kernel-only node)") ok() proc ensureChannels*(self: LogosDelivery): Result[void, string] = - ## Fails if the node has no reliable channel manager (a kernel-only / fleet node). + ## Fails if the node has no reliable channel manager (a kernel-only node). if self.isNil() or self.reliableChannelManager.isNil(): - return err("node has no reliable channel manager (kernel-only/fleet node)") + return err("node has no reliable channel manager (kernel-only node)") ok() # Compile-time check that each concrete type satisfies its API concept. diff --git a/logos_delivery/messaging/rest_api/client.nim b/logos_delivery/messaging/rest_api/client.nim new file mode 100644 index 000000000..c618849ca --- /dev/null +++ b/logos_delivery/messaging/rest_api/client.nim @@ -0,0 +1,63 @@ +{.push raises: [].} + +import chronicles, json_serialization, presto/[route, client, common] +import + logos_delivery/waku/rest_api/endpoint/serdes, + logos_delivery/waku/rest_api/endpoint/rest_serdes, + logos_delivery/api/types, + ./types + +export types + +logScope: + topics = "messaging rest client" + +proc encodeBytes*( + value: seq[ContentTopic], contentType: string +): RestResult[seq[byte]] = + return encodeBytesOf(value, contentType) + +proc encodeBytes*( + value: MessagingPostMessageRequest, contentType: string +): RestResult[seq[byte]] = + return encodeBytesOf(value, contentType) + +proc messagingPostSubscriptionsV1*( + body: seq[ContentTopic] +): RestResponse[string] {. + rest, endpoint: "/messaging/v1/subscriptions", meth: HttpMethod.MethodPost +.} + +proc messagingDeleteSubscriptionsV1*( + body: seq[ContentTopic] +): RestResponse[string] {. + rest, endpoint: "/messaging/v1/subscriptions", meth: HttpMethod.MethodDelete +.} + +proc messagingPostMessagesV1*( + body: MessagingPostMessageRequest +): RestResponse[MessagingSendResponse] {. + rest, endpoint: "/messaging/v1/messages", meth: HttpMethod.MethodPost +.} + +proc messagingGetSendEventsV1*(): RestResponse[seq[SendStatus]] {. + rest, endpoint: "/messaging/v1/events/send", meth: HttpMethod.MethodGet +.} + +proc messagingGetSendEventsByIdV1*( + requestId: string +): RestResponse[SendStatus] {. + rest, endpoint: "/messaging/v1/events/send/{requestId}", meth: HttpMethod.MethodGet +.} + +# Raw variant: the typed client above cannot decode a non-2xx text error body, +# so status-code assertions (e.g. 404) use this string form. +proc messagingGetSendEventsByIdRawV1*( + requestId: string +): RestResponse[string] {. + rest, endpoint: "/messaging/v1/events/send/{requestId}", meth: HttpMethod.MethodGet +.} + +proc messagingGetReceivedMessagesV1*(): RestResponse[seq[ReceivedMessageRecord]] {. + rest, endpoint: "/messaging/v1/events/received", meth: HttpMethod.MethodGet +.} diff --git a/logos_delivery/messaging/rest_api/event_cache.nim b/logos_delivery/messaging/rest_api/event_cache.nim new file mode 100644 index 000000000..3326ade88 --- /dev/null +++ b/logos_delivery/messaging/rest_api/event_cache.nim @@ -0,0 +1,115 @@ +## In-memory cache backing the messaging event REST endpoints. +## +## REST is a poll-based client/server surface, so the interactive MessagingClient +## events must be buffered here for later observation: +## * send-related events (sent / propagated / error) grouped by request id +## * received messages +## +## Both surfaces are evict-after-poll: a GET returns the buffered data and clears +## it. A generous overflow cap bounds memory if nobody polls. All access happens +## on the single chronos event loop (broker listeners + REST handlers), so the +## synchronous (no-await) ops below need no locking. + +{.push raises: [].} + +import std/[tables, deques, options] +import results +import logos_delivery/waku/waku_core/time +import ./types + +const + DefaultMaxReceived* = 50 ## Received messages kept between polls (spec default). + DefaultMaxSendRequests* = 10_000 + ## Overflow guard: distinct request ids retained between polls. + +type MessagingEventCache* = ref object + # Send events grouped by request id, with FIFO insertion order for overflow + # eviction. Cleared on poll. + sendByReqId: Table[string, SendStatus] + sendOrder: Deque[string] + maxSendRequests: int + # Received messages, bounded ring. Cleared on poll. + received: Deque[ReceivedMessageRecord] + maxReceived: int + +proc new*( + T: type MessagingEventCache, + maxReceived = DefaultMaxReceived, + maxSendRequests = DefaultMaxSendRequests, +): MessagingEventCache = + MessagingEventCache( + sendByReqId: initTable[string, SendStatus](), + sendOrder: initDeque[string](), + maxSendRequests: maxSendRequests, + received: initDeque[ReceivedMessageRecord](), + maxReceived: maxReceived, + ) + +proc recordSend*( + self: MessagingEventCache, + requestId: string, + messageHash: string, + kind: SendEventKind, + error = "", +) = + ## Append a send event to its request id's timeline, creating the entry (and + ## evicting the oldest request id past the overflow cap) as needed. + let record = SendEventRecord( + kind: kind, + messageHash: messageHash, + error: error, + timestamp: getNowInNanosecondTime(), + ) + + if not self.sendByReqId.hasKey(requestId): + self.sendByReqId[requestId] = SendStatus(requestId: requestId, events: @[]) + self.sendOrder.addLast(requestId) + + while self.sendOrder.len > self.maxSendRequests: + let evicted = self.sendOrder.popFirst() + self.sendByReqId.del(evicted) + + self.sendByReqId.withValue(requestId, status): + status[].events.add(record) + +proc recordReceived*( + self: MessagingEventCache, messageHash: string, message: RelayWakuMessage +) = + ## Buffer a received message, dropping the oldest past the ring capacity. + self.received.addLast( + ReceivedMessageRecord(messageHash: messageHash, message: message) + ) + + while self.received.len > self.maxReceived: + discard self.received.popFirst() + +proc pollAllSend*(self: MessagingEventCache): seq[SendStatus] = + ## Return all buffered send statuses and clear the store (evict-after-poll). + for reqId in self.sendOrder: + self.sendByReqId.withValue(reqId, status): + result.add(status[]) + self.sendByReqId.clear() + self.sendOrder.clear() + +proc pollSend*(self: MessagingEventCache, requestId: string): Opt[SendStatus] = + ## Return one request id's send status and remove it (evict-after-poll). + var status: SendStatus + if not self.sendByReqId.pop(requestId, status): + return Opt.none(SendStatus) + + # Deque has no random removal; rebuild order without the polled id. + var rebuilt = initDeque[string]() + for reqId in self.sendOrder: + if reqId != requestId: + rebuilt.addLast(reqId) + self.sendOrder = rebuilt + + return Opt.some(status) + +proc pollReceived*(self: MessagingEventCache): seq[ReceivedMessageRecord] = + ## Return buffered received messages (oldest first) and clear (evict-after-poll). + for record in self.received: + result.add(record) + self.received.clear() + +{.pop.} diff --git a/logos_delivery/messaging/rest_api/handlers.nim b/logos_delivery/messaging/rest_api/handlers.nim new file mode 100644 index 000000000..ea3d82001 --- /dev/null +++ b/logos_delivery/messaging/rest_api/handlers.nim @@ -0,0 +1,178 @@ +{.push raises: [].} + +import chronos, chronicles, results, json_serialization, json_serialization/std/options +import presto/[route, common] +import + logos_delivery/waku/waku, + logos_delivery/waku/rest_api/endpoint/serdes, + logos_delivery/waku/rest_api/endpoint/responses, + logos_delivery/waku/rest_api/endpoint/rest_serdes, + logos_delivery/messaging/messaging_client, + logos_delivery/messaging/api/subscription, + logos_delivery/messaging/api/send, + logos_delivery/api/types, + logos_delivery/api/events/messaging_client_events, + ./types, + ./event_cache + +export types + +logScope: + topics = "messaging rest api" + +#### Routes + +const ROUTE_MESSAGING_SUBSCRIPTIONSV1* = "/messaging/v1/subscriptions" +const ROUTE_MESSAGING_MESSAGESV1* = "/messaging/v1/messages" +const ROUTE_MESSAGING_EVENTS_SENDV1* = "/messaging/v1/events/send" +const ROUTE_MESSAGING_EVENTS_SEND_BY_IDV1* = "/messaging/v1/events/send/{requestId}" +const ROUTE_MESSAGING_EVENTS_RECEIVEDV1* = "/messaging/v1/events/received" + +proc installEventListeners(brokerCtx: BrokerContext, cache: MessagingEventCache) = + ## Buffers the MessagingClient events into `cache` so the poll-based REST + ## endpoints can observe them. Listeners live for the node's lifetime (the + ## captured `cache` keeps them and their data alive); no teardown is wired. + discard MessageSentEvent.listen( + brokerCtx, + proc(evt: MessageSentEvent): Future[void] {.async: (raises: []).} = + cache.recordSend($evt.requestId, evt.messageHash, SendEventKind.Sent), + ) + + discard MessagePropagatedEvent.listen( + brokerCtx, + proc(evt: MessagePropagatedEvent): Future[void] {.async: (raises: []).} = + cache.recordSend($evt.requestId, evt.messageHash, SendEventKind.Propagated), + ) + + discard MessageErrorEvent.listen( + brokerCtx, + proc(evt: MessageErrorEvent): Future[void] {.async: (raises: []).} = + cache.recordSend($evt.requestId, evt.messageHash, SendEventKind.Error, evt.error), + ) + + discard MessageReceivedEvent.listen( + brokerCtx, + proc(evt: MessageReceivedEvent): Future[void] {.async: (raises: []).} = + cache.recordReceived(evt.messageHash, toRelayWakuMessage(evt.message)), + ) + +proc installMessagingApiHandlers*(router: var RestRouter, client: MessagingClient) = + ## Mounts the MessagingClient subscribe / unsubscribe / send operations as + ## REST endpoints onto the given (kernel-owned) router. Subscriptions are + ## keyed by content topic, matching the messaging layer's content-topic API. + + # Event observability: buffer send/received events for the poll-based GETs. + let eventCache = MessagingEventCache.new() + installEventListeners(client.waku.brokerCtx, eventCache) + + router.api(MethodOptions, ROUTE_MESSAGING_SUBSCRIPTIONSV1) do() -> RestApiResponse: + return RestApiResponse.ok() + + router.api(MethodPost, ROUTE_MESSAGING_SUBSCRIPTIONSV1) do( + contentBody: Option[ContentBody] + ) -> RestApiResponse: + ## Subscribes the messaging client to a list of content topics. + let req: seq[ContentTopic] = decodeRequestBody[seq[ContentTopic]](contentBody).valueOr: + return error + + for contentTopic in req: + (await client.subscribe(contentTopic)).isOkOr: + let errorMsg = "Subscribe failed: " & error + error "messaging SUBSCRIBE failed", error = errorMsg + return RestApiResponse.internalServerError(errorMsg) + + return RestApiResponse.ok() + + router.api(MethodDelete, ROUTE_MESSAGING_SUBSCRIPTIONSV1) do( + contentBody: Option[ContentBody] + ) -> RestApiResponse: + ## Unsubscribes the messaging client from a list of content topics. + let req: seq[ContentTopic] = decodeRequestBody[seq[ContentTopic]](contentBody).valueOr: + return error + + for contentTopic in req: + client.unsubscribe(contentTopic).isOkOr: + let errorMsg = "Unsubscribe failed: " & error + error "messaging UNSUBSCRIBE failed", error = errorMsg + return RestApiResponse.internalServerError(errorMsg) + + return RestApiResponse.ok() + + router.api(MethodOptions, ROUTE_MESSAGING_MESSAGESV1) do() -> RestApiResponse: + return RestApiResponse.ok() + + router.api(MethodPost, ROUTE_MESSAGING_MESSAGESV1) do( + contentBody: Option[ContentBody] + ) -> RestApiResponse: + ## Sends a message through the messaging client, returning the request id. + let req: MessagingJsonEnvelope = decodeRequestBody[MessagingJsonEnvelope]( + contentBody + ).valueOr: + return error + + let envelope = req.toMessageEnvelope().valueOr: + return RestApiResponse.badRequest("Invalid message: " & error) + + let requestId = (await client.send(envelope)).valueOr: + error "messaging SEND failed", error = error + return RestApiResponse.internalServerError("Send failed: " & error) + + let data = MessagingSendResponse(requestId: $requestId) + return RestApiResponse.jsonResponse(data, status = Http200).valueOr: + error "An error occurred while building the json response", error = error + return RestApiResponse.internalServerError($error) + + #### Event observability endpoints (poll-based, evict-after-poll) + + router.api(MethodOptions, ROUTE_MESSAGING_EVENTS_SENDV1) do() -> RestApiResponse: + return RestApiResponse.ok() + + router.api(MethodGet, ROUTE_MESSAGING_EVENTS_SENDV1) do() -> RestApiResponse: + ## Returns all buffered send events grouped by request id, then clears them. + let data = eventCache.pollAllSend() + return RestApiResponse.jsonResponse(data, status = Http200).valueOr: + error "An error occurred while building the json response", error = error + return RestApiResponse.internalServerError($error) + + router.api(MethodOptions, ROUTE_MESSAGING_EVENTS_SEND_BY_IDV1) do( + requestId: string + ) -> RestApiResponse: + return RestApiResponse.ok() + + router.api(MethodGet, ROUTE_MESSAGING_EVENTS_SEND_BY_IDV1) do( + requestId: string + ) -> RestApiResponse: + ## Returns the buffered send events for one request id, then removes them. + let reqId = requestId.valueOr: + return RestApiResponse.badRequest("Invalid requestId") + + let status = eventCache.pollSend(reqId).valueOr: + return RestApiResponse.notFound("No send events for requestId: " & reqId) + + return RestApiResponse.jsonResponse(status, status = Http200).valueOr: + error "An error occurred while building the json response", error = error + return RestApiResponse.internalServerError($error) + + router.api(MethodOptions, ROUTE_MESSAGING_EVENTS_RECEIVEDV1) do() -> RestApiResponse: + return RestApiResponse.ok() + + router.api(MethodGet, ROUTE_MESSAGING_EVENTS_RECEIVEDV1) do() -> RestApiResponse: + ## Returns buffered received messages (up to the cache capacity, oldest + ## first), then clears them — optimized for polling. + let data = eventCache.pollReceived() + return RestApiResponse.jsonResponse(data, status = Http200).valueOr: + error "An error occurred while building the json response", error = error + return RestApiResponse.internalServerError($error) + +proc mountRestApi*(client: MessagingClient) = + ## Mounts the messaging REST endpoints onto the kernel-owned REST router, if + ## the REST server is enabled. Called by the `LogosDelivery` concentrator + ## after the messaging layer has started. Lives here (not in the core + ## `messaging_client` module) so the core need not depend on the REST layer + ## above it — that would form an import cycle. + if not client.waku.restServer.isNil(): + # The BTree route table is ref-backed, so mutating the copied router persists + # (same pattern as the waku REST builder). + var router = client.waku.restServer.router + installMessagingApiHandlers(router, client) + info "Mounted messaging REST API endpoints" diff --git a/logos_delivery/messaging/rest_api/types.nim b/logos_delivery/messaging/rest_api/types.nim new file mode 100644 index 000000000..0688db266 --- /dev/null +++ b/logos_delivery/messaging/rest_api/types.nim @@ -0,0 +1,276 @@ +{.push raises: [].} + +import + std/[sets, strformat], + chronicles, + results, + json_serialization, + json_serialization/pkg/results, + presto/[route, client, common] +import + logos_delivery/waku/common/base64, + logos_delivery/waku/rest_api/endpoint/serdes, + logos_delivery/waku/rest_api/endpoint/relay/types as relay_types, + logos_delivery/api/types + +export types, relay_types + +#### Types + +type MessagingJsonEnvelope* = object + ## REST wire (JSON) representation of the messaging API's `MessageEnvelope`. + ## `payload` / `meta` are base64. Fields mirror `MessageEnvelope` exactly. + payload*: Base64String + contentTopic*: ContentTopic + ephemeral*: Opt[bool] + meta*: Opt[Base64String] + +type MessagingPostMessageRequest* = MessagingJsonEnvelope + +type MessagingSendResponse* = object + ## Returned by the send endpoint on success; correlates with + ## `MessageSentEvent` / `MessageErrorEvent`. + requestId*: string + +#### Type conversion + +proc toMessageEnvelope*(msg: MessagingJsonEnvelope): Result[MessageEnvelope, string] = + let + payload = ?msg.payload.decode() + meta = ?msg.meta.get(Base64String("")).decode() + + return ok( + MessageEnvelope( + contentTopic: msg.contentTopic, + payload: payload, + ephemeral: msg.ephemeral.get(false), + meta: meta, + ) + ) + +#### Serialization and deserialization + +proc writeValue*( + writer: var JsonWriter[RestJson], value: MessagingJsonEnvelope +) {.raises: [IOError].} = + writer.beginRecord() + writer.writeField("payload", value.payload) + writer.writeField("contentTopic", value.contentTopic) + if value.ephemeral.isSome(): + writer.writeField("ephemeral", value.ephemeral.get()) + if value.meta.isSome(): + writer.writeField("meta", value.meta.get()) + writer.endRecord() + +proc readValue*( + reader: var JsonReader[RestJson], value: var MessagingJsonEnvelope +) {.raises: [SerializationError, IOError].} = + var + payload = Opt.none(Base64String) + contentTopic = Opt.none(ContentTopic) + ephemeral = Opt.none(bool) + meta = Opt.none(Base64String) + + var keys = initHashSet[string]() + for fieldName in readObjectFields(reader): + # Check for repeated keys + if keys.containsOrIncl(fieldName): + let err = + try: + fmt"Multiple `{fieldName}` fields found" + except CatchableError: + "Multiple fields with the same name found" + reader.raiseUnexpectedField(err, "MessagingJsonEnvelope") + + case fieldName + of "payload": + payload = Opt.some(reader.readValue(Base64String)) + of "contentTopic": + contentTopic = Opt.some(reader.readValue(ContentTopic)) + of "ephemeral": + ephemeral = Opt.some(reader.readValue(bool)) + of "meta": + meta = Opt.some(reader.readValue(Base64String)) + else: + unrecognizedFieldWarning(value) + + if payload.isNone() or isEmptyOrWhitespace(string(payload.get())): + reader.raiseUnexpectedValue("Field `payload` is missing or empty") + + if contentTopic.isNone() or contentTopic.get().isEmptyOrWhitespace(): + reader.raiseUnexpectedValue("Field `contentTopic` is missing or empty") + + value = MessagingJsonEnvelope( + payload: payload.get(), + contentTopic: contentTopic.get(), + ephemeral: ephemeral, + meta: meta, + ) + +proc writeValue*( + writer: var JsonWriter[RestJson], value: MessagingSendResponse +) {.raises: [IOError].} = + writer.beginRecord() + writer.writeField("requestId", value.requestId) + writer.endRecord() + +proc readValue*( + reader: var JsonReader[RestJson], value: var MessagingSendResponse +) {.raises: [SerializationError, IOError].} = + var requestId = Opt.none(string) + + var keys = initHashSet[string]() + for fieldName in readObjectFields(reader): + if keys.containsOrIncl(fieldName): + let err = + try: + fmt"Multiple `{fieldName}` fields found" + except CatchableError: + "Multiple fields with the same name found" + reader.raiseUnexpectedField(err, "MessagingSendResponse") + + case fieldName + of "requestId": + requestId = Opt.some(reader.readValue(string)) + else: + unrecognizedFieldWarning(value) + + if requestId.isNone(): + reader.raiseUnexpectedValue("Field `requestId` is missing") + + value = MessagingSendResponse(requestId: requestId.get()) + +#### Event observability DTOs +## +## Send-related events (sent / propagated / error) are grouped per request id. +## Received messages carry the full `WakuMessage` (serialized as +## `RelayWakuMessage`), matching the nim `MessageReceivedEvent`. Both surfaces +## are populated by the broker listeners installed in the messaging handlers. + +type + SendEventKind* {.pure.} = enum + Sent = "sent" + Propagated = "propagated" + Error = "error" + + SendEventRecord* = object + kind*: SendEventKind + messageHash*: string + error*: string ## populated only for `Error` + timestamp*: int64 ## nanoseconds, stamped when cached + + SendStatus* = object ## All send events observed so far for a single request id. + requestId*: string + events*: seq[SendEventRecord] + + ReceivedMessageRecord* = object + messageHash*: string + message*: RelayWakuMessage ## the received WakuMessage, full fidelity + +#### Event DTO serialization + +proc writeValue*( + writer: var JsonWriter[RestJson], value: SendEventRecord +) {.raises: [IOError].} = + writer.beginRecord() + writer.writeField("kind", $value.kind) + writer.writeField("messageHash", value.messageHash) + if value.error.len > 0: + writer.writeField("error", value.error) + writer.writeField("timestamp", value.timestamp) + writer.endRecord() + +proc writeValue*( + writer: var JsonWriter[RestJson], value: SendStatus +) {.raises: [IOError].} = + writer.beginRecord() + writer.writeField("requestId", value.requestId) + writer.writeField("events", value.events) + writer.endRecord() + +proc writeValue*( + writer: var JsonWriter[RestJson], value: ReceivedMessageRecord +) {.raises: [IOError].} = + writer.beginRecord() + writer.writeField("messageHash", value.messageHash) + writer.writeField("message", value.message) + writer.endRecord() + +proc readValue*( + reader: var JsonReader[RestJson], value: var SendEventKind +) {.raises: [SerializationError, IOError].} = + let s = reader.readValue(string) + case s + of "sent": + value = SendEventKind.Sent + of "propagated": + value = SendEventKind.Propagated + of "error": + value = SendEventKind.Error + else: + reader.raiseUnexpectedValue("Invalid send event kind: " & s) + +proc readValue*( + reader: var JsonReader[RestJson], value: var SendEventRecord +) {.raises: [SerializationError, IOError].} = + var + kind = Opt.none(SendEventKind) + messageHash = "" + error = "" + timestamp = int64(0) + + for fieldName in readObjectFields(reader): + case fieldName + of "kind": + kind = Opt.some(reader.readValue(SendEventKind)) + of "messageHash": + messageHash = reader.readValue(string) + of "error": + error = reader.readValue(string) + of "timestamp": + timestamp = reader.readValue(int64) + else: + unrecognizedFieldWarning(value) + + if kind.isNone(): + reader.raiseUnexpectedValue("Field `kind` is missing") + + value = SendEventRecord( + kind: kind.get(), messageHash: messageHash, error: error, timestamp: timestamp + ) + +proc readValue*( + reader: var JsonReader[RestJson], value: var SendStatus +) {.raises: [SerializationError, IOError].} = + var + requestId = "" + events: seq[SendEventRecord] = @[] + + for fieldName in readObjectFields(reader): + case fieldName + of "requestId": + requestId = reader.readValue(string) + of "events": + events = reader.readValue(seq[SendEventRecord]) + else: + unrecognizedFieldWarning(value) + + value = SendStatus(requestId: requestId, events: events) + +proc readValue*( + reader: var JsonReader[RestJson], value: var ReceivedMessageRecord +) {.raises: [SerializationError, IOError].} = + var + messageHash = "" + message = RelayWakuMessage() + + for fieldName in readObjectFields(reader): + case fieldName + of "messageHash": + messageHash = reader.readValue(string) + of "message": + message = reader.readValue(RelayWakuMessage) + else: + unrecognizedFieldWarning(value) + + value = ReceivedMessageRecord(messageHash: messageHash, message: message) diff --git a/tests/api/test_all.nim b/tests/api/test_all.nim index 108551bc2..3d7bdee4b 100644 --- a/tests/api/test_all.nim +++ b/tests/api/test_all.nim @@ -7,4 +7,6 @@ import ./test_api_send, ./test_api_subscription, ./test_api_receive, - ./test_api_health + ./test_api_health, + ./test_messaging_rest, + ./test_entry_layer diff --git a/tests/api/test_conf.nim b/tests/api/test_conf.nim index 06a115897..7f4f23315 100644 --- a/tests/api/test_conf.nim +++ b/tests/api/test_conf.nim @@ -246,9 +246,9 @@ suite "parseLogosDeliveryConf - JSON parsing": kc.storeMessageRetentionPolicy == "time:3600" kc.storeMaxNumDbConnections == 7 - test "fleet mode parses a raw kernelConf": + test "kernel entry layer parses a raw kernelConf": let lc = parseLogosDeliveryConf( - """{"mode": "fleet", "kernelConf": {"relay": false, "maxMessageSize": "150KiB"}}""" + """{"entrylayer": "kernel", "kernelConf": {"relay": false, "maxMessageSize": "150KiB"}}""" ).valueOr: raiseAssert error check: @@ -259,24 +259,25 @@ suite "parseLogosDeliveryConf - JSON parsing": kc.relay == false kc.maxMessageSize == "150KiB" - test "fleet mode requires a kernelConf": - check parseLogosDeliveryConf("""{"mode": "fleet"}""").isErr() + test "kernel entry layer requires a kernelConf": + check parseLogosDeliveryConf("""{"entrylayer": "kernel"}""").isErr() - test "fleet mode rejects anything besides kernelConf": + test "kernel entry layer rejects anything besides kernelConf": check parseLogosDeliveryConf( - """{"mode": "fleet", "kernelConf": {}, "messagingOverrides": {"clusterId": 1}}""" + """{"entrylayer": "kernel", "kernelConf": {}, "messagingOverrides": {"clusterId": 1}}""" ) .isErr() - # a preset for a fleet node goes inside kernelConf (WakuNodeConf.preset); fleet - # treats kernelConf as a finished object and overlays nothing onto it + # a preset for a kernel node goes inside kernelConf (WakuNodeConf.preset); the + # kernel entry layer treats kernelConf as a finished object and overlays nothing check parseLogosDeliveryConf( - """{"mode": "fleet", "kernelConf": {}, "preset": "twn"}""" + """{"entrylayer": "kernel", "kernelConf": {}, "preset": "twn"}""" ) .isErr() - test "kernelConf is rejected outside fleet mode": - # kernelConf is a fleet-only wrapper. Under Core/Edge it is neither consumed by the - # structured path nor a flat kernel field, so it must surface as an unknown key. + test "kernelConf is rejected outside the kernel entry layer": + # kernelConf is a kernel-entry-layer wrapper. Under messaging/channels it is + # neither consumed by the structured path nor a flat kernel field, so it must + # surface as an unknown key. check parseLogosDeliveryConf("""{"mode": "core", "kernelConf": {}}""").isErr() suite "LogosDelivery.new - construction (the app-dev entry)": @@ -285,9 +286,9 @@ suite "LogosDelivery.new - construction (the app-dev entry)": lockNewGlobalBrokerContext: node = ( await LogosDelivery.new( - LogosDeliveryMode.Core, - "", - MessagingClientConf( + mode = LogosDeliveryMode.Core, + preset = "", + messagingOverrides = MessagingClientConf( clusterId: Opt.some(3'u16), numShardsInCluster: Opt.some(1'u16), listenIpv4: Opt.some(parseIpAddress("0.0.0.0")), @@ -307,9 +308,10 @@ suite "LogosDelivery.new - construction (the app-dev entry)": lockNewGlobalBrokerContext: node = ( await LogosDelivery.new( - LogosDeliveryMode.Core, - "logostest", - MessagingClientConf(listenIpv4: Opt.some(parseIpAddress("0.0.0.0"))), + mode = LogosDeliveryMode.Core, + preset = "logostest", + messagingOverrides = + MessagingClientConf(listenIpv4: Opt.some(parseIpAddress("0.0.0.0"))), ) ).valueOr: raiseAssert error @@ -327,7 +329,7 @@ suite "MessagingClientConf - store override": kc.relay == false # protocols are owned by the mode, not overridable suite "LogosDelivery.new - raw kernel construction": - asyncTest "a fleet node mounts the kernel only; start/stop tolerate the nil layers": + asyncTest "a kernel-only node mounts the kernel only; start/stop tolerate the nil layers": let kernel = MessagingClientConf(listenIpv4: Opt.some(parseIpAddress("0.0.0.0"))).toWakuNodeConf( LogosDeliveryMode.Core ).valueOr: diff --git a/tests/api/test_entry_layer.nim b/tests/api/test_entry_layer.nim new file mode 100644 index 000000000..710a16653 --- /dev/null +++ b/tests/api/test_entry_layer.nim @@ -0,0 +1,128 @@ +{.used.} + +import std/[options, net] +import chronos, testutils/unittests, presto, presto/client as presto_client +import brokers/broker_context +import logos_delivery +import + logos_delivery/api/conf/logos_delivery_conf, + logos_delivery/messaging/rest_api/client as messaging_rest_client, + logos_delivery/waku/rest_api/endpoint/client +import tools/confutils/cli_args +import ../testlib/testasync + +## Validates the layer-selection invariant of `LogosDelivery.new(WakuNodeConf)`: +## `messagingClient` (and `reliableChannelManager`) are instantiated only for the +## entry layers that call for them. +## +## kernel -> waku only +## messaging -> waku + messagingClient +## channels -> waku + messagingClient + reliableChannelManager + +proc nodeConf(entryLayer: EntryLayer, rest = false): WakuNodeConf = + var conf = defaultWakuNodeConf().valueOr: + raiseAssert error + conf.entryLayer = entryLayer + conf.mode = LogosDeliveryMode.Core + conf.listenAddress = parseIpAddress("0.0.0.0") + conf.tcpPort = Port(0) + conf.discv5UdpPort = Port(0) + conf.clusterId = Opt.some(3'u16) + conf.numShardsInNetwork = 1 + conf.rest = rest + conf.restAddress = parseIpAddress("127.0.0.1") + conf.restPort = 0'u16 # bind to an ephemeral port + return conf + +proc restClientFor(node: LogosDelivery): RestClientRef = + let boundPort = node.waku.restServer.httpServer.address.port + newRestHttpClient(initTAddress(parseIpAddress("127.0.0.1"), boundPort)) + +suite "LogosDelivery - entry layer selection": + asyncTest "kernel: waku only, no messaging / channels": + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(nodeConf(EntryLayer.kernel))).valueOr: + raiseAssert error + check: + not node.waku.isNil() + node.messagingClient.isNil() + node.reliableChannelManager.isNil() + node.ensureMessaging().isErr() + node.ensureChannels().isErr() + (await node.stop()).isOkOr: + raiseAssert "stop failed: " & error + + asyncTest "messaging: waku + messagingClient, no channels": + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(nodeConf(EntryLayer.messaging))).valueOr: + raiseAssert error + check: + not node.waku.isNil() + not node.messagingClient.isNil() + node.reliableChannelManager.isNil() + node.ensureMessaging().isOk() + node.ensureChannels().isErr() + (await node.stop()).isOkOr: + raiseAssert "stop failed: " & error + + asyncTest "channels: full stack": + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(nodeConf(EntryLayer.channels))).valueOr: + raiseAssert error + check: + not node.waku.isNil() + not node.messagingClient.isNil() + not node.reliableChannelManager.isNil() + node.ensureMessaging().isOk() + node.ensureChannels().isOk() + (await node.stop()).isOkOr: + raiseAssert "stop failed: " & error + + asyncTest "messaging + rest: messaging REST endpoints are installed and working": + ## entry-layer=messaging, mode=Core, rest=true -> `start` mounts the messaging + ## REST endpoints; they respond over HTTP. + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(nodeConf(EntryLayer.messaging, rest = true))).valueOr: + raiseAssert error + (await node.start()).isOkOr: + raiseAssert "start failed: " & error + + check not node.messagingClient.isNil() + + let client = restClientFor(node) + + # A command endpoint and an observability endpoint both respond -> the + # handlers were installed onto the kernel router. + let subResp = + await client.messagingPostSubscriptionsV1(@["/test/1/entry-layer/proto"]) + check subResp.status == 200 + + let sendEventsResp = await client.messagingGetSendEventsV1() + check sendEventsResp.status == 200 + + (await node.stop()).isOkOr: + raiseAssert "stop failed: " & error + + asyncTest "kernel + rest: messaging REST endpoints are NOT installed": + ## Gating check: a kernel-only node still starts a REST server, but the + ## messaging endpoints must be absent (no messaging client to mount them). + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(nodeConf(EntryLayer.kernel, rest = true))).valueOr: + raiseAssert error + (await node.start()).isOkOr: + raiseAssert "start failed: " & error + + check node.messagingClient.isNil() + + let client = restClientFor(node) + let subResp = + await client.messagingPostSubscriptionsV1(@["/test/1/entry-layer/proto"]) + check subResp.status == 404 # route not mounted + + (await node.stop()).isOkOr: + raiseAssert "stop failed: " & error diff --git a/tests/api/test_messaging_rest.nim b/tests/api/test_messaging_rest.nim new file mode 100644 index 000000000..fcad51d44 --- /dev/null +++ b/tests/api/test_messaging_rest.nim @@ -0,0 +1,174 @@ +{.used.} + +import + std/[options, net, sequtils], + chronos, + testutils/unittests, + presto, + presto/client as presto_client, + libp2p/crypto/crypto +import brokers/broker_context +import logos_delivery +import + logos_delivery/api/conf/logos_delivery_conf, + logos_delivery/messaging/rest_api/client as messaging_rest_client, + logos_delivery/waku/rest_api/endpoint/client, + logos_delivery/waku/common/base64 +import tools/confutils/cli_args +import ../testlib/[wakucore, testasync] + +## Integration test for the messaging REST endpoints and their event cache. +## +## A full `LogosDelivery` node is started with REST enabled (so `start` mounts +## the messaging handlers + event listeners), then driven through the generated +## `client.nim` stubs. Event observability is exercised deterministically by +## emitting the MessagingClient events directly on the node's broker context — +## the same context the send/recv services emit on — so we do not depend on real +## network delivery. + +proc restNodeConf(): WakuNodeConf = + var conf = defaultWakuNodeConf().valueOr: + raiseAssert error + conf.entryLayer = EntryLayer.messaging + conf.mode = LogosDeliveryMode.Core + conf.listenAddress = parseIpAddress("0.0.0.0") + conf.tcpPort = Port(0) + conf.discv5UdpPort = Port(0) + conf.clusterId = Opt.some(3'u16) + conf.numShardsInNetwork = 1 + conf.rest = true + conf.restAddress = parseIpAddress("127.0.0.1") + conf.restPort = 0'u16 # bind to an ephemeral port + return conf + +proc restClientFor(node: LogosDelivery): RestClientRef = + let boundPort = node.waku.restServer.httpServer.address.port + newRestHttpClient(initTAddress(parseIpAddress("127.0.0.1"), boundPort)) + +const settleDelay = 200.milliseconds + ## Event emit + listener run are asyncSpawned; give them a turn before polling. + +suite "Messaging REST API": + asyncTest "subscribe / unsubscribe / send endpoints respond": + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(restNodeConf())).valueOr: + raiseAssert error + (await node.start()).isOkOr: + raiseAssert "Failed to start node: " & error + let client = restClientFor(node) + + let contentTopic = "/test/1/messaging-rest/proto" + + let subResp = await client.messagingPostSubscriptionsV1(@[contentTopic]) + check subResp.status == 200 + + let msg = MessagingJsonEnvelope( + payload: base64.encode("hello rest"), + contentTopic: contentTopic, + ephemeral: Opt.none(bool), + meta: Opt.none(Base64String), + ) + let sendResp = await client.messagingPostMessagesV1(msg) + check: + sendResp.status == 200 + sendResp.data.requestId.len > 0 + + let unsubResp = await client.messagingDeleteSubscriptionsV1(@[contentTopic]) + check unsubResp.status == 200 + + (await node.stop()).isOkOr: + raiseAssert "Failed to stop node: " & error + + asyncTest "send events are grouped by requestId and evict after poll": + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(restNodeConf())).valueOr: + raiseAssert error + (await node.start()).isOkOr: + raiseAssert "Failed to start node: " & error + let client = restClientFor(node) + let brokerCtx = node.waku.brokerCtx + + let reqA = RequestId("req-A") + let reqB = RequestId("req-B") + + # reqA sees the full lifecycle; reqB only an error. + MessagePropagatedEvent.emit( + brokerCtx, MessagePropagatedEvent(requestId: reqA, messageHash: "0xaa") + ) + MessageSentEvent.emit( + brokerCtx, MessageSentEvent(requestId: reqA, messageHash: "0xaa") + ) + MessageErrorEvent.emit( + brokerCtx, MessageErrorEvent(requestId: reqB, messageHash: "0xbb", error: "boom") + ) + await sleepAsync(settleDelay) + + # GET by id returns only reqA and removes it. + let byIdResp = await client.messagingGetSendEventsByIdV1($reqA) + check: + byIdResp.status == 200 + byIdResp.data.requestId == $reqA + byIdResp.data.events.len == 2 + byIdResp.data.events.anyIt(it.kind == SendEventKind.Sent) + byIdResp.data.events.anyIt(it.kind == SendEventKind.Propagated) + + # Unknown / already-polled id → 404 (raw string client so the text error + # body decodes; the typed client would raise on a non-2xx body). + let missingResp = await client.messagingGetSendEventsByIdRawV1($reqA) + check missingResp.status == 404 + + # GET all now returns only reqB (with its error), then clears. + let allResp = await client.messagingGetSendEventsV1() + check: + allResp.status == 200 + allResp.data.len == 1 + allResp.data[0].requestId == $reqB + allResp.data[0].events.len == 1 + allResp.data[0].events[0].kind == SendEventKind.Error + allResp.data[0].events[0].error == "boom" + + let emptyResp = await client.messagingGetSendEventsV1() + check: + emptyResp.status == 200 + emptyResp.data.len == 0 + + (await node.stop()).isOkOr: + raiseAssert "Failed to stop node: " & error + + asyncTest "received messages are observable, capped, and evict after poll": + var node: LogosDelivery + lockNewGlobalBrokerContext: + node = (await LogosDelivery.new(restNodeConf())).valueOr: + raiseAssert error + (await node.start()).isOkOr: + raiseAssert "Failed to start node: " & error + let client = restClientFor(node) + let brokerCtx = node.waku.brokerCtx + + # Emit more than the cache capacity (50); oldest must be dropped. + const total = 55 + for i in 0 ..< total: + let wm = + fakeWakuMessage(payload = "msg-" & $i, contentTopic = "/test/1/recv/proto") + MessageReceivedEvent.emit( + brokerCtx, MessageReceivedEvent(messageHash: "0x" & $i, message: wm) + ) + await sleepAsync(settleDelay) + + let resp = await client.messagingGetReceivedMessagesV1() + check: + resp.status == 200 + resp.data.len == 50 # capped at DefaultMaxReceived + # oldest (0..4) evicted, newest retained, oldest-first ordering + resp.data[0].messageHash == "0x5" + resp.data[^1].messageHash == "0x" & $(total - 1) + + let emptyResp = await client.messagingGetReceivedMessagesV1() + check: + emptyResp.status == 200 + emptyResp.data.len == 0 + + (await node.stop()).isOkOr: + raiseAssert "Failed to stop node: " & error diff --git a/tools/confutils/cli_args.nim b/tools/confutils/cli_args.nim index 1774f1445..1cb6292d2 100644 --- a/tools/confutils/cli_args.nim +++ b/tools/confutils/cli_args.nim @@ -21,6 +21,7 @@ import json import + logos_delivery/api/conf/modes, logos_delivery/waku/factory/[waku_conf, conf_builder/conf_builder, networks_config], logos_delivery/waku/common/[logging], logos_delivery/waku/[ @@ -37,7 +38,7 @@ import ./envvar as confEnvvarDefs, ./envvar_net as confEnvvarNet export confTomlDefs, confTomlNet, confEnvvarDefs, confEnvvarNet, ProtectedShard, - DefaultMaxWakuMessageSizeStr, DefaultAgentString + DefaultMaxWakuMessageSizeStr, DefaultAgentString, modes logScope: topics = "waku cli args" @@ -168,6 +169,20 @@ type WakuNodeConf* = object name: "preset" .}: string + entryLayer* {. + desc: + "Top API layer to run: kernel (transport only), messaging, or channels (messaging + reliable channels).", + defaultValue: EntryLayer.channels, + name: "entry-layer" + .}: EntryLayer + + mode* {. + desc: + "Kernel operating mode: Edge (client-only) or Core (full service node). Applied only for --entry-layer=messaging|channels; ignored for kernel.", + defaultValue: LogosDeliveryMode.Core, + name: "mode" + .}: LogosDeliveryMode + # Opt-typed; desc states the default since the CLI can't auto-show it for Opt.none(). clusterId* {. desc: static(