diff --git a/bitchat/Services/MessageRouter.swift b/bitchat/Services/MessageRouter.swift index 4195658d..6c624c33 100644 --- a/bitchat/Services/MessageRouter.swift +++ b/bitchat/Services/MessageRouter.swift @@ -107,12 +107,20 @@ final class MessageRouter { private var bridgeDepositsInFlight = Set() private var outbox: [PeerID: [QueuedMessage]] = [:] + /// IDs whose latest router-owned transmission used an already-established + /// secure session and still awaits an ack. This deliberately excludes + /// messages handed to BLE while a handshake is pending: BLE owns those + /// sends and drains its queue after authentication, so retrying them here + /// would duplicate every normal first-handshake DM. + private var securelyTransmittedMessageIDs = Set() // Outbox limits to prevent unbounded memory growth private static let maxMessagesPerPeer = 100 private static let messageTTLSeconds: TimeInterval = 24 * 60 * 60 // 24 hours - // Bound resends of messages sent on a weak reachability signal that never - // get a delivery ack (e.g. peer on an old client that doesn't ack). + // Bound actual sends that never receive an ack, whether they used weak + // reachability or an apparently secure session that keeps being replaced. + // Connected pre-handshake sends are transport-owned and do not burn this + // cap because BLE queues/drains them itself. private static let maxSendAttempts = 8 // Redundant couriers improve delivery odds; receivers dedup by message ID. private static let maxCouriersPerMessage = 3 @@ -181,15 +189,28 @@ final class MessageRouter { // MARK: - Message Sending func sendPrivate(_ content: String, to peerID: PeerID, recipientNickname: String, messageID: String) { + let message = QueuedMessage( + content: content, + nickname: recipientNickname, + messageID: messageID, + timestamp: now(), + sendAttempts: 1 + ) + if let transport = connectedTransport(for: peerID), transport.canDeliverSecurely(to: peerID) { - // A live link that can complete an encrypted delivery is a - // strong delivery signal; trust it outright. + // Even an established Noise session can be stale after the peer + // restarts or replaces its app. Persist before handing the packet + // to the transport so a fast ack cannot race ahead of retention, + // then keep the copy until a delivery/read ack clears it. A + // replacement handshake will retry this same message ID, which + // receivers deduplicate. + enqueue(message, for: peerID) + securelyTransmittedMessageIDs.insert(messageID) SecureLogger.debug("Routing PM via \(type(of: transport)) (connected) to \(peerID.id.prefix(8))… id=\(messageID.prefix(8))…", category: .session) transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID) return } - let message = QueuedMessage(content: content, nickname: recipientNickname, messageID: messageID, timestamp: now(), sendAttempts: 1) if let transport = connectedTransport(for: peerID) { // "Connected" without an established secure session is forgeable: // link bindings heal on signature-verified "direct" announces, but @@ -207,8 +228,9 @@ final class MessageRouter { // deposit is cleared on ack. Don't "optimize" the courier call // away. SecureLogger.debug("Routing PM via \(type(of: transport)) (connected, no secure session) to \(peerID.id.prefix(8))… id=\(messageID.prefix(8))…", category: .session) - transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID) enqueue(message, for: peerID) + securelyTransmittedMessageIDs.remove(messageID) + transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID) attemptCourierDeposit(messageID: messageID, for: peerID) return } @@ -219,8 +241,8 @@ final class MessageRouter { // Send now, but retain a copy until a delivery/read ack clears it; // receivers dedup resends by message ID. SecureLogger.debug("Routing PM via \(type(of: transport)) (reachable) to \(peerID.id.prefix(8))… id=\(messageID.prefix(8))…", category: .session) - transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID) enqueue(message, for: peerID) + transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID) // "Reachable" without prompt delivery means the send only joined // a queue (Nostr with relays down): also hand a sealed copy to // any connected couriers rather than waiting for internet that @@ -366,15 +388,41 @@ final class MessageRouter { // MARK: - Outbox Management - /// A delivery or read ack confirms receipt; stop retaining the message. + /// A locally trusted delivery transition confirms receipt; stop retaining + /// every copy of the message. Authenticated remote receipts must use the + /// peer-bound overload below instead. func markDelivered(_ messageID: String) { + clearRetainedMessage(messageID, allowedPeerIDs: nil) + } + + /// Stops retaining a message only for the authenticated conversation + /// aliases that produced the accepted receipt. A peer that learns another + /// conversation's message ID cannot use it to clear that conversation's + /// retry state. + func markDelivered(_ messageID: String, from peerIDs: Set) { + guard !peerIDs.isEmpty else { return } + clearRetainedMessage(messageID, allowedPeerIDs: peerIDs) + } + + private func clearRetainedMessage( + _ messageID: String, + allowedPeerIDs: Set? + ) { var cleared = false for (peerID, queue) in outbox { + if let allowedPeerIDs, !allowedPeerIDs.contains(peerID) { + continue + } let filtered = queue.filter { $0.messageID != messageID } guard filtered.count != queue.count else { continue } outbox[peerID] = filtered.isEmpty ? nil : filtered cleared = true } + if !outbox.values.contains(where: { queue in + queue.contains { $0.messageID == messageID } + }) { + securelyTransmittedMessageIDs.remove(messageID) + } // The durable snapshot may still be hidden by protected data. Record // the ack even when this cold-load view cannot find the message, then // persist the current view so the store retains a removal tombstone. @@ -401,6 +449,11 @@ final class MessageRouter { outbox[peerID] = filtered.isEmpty ? nil : filtered cleared = true } + if !outbox.values.contains(where: { queue in + queue.contains { $0.messageID == messageID } + }) { + securelyTransmittedMessageIDs.remove(messageID) + } // Preserve the scoped ack even when protected data hides the durable // queue during a cold launch. outboxStore?.recordRemoval(messageID: messageID, for: peerIDs) @@ -434,7 +487,34 @@ final class MessageRouter { persistOutbox() } + @discardableResult + private func removeQueuedMessage(_ messageID: String, for peerID: PeerID) -> Bool { + guard let queue = outbox[peerID], + let index = queue.firstIndex(where: { $0.messageID == messageID }) else { + return false + } + var updated = queue + updated.remove(at: index) + outbox[peerID] = updated.isEmpty ? nil : updated + return true + } + + @discardableResult + private func incrementSendAttemptsIfQueued(_ messageID: String, for peerID: PeerID) -> Bool { + guard var queue = outbox[peerID], + let index = queue.firstIndex(where: { $0.messageID == messageID }) else { + // A synchronous delivery/read ack may have cleared the retained + // copy while `sendPrivateMessage` was on the stack. Never + // resurrect it from the flush snapshot. + return false + } + queue[index].sendAttempts += 1 + outbox[peerID] = queue + return true + } + private func dropMessage(_ messageID: String, for peerID: PeerID) { + securelyTransmittedMessageIDs.remove(messageID) metrics?.record(.outboxDropped) onMessageDropped?(messageID, peerID) } @@ -478,6 +558,7 @@ final class MessageRouter { /// Panic wipe: forget queued mail on disk and in memory. func wipeOutbox() { outbox.removeAll() + securelyTransmittedMessageIDs.removeAll() outboxStore?.wipe() } @@ -505,26 +586,156 @@ final class MessageRouter { } } + /// Retries only messages that the router previously transmitted through + /// an already-established secure session and that still await an ack. + /// + /// A peer restart can leave that local session looking usable until the + /// replacement handshake arrives; the first ciphertext is then + /// undecryptable remotely. Normal pre-handshake sends are intentionally + /// absent from `securelyTransmittedMessageIDs` because BLE already queues + /// and drains them when authentication completes. + func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) { + typealias Candidate = ( + peerID: PeerID, + message: QueuedMessage, + aliasOrder: Int, + queueOrder: Int + ) + + var visitedPeerIDs = Set() + var retriedMessageIDs = Set() + var outboxChanged = false + let currentDate = now() + var candidates: [Candidate] = [] + + for (aliasOrder, peerID) in peerIDAliases.enumerated() { + guard visitedPeerIDs.insert(peerID).inserted else { continue } + guard let queued = outbox[peerID], !queued.isEmpty, + let transport = connectedTransport(for: peerID), + transport.canDeliverSecurely(to: peerID) else { + continue + } + + for (queueOrder, message) in queued.enumerated() { + guard securelyTransmittedMessageIDs.contains(message.messageID) else { continue } + candidates.append(( + peerID: peerID, + message: message, + aliasOrder: aliasOrder, + queueOrder: queueOrder + )) + } + } + + // Conversation migration can leave retained messages split across the + // ephemeral and stable outbox keys. Merge both queues into one + // chronological stream so callback alias order cannot send newer mail + // ahead of older mail. + candidates.sort { lhs, rhs in + if lhs.message.timestamp != rhs.message.timestamp { + return lhs.message.timestamp < rhs.message.timestamp + } + if lhs.aliasOrder != rhs.aliasOrder { + return lhs.aliasOrder < rhs.aliasOrder + } + if lhs.queueOrder != rhs.queueOrder { + return lhs.queueOrder < rhs.queueOrder + } + return lhs.message.messageID < rhs.message.messageID + } + + for candidate in candidates { + let peerID = candidate.peerID + let message = candidate.message + guard retriedMessageIDs.insert(message.messageID).inserted, + securelyTransmittedMessageIDs.contains(message.messageID), + queuedMessage(message.messageID, for: peerID) != nil, + let transport = connectedTransport(for: peerID), + transport.canDeliverSecurely(to: peerID) else { + continue + } + + if currentDate.timeIntervalSince(message.timestamp) > Self.messageTTLSeconds { + if removeQueuedMessage(message.messageID, for: peerID) { + dropMessage(message.messageID, for: peerID) + outboxChanged = true + } + continue + } + + guard message.sendAttempts < Self.maxSendAttempts else { + SecureLogger.warning( + "📤 Dropping unacked PM for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))… after \(message.sendAttempts) secure attempts", + category: .session + ) + if removeQueuedMessage(message.messageID, for: peerID) { + dropMessage(message.messageID, for: peerID) + outboxChanged = true + } + continue + } + + SecureLogger.debug( + "Auth retry -> \(type(of: transport)) for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))…", + category: .session + ) + transport.sendPrivateMessage( + message.content, + to: peerID, + recipientNickname: message.nickname, + messageID: message.messageID + ) + metrics?.record(.outboxResent) + outboxChanged = incrementSendAttemptsIfQueued(message.messageID, for: peerID) || outboxChanged + } + + if outboxChanged { + persistOutbox() + } + } + func flushOutbox(for peerID: PeerID) { guard let queued = outbox[peerID], !queued.isEmpty else { return } SecureLogger.debug("Flushing outbox for \(peerID.id.prefix(8))… count=\(queued.count)", category: .session) let now = now() - var remaining: [QueuedMessage] = [] + var outboxChanged = false for message in queued { + // A synchronous ack from an earlier send in this flush may have + // removed an entry from the live outbox. The snapshot is only an + // iteration order; never use it to recreate removed messages. + guard queuedMessage(message.messageID, for: peerID) != nil else { continue } + // Skip expired messages (TTL exceeded) if now.timeIntervalSince(message.timestamp) > Self.messageTTLSeconds { SecureLogger.debug("⏰ Expired queued message for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))… (age: \(Int(now.timeIntervalSince(message.timestamp)))s)", category: .session) - dropMessage(message.messageID, for: peerID) + if removeQueuedMessage(message.messageID, for: peerID) { + dropMessage(message.messageID, for: peerID) + outboxChanged = true + } continue } if let transport = connectedTransport(for: peerID), transport.canDeliverSecurely(to: peerID) { - // Live link with a secure session: send and stop retaining. + // A secure session is meaningful enough to retry, but not + // proof that this particular ciphertext reached the peer: the + // remote app may have restarted while our old session still + // looked established. Retain until an ack, while bounding + // actual secure transmissions for peers that never ack. + guard message.sendAttempts < Self.maxSendAttempts else { + SecureLogger.warning("📤 Dropping unacked PM for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))… after \(message.sendAttempts) attempts", category: .session) + if removeQueuedMessage(message.messageID, for: peerID) { + dropMessage(message.messageID, for: peerID) + outboxChanged = true + } + continue + } SecureLogger.debug("Outbox -> \(type(of: transport)) (connected) for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))…", category: .session) + securelyTransmittedMessageIDs.insert(message.messageID) transport.sendPrivateMessage(message.content, to: peerID, recipientNickname: message.nickname, messageID: message.messageID) metrics?.record(.outboxResent) + outboxChanged = incrementSendAttemptsIfQueued(message.messageID, for: peerID) || outboxChanged } else if let transport = connectedTransport(for: peerID) { // "Connected" without a secure session — possibly a stolen // binding from a replayed announce: send (a genuine link @@ -537,9 +748,9 @@ final class MessageRouter { // preserve. Retention stays bounded by the 24h outbox TTL // and the per-peer FIFO cap. SecureLogger.debug("Outbox -> \(type(of: transport)) (connected, no secure session) for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))…", category: .session) + securelyTransmittedMessageIDs.remove(message.messageID) transport.sendPrivateMessage(message.content, to: peerID, recipientNickname: message.nickname, messageID: message.messageID) metrics?.record(.outboxResent) - remaining.append(message) } else if let transport = reachableTransport(for: peerID) { // Reachability without a connection is a freshness heuristic, // so the send can silently go nowhere: send but keep retaining @@ -547,26 +758,22 @@ final class MessageRouter { // that never ack. guard message.sendAttempts < Self.maxSendAttempts else { SecureLogger.warning("📤 Dropping unacked PM for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))… after \(message.sendAttempts) attempts", category: .session) - dropMessage(message.messageID, for: peerID) + if removeQueuedMessage(message.messageID, for: peerID) { + dropMessage(message.messageID, for: peerID) + outboxChanged = true + } continue } SecureLogger.debug("Outbox -> \(type(of: transport)) (reachable) for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))…", category: .session) transport.sendPrivateMessage(message.content, to: peerID, recipientNickname: message.nickname, messageID: message.messageID) metrics?.record(.outboxResent) - var retained = message - retained.sendAttempts += 1 - remaining.append(retained) - } else { - remaining.append(message) + outboxChanged = incrementSendAttemptsIfQueued(message.messageID, for: peerID) || outboxChanged } } - if remaining.isEmpty { - outbox.removeValue(forKey: peerID) - } else { - outbox[peerID] = remaining + if outboxChanged { + persistOutbox() } - persistOutbox() } func flushAllOutbox() { diff --git a/bitchat/ViewModels/ChatDeliveryCoordinator.swift b/bitchat/ViewModels/ChatDeliveryCoordinator.swift index 20049439..c071cef1 100644 --- a/bitchat/ViewModels/ChatDeliveryCoordinator.swift +++ b/bitchat/ViewModels/ChatDeliveryCoordinator.swift @@ -33,6 +33,12 @@ protocol ChatDeliveryContext: AnyObject { func notifyUIChanged() /// Confirms receipt so the message router stops retaining the message for resend. func markMessageDelivered(_ messageID: String) + /// Peer-bound form for authenticated remote receipts. Only the supplied + /// conversation aliases may have retained state terminalized. + func markMessageDelivered(_ messageID: String, from peerIDs: Set) + /// Returns true only when `messageID` is one of our outgoing messages in + /// at least one of the authenticated peer's direct-conversation aliases. + func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set) -> Bool } extension ChatViewModel: ChatDeliveryContext { @@ -59,6 +65,18 @@ extension ChatViewModel: ChatDeliveryContext { messageID: messageID ) } + + func markMessageDelivered(_ messageID: String, from peerIDs: Set) { + messageRouter.markDelivered(messageID, from: peerIDs) + } + + func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set) -> Bool { + peerIDs.contains { peerID in + privateMessages(for: peerID).contains { message in + message.id == messageID && message.senderPeerID == myPeerID + } + } + } } /// Thin mapper from delivery events (read receipts, transport delivery @@ -86,9 +104,10 @@ final class ChatDeliveryCoordinator { @MainActor func didReceiveReadReceipt(_ receipt: ReadReceipt) { - updateMessageDeliveryStatus( + updateAcknowledgedMessageDeliveryStatus( receipt.originalMessageID, - status: .read(by: receipt.readerNickname, at: receipt.timestamp) + status: .read(by: receipt.readerNickname, at: receipt.timestamp), + from: [receipt.readerID] ) } @@ -105,17 +124,43 @@ final class ChatDeliveryCoordinator { @MainActor @discardableResult func updateMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { + guard context.setDeliveryStatus(status, forMessageID: messageID) else { + return false + } switch status { case .delivered, .read: - // Confirmed receipt — stop retaining the message for resend. + // Terminalize only after the store accepted the transition. context.markMessageDelivered(messageID) default: break } + context.notifyUIChanged() + return true + } - guard context.setDeliveryStatus(status, forMessageID: messageID) else { + /// Applies an authenticated remote delivery/read receipt only when it + /// belongs to one of our outgoing messages in that peer's conversation. + /// Retry state is cleared after, never before, the status transition is + /// accepted by the store. + @MainActor + @discardableResult + func updateAcknowledgedMessageDeliveryStatus( + _ messageID: String, + status: DeliveryStatus, + from peerIDAliases: Set + ) -> Bool { + switch status { + case .delivered, .read: + break + default: return false } + guard !peerIDAliases.isEmpty, + context.isOutgoingPrivateMessage(messageID, toAny: peerIDAliases), + context.setDeliveryStatus(status, forMessageID: messageID) else { + return false + } + context.markMessageDelivered(messageID, from: peerIDAliases) context.notifyUIChanged() return true } diff --git a/bitchat/ViewModels/ChatTransportEventCoordinator.swift b/bitchat/ViewModels/ChatTransportEventCoordinator.swift index 72641812..ee2d8b6c 100644 --- a/bitchat/ViewModels/ChatTransportEventCoordinator.swift +++ b/bitchat/ViewModels/ChatTransportEventCoordinator.swift @@ -62,10 +62,15 @@ protocol ChatTransportEventContext: AnyObject { func sendMeshDeliveryAck(for messageID: String, to peerID: PeerID) // MARK: Delivery status - /// Applies the status to every known location of the message. - /// Returns `false` when no message with that ID was updated. + /// Applies an authenticated receipt to the message only when it belongs + /// to the supplied peer conversation aliases. Returns `false` for an + /// unknown ID, wrong peer, or rejected status transition. @discardableResult - func applyMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool + func applyAcknowledgedMessageDeliveryStatus( + _ messageID: String, + status: DeliveryStatus, + from peerIDAliases: Set + ) -> Bool func deliveryStatus(for messageID: String) -> DeliveryStatus? // MARK: Verification payloads @@ -122,8 +127,16 @@ extension ChatViewModel: ChatTransportEventContext { } @discardableResult - func applyMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { - deliveryCoordinator.updateMessageDeliveryStatus(messageID, status: status) + func applyAcknowledgedMessageDeliveryStatus( + _ messageID: String, + status: DeliveryStatus, + from peerIDAliases: Set + ) -> Bool { + deliveryCoordinator.updateAcknowledgedMessageDeliveryStatus( + messageID, + status: status, + from: peerIDAliases + ) } func deliveryStatus(for messageID: String) -> DeliveryStatus? { @@ -446,9 +459,10 @@ private extension ChatTransportEventCoordinator { guard let messageID = String(data: payload, encoding: .utf8) else { return } let name = deliveryStatusName(for: peerID, in: context) - let didUpdate = context.applyMessageDeliveryStatus( + let didUpdate = context.applyAcknowledgedMessageDeliveryStatus( messageID, - status: .delivered(to: name, at: Date()) + status: .delivered(to: name, at: Date()), + from: receiptPeerAliases(for: peerID, in: context) ) if !didUpdate { @@ -463,9 +477,10 @@ private extension ChatTransportEventCoordinator { guard let messageID = String(data: payload, encoding: .utf8) else { return } let name = deliveryStatusName(for: peerID, in: context) - let didUpdate = context.applyMessageDeliveryStatus( + let didUpdate = context.applyAcknowledgedMessageDeliveryStatus( messageID, - status: .read(by: name, at: Date()) + status: .read(by: name, at: Date()), + from: receiptPeerAliases(for: peerID, in: context) ) if !didUpdate { @@ -503,4 +518,21 @@ private extension ChatTransportEventCoordinator { func deliveryStatusName(for peerID: PeerID, in context: any ChatTransportEventContext) -> String { context.unifiedPeer(for: peerID)?.nickname ?? context.resolveNickname(for: peerID) } + + @MainActor + func receiptPeerAliases( + for peerID: PeerID, + in context: any ChatTransportEventContext + ) -> Set { + var aliases: Set = [peerID] + // The active authenticated Noise key is authoritative. A cached + // ephemeral→stable mapping can predate an identity replacement, so + // use it only when the live session cannot provide its static key. + if let keyData = context.noiseSessionPublicKeyData(for: peerID) { + aliases.insert(PeerID(hexData: keyData)) + } else if let stablePeerID = context.cachedStablePeerID(for: peerID) { + aliases.insert(stablePeerID) + } + return aliases + } } diff --git a/bitchat/ViewModels/ChatVerificationCoordinator.swift b/bitchat/ViewModels/ChatVerificationCoordinator.swift index e767a5e9..2de291ff 100644 --- a/bitchat/ViewModels/ChatVerificationCoordinator.swift +++ b/bitchat/ViewModels/ChatVerificationCoordinator.swift @@ -59,6 +59,10 @@ protocol ChatVerificationContext: AnyObject { func hasEstablishedNoiseSession(with peerID: PeerID) -> Bool func triggerHandshake(with peerID: PeerID) func privateMediaPeerDidAuthenticate(_ peerID: PeerID) + /// Retries only private messages previously transmitted through a secure + /// session and still pending an ack. Both ephemeral and stable aliases + /// are supplied because either can own the outbox entry. + func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) func sendVerifyResponse(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) @@ -121,6 +125,10 @@ extension ChatViewModel: ChatVerificationContext { mediaTransferCoordinator.peerDidAuthenticate(peerID.toShort()) } + func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) { + messageRouter.retrySecurePrivateMessagesAfterAuthentication(for: peerIDAliases) + } + func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) { meshService.sendVerifyChallenge(to: peerID, noiseKeyHex: noiseKeyHex, nonceA: nonceA) } @@ -216,16 +224,37 @@ final class ChatVerificationCoordinator { self.context.invalidateEncryptionCache(for: peerID) - if self.context.cachedStablePeerID(for: peerID) == nil, - let keyData = self.context.noiseSessionPublicKeyData(for: peerID) { + var authenticatedStablePeerID: PeerID? + if let keyData = self.context.noiseSessionPublicKeyData(for: peerID) { let stablePeerID = PeerID(hexData: keyData) - self.context.cacheStablePeerID(stablePeerID, for: peerID) + authenticatedStablePeerID = stablePeerID + if self.context.cachedStablePeerID(for: peerID) != stablePeerID { + // The freshly authenticated Noise key outranks a + // stale announce-derived alias. + self.context.cacheStablePeerID(stablePeerID, for: peerID) + } SecureLogger.debug( "🗺️ Mapped short peerID to Noise key for header continuity: \(peerID) -> \(stablePeerID.id.prefix(8))…", category: .session ) } + // A locally established session may have belonged to the + // peer's previous app process. The first ciphertext sent + // into that stale session is retained by MessageRouter; + // retry it now that this newly authenticated/replacement + // session can actually decrypt it. + var peerIDAliases = [peerID] + if let stablePeerID = authenticatedStablePeerID + ?? self.context.cachedStablePeerID(for: peerID), + stablePeerID != peerID { + // Conversations can migrate from the ephemeral BLE ID + // to the authenticated Noise-key ID. Retry both aliases + // because either may own the retained outbox entry. + peerIDAliases.append(stablePeerID) + } + self.context.retrySecurePrivateMessagesAfterAuthentication(for: peerIDAliases) + if var pending = self.pendingQRVerifications[peerID], pending.sent == false { self.context.sendVerifyChallenge( to: peerID, diff --git a/bitchatTests/ChatTransportEventCoordinatorContextTests.swift b/bitchatTests/ChatTransportEventCoordinatorContextTests.swift index 31bea94a..aa02b255 100644 --- a/bitchatTests/ChatTransportEventCoordinatorContextTests.swift +++ b/bitchatTests/ChatTransportEventCoordinatorContextTests.swift @@ -133,11 +133,17 @@ private final class MockChatTransportEventContext: ChatTransportEventContext { // Delivery status var applyMessageDeliveryStatusResult = true var deliveryStatusesByMessageID: [String: DeliveryStatus] = [:] - private(set) var appliedDeliveryStatuses: [(messageID: String, status: DeliveryStatus)] = [] + private(set) var appliedDeliveryStatuses: [ + (messageID: String, status: DeliveryStatus, peerIDAliases: Set) + ] = [] @discardableResult - func applyMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { - appliedDeliveryStatuses.append((messageID, status)) + func applyAcknowledgedMessageDeliveryStatus( + _ messageID: String, + status: DeliveryStatus, + from peerIDAliases: Set + ) -> Bool { + appliedDeliveryStatuses.append((messageID, status, peerIDAliases)) return applyMessageDeliveryStatusResult } @@ -417,6 +423,10 @@ struct ChatTransportEventCoordinatorContextTests { let peerID = PeerID(str: "99aabbccddeeff00") let noiseKey = Data(repeating: 0x44, count: 32) context.peersByID[peerID] = BitchatPeer(peerID: peerID, noisePublicKey: noiseKey, nickname: "alice") + let stablePeerID = PeerID(hexData: noiseKey) + let staleStablePeerID = PeerID(hexData: Data(repeating: 0x55, count: 32)) + context.cacheStablePeerID(staleStablePeerID, for: peerID) + context.noiseSessionKeysByPeerID[peerID] = noiseKey // Inbound private message: decoded, handled, and delivery-acked. let packet = PrivateMessagePacket(messageID: "pm-1", content: "hi there") @@ -438,6 +448,8 @@ struct ChatTransportEventCoordinatorContextTests { await drainMainActorTasks() #expect(context.appliedDeliveryStatuses.count == 2) #expect(context.appliedDeliveryStatuses[0].messageID == "m-1") + #expect(context.appliedDeliveryStatuses[0].peerIDAliases == [peerID, stablePeerID]) + #expect(!context.appliedDeliveryStatuses[0].peerIDAliases.contains(staleStablePeerID)) if case .delivered(let to, _) = context.appliedDeliveryStatuses[0].status { #expect(to == "alice") } else { diff --git a/bitchatTests/ChatVerificationCoordinatorContextTests.swift b/bitchatTests/ChatVerificationCoordinatorContextTests.swift index 59d04b9a..83cc7296 100644 --- a/bitchatTests/ChatVerificationCoordinatorContextTests.swift +++ b/bitchatTests/ChatVerificationCoordinatorContextTests.swift @@ -98,6 +98,7 @@ private final class MockChatVerificationContext: ChatVerificationContext { private(set) var installedCallbacks: (onPeerAuthenticated: (PeerID, String) -> Void, onHandshakeRequired: (PeerID) -> Void)? private(set) var triggeredHandshakes: [PeerID] = [] private(set) var privateMediaAuthenticatedPeers: [PeerID] = [] + private(set) var securePrivateMessageRetryAliases: [[PeerID]] = [] private(set) var sentChallenges: [(peerID: PeerID, noiseKeyHex: String, nonceA: Data)] = [] private(set) var sentResponses: [(peerID: PeerID, noiseKeyHex: String, nonceA: Data)] = [] @@ -118,6 +119,10 @@ private final class MockChatVerificationContext: ChatVerificationContext { privateMediaAuthenticatedPeers.append(peerID) } + func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) { + securePrivateMessageRetryAliases.append(peerIDAliases) + } + func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) { sentChallenges.append((peerID, noiseKeyHex, nonceA)) } @@ -269,6 +274,10 @@ struct ChatVerificationCoordinatorContextTests { let peerID = PeerID(str: "1122334455667788") let noiseKey = Data(repeating: 0x33, count: 32) context.noiseSessionKeysByPeerID[peerID] = noiseKey + context.cacheStablePeerID( + PeerID(hexData: Data(repeating: 0x44, count: 32)), + for: peerID + ) context.verifiedFingerprints = ["fp-verified"] coordinator.setupNoiseCallbacks() @@ -279,9 +288,11 @@ struct ChatVerificationCoordinatorContextTests { callbacks?.onPeerAuthenticated(peerID, "fp-verified") await waitForMainQueue() #expect(context.encryptionStatuses[peerID] == .noiseVerified) - #expect(context.stablePeerIDCache[peerID] == PeerID(hexData: noiseKey)) + let stablePeerID = PeerID(hexData: noiseKey) + #expect(context.stablePeerIDCache[peerID] == stablePeerID) #expect(context.invalidatedEncryptionCachePeers.contains(peerID)) #expect(context.privateMediaAuthenticatedPeers == [peerID]) + #expect(context.securePrivateMessageRetryAliases == [[peerID, stablePeerID]]) // Handshake required -> handshaking status. callbacks?.onHandshakeRequired(peerID) diff --git a/bitchatTests/ChatViewModelDeliveryStatusTests.swift b/bitchatTests/ChatViewModelDeliveryStatusTests.swift index fca755f2..1c898d13 100644 --- a/bitchatTests/ChatViewModelDeliveryStatusTests.swift +++ b/bitchatTests/ChatViewModelDeliveryStatusTests.swift @@ -298,6 +298,65 @@ struct ChatViewModelDeliveryStatusTests { }()) } + @Test @MainActor + func authenticatedNoiseAckCannotClearAnotherPeersRetryState() async { + let (viewModel, transport) = makeTestableViewModel() + let intendedPeer = PeerID(str: "0102030405060708") + let otherPeer = PeerID(str: "1112131415161718") + let messageID = "noise-peer-bound-ack" + + viewModel.seedPrivateChat( + [ + BitchatMessage( + id: messageID, + sender: viewModel.nickname, + content: "Keep retrying", + timestamp: Date(), + isRelay: false, + isPrivate: true, + recipientNickname: "Intended", + senderPeerID: viewModel.myPeerID, + deliveryStatus: .sent + ) + ], + for: intendedPeer + ) + transport.reachablePeers.insert(intendedPeer) + viewModel.messageRouter.sendPrivate( + "Keep retrying", + to: intendedPeer, + recipientNickname: "Intended", + messageID: messageID + ) + #expect(transport.sentPrivateMessages.count == 1) + + // This models a decrypted Noise receipt: the transport-authenticated + // peer is authoritative, not the attacker-controlled message ID. + viewModel.didReceiveNoisePayload( + from: otherPeer, + type: .delivered, + payload: Data(messageID.utf8), + timestamp: Date() + ) + for _ in 0..<10 { await Task.yield() } + + #expect(isSent(viewModel.conversations.deliveryStatus(forMessageID: messageID))) + viewModel.messageRouter.flushOutbox(for: intendedPeer) + #expect(transport.sentPrivateMessages.count == 2) + + viewModel.didReceiveNoisePayload( + from: intendedPeer, + type: .delivered, + payload: Data(messageID.utf8), + timestamp: Date() + ) + for _ in 0..<10 { await Task.yield() } + + #expect(isDelivered(viewModel.conversations.deliveryStatus(forMessageID: messageID))) + viewModel.messageRouter.flushOutbox(for: intendedPeer) + #expect(transport.sentPrivateMessages.count == 2) + } + @Test @MainActor func cleanupOldReadReceipts_removesReceiptIDsWithoutMessages() async { let (viewModel, transport) = makeTestableViewModel() @@ -577,6 +636,7 @@ private final class MockChatDeliveryContext: ChatDeliveryContext { var isStartupPhase = false private(set) var notifyUIChangedCount = 0 private(set) var markedDeliveredMessageIDs: [String] = [] + private(set) var peerBoundDeliveredMessages: [(messageID: String, peerIDs: Set)] = [] @discardableResult func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool { @@ -604,6 +664,24 @@ private final class MockChatDeliveryContext: ChatDeliveryContext { func markMessageDelivered(_ messageID: String) { markedDeliveredMessageIDs.append(messageID) } + + func markMessageDelivered(_ messageID: String, from peerIDs: Set) { + peerBoundDeliveredMessages.append((messageID, peerIDs)) + } + + func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set) -> Bool { + peerIDs.contains { peerID in + contextMessages(for: peerID).contains { message in + message.id == messageID && message.senderPeerID == localPeerID + } + } + } + + private let localPeerID = PeerID(str: "aabbccddeeff0011") + + private func contextMessages(for peerID: PeerID) -> [BitchatMessage] { + store.conversationsByID[.directPeer(peerID)]?.messages ?? [] + } } @MainActor @@ -685,7 +763,63 @@ struct ChatDeliveryCoordinatorContextTests { #expect(isRead(coordinator.deliveryStatus(for: messageID))) #expect(context.notifyUIChangedCount == 1) - #expect(context.markedDeliveredMessageIDs == [messageID]) + #expect(context.markedDeliveredMessageIDs.isEmpty) + #expect(context.peerBoundDeliveredMessages.count == 1) + #expect(context.peerBoundDeliveredMessages[0].messageID == messageID) + #expect(context.peerBoundDeliveredMessages[0].peerIDs == [peerID]) + } + + @Test @MainActor + func receiptFromWrongPeerDoesNotUpdateOrTerminalizeOutgoingMessage() async { + let context = MockChatDeliveryContext() + let coordinator = ChatDeliveryCoordinator(context: context) + let intendedPeer = PeerID(str: "0102030405060708") + let otherPeer = PeerID(str: "1112131415161718") + let messageID = "wrong-peer-receipt" + context.store.append( + makePrivateMessage(id: messageID, status: .sent), + to: .directPeer(intendedPeer) + ) + + coordinator.didReceiveReadReceipt( + ReadReceipt( + originalMessageID: messageID, + readerID: otherPeer, + readerNickname: "Other" + ) + ) + + #expect(isSent(coordinator.deliveryStatus(for: messageID))) + #expect(context.peerBoundDeliveredMessages.isEmpty) + #expect(context.markedDeliveredMessageIDs.isEmpty) + #expect(context.notifyUIChangedCount == 0) + } + + @Test @MainActor + func rejectedStaleReceiptDoesNotTerminalizeRetryState() async { + let context = MockChatDeliveryContext() + let coordinator = ChatDeliveryCoordinator(context: context) + let peerID = PeerID(str: "0102030405060708") + let messageID = "stale-receipt" + context.store.append( + makePrivateMessage( + id: messageID, + status: .read(by: "Peer", at: Date()) + ), + to: .directPeer(peerID) + ) + + let didUpdate = coordinator.updateAcknowledgedMessageDeliveryStatus( + messageID, + status: .delivered(to: "Peer", at: Date()), + from: [peerID] + ) + + #expect(!didUpdate) + #expect(isRead(coordinator.deliveryStatus(for: messageID))) + #expect(context.peerBoundDeliveredMessages.isEmpty) + #expect(context.markedDeliveredMessageIDs.isEmpty) + #expect(context.notifyUIChangedCount == 0) } @Test @MainActor diff --git a/bitchatTests/Mocks/MockTransport.swift b/bitchatTests/Mocks/MockTransport.swift index 2efe9f15..dcf8deee 100644 --- a/bitchatTests/Mocks/MockTransport.swift +++ b/bitchatTests/Mocks/MockTransport.swift @@ -70,6 +70,9 @@ final class MockTransport: Transport, PrivateMediaDeletionPersisting { private var pendingDeletedPrivateMediaCompletions: [ @MainActor (Bool) -> Void ] = [] + /// Optional synchronous hook for send-ordering tests (for example, an ack + /// arriving before the router's send call returns). + var onSendPrivateMessage: (@MainActor (_ messageID: String) -> Void)? private let mockKeychain = MockKeychain() // MARK: - Transport Protocol Implementation @@ -174,6 +177,11 @@ final class MockTransport: Transport, PrivateMediaDeletionPersisting { func sendPrivateMessage(_ content: String, to peerID: PeerID, recipientNickname: String, messageID: String) { sentPrivateMessages.append((content, peerID, recipientNickname, messageID)) + if let onSendPrivateMessage { + MainActor.assumeIsolated { + onSendPrivateMessage(messageID) + } + } } func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) { diff --git a/bitchatTests/Performance/PerformanceBaselineTests.swift b/bitchatTests/Performance/PerformanceBaselineTests.swift index 0cfb6f92..b102a69b 100644 --- a/bitchatTests/Performance/PerformanceBaselineTests.swift +++ b/bitchatTests/Performance/PerformanceBaselineTests.swift @@ -690,6 +690,10 @@ private final class PerfDeliveryContext: ChatDeliveryContext { func notifyUIChanged() {} func markMessageDelivered(_ messageID: String) {} + func markMessageDelivered(_ messageID: String, from peerIDs: Set) {} + func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set) -> Bool { + true + } @discardableResult func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool { diff --git a/bitchatTests/Services/MessageRouterTests.swift b/bitchatTests/Services/MessageRouterTests.swift index a45b0bcd..ccda1a92 100644 --- a/bitchatTests/Services/MessageRouterTests.swift +++ b/bitchatTests/Services/MessageRouterTests.swift @@ -79,18 +79,179 @@ struct MessageRouterTests { } @Test @MainActor - func sendPrivate_connectedSendIsNotRetained() async { + func peerBoundDeliveryAckCannotClearAnotherPeersRetainedMessage() async { + let intendedPeer = PeerID(str: "0000000000000023") + let otherPeer = PeerID(str: "0000000000000024") + let transport = MockTransport() + transport.reachablePeers = [intendedPeer, otherPeer] + + let router = MessageRouter(transports: [transport]) + router.sendPrivate( + "Secret", + to: intendedPeer, + recipientNickname: "Intended", + messageID: "peer-bound-ack" + ) + #expect(transport.sentPrivateMessages.count == 1) + + // Even a receipt arriving over another authenticated conversation + // must not terminalize the intended peer's retained retry. + router.markDelivered("peer-bound-ack", from: [otherPeer]) + router.flushOutbox(for: intendedPeer) + #expect(transport.sentPrivateMessages.count == 2) + + router.markDelivered("peer-bound-ack", from: [intendedPeer]) + router.flushOutbox(for: intendedPeer) + #expect(transport.sentPrivateMessages.count == 2) + } + + @Test @MainActor + func sendPrivate_connectedSecureSendRetainsUntilDeliveryAck() async { let peerID = PeerID(str: "0000000000000007") let transport = MockTransport() transport.connectedPeers.insert(peerID) transport.reachablePeers.insert(peerID) + transport.securePeers = [peerID] let router = MessageRouter(transports: [transport]) router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "m7") #expect(transport.sentPrivateMessages.count == 1) - router.flushOutbox(for: peerID) + // A newly authenticated/replacement session retries the retained + // message instead of losing the first ciphertext to a stale session. + router.retrySecurePrivateMessagesAfterAuthentication(for: [peerID]) + #expect(transport.sentPrivateMessages.count == 2) + + router.markDelivered("m7") + router.retrySecurePrivateMessagesAfterAuthentication(for: [peerID]) + #expect(transport.sentPrivateMessages.count == 2) + } + + @Test @MainActor + func authenticationRetry_matchesStableOutboxAliasWithoutDoubleSending() async { + let shortPeerID = PeerID(str: "0000000000000019") + let stablePeerID = PeerID(hexData: Data(repeating: 0x19, count: 32)) + let transport = MockTransport() + transport.connectedPeers.insert(stablePeerID) + transport.securePeers = [stablePeerID] + + let router = MessageRouter(transports: [transport]) + router.sendPrivate("Hello", to: stablePeerID, recipientNickname: "Peer", messageID: "alias-retry") + + router.retrySecurePrivateMessagesAfterAuthentication(for: [shortPeerID, stablePeerID, stablePeerID]) + + #expect(transport.sentPrivateMessages.map(\.messageID) == ["alias-retry", "alias-retry"]) + #expect(transport.sentPrivateMessages.allSatisfy { $0.peerID == stablePeerID }) + } + + @Test @MainActor + func authenticationRetry_preservesFIFOAcrossSplitAliases() async { + let shortPeerID = PeerID(str: "0000000000000022") + let stablePeerID = PeerID(hexData: Data(repeating: 0x22, count: 32)) + let transport = MockTransport() + transport.connectedPeers = [shortPeerID, stablePeerID] + transport.securePeers = [shortPeerID, stablePeerID] + let clock = MutableTestClock() + let router = MessageRouter(transports: [transport], now: { clock.now }) + + // The older message lives under the stable key, even though the auth + // callback supplies the ephemeral alias first. + router.sendPrivate("Older", to: stablePeerID, recipientNickname: "Peer", messageID: "fifo-old") + clock.now = clock.now.addingTimeInterval(1) + router.sendPrivate("Newer", to: shortPeerID, recipientNickname: "Peer", messageID: "fifo-new") + transport.resetRecordings() + + router.retrySecurePrivateMessagesAfterAuthentication(for: [shortPeerID, stablePeerID]) + + #expect(transport.sentPrivateMessages.map(\.messageID) == ["fifo-old", "fifo-new"]) + #expect(transport.sentPrivateMessages.map(\.peerID) == [stablePeerID, shortPeerID]) + } + + @Test @MainActor + func authenticationRetry_doesNotDuplicateNormalPendingHandshakeSend() async { + let peerID = PeerID(str: "0000000000000020") + let transport = MockTransport() + transport.connectedPeers.insert(peerID) + transport.securePeers = [] + + let router = MessageRouter(transports: [transport]) + router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "normal-handshake") #expect(transport.sentPrivateMessages.count == 1) + + // BLE owns this pending send and drains it after authentication. Once + // the session becomes secure, the router's targeted auth retry must + // stay silent instead of producing a second copy. + transport.securePeers = [peerID] + router.retrySecurePrivateMessagesAfterAuthentication(for: [peerID]) + #expect(transport.sentPrivateMessages.count == 1) + + router.markDelivered("normal-handshake") + } + + @Test @MainActor + func authenticationRetry_doesNotDuplicateMessageRequeuedByBLEForHandshake() async { + let peerID = PeerID(str: "0000000000000021") + let transport = MockTransport() + transport.connectedPeers.insert(peerID) + transport.securePeers = [peerID] + + let router = MessageRouter(transports: [transport]) + router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "session-lost") + #expect(transport.sentPrivateMessages.count == 1) + + // The session disappears before a normal outbox flush. That send is + // now owned by BLE's pending-handshake queue, so it clears the + // router's secure-auth retry marker. + transport.securePeers = [] + router.flushOutbox(for: peerID) + #expect(transport.sentPrivateMessages.count == 2) + + transport.securePeers = [peerID] + router.retrySecurePrivateMessagesAfterAuthentication(for: [peerID]) + #expect(transport.sentPrivateMessages.count == 2) + + router.markDelivered("session-lost") + } + + @Test @MainActor + func sendPrivate_fastDeliveryAckCannotRaceAheadOfRetention() async { + let peerID = PeerID(str: "0000000000000017") + let transport = MockTransport() + transport.connectedPeers.insert(peerID) + transport.securePeers = [peerID] + + let router = MessageRouter(transports: [transport]) + transport.onSendPrivateMessage = { messageID in + router.markDelivered(messageID) + } + + router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "fast-ack") + #expect(transport.sentPrivateMessages.map(\.messageID) == ["fast-ack"]) + + transport.onSendPrivateMessage = nil + router.flushOutbox(for: peerID) + #expect(transport.sentPrivateMessages.map(\.messageID) == ["fast-ack"]) + } + + @Test @MainActor + func flushOutbox_synchronousAckDoesNotResurrectSnapshotEntry() async { + let peerID = PeerID(str: "0000000000000018") + let transport = MockTransport() + transport.connectedPeers.insert(peerID) + transport.securePeers = [peerID] + + let router = MessageRouter(transports: [transport]) + router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "flush-fast-ack") + transport.onSendPrivateMessage = { messageID in + router.markDelivered(messageID) + } + + router.flushOutbox(for: peerID) + #expect(transport.sentPrivateMessages.count == 2) + + transport.onSendPrivateMessage = nil + router.flushOutbox(for: peerID) + #expect(transport.sentPrivateMessages.count == 2) } @Test @MainActor @@ -392,10 +553,31 @@ struct MessageRouterTests { #expect(transport.sentPrivateMessages.count == 11) } - /// With an established secure session the connected fast-path stays - /// exactly as before: trusted outright, no retained copy, no courier. @Test @MainActor - func sendPrivate_connectedWithSecureSessionIsTrustedOutright() async { + func authenticationRetry_capsActualSecureTransmissions() async { + let peerID = PeerID(str: "00000000000000ad") + let transport = MockTransport() + transport.connectedPeers.insert(peerID) + transport.securePeers = [peerID] + + let router = MessageRouter(transports: [transport]) + var dropped: [String] = [] + router.onMessageDropped = { messageID, _ in dropped.append(messageID) } + + router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "secure-retry") + for _ in 0..<10 { + router.retrySecurePrivateMessagesAfterAuthentication(for: [peerID]) + } + + #expect(dropped == ["secure-retry"]) + #expect(transport.sentPrivateMessages.count == 8) + } + + /// With an established secure session the connected fast-path sends + /// immediately and never leaks to couriers, but retains a local encrypted + /// outbox copy until the peer confirms receipt. + @Test @MainActor + func sendPrivate_connectedWithSecureSessionRetainsLocallyWithoutCourier() async { let peerID = PeerID(str: "00000000000000ab") let peerKey = Data(repeating: 0xAB, count: 32) let courier = PeerID(str: "00000000000000cc") @@ -415,7 +597,11 @@ struct MessageRouterTests { #expect(transport.sentPrivateMessages.map(\.messageID) == ["cs2"]) #expect(transport.sentCourierMessages.isEmpty) router.flushOutbox(for: peerID) - #expect(transport.sentPrivateMessages.count == 1) + #expect(transport.sentPrivateMessages.count == 2) + #expect(transport.sentCourierMessages.isEmpty) + router.markDelivered("cs2") + router.flushOutbox(for: peerID) + #expect(transport.sentPrivateMessages.count == 2) } @Test @MainActor @@ -853,12 +1039,14 @@ struct MessageRouterTests { protectedDataUnavailable = false restoredStore.retryDeferredLoad() // captures unseen durable + known wake - // Secure direct flush removes the wake message before recovery's - // MainActor merge. It must remain removed, while the unseen durable - // message still arrives through the pending recovery claim. + // A secure direct retry followed by its delivery ack removes the wake + // message before recovery's MainActor merge. It must remain removed, + // while the unseen durable message still arrives through the pending + // recovery claim. transport.connectedPeers.insert(peerID) transport.securePeers = [peerID] router.flushOutbox(for: peerID) + router.markDelivered("recovery-gap-known") await Task.yield() await Task.yield() @@ -929,7 +1117,8 @@ struct MessageRouterTests { restoredStore.retryDeferredLoad() // persists D+W and queues recovery transport.connectedPeers.insert(peerID) transport.securePeers = [peerID] - router.flushOutbox(for: peerID) // removes W before queued callback + router.flushOutbox(for: peerID) + router.markDelivered("recovery-write-failure-known") // removes W before queued callback // The gap save may remove W, but it must leave unseen D durable until // MessageRouter receives the pending recovery callback.