mirror of
https://github.com/logos-messaging/logos-delivery.git
synced 2026-07-29 07:53:15 +00:00
route reliable channel send through waku.send
Reviewer flagged that constructing MessageEnvelope/DeliveryTask, minting the messagingReqId, poking deliveryTask.msg.meta, and calling sendService.send directly all reach past the messaging API into its internals — a layer-boundary violation. Replace that block in onReadyToSend with `await self.waku.send(envelope)` where `envelope.meta` carries the wire-format marker. The channel now holds a `Waku` instead of a `DeliveryService`. waku.send already handles requestId generation, subscription, and DeliveryTask construction. `waku.send` isn't `(raises: [])` so the await is wrapped in try/except to keep the listener annotation honest. Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
parent
6f1c0ffe18
commit
3f1a45c73c
@ -20,7 +20,8 @@ import bearssl/rand
|
|||||||
import stew/byteutils
|
import stew/byteutils
|
||||||
import libp2p/crypto/crypto as libp2p_crypto
|
import libp2p/crypto/crypto as libp2p_crypto
|
||||||
|
|
||||||
import waku/node/delivery_service/delivery_service
|
import waku/api/api
|
||||||
|
import waku/factory/waku as waku_factory
|
||||||
import waku/node/delivery_service/send_service
|
import waku/node/delivery_service/send_service
|
||||||
import waku/waku_core/topics
|
import waku/waku_core/topics
|
||||||
|
|
||||||
@ -31,8 +32,8 @@ import ./rate_limit_manager/rate_limit_manager
|
|||||||
import ./encryption/encryption
|
import ./encryption/encryption
|
||||||
|
|
||||||
export
|
export
|
||||||
delivery_service, send_service, events, segmentation, scalable_data_sync,
|
api, waku_factory, events, segmentation, scalable_data_sync, rate_limit_manager,
|
||||||
rate_limit_manager, encryption
|
encryption
|
||||||
|
|
||||||
const LipWireReliableChannelVersion* = "RELIABLE-CHANNEL-API/1"
|
const LipWireReliableChannelVersion* = "RELIABLE-CHANNEL-API/1"
|
||||||
## Wire-format spec marker for the Reliable Channel layer, as defined
|
## Wire-format spec marker for the Reliable Channel layer, as defined
|
||||||
@ -80,7 +81,7 @@ type
|
|||||||
## Spec-defined public type. Fields are private so callers cannot
|
## Spec-defined public type. Fields are private so callers cannot
|
||||||
## mutate internals and break invariants. Getters are added below
|
## mutate internals and break invariants. Getters are added below
|
||||||
## for the few values consumers may need.
|
## for the few values consumers may need.
|
||||||
deliveryService: DeliveryService
|
waku: Waku
|
||||||
channelId: ChannelId
|
channelId: ChannelId
|
||||||
contentTopic: ContentTopic
|
contentTopic: ContentTopic
|
||||||
senderId: SdsParticipantID
|
senderId: SdsParticipantID
|
||||||
@ -224,38 +225,41 @@ proc onReadyToSend(
|
|||||||
continue
|
continue
|
||||||
let wireBytes = seq[byte](encrypted)
|
let wireBytes = seq[byte](encrypted)
|
||||||
|
|
||||||
|
## The `meta` field carries the Reliable Channel wire-format spec
|
||||||
|
## marker so the ingress side of any peer can route this WakuMessage
|
||||||
|
## to its Reliable Channel layer.
|
||||||
let envelope = MessageEnvelope(
|
let envelope = MessageEnvelope(
|
||||||
contentTopic: self.contentTopic, payload: wireBytes, ephemeral: isEphemeral
|
contentTopic: self.contentTopic,
|
||||||
|
payload: wireBytes,
|
||||||
|
ephemeral: isEphemeral,
|
||||||
|
meta: LipWireReliableChannelVersion.toBytes(),
|
||||||
)
|
)
|
||||||
|
|
||||||
let messagingReqId = RequestId.new(self.rng)
|
## `waku.send` is not annotated `(raises: [])`, but this listener is.
|
||||||
let deliveryTask = DeliveryTask.new(
|
## Convert any raise to a Result error so the state machine handles
|
||||||
messagingReqId, envelope, self.brokerCtx
|
## both failure modes (Result.err and exception) through one path.
|
||||||
).valueOr:
|
let sendRes =
|
||||||
|
try:
|
||||||
|
await self.waku.send(envelope)
|
||||||
|
except CatchableError as e:
|
||||||
|
Result[RequestId, string].err("waku send raised: " & e.msg)
|
||||||
|
|
||||||
|
let messagingReqId = sendRes.valueOr:
|
||||||
MessageErrorEvent.emit(
|
MessageErrorEvent.emit(
|
||||||
self.brokerCtx,
|
self.brokerCtx,
|
||||||
MessageErrorEvent(
|
MessageErrorEvent(
|
||||||
requestId: channelReqId,
|
requestId: channelReqId,
|
||||||
messageHash: "",
|
messageHash: "",
|
||||||
error: "delivery task setup failed: " & error,
|
error: "waku send failed: " & error,
|
||||||
),
|
),
|
||||||
)
|
)
|
||||||
self.pendingMessagingRequests[idx].messagingReqId = some(messagingReqId)
|
|
||||||
self.pendingMessagingRequests[idx].segmentSendState = SegmentSendState.Failed
|
self.pendingMessagingRequests[idx].segmentSendState = SegmentSendState.Failed
|
||||||
self.pruneCompletedChannelReqs()
|
self.pruneCompletedChannelReqs()
|
||||||
idx.inc()
|
idx.inc()
|
||||||
continue
|
continue
|
||||||
|
|
||||||
## Stamp the Reliable Channel wire-format spec marker so the ingress
|
|
||||||
## side of any peer can route this WakuMessage to its Reliable
|
|
||||||
## Channel layer. Done on the constructed WakuMessage rather than
|
|
||||||
## via the envelope because `MessageEnvelope` does not expose a
|
|
||||||
## `meta` field.
|
|
||||||
deliveryTask.msg.meta = LipWireReliableChannelVersion.toBytes()
|
|
||||||
|
|
||||||
self.pendingMessagingRequests[idx].messagingReqId = some(messagingReqId)
|
self.pendingMessagingRequests[idx].messagingReqId = some(messagingReqId)
|
||||||
self.pendingMessagingRequests[idx].segmentSendState = SegmentSendState.InFlight
|
self.pendingMessagingRequests[idx].segmentSendState = SegmentSendState.InFlight
|
||||||
asyncSpawn self.deliveryService.sendService.send(deliveryTask)
|
|
||||||
self.requestIds.mgetOrPut(channelReqId, @[]).add(messagingReqId)
|
self.requestIds.mgetOrPut(channelReqId, @[]).add(messagingReqId)
|
||||||
idx.inc()
|
idx.inc()
|
||||||
|
|
||||||
@ -352,7 +356,7 @@ proc onMessageReceived(
|
|||||||
|
|
||||||
proc new*(
|
proc new*(
|
||||||
T: type ReliableChannel,
|
T: type ReliableChannel,
|
||||||
deliveryService: DeliveryService,
|
waku: Waku,
|
||||||
channelId: ChannelId,
|
channelId: ChannelId,
|
||||||
contentTopic: ContentTopic,
|
contentTopic: ContentTopic,
|
||||||
senderId: SdsParticipantID,
|
senderId: SdsParticipantID,
|
||||||
@ -368,7 +372,7 @@ proc new*(
|
|||||||
## `Decrypt` request brokers, so the channel keeps no per-instance
|
## `Decrypt` request brokers, so the channel keeps no per-instance
|
||||||
## encryption state either.
|
## encryption state either.
|
||||||
let chn = T(
|
let chn = T(
|
||||||
deliveryService: deliveryService,
|
waku: waku,
|
||||||
channelId: channelId,
|
channelId: channelId,
|
||||||
contentTopic: contentTopic,
|
contentTopic: contentTopic,
|
||||||
senderId: senderId,
|
senderId: senderId,
|
||||||
|
|||||||
@ -14,6 +14,7 @@ import waku/api/api
|
|||||||
import waku/api/api_conf
|
import waku/api/api_conf
|
||||||
import waku/events/message_events as waku_message_events
|
import waku/events/message_events as waku_message_events
|
||||||
import waku/factory/waku as waku_factory
|
import waku/factory/waku as waku_factory
|
||||||
|
import waku/node/delivery_service/delivery_service
|
||||||
import waku/waku_core/topics
|
import waku/waku_core/topics
|
||||||
|
|
||||||
import ./reliable_channel
|
import ./reliable_channel
|
||||||
@ -23,11 +24,10 @@ export reliable_channel
|
|||||||
|
|
||||||
type ReliableChannelManager* = ref object
|
type ReliableChannelManager* = ref object
|
||||||
channels: Table[ChannelId, ReliableChannel]
|
channels: Table[ChannelId, ReliableChannel]
|
||||||
deliveryService: DeliveryService
|
waku: Waku
|
||||||
## Owned by the manager. The ownership chain is
|
## Owned by the manager. The channel layer reaches the messaging
|
||||||
## ReliableChannelManager -> DeliveryService -> Waku -> WakuNode.
|
## API through `waku.send(envelope)`; constructing DeliveryTasks
|
||||||
## Hidden so callers can't substitute their own and bypass the
|
## directly would breach the layer boundary.
|
||||||
## manager's pipeline.
|
|
||||||
brokerCtx: BrokerContext
|
brokerCtx: BrokerContext
|
||||||
|
|
||||||
proc new*(
|
proc new*(
|
||||||
@ -38,14 +38,14 @@ proc new*(
|
|||||||
## TODO !! The proper ownership chain is:
|
## TODO !! The proper ownership chain is:
|
||||||
## ReliableChannelManager -> DeliveryService (MessagingClient) -> Waku (Kernel/Protocols) -> WakuNode,
|
## ReliableChannelManager -> DeliveryService (MessagingClient) -> Waku (Kernel/Protocols) -> WakuNode,
|
||||||
## and this will be implemented in the future. For now, `createNode`
|
## and this will be implemented in the future. For now, `createNode`
|
||||||
## is called here to get a DeliveryService instance, and the WakuNode is immediately discarded.
|
## is called here to get a Waku instance, and the WakuNode is immediately discarded.
|
||||||
## This is a temporary workaround to get the API
|
## This is a temporary workaround to get the API
|
||||||
|
|
||||||
let waku = ?(await createNode(conf))
|
let waku = ?(await createNode(conf))
|
||||||
|
|
||||||
let manager = T(
|
let manager = T(
|
||||||
channels: initTable[ChannelId, ReliableChannel](),
|
channels: initTable[ChannelId, ReliableChannel](),
|
||||||
deliveryService: waku.deliveryService,
|
waku: waku,
|
||||||
brokerCtx: brokerCtx,
|
brokerCtx: brokerCtx,
|
||||||
)
|
)
|
||||||
|
|
||||||
@ -55,11 +55,11 @@ proc start*(self: ReliableChannelManager): Result[void, string] =
|
|||||||
## Bring the owned DeliveryService up. Separated from `new` so callers
|
## Bring the owned DeliveryService up. Separated from `new` so callers
|
||||||
## can register encryption providers / create channels before traffic
|
## can register encryption providers / create channels before traffic
|
||||||
## starts flowing.
|
## starts flowing.
|
||||||
self.deliveryService.startDeliveryService()
|
self.waku.deliveryService.startDeliveryService()
|
||||||
|
|
||||||
proc stop*(self: ReliableChannelManager) {.async.} =
|
proc stop*(self: ReliableChannelManager) {.async.} =
|
||||||
if not self.deliveryService.isNil():
|
if not self.waku.isNil():
|
||||||
await self.deliveryService.stopDeliveryService()
|
await self.waku.deliveryService.stopDeliveryService()
|
||||||
|
|
||||||
proc getChannelForTest*(
|
proc getChannelForTest*(
|
||||||
self: ReliableChannelManager, channelId: ChannelId
|
self: ReliableChannelManager, channelId: ChannelId
|
||||||
@ -103,7 +103,7 @@ proc createReliableChannel*(
|
|||||||
)
|
)
|
||||||
|
|
||||||
let chn = ReliableChannel.new(
|
let chn = ReliableChannel.new(
|
||||||
deliveryService = self.deliveryService,
|
waku = self.waku,
|
||||||
channelId = channelId,
|
channelId = channelId,
|
||||||
contentTopic = contentTopic,
|
contentTopic = contentTopic,
|
||||||
senderId = senderId,
|
senderId = senderId,
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user