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 <noreply@anthropic.com>

* 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 <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com>
This commit is contained in:
NagyZoltanPeter 2026-07-17 18:14:52 +02:00 committed by GitHub
parent ce918b0819
commit 9827be5990
No known key found for this signature in database
GPG Key ID: B5690EEEBB952194
19 changed files with 1226 additions and 59 deletions

View File

@ -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 }}

View File

@ -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 <task>` 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

View File

@ -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()

View File

@ -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 = "../.."

View File

@ -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:../.."

View File

@ -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),
)
)

View File

@ -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.}

View File

@ -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()

View File

@ -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

View File

@ -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.

View File

@ -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
.}

View File

@ -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.}

View File

@ -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"

View File

@ -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)

View File

@ -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

View File

@ -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:

View File

@ -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

View File

@ -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

View File

@ -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(