Retry private messages after Noise session replacement

This commit is contained in:
jack
2026-07-25 17:36:43 +02:00
committed by jack
parent 37e09393fa
commit ced9ac7610
10 changed files with 724 additions and 54 deletions
+230 -23
View File
@@ -107,12 +107,20 @@ final class MessageRouter {
private var bridgeDepositsInFlight = Set<String>() private var bridgeDepositsInFlight = Set<String>()
private var outbox: [PeerID: [QueuedMessage]] = [:] 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<String>()
// Outbox limits to prevent unbounded memory growth // Outbox limits to prevent unbounded memory growth
private static let maxMessagesPerPeer = 100 private static let maxMessagesPerPeer = 100
private static let messageTTLSeconds: TimeInterval = 24 * 60 * 60 // 24 hours private static let messageTTLSeconds: TimeInterval = 24 * 60 * 60 // 24 hours
// Bound resends of messages sent on a weak reachability signal that never // Bound actual sends that never receive an ack, whether they used weak
// get a delivery ack (e.g. peer on an old client that doesn't ack). // 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 private static let maxSendAttempts = 8
// Redundant couriers improve delivery odds; receivers dedup by message ID. // Redundant couriers improve delivery odds; receivers dedup by message ID.
private static let maxCouriersPerMessage = 3 private static let maxCouriersPerMessage = 3
@@ -171,15 +179,28 @@ final class MessageRouter {
// MARK: - Message Sending // MARK: - Message Sending
func sendPrivate(_ content: String, to peerID: PeerID, recipientNickname: String, messageID: String) { 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) { if let transport = connectedTransport(for: peerID), transport.canDeliverSecurely(to: peerID) {
// A live link that can complete an encrypted delivery is a // Even an established Noise session can be stale after the peer
// strong delivery signal; trust it outright. // 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) 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) transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID)
return return
} }
let message = QueuedMessage(content: content, nickname: recipientNickname, messageID: messageID, timestamp: now(), sendAttempts: 1)
if let transport = connectedTransport(for: peerID) { if let transport = connectedTransport(for: peerID) {
// "Connected" without an established secure session is forgeable: // "Connected" without an established secure session is forgeable:
// link bindings heal on signature-verified "direct" announces, but // link bindings heal on signature-verified "direct" announces, but
@@ -197,8 +218,9 @@ final class MessageRouter {
// deposit is cleared on ack. Don't "optimize" the courier call // deposit is cleared on ack. Don't "optimize" the courier call
// away. // away.
SecureLogger.debug("Routing PM via \(type(of: transport)) (connected, no secure session) to \(peerID.id.prefix(8))… id=\(messageID.prefix(8))", category: .session) 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) enqueue(message, for: peerID)
securelyTransmittedMessageIDs.remove(messageID)
transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID)
attemptCourierDeposit(messageID: messageID, for: peerID) attemptCourierDeposit(messageID: messageID, for: peerID)
return return
} }
@@ -209,8 +231,8 @@ final class MessageRouter {
// Send now, but retain a copy until a delivery/read ack clears it; // Send now, but retain a copy until a delivery/read ack clears it;
// receivers dedup resends by message ID. // 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) 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) enqueue(message, for: peerID)
transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID)
// "Reachable" without prompt delivery means the send only joined // "Reachable" without prompt delivery means the send only joined
// a queue (Nostr with relays down): also hand a sealed copy to // a queue (Nostr with relays down): also hand a sealed copy to
// any connected couriers rather than waiting for internet that // any connected couriers rather than waiting for internet that
@@ -356,15 +378,41 @@ final class MessageRouter {
// MARK: - Outbox Management // 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) { 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<PeerID>) {
guard !peerIDs.isEmpty else { return }
clearRetainedMessage(messageID, allowedPeerIDs: peerIDs)
}
private func clearRetainedMessage(
_ messageID: String,
allowedPeerIDs: Set<PeerID>?
) {
var cleared = false var cleared = false
for (peerID, queue) in outbox { for (peerID, queue) in outbox {
if let allowedPeerIDs, !allowedPeerIDs.contains(peerID) {
continue
}
let filtered = queue.filter { $0.messageID != messageID } let filtered = queue.filter { $0.messageID != messageID }
guard filtered.count != queue.count else { continue } guard filtered.count != queue.count else { continue }
outbox[peerID] = filtered.isEmpty ? nil : filtered outbox[peerID] = filtered.isEmpty ? nil : filtered
cleared = true 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 durable snapshot may still be hidden by protected data. Record
// the ack even when this cold-load view cannot find the message, then // the ack even when this cold-load view cannot find the message, then
// persist the current view so the store retains a removal tombstone. // persist the current view so the store retains a removal tombstone.
@@ -391,6 +439,11 @@ final class MessageRouter {
outbox[peerID] = filtered.isEmpty ? nil : filtered outbox[peerID] = filtered.isEmpty ? nil : filtered
cleared = true 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 // Preserve the scoped ack even when protected data hides the durable
// queue during a cold launch. // queue during a cold launch.
outboxStore?.recordRemoval(messageID: messageID, for: peerIDs) outboxStore?.recordRemoval(messageID: messageID, for: peerIDs)
@@ -424,7 +477,34 @@ final class MessageRouter {
persistOutbox() 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) { private func dropMessage(_ messageID: String, for peerID: PeerID) {
securelyTransmittedMessageIDs.remove(messageID)
metrics?.record(.outboxDropped) metrics?.record(.outboxDropped)
onMessageDropped?(messageID, peerID) onMessageDropped?(messageID, peerID)
} }
@@ -468,6 +548,7 @@ final class MessageRouter {
/// Panic wipe: forget queued mail on disk and in memory. /// Panic wipe: forget queued mail on disk and in memory.
func wipeOutbox() { func wipeOutbox() {
outbox.removeAll() outbox.removeAll()
securelyTransmittedMessageIDs.removeAll()
outboxStore?.wipe() outboxStore?.wipe()
} }
@@ -495,26 +576,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<PeerID>()
var retriedMessageIDs = Set<String>()
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) { func flushOutbox(for peerID: PeerID) {
guard let queued = outbox[peerID], !queued.isEmpty else { return } guard let queued = outbox[peerID], !queued.isEmpty else { return }
SecureLogger.debug("Flushing outbox for \(peerID.id.prefix(8))… count=\(queued.count)", category: .session) SecureLogger.debug("Flushing outbox for \(peerID.id.prefix(8))… count=\(queued.count)", category: .session)
let now = now() let now = now()
var remaining: [QueuedMessage] = [] var outboxChanged = false
for message in queued { 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) // Skip expired messages (TTL exceeded)
if now.timeIntervalSince(message.timestamp) > Self.messageTTLSeconds { 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) 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 continue
} }
if let transport = connectedTransport(for: peerID), transport.canDeliverSecurely(to: peerID) { 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) 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) transport.sendPrivateMessage(message.content, to: peerID, recipientNickname: message.nickname, messageID: message.messageID)
metrics?.record(.outboxResent) metrics?.record(.outboxResent)
outboxChanged = incrementSendAttemptsIfQueued(message.messageID, for: peerID) || outboxChanged
} else if let transport = connectedTransport(for: peerID) { } else if let transport = connectedTransport(for: peerID) {
// "Connected" without a secure session possibly a stolen // "Connected" without a secure session possibly a stolen
// binding from a replayed announce: send (a genuine link // binding from a replayed announce: send (a genuine link
@@ -527,9 +738,9 @@ final class MessageRouter {
// preserve. Retention stays bounded by the 24h outbox TTL // preserve. Retention stays bounded by the 24h outbox TTL
// and the per-peer FIFO cap. // 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) 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) transport.sendPrivateMessage(message.content, to: peerID, recipientNickname: message.nickname, messageID: message.messageID)
metrics?.record(.outboxResent) metrics?.record(.outboxResent)
remaining.append(message)
} else if let transport = reachableTransport(for: peerID) { } else if let transport = reachableTransport(for: peerID) {
// Reachability without a connection is a freshness heuristic, // Reachability without a connection is a freshness heuristic,
// so the send can silently go nowhere: send but keep retaining // so the send can silently go nowhere: send but keep retaining
@@ -537,26 +748,22 @@ final class MessageRouter {
// that never ack. // that never ack.
guard message.sendAttempts < Self.maxSendAttempts else { 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) 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 continue
} }
SecureLogger.debug("Outbox -> \(type(of: transport)) (reachable) for \(peerID.id.prefix(8))… id=\(message.messageID.prefix(8))", category: .session) 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) transport.sendPrivateMessage(message.content, to: peerID, recipientNickname: message.nickname, messageID: message.messageID)
metrics?.record(.outboxResent) metrics?.record(.outboxResent)
var retained = message outboxChanged = incrementSendAttemptsIfQueued(message.messageID, for: peerID) || outboxChanged
retained.sendAttempts += 1
remaining.append(retained)
} else {
remaining.append(message)
} }
} }
if remaining.isEmpty { if outboxChanged {
outbox.removeValue(forKey: peerID) persistOutbox()
} else {
outbox[peerID] = remaining
} }
persistOutbox()
} }
func flushAllOutbox() { func flushAllOutbox() {
@@ -33,6 +33,12 @@ protocol ChatDeliveryContext: AnyObject {
func notifyUIChanged() func notifyUIChanged()
/// Confirms receipt so the message router stops retaining the message for resend. /// Confirms receipt so the message router stops retaining the message for resend.
func markMessageDelivered(_ messageID: String) 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<PeerID>)
/// 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<PeerID>) -> Bool
} }
extension ChatViewModel: ChatDeliveryContext { extension ChatViewModel: ChatDeliveryContext {
@@ -56,6 +62,18 @@ extension ChatViewModel: ChatDeliveryContext {
func markMessageDelivered(_ messageID: String) { func markMessageDelivered(_ messageID: String) {
messageRouter.markDelivered(messageID) messageRouter.markDelivered(messageID)
} }
func markMessageDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {
messageRouter.markDelivered(messageID, from: peerIDs)
}
func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set<PeerID>) -> 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 /// Thin mapper from delivery events (read receipts, transport delivery
@@ -83,9 +101,10 @@ final class ChatDeliveryCoordinator {
@MainActor @MainActor
func didReceiveReadReceipt(_ receipt: ReadReceipt) { func didReceiveReadReceipt(_ receipt: ReadReceipt) {
updateMessageDeliveryStatus( updateAcknowledgedMessageDeliveryStatus(
receipt.originalMessageID, receipt.originalMessageID,
status: .read(by: receipt.readerNickname, at: receipt.timestamp) status: .read(by: receipt.readerNickname, at: receipt.timestamp),
from: [receipt.readerID]
) )
} }
@@ -102,17 +121,43 @@ final class ChatDeliveryCoordinator {
@MainActor @MainActor
@discardableResult @discardableResult
func updateMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { func updateMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool {
guard context.setDeliveryStatus(status, forMessageID: messageID) else {
return false
}
switch status { switch status {
case .delivered, .read: case .delivered, .read:
// Confirmed receipt stop retaining the message for resend. // Terminalize only after the store accepted the transition.
context.markMessageDelivered(messageID) context.markMessageDelivered(messageID)
default: default:
break 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<PeerID>
) -> Bool {
switch status {
case .delivered, .read:
break
default:
return false 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() context.notifyUIChanged()
return true return true
} }
@@ -62,10 +62,15 @@ protocol ChatTransportEventContext: AnyObject {
func sendMeshDeliveryAck(for messageID: String, to peerID: PeerID) func sendMeshDeliveryAck(for messageID: String, to peerID: PeerID)
// MARK: Delivery status // MARK: Delivery status
/// Applies the status to every known location of the message. /// Applies an authenticated receipt to the message only when it belongs
/// Returns `false` when no message with that ID was updated. /// to the supplied peer conversation aliases. Returns `false` for an
/// unknown ID, wrong peer, or rejected status transition.
@discardableResult @discardableResult
func applyMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool func applyAcknowledgedMessageDeliveryStatus(
_ messageID: String,
status: DeliveryStatus,
from peerIDAliases: Set<PeerID>
) -> Bool
func deliveryStatus(for messageID: String) -> DeliveryStatus? func deliveryStatus(for messageID: String) -> DeliveryStatus?
// MARK: Verification payloads // MARK: Verification payloads
@@ -122,8 +127,16 @@ extension ChatViewModel: ChatTransportEventContext {
} }
@discardableResult @discardableResult
func applyMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { func applyAcknowledgedMessageDeliveryStatus(
deliveryCoordinator.updateMessageDeliveryStatus(messageID, status: status) _ messageID: String,
status: DeliveryStatus,
from peerIDAliases: Set<PeerID>
) -> Bool {
deliveryCoordinator.updateAcknowledgedMessageDeliveryStatus(
messageID,
status: status,
from: peerIDAliases
)
} }
func deliveryStatus(for messageID: String) -> DeliveryStatus? { func deliveryStatus(for messageID: String) -> DeliveryStatus? {
@@ -364,9 +377,10 @@ private extension ChatTransportEventCoordinator {
guard let messageID = String(data: payload, encoding: .utf8) else { return } guard let messageID = String(data: payload, encoding: .utf8) else { return }
let name = deliveryStatusName(for: peerID, in: context) let name = deliveryStatusName(for: peerID, in: context)
let didUpdate = context.applyMessageDeliveryStatus( let didUpdate = context.applyAcknowledgedMessageDeliveryStatus(
messageID, messageID,
status: .delivered(to: name, at: Date()) status: .delivered(to: name, at: Date()),
from: receiptPeerAliases(for: peerID, in: context)
) )
if !didUpdate { if !didUpdate {
@@ -381,9 +395,10 @@ private extension ChatTransportEventCoordinator {
guard let messageID = String(data: payload, encoding: .utf8) else { return } guard let messageID = String(data: payload, encoding: .utf8) else { return }
let name = deliveryStatusName(for: peerID, in: context) let name = deliveryStatusName(for: peerID, in: context)
let didUpdate = context.applyMessageDeliveryStatus( let didUpdate = context.applyAcknowledgedMessageDeliveryStatus(
messageID, messageID,
status: .read(by: name, at: Date()) status: .read(by: name, at: Date()),
from: receiptPeerAliases(for: peerID, in: context)
) )
if !didUpdate { if !didUpdate {
@@ -414,4 +429,21 @@ private extension ChatTransportEventCoordinator {
func deliveryStatusName(for peerID: PeerID, in context: any ChatTransportEventContext) -> String { func deliveryStatusName(for peerID: PeerID, in context: any ChatTransportEventContext) -> String {
context.unifiedPeer(for: peerID)?.nickname ?? context.resolveNickname(for: peerID) context.unifiedPeer(for: peerID)?.nickname ?? context.resolveNickname(for: peerID)
} }
@MainActor
func receiptPeerAliases(
for peerID: PeerID,
in context: any ChatTransportEventContext
) -> Set<PeerID> {
var aliases: Set<PeerID> = [peerID]
// The active authenticated Noise key is authoritative. A cached
// ephemeralstable 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
}
} }
@@ -58,6 +58,10 @@ protocol ChatVerificationContext: AnyObject {
func noiseStaticPublicKeyData() -> Data func noiseStaticPublicKeyData() -> Data
func hasEstablishedNoiseSession(with peerID: PeerID) -> Bool func hasEstablishedNoiseSession(with peerID: PeerID) -> Bool
func triggerHandshake(with peerID: PeerID) func triggerHandshake(with 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 sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data)
func sendVerifyResponse(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) func sendVerifyResponse(to peerID: PeerID, noiseKeyHex: String, nonceA: Data)
@@ -116,6 +120,10 @@ extension ChatViewModel: ChatVerificationContext {
meshService.noiseStaticPublicKeyData() meshService.noiseStaticPublicKeyData()
} }
func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) {
messageRouter.retrySecurePrivateMessagesAfterAuthentication(for: peerIDAliases)
}
func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) { func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) {
meshService.sendVerifyChallenge(to: peerID, noiseKeyHex: noiseKeyHex, nonceA: nonceA) meshService.sendVerifyChallenge(to: peerID, noiseKeyHex: noiseKeyHex, nonceA: nonceA)
} }
@@ -206,16 +214,37 @@ final class ChatVerificationCoordinator {
self.context.invalidateEncryptionCache(for: peerID) self.context.invalidateEncryptionCache(for: peerID)
if self.context.cachedStablePeerID(for: peerID) == nil, var authenticatedStablePeerID: PeerID?
let keyData = self.context.noiseSessionPublicKeyData(for: peerID) { if let keyData = self.context.noiseSessionPublicKeyData(for: peerID) {
let stablePeerID = PeerID(hexData: keyData) 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( SecureLogger.debug(
"🗺️ Mapped short peerID to Noise key for header continuity: \(peerID) -> \(stablePeerID.id.prefix(8))", "🗺️ Mapped short peerID to Noise key for header continuity: \(peerID) -> \(stablePeerID.id.prefix(8))",
category: .session 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 { if var pending = self.pendingQRVerifications[peerID], pending.sent == false {
self.context.sendVerifyChallenge( self.context.sendVerifyChallenge(
to: peerID, to: peerID,
@@ -133,11 +133,17 @@ private final class MockChatTransportEventContext: ChatTransportEventContext {
// Delivery status // Delivery status
var applyMessageDeliveryStatusResult = true var applyMessageDeliveryStatusResult = true
var deliveryStatusesByMessageID: [String: DeliveryStatus] = [:] var deliveryStatusesByMessageID: [String: DeliveryStatus] = [:]
private(set) var appliedDeliveryStatuses: [(messageID: String, status: DeliveryStatus)] = [] private(set) var appliedDeliveryStatuses: [
(messageID: String, status: DeliveryStatus, peerIDAliases: Set<PeerID>)
] = []
@discardableResult @discardableResult
func applyMessageDeliveryStatus(_ messageID: String, status: DeliveryStatus) -> Bool { func applyAcknowledgedMessageDeliveryStatus(
appliedDeliveryStatuses.append((messageID, status)) _ messageID: String,
status: DeliveryStatus,
from peerIDAliases: Set<PeerID>
) -> Bool {
appliedDeliveryStatuses.append((messageID, status, peerIDAliases))
return applyMessageDeliveryStatusResult return applyMessageDeliveryStatusResult
} }
@@ -332,6 +338,10 @@ struct ChatTransportEventCoordinatorContextTests {
let peerID = PeerID(str: "99aabbccddeeff00") let peerID = PeerID(str: "99aabbccddeeff00")
let noiseKey = Data(repeating: 0x44, count: 32) let noiseKey = Data(repeating: 0x44, count: 32)
context.peersByID[peerID] = BitchatPeer(peerID: peerID, noisePublicKey: noiseKey, nickname: "alice") 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. // Inbound private message: decoded, handled, and delivery-acked.
let packet = PrivateMessagePacket(messageID: "pm-1", content: "hi there") let packet = PrivateMessagePacket(messageID: "pm-1", content: "hi there")
@@ -353,6 +363,8 @@ struct ChatTransportEventCoordinatorContextTests {
await drainMainActorTasks() await drainMainActorTasks()
#expect(context.appliedDeliveryStatuses.count == 2) #expect(context.appliedDeliveryStatuses.count == 2)
#expect(context.appliedDeliveryStatuses[0].messageID == "m-1") #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 { if case .delivered(let to, _) = context.appliedDeliveryStatuses[0].status {
#expect(to == "alice") #expect(to == "alice")
} else { } else {
@@ -97,6 +97,7 @@ private final class MockChatVerificationContext: ChatVerificationContext {
var noiseSessionKeysByPeerID: [PeerID: Data] = [:] var noiseSessionKeysByPeerID: [PeerID: Data] = [:]
private(set) var installedCallbacks: (onPeerAuthenticated: (PeerID, String) -> Void, onHandshakeRequired: (PeerID) -> Void)? private(set) var installedCallbacks: (onPeerAuthenticated: (PeerID, String) -> Void, onHandshakeRequired: (PeerID) -> Void)?
private(set) var triggeredHandshakes: [PeerID] = [] private(set) var triggeredHandshakes: [PeerID] = []
private(set) var securePrivateMessageRetryAliases: [[PeerID]] = []
private(set) var sentChallenges: [(peerID: PeerID, noiseKeyHex: String, nonceA: Data)] = [] private(set) var sentChallenges: [(peerID: PeerID, noiseKeyHex: String, nonceA: Data)] = []
private(set) var sentResponses: [(peerID: PeerID, noiseKeyHex: String, nonceA: Data)] = [] private(set) var sentResponses: [(peerID: PeerID, noiseKeyHex: String, nonceA: Data)] = []
@@ -113,6 +114,9 @@ private final class MockChatVerificationContext: ChatVerificationContext {
establishedNoiseSessions.contains(peerID) establishedNoiseSessions.contains(peerID)
} }
func triggerHandshake(with peerID: PeerID) { triggeredHandshakes.append(peerID) } func triggerHandshake(with peerID: PeerID) { triggeredHandshakes.append(peerID) }
func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) {
securePrivateMessageRetryAliases.append(peerIDAliases)
}
func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) { func sendVerifyChallenge(to peerID: PeerID, noiseKeyHex: String, nonceA: Data) {
sentChallenges.append((peerID, noiseKeyHex, nonceA)) sentChallenges.append((peerID, noiseKeyHex, nonceA))
@@ -265,6 +269,10 @@ struct ChatVerificationCoordinatorContextTests {
let peerID = PeerID(str: "1122334455667788") let peerID = PeerID(str: "1122334455667788")
let noiseKey = Data(repeating: 0x33, count: 32) let noiseKey = Data(repeating: 0x33, count: 32)
context.noiseSessionKeysByPeerID[peerID] = noiseKey context.noiseSessionKeysByPeerID[peerID] = noiseKey
context.cacheStablePeerID(
PeerID(hexData: Data(repeating: 0x44, count: 32)),
for: peerID
)
context.verifiedFingerprints = ["fp-verified"] context.verifiedFingerprints = ["fp-verified"]
coordinator.setupNoiseCallbacks() coordinator.setupNoiseCallbacks()
@@ -275,8 +283,10 @@ struct ChatVerificationCoordinatorContextTests {
callbacks?.onPeerAuthenticated(peerID, "fp-verified") callbacks?.onPeerAuthenticated(peerID, "fp-verified")
await waitForMainQueue() await waitForMainQueue()
#expect(context.encryptionStatuses[peerID] == .noiseVerified) #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.invalidatedEncryptionCachePeers.contains(peerID))
#expect(context.securePrivateMessageRetryAliases == [[peerID, stablePeerID]])
// Handshake required -> handshaking status. // Handshake required -> handshaking status.
callbacks?.onHandshakeRequired(peerID) callbacks?.onHandshakeRequired(peerID)
@@ -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 @Test @MainActor
func cleanupOldReadReceipts_removesReceiptIDsWithoutMessages() async { func cleanupOldReadReceipts_removesReceiptIDsWithoutMessages() async {
let (viewModel, transport) = makeTestableViewModel() let (viewModel, transport) = makeTestableViewModel()
@@ -577,6 +636,7 @@ private final class MockChatDeliveryContext: ChatDeliveryContext {
var isStartupPhase = false var isStartupPhase = false
private(set) var notifyUIChangedCount = 0 private(set) var notifyUIChangedCount = 0
private(set) var markedDeliveredMessageIDs: [String] = [] private(set) var markedDeliveredMessageIDs: [String] = []
private(set) var peerBoundDeliveredMessages: [(messageID: String, peerIDs: Set<PeerID>)] = []
@discardableResult @discardableResult
func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool { func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool {
@@ -604,6 +664,24 @@ private final class MockChatDeliveryContext: ChatDeliveryContext {
func markMessageDelivered(_ messageID: String) { func markMessageDelivered(_ messageID: String) {
markedDeliveredMessageIDs.append(messageID) markedDeliveredMessageIDs.append(messageID)
} }
func markMessageDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {
peerBoundDeliveredMessages.append((messageID, peerIDs))
}
func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set<PeerID>) -> 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 @MainActor
@@ -685,7 +763,63 @@ struct ChatDeliveryCoordinatorContextTests {
#expect(isRead(coordinator.deliveryStatus(for: messageID))) #expect(isRead(coordinator.deliveryStatus(for: messageID)))
#expect(context.notifyUIChangedCount == 1) #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 @Test @MainActor
+8
View File
@@ -58,6 +58,9 @@ final class MockTransport: Transport {
var peerNicknames: [PeerID: String] = [:] var peerNicknames: [PeerID: String] = [:]
var peerFingerprints: [PeerID: String] = [:] var peerFingerprints: [PeerID: String] = [:]
var peerNoiseStates: [PeerID: LazyHandshakeState] = [:] var peerNoiseStates: [PeerID: LazyHandshakeState] = [:]
/// 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() private let mockKeychain = MockKeychain()
// MARK: - Transport Protocol Implementation // MARK: - Transport Protocol Implementation
@@ -162,6 +165,11 @@ final class MockTransport: Transport {
func sendPrivateMessage(_ content: String, to peerID: PeerID, recipientNickname: String, messageID: String) { func sendPrivateMessage(_ content: String, to peerID: PeerID, recipientNickname: String, messageID: String) {
sentPrivateMessages.append((content, peerID, recipientNickname, messageID)) sentPrivateMessages.append((content, peerID, recipientNickname, messageID))
if let onSendPrivateMessage {
MainActor.assumeIsolated {
onSendPrivateMessage(messageID)
}
}
} }
func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) { func sendReadReceipt(_ receipt: ReadReceipt, to peerID: PeerID) {
@@ -690,6 +690,10 @@ private final class PerfDeliveryContext: ChatDeliveryContext {
func notifyUIChanged() {} func notifyUIChanged() {}
func markMessageDelivered(_ messageID: String) {} func markMessageDelivered(_ messageID: String) {}
func markMessageDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {}
func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set<PeerID>) -> Bool {
true
}
@discardableResult @discardableResult
func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool { func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool {
+199 -10
View File
@@ -79,18 +79,179 @@ struct MessageRouterTests {
} }
@Test @MainActor @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 peerID = PeerID(str: "0000000000000007")
let transport = MockTransport() let transport = MockTransport()
transport.connectedPeers.insert(peerID) transport.connectedPeers.insert(peerID)
transport.reachablePeers.insert(peerID) transport.reachablePeers.insert(peerID)
transport.securePeers = [peerID]
let router = MessageRouter(transports: [transport]) let router = MessageRouter(transports: [transport])
router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "m7") router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "m7")
#expect(transport.sentPrivateMessages.count == 1) #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) #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 @Test @MainActor
@@ -392,10 +553,31 @@ struct MessageRouterTests {
#expect(transport.sentPrivateMessages.count == 11) #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 @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 peerID = PeerID(str: "00000000000000ab")
let peerKey = Data(repeating: 0xAB, count: 32) let peerKey = Data(repeating: 0xAB, count: 32)
let courier = PeerID(str: "00000000000000cc") let courier = PeerID(str: "00000000000000cc")
@@ -415,7 +597,11 @@ struct MessageRouterTests {
#expect(transport.sentPrivateMessages.map(\.messageID) == ["cs2"]) #expect(transport.sentPrivateMessages.map(\.messageID) == ["cs2"])
#expect(transport.sentCourierMessages.isEmpty) #expect(transport.sentCourierMessages.isEmpty)
router.flushOutbox(for: peerID) 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 @Test @MainActor
@@ -853,12 +1039,14 @@ struct MessageRouterTests {
protectedDataUnavailable = false protectedDataUnavailable = false
restoredStore.retryDeferredLoad() // captures unseen durable + known wake restoredStore.retryDeferredLoad() // captures unseen durable + known wake
// Secure direct flush removes the wake message before recovery's // A secure direct retry followed by its delivery ack removes the wake
// MainActor merge. It must remain removed, while the unseen durable // message before recovery's MainActor merge. It must remain removed,
// message still arrives through the pending recovery claim. // while the unseen durable message still arrives through the pending
// recovery claim.
transport.connectedPeers.insert(peerID) transport.connectedPeers.insert(peerID)
transport.securePeers = [peerID] transport.securePeers = [peerID]
router.flushOutbox(for: peerID) router.flushOutbox(for: peerID)
router.markDelivered("recovery-gap-known")
await Task.yield() await Task.yield()
await Task.yield() await Task.yield()
@@ -929,7 +1117,8 @@ struct MessageRouterTests {
restoredStore.retryDeferredLoad() // persists D+W and queues recovery restoredStore.retryDeferredLoad() // persists D+W and queues recovery
transport.connectedPeers.insert(peerID) transport.connectedPeers.insert(peerID)
transport.securePeers = [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 // The gap save may remove W, but it must leave unseen D durable until
// MessageRouter receives the pending recovery callback. // MessageRouter receives the pending recovery callback.