From cb9c7f6e098a7450b61bc5a025f70ad0a06a5d4a Mon Sep 17 00:00:00 2001 From: NagyZoltanPeter <113987313+NagyZoltanPeter@users.noreply.github.com> Date: Thu, 16 Jul 2026 17:55:05 +0200 Subject: [PATCH] Refactor Messaging REST API to better match Messaging Send and Receive APIs --- .../messaging/rest_api/event_cache.nim | 6 +- .../messaging/rest_api/handlers.nim | 6 +- logos_delivery/messaging/rest_api/types.nim | 63 +++++++------------ tests/api/test_messaging_rest.nim | 2 +- 4 files changed, 30 insertions(+), 47 deletions(-) diff --git a/logos_delivery/messaging/rest_api/event_cache.nim b/logos_delivery/messaging/rest_api/event_cache.nim index 80a6cb47a..3326ade88 100644 --- a/logos_delivery/messaging/rest_api/event_cache.nim +++ b/logos_delivery/messaging/rest_api/event_cache.nim @@ -73,13 +73,11 @@ proc recordSend*( status[].events.add(record) proc recordReceived*( - self: MessagingEventCache, messageHash: string, message: MessagingMessage + 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, timestamp: getNowInNanosecondTime() - ) + ReceivedMessageRecord(messageHash: messageHash, message: message) ) while self.received.len > self.maxReceived: diff --git a/logos_delivery/messaging/rest_api/handlers.nim b/logos_delivery/messaging/rest_api/handlers.nim index 3bfdf8ca4..ea3d82001 100644 --- a/logos_delivery/messaging/rest_api/handlers.nim +++ b/logos_delivery/messaging/rest_api/handlers.nim @@ -53,7 +53,7 @@ proc installEventListeners(brokerCtx: BrokerContext, cache: MessagingEventCache) discard MessageReceivedEvent.listen( brokerCtx, proc(evt: MessageReceivedEvent): Future[void] {.async: (raises: []).} = - cache.recordReceived(evt.messageHash, toMessagingMessage(evt.message)), + cache.recordReceived(evt.messageHash, toRelayWakuMessage(evt.message)), ) proc installMessagingApiHandlers*(router: var RestRouter, client: MessagingClient) = @@ -105,7 +105,9 @@ proc installMessagingApiHandlers*(router: var RestRouter, client: MessagingClien contentBody: Option[ContentBody] ) -> RestApiResponse: ## Sends a message through the messaging client, returning the request id. - let req: MessagingMessage = decodeRequestBody[MessagingMessage](contentBody).valueOr: + let req: MessagingJsonEnvelope = decodeRequestBody[MessagingJsonEnvelope]( + contentBody + ).valueOr: return error let envelope = req.toMessageEnvelope().valueOr: diff --git a/logos_delivery/messaging/rest_api/types.nim b/logos_delivery/messaging/rest_api/types.nim index 31067c48e..c6845f183 100644 --- a/logos_delivery/messaging/rest_api/types.nim +++ b/logos_delivery/messaging/rest_api/types.nim @@ -10,30 +10,31 @@ import 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 +export types, relay_types #### Types -type MessagingMessage* = object - ## REST wire representation of a `MessageEnvelope`. `payload` is base64. +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*: Option[bool] meta*: Option[Base64String] -type - MessagingPostMessageRequest* = MessagingMessage +type MessagingPostMessageRequest* = MessagingJsonEnvelope - MessagingSendResponse* = object - ## Returned by the send endpoint; correlates with `MessageSentEvent` / - ## `MessageErrorEvent`. - requestId*: string +type MessagingSendResponse* = object + ## Returned by the send endpoint on success; correlates with + ## `MessageSentEvent` / `MessageErrorEvent`. + requestId*: string #### Type conversion -proc toMessageEnvelope*(msg: MessagingMessage): Result[MessageEnvelope, string] = +proc toMessageEnvelope*(msg: MessagingJsonEnvelope): Result[MessageEnvelope, string] = let payload = ?msg.payload.decode() meta = ?msg.meta.get(Base64String("")).decode() @@ -50,7 +51,7 @@ proc toMessageEnvelope*(msg: MessagingMessage): Result[MessageEnvelope, string] #### Serialization and deserialization proc writeValue*( - writer: var JsonWriter[RestJson], value: MessagingMessage + writer: var JsonWriter[RestJson], value: MessagingJsonEnvelope ) {.raises: [IOError].} = writer.beginRecord() writer.writeField("payload", value.payload) @@ -62,7 +63,7 @@ proc writeValue*( writer.endRecord() proc readValue*( - reader: var JsonReader[RestJson], value: var MessagingMessage + reader: var JsonReader[RestJson], value: var MessagingJsonEnvelope ) {.raises: [SerializationError, IOError].} = var payload = none(Base64String) @@ -79,7 +80,7 @@ proc readValue*( fmt"Multiple `{fieldName}` fields found" except CatchableError: "Multiple fields with the same name found" - reader.raiseUnexpectedField(err, "MessagingMessage") + reader.raiseUnexpectedField(err, "MessagingJsonEnvelope") case fieldName of "payload": @@ -99,7 +100,7 @@ proc readValue*( if contentTopic.isNone() or contentTopic.get().isEmptyOrWhitespace(): reader.raiseUnexpectedValue("Field `contentTopic` is missing or empty") - value = MessagingMessage( + value = MessagingJsonEnvelope( payload: payload.get(), contentTopic: contentTopic.get(), ephemeral: ephemeral, @@ -141,9 +142,10 @@ proc readValue*( #### Event observability DTOs ## -## Send-related events (sent / propagated / error) are grouped per request id; -## received messages are cached for polling. Both are populated by the broker -## listeners installed in the messaging REST handlers. +## 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 @@ -163,20 +165,7 @@ type ReceivedMessageRecord* = object messageHash*: string - message*: MessagingMessage - timestamp*: int64 ## nanoseconds, stamped when cached - -proc toMessagingMessage*(msg: WakuMessage): MessagingMessage = - MessagingMessage( - payload: base64.encode(msg.payload), - contentTopic: msg.contentTopic, - ephemeral: some(msg.ephemeral), - meta: - if msg.meta.len > 0: - some(base64.encode(msg.meta)) - else: - none(Base64String), - ) + message*: RelayWakuMessage ## the received WakuMessage, full fidelity #### Event DTO serialization @@ -205,7 +194,6 @@ proc writeValue*( writer.beginRecord() writer.writeField("messageHash", value.messageHash) writer.writeField("message", value.message) - writer.writeField("timestamp", value.timestamp) writer.endRecord() proc readValue*( @@ -274,20 +262,15 @@ proc readValue*( ) {.raises: [SerializationError, IOError].} = var messageHash = "" - message = MessagingMessage() - timestamp = int64(0) + message = RelayWakuMessage() for fieldName in readObjectFields(reader): case fieldName of "messageHash": messageHash = reader.readValue(string) of "message": - message = reader.readValue(MessagingMessage) - of "timestamp": - timestamp = reader.readValue(int64) + message = reader.readValue(RelayWakuMessage) else: unrecognizedFieldWarning(value) - value = ReceivedMessageRecord( - messageHash: messageHash, message: message, timestamp: timestamp - ) + value = ReceivedMessageRecord(messageHash: messageHash, message: message) diff --git a/tests/api/test_messaging_rest.nim b/tests/api/test_messaging_rest.nim index 91ebaf349..ccd7ad46b 100644 --- a/tests/api/test_messaging_rest.nim +++ b/tests/api/test_messaging_rest.nim @@ -63,7 +63,7 @@ suite "Messaging REST API": let subResp = await client.messagingPostSubscriptionsV1(@[contentTopic]) check subResp.status == 200 - let msg = MessagingMessage( + let msg = MessagingJsonEnvelope( payload: base64.encode("hello rest"), contentTopic: contentTopic, ephemeral: none(bool),