mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-26 19:25:23 +00:00
Retry private messages after Noise session replacement
This commit is contained in:
@@ -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 {
|
||||||
@@ -59,6 +65,18 @@ extension ChatViewModel: ChatDeliveryContext {
|
|||||||
messageID: messageID
|
messageID: 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
|
||||||
@@ -86,9 +104,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]
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -105,17 +124,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? {
|
||||||
@@ -446,9 +459,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 {
|
||||||
@@ -463,9 +477,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 {
|
||||||
@@ -503,4 +518,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
|
||||||
|
// 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
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -59,6 +59,10 @@ protocol ChatVerificationContext: AnyObject {
|
|||||||
func hasEstablishedNoiseSession(with peerID: PeerID) -> Bool
|
func hasEstablishedNoiseSession(with peerID: PeerID) -> Bool
|
||||||
func triggerHandshake(with peerID: PeerID)
|
func triggerHandshake(with peerID: PeerID)
|
||||||
func privateMediaPeerDidAuthenticate(_ 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 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)
|
||||||
|
|
||||||
@@ -121,6 +125,10 @@ extension ChatViewModel: ChatVerificationContext {
|
|||||||
mediaTransferCoordinator.peerDidAuthenticate(peerID.toShort())
|
mediaTransferCoordinator.peerDidAuthenticate(peerID.toShort())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
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)
|
||||||
}
|
}
|
||||||
@@ -216,16 +224,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
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -417,6 +423,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")
|
||||||
@@ -438,6 +448,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 {
|
||||||
|
|||||||
@@ -98,6 +98,7 @@ private final class MockChatVerificationContext: ChatVerificationContext {
|
|||||||
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 privateMediaAuthenticatedPeers: [PeerID] = []
|
private(set) var privateMediaAuthenticatedPeers: [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)] = []
|
||||||
|
|
||||||
@@ -118,6 +119,10 @@ private final class MockChatVerificationContext: ChatVerificationContext {
|
|||||||
privateMediaAuthenticatedPeers.append(peerID)
|
privateMediaAuthenticatedPeers.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))
|
||||||
}
|
}
|
||||||
@@ -269,6 +274,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()
|
||||||
@@ -279,9 +288,11 @@ 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.privateMediaAuthenticatedPeers == [peerID])
|
#expect(context.privateMediaAuthenticatedPeers == [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
|
||||||
|
|||||||
@@ -70,6 +70,9 @@ final class MockTransport: Transport, PrivateMediaDeletionPersisting {
|
|||||||
private var pendingDeletedPrivateMediaCompletions: [
|
private var pendingDeletedPrivateMediaCompletions: [
|
||||||
@MainActor (Bool) -> Void
|
@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()
|
private let mockKeychain = MockKeychain()
|
||||||
|
|
||||||
// MARK: - Transport Protocol Implementation
|
// MARK: - Transport Protocol Implementation
|
||||||
@@ -174,6 +177,11 @@ final class MockTransport: Transport, PrivateMediaDeletionPersisting {
|
|||||||
|
|
||||||
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) {
|
||||||
|
|||||||
@@ -749,6 +749,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 {
|
||||||
|
|||||||
@@ -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.
|
||||||
|
|||||||
Reference in New Issue
Block a user