Fixes after rebase

This commit is contained in:
NagyZoltanPeter 2026-06-19 11:58:13 +02:00
parent 3006eefdbe
commit 4af5ee86fd
No known key found for this signature in database
GPG Key ID: 3E1F97CF4A7B6F42
5 changed files with 171 additions and 129 deletions

View File

@ -13,8 +13,7 @@
{.push raises: [].}
import bearssl/rand, std/times, chronos
import stew/byteutils
import libp2p/crypto/crypto, std/times, chronos
import std/hashes
import
@ -77,7 +76,7 @@ type
## ===== helpers =====
##
proc new*(T: typedesc[RequestId], rng: ref HmacDrbgContext): T =
proc new*(T: typedesc[RequestId], rng: crypto.Rng): T =
## Generate a new RequestId using the provided RNG.
RequestId(request_utils.generateRequestId(rng))
@ -91,10 +90,8 @@ proc hash*(r: RequestId): Hash =
## Allows `RequestId` to be used as a `Table` key.
hash(string(r))
proc generateRequestId*(rng: ref HmacDrbgContext): RequestId =
var bytes: array[10, byte]
hmacDrbgGenerate(rng[], bytes)
return RequestId(byteutils.toHex(bytes))
proc generateRequestId*(rng: crypto.Rng): RequestId =
RequestId(request_utils.generateRequestId(rng))
proc init*(
T: type MessageEnvelope,

View File

@ -22,7 +22,7 @@ import stew/byteutils
import libp2p/crypto/crypto as libp2p_crypto
import logos_delivery/api/messaging_client_interface as mci
import logos_delivery/waku/api/types
import logos_delivery/api/types
import logos_delivery/waku/waku_core/topics
import ./events
@ -359,7 +359,7 @@ proc dispatchRepair(self: ReliableChannel, wire: seq[byte]) {.async: (raises: []
let sendRes =
try:
await self.sendHandler(envelope)
await self.messagingClient.send(envelope)
except CatchableError as e:
Result[RequestId, string].err("messaging send raised: " & e.msg)
if sendRes.isErr():
@ -425,9 +425,9 @@ proc new*(
## `Decrypt` request brokers, so the channel keeps no per-instance
## encryption state either.
##
## `sendHandler` is the egress dispatch. The owning `ReliableChannelManager`
## typically constructs it as a closure over `MessagingClient.send`. Tests
## pass a fake to drive the send state machine without touching the network.
## `messagingClient` is the egress dispatch — the channel calls
## `messagingClient.send` to transmit. Tests pass a fake `MessagingClientInterface`
## to drive the send state machine without touching the network.
let chn = T(
messagingClient: messagingClient,
channelId: channelId,

View File

@ -140,7 +140,7 @@ BrokerImplement ReliableChannelManager of ReliableChannelManagerInterface:
let chn = self.channels.getOrDefault(channelId)
if chn.isNil():
return err("unknown channel: " & channelId)
return chn.send(appPayload, ephemeral)
return await chn.send(appPayload, ephemeral)
## Inbound messages are not handed to the manager by direct call. Each
## `ReliableChannel` installs its own `MessageReceivedEvent` listener

View File

@ -39,7 +39,7 @@ type RecvService* = ref object of RootObj
brokerCtx: BrokerContext
node: WakuNode
seenMsgListener: MessageSeenEventListener
connStatusListener: EventConnectionStatusChangeListener
connStatusListener: ConnectionStatusChangeEventListener
recentReceivedMsgs: seq[RecvMessage]
@ -200,17 +200,17 @@ proc startRecvService*(self: RecvService) =
error "Failed to set MessageSeenEvent listener", error = error
quit(QuitFailure)
self.connStatusListener = EventConnectionStatusChange.listen(
self.connStatusListener = ConnectionStatusChangeEvent.listen(
self.brokerCtx,
proc(event: EventConnectionStatusChange) {.async: (raises: []).} =
proc(event: ConnectionStatusChangeEvent) {.async: (raises: []).} =
self.onConnectionStatusChange(event.connectionStatus),
).valueOr:
error "Failed to set EventConnectionStatusChange listener", error = error
error "Failed to set ConnectionStatusChangeEvent listener", error = error
quit(QuitFailure)
proc stopRecvService*(self: RecvService) {.async.} =
await MessageSeenEvent.dropListener(self.brokerCtx, self.seenMsgListener)
await EventConnectionStatusChange.dropListener(
await ConnectionStatusChangeEvent.dropListener(
self.brokerCtx, self.connStatusListener
)
if not self.backfillHandler.isNil():

View File

@ -24,8 +24,6 @@ import sds
import snapshot_codec
import logos_delivery/api/messaging_client_interface
# for the `Send` request broker — send-side tests mock its provider to
# intercept `messagingClient.send` without touching the network.
const TestTimeout = chronos.seconds(15)
@ -177,7 +175,7 @@ suite "Reliable Channel - ingress":
suite "Reliable Channel - send state machine":
asyncTest "MessageSentEvent finalises the channelReqId as Sent":
## Drives the real send pipeline (`send` -> segmentation -> SDS ->
## rate_limit -> encrypt -> dispatch) via a fake `SendHandler` that
## rate_limit -> encrypt -> dispatch) via a mocked `Send` provider that
## returns a canned `RequestId` instead of hitting the network.
## Emitting the delivery-layer `MessageSentEvent` must drive the
## channel-level state machine through `Confirmed` and produce a
@ -204,13 +202,14 @@ suite "Reliable Channel - send state machine":
## Mock the messaging egress: override the `Send` provider the mounted
## MessagingClient installed, so the channel's `messagingClient.send`
## returns a canned RequestId without hitting the network.
Send.clearProvider(brokerCtx)
Send.setProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
sendCalls.inc
return ok(fakeMsgReqId),
).expect("set Send provider")
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
sendCalls.inc
return ok(fakeMsgReqId),
)
.expect("set Send provider")
discard (
await manager.createReliableChannel(
@ -253,7 +252,7 @@ suite "Reliable Channel - send state machine":
## Two `send()` calls -> two independent `channelReqId`s, each with
## one segment under the current segmentation skeleton
## (`performSegmentation` always emits exactly one segment). The
## fake `SendHandler` returns distinct `messagingReqId`s; finalising
## mocked `Send` provider returns distinct `messagingReqId`s; finalising
## the first emits `ChannelMessageSentEvent` for its `channelReqId`,
## finalising the second as a failure emits `ChannelMessageErrorEvent`
## for the other.
@ -275,14 +274,15 @@ suite "Reliable Channel - send state machine":
var msgReqIds: seq[RequestId]
## Mock the messaging egress (see the single-send case above).
Send.clearProvider(brokerCtx)
Send.setProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
let id = RequestId("fake-msg-req-" & $(msgReqIds.len + 1))
msgReqIds.add(id)
return ok(id),
).expect("set Send provider")
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
let id = RequestId("fake-msg-req-" & $(msgReqIds.len + 1))
msgReqIds.add(id)
return ok(id),
)
.expect("set Send provider")
discard (
await manager.createReliableChannel(
@ -311,10 +311,12 @@ suite "Reliable Channel - send state machine":
)
.expect("listen ChannelMessageErrorEvent")
let channelReqId1 =
(await manager.sendOnChannel(channelId, "first".toBytes(), false)).expect("send 1")
let channelReqId2 =
(await manager.sendOnChannel(channelId, "second".toBytes(), false)).expect("send 2")
let channelReqId1 = (
await manager.sendOnChannel(channelId, "first".toBytes(), false)
).expect("send 1")
let channelReqId2 = (
await manager.sendOnChannel(channelId, "second".toBytes(), false)
).expect("send 2")
let dispatchDeadline = Moment.now() + 1.seconds
while Moment.now() < dispatchDeadline and msgReqIds.len < 2:
@ -353,11 +355,11 @@ suite "Reliable Channel - send state machine":
## `channelReqId` cannot be produced through the real pipeline.
## Implement once segmentation does real chunking: send a payload
## larger than `DefaultSegmentSizeBytes`, capture the N
## `messagingReqId`s from a fake `SendHandler`, finalise some, and
## `messagingReqId`s from the mocked `Send` provider, finalise some, and
## assert prune only fires once every sibling is final.
skip()
asyncTest "sibling MessageSentEvent during sendHandler await does not corrupt state":
asyncTest "sibling MessageSentEvent during send await does not corrupt state":
## Regression test for the prune-during-await race
## (PR #3914 review comment r3324891059). Locks in that a sibling
## `MessageSentEvent` firing while `onReadyToSend` is paused at an
@ -385,23 +387,24 @@ suite "Reliable Channel - send state machine":
## event and then yields, so the listener task runs while the second
## segment is still mid-`await` in `onReadyToSend` — the exact race
## window the regression test targets.
Send.clearProvider(brokerCtx)
Send.setProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
let id = RequestId("race-msg-req-" & $(msgReqIds.len + 1))
msgReqIds.add(id)
if msgReqIds.len == 2:
waku_message_events.MessageSentEvent.emit(
brokerCtx,
waku_message_events.MessageSentEvent(
requestId: msgReqIds[0], messageHash: ""
),
)
await sleepAsync(50.milliseconds)
sendsReturned.inc()
return ok(id),
).expect("set Send provider")
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
let id = RequestId("race-msg-req-" & $(msgReqIds.len + 1))
msgReqIds.add(id)
if msgReqIds.len == 2:
waku_message_events.MessageSentEvent.emit(
brokerCtx,
waku_message_events.MessageSentEvent(
requestId: msgReqIds[0], messageHash: ""
),
)
await sleepAsync(50.milliseconds)
sendsReturned.inc()
return ok(id),
)
.expect("set Send provider")
discard (
await manager.createReliableChannel(
@ -423,8 +426,9 @@ suite "Reliable Channel - send state machine":
)
.expect("listen ChannelMessageSentEvent")
let channelReqId1 =
(await manager.sendOnChannel(channelId, "first".toBytes(), false)).expect("send 1")
let channelReqId1 = (
await manager.sendOnChannel(channelId, "first".toBytes(), false)
).expect("send 1")
## Drain the first segment fully before queueing the second, so
## the rate-limit FIFO between sibling sends isn't itself under
@ -434,8 +438,9 @@ suite "Reliable Channel - send state machine":
await sleepAsync(5.milliseconds)
check msgReqIds.len == 1
let channelReqId2 =
(await manager.sendOnChannel(channelId, "second".toBytes(), false)).expect("send 2")
let channelReqId2 = (
await manager.sendOnChannel(channelId, "second".toBytes(), false)
).expect("send 2")
## Wait until `sendOnChannel(m2)` has fully returned and yield once
## more so `onReadyToSend`'s post-await continuation gets a chance
@ -482,7 +487,9 @@ suite "Reliable Channel - SDS persistence":
var waku: Waku
var manager: ReliableChannelManager
var brokerCtx: BrokerContext
lockNewGlobalBrokerContext:
brokerCtx = globalBrokerContext()
waku = (await createNode(createApiNodeConf())).expect("createNode")
waku.mountMessagingClient().expect("mountMessagingClient")
waku.mountReliableChannelManager().expect("mountReliableChannelManager")
@ -490,18 +497,25 @@ suite "Reliable Channel - SDS persistence":
setNoopEncryption()
let fakeSend: SendHandler = proc(
env: MessageEnvelope
): Future[Result[RequestId, string]] {.async: (raises: [CatchableError]), gcsafe.} =
return ok(RequestId("persist-msg-req-1"))
discard manager
.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local"), sendHandler = fakeSend
## Mock the messaging egress so the channel's `messagingClient.send`
## returns a canned RequestId without hitting the network.
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
return ok(RequestId("persist-msg-req-1")),
)
.expect("createReliableChannel")
.expect("set Send provider")
discard (await manager.sendOnChannel(channelId, "persist me".toBytes(), false)).expect("send")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
discard (await manager.sendOnChannel(channelId, "persist me".toBytes(), false)).expect(
"send"
)
## Same handle the channel layer writes through (`openJob` is idempotent).
let job = persistency.openJob("sds").expect("openJob sds")
@ -562,9 +576,11 @@ suite "Reliable Channel - SDS lifecycle":
setNoopEncryption()
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var deliveries: seq[seq[byte]]
discard ChannelMessageReceivedEvent
@ -635,9 +651,11 @@ suite "Reliable Channel - SDS lifecycle":
setNoopEncryption()
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var deliveryCount = 0
discard ChannelMessageReceivedEvent
@ -694,9 +712,11 @@ suite "Reliable Channel - SDS lifecycle":
setNoopEncryption()
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var fired = false
discard ChannelMessageReceivedEvent
@ -756,9 +776,11 @@ suite "Reliable Channel - SDS lifecycle":
setNoopEncryption()
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var deliveryCount = 0
discard ChannelMessageReceivedEvent
@ -804,9 +826,11 @@ suite "Reliable Channel - SDS lifecycle":
(await manager.closeChannel(channelId)).expect("closeChannel")
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("re-createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("re-createReliableChannel")
## Replay the same envelope. Only a restored history suppresses it.
waku_message_events.MessageReceivedEvent.emit(
@ -841,18 +865,22 @@ suite "Reliable Channel - SDS protocol semantics":
setNoopEncryption()
var capturedWires: seq[seq[byte]]
let fakeSend: SendHandler = proc(
env: MessageEnvelope
): Future[Result[RequestId, string]] {.async: (raises: [CatchableError]), gcsafe.} =
## Noop encryption is identity, so the envelope payload IS the SDS wire.
capturedWires.add(env.payload)
return ok(RequestId("semantics-req-" & $capturedWires.len))
discard manager
.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local"), sendHandler = fakeSend
## Mock the messaging egress; noop encryption is identity, so the
## envelope payload IS the SDS wire.
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
capturedWires.add(env.payload)
return ok(RequestId("semantics-req-" & $capturedWires.len)),
)
.expect("createReliableChannel")
.expect("set Send provider")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
let remotePeer =
ReliabilityManager.new(SdsParticipantID("remote"), ReliabilityConfig.init())
@ -871,7 +899,8 @@ suite "Reliable Channel - SDS protocol semantics":
)
await sleepAsync(100.milliseconds)
discard (await manager.sendOnChannel(channelId, "reply".toBytes(), false)).expect("send")
discard
(await manager.sendOnChannel(channelId, "reply".toBytes(), false)).expect("send")
var deadline = Moment.now() + 2.seconds
while Moment.now() < deadline and capturedWires.len < 1:
await sleepAsync(5.milliseconds)
@ -911,19 +940,25 @@ suite "Reliable Channel - SDS protocol semantics":
setNoopEncryption()
var capturedWires: seq[seq[byte]]
let fakeSend: SendHandler = proc(
env: MessageEnvelope
): Future[Result[RequestId, string]] {.async: (raises: [CatchableError]), gcsafe.} =
capturedWires.add(env.payload)
return ok(RequestId("ack-req-" & $capturedWires.len))
discard manager
.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local"), sendHandler = fakeSend
## Mock the messaging egress (see other send-side cases).
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
capturedWires.add(env.payload)
return ok(RequestId("ack-req-" & $capturedWires.len)),
)
.expect("createReliableChannel")
.expect("set Send provider")
discard (await manager.sendOnChannel(channelId, "needs ack".toBytes(), false)).expect("send")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
discard (await manager.sendOnChannel(channelId, "needs ack".toBytes(), false)).expect(
"send"
)
var deadline = Moment.now() + 2.seconds
while Moment.now() < deadline and capturedWires.len < 1:
await sleepAsync(5.milliseconds)
@ -1000,9 +1035,11 @@ suite "Reliable Channel - SDS protocol semantics":
setNoopEncryption()
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var deliveries: seq[seq[byte]]
discard ChannelMessageReceivedEvent
@ -1077,9 +1114,11 @@ suite "Reliable Channel - SDS protocol semantics":
setNoopEncryption()
discard manager
.createReliableChannel(channelId, contentTopic, SdsParticipantID("local"))
.expect("createReliableChannel")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var deliveryCount = 0
discard ChannelMessageReceivedEvent
@ -1153,17 +1192,21 @@ suite "Reliable Channel - SDS protocol semantics":
setNoopEncryption()
var capturedWires: seq[seq[byte]]
let fakeSend: SendHandler = proc(
env: MessageEnvelope
): Future[Result[RequestId, string]] {.async: (raises: [CatchableError]), gcsafe.} =
capturedWires.add(env.payload)
return ok(RequestId("unique-req-" & $capturedWires.len))
discard manager
.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local"), sendHandler = fakeSend
## Mock the messaging egress (see other send-side cases).
Send
.replaceProvider(
brokerCtx,
proc(env: MessageEnvelope): Future[Result[RequestId, string]] {.async.} =
capturedWires.add(env.payload)
return ok(RequestId("unique-req-" & $capturedWires.len)),
)
.expect("createReliableChannel")
.expect("set Send provider")
discard (
await manager.createReliableChannel(
channelId, contentTopic, SdsParticipantID("local")
)
).expect("createReliableChannel")
var deliveries: seq[seq[byte]]
discard ChannelMessageReceivedEvent
@ -1218,7 +1261,9 @@ suite "Reliable Channel - SDS protocol semantics":
waku.mountReliableChannelManager().expect("mountReliableChannelManager")
manager = waku.reliableChannelManager
check (await manager.sendOnChannel(ChannelId("no-such-channel"), "x".toBytes(), false)).isErr()
check (
await manager.sendOnChannel(ChannelId("no-such-channel"), "x".toBytes(), false)
).isErr()
check (await manager.closeChannel(ChannelId("no-such-channel"))).isErr()
(await waku.stop()).expect("stop")