mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-26 13:05:21 +00:00
Scope authenticated delivery cleanup
This commit is contained in:
@@ -426,6 +426,40 @@ final class ConversationStore: ObservableObject {
|
|||||||
@discardableResult
|
@discardableResult
|
||||||
func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool {
|
func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool {
|
||||||
guard let ids = conversationIDsByMessageID[messageID] else { return false }
|
guard let ids = conversationIDsByMessageID[messageID] else { return false }
|
||||||
|
return applyDeliveryStatus(status, forMessageID: messageID, among: ids)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Applies an authenticated delivery/read receipt only to the supplied
|
||||||
|
/// direct-conversation aliases. A colliding message ID in another peer's
|
||||||
|
/// conversation (or a public timeline) must not inherit the receipt.
|
||||||
|
///
|
||||||
|
/// Stable and ephemeral aliases can temporarily hold the same message
|
||||||
|
/// instance during handoff. The shared helper republishes every targeted
|
||||||
|
/// alias even when the first mutation already changed that instance.
|
||||||
|
@discardableResult
|
||||||
|
func setDeliveryStatus(
|
||||||
|
_ status: DeliveryStatus,
|
||||||
|
forMessageID messageID: String,
|
||||||
|
inDirectPeerAliases peerIDs: Set<PeerID>
|
||||||
|
) -> Bool {
|
||||||
|
guard !peerIDs.isEmpty,
|
||||||
|
let indexedIDs = conversationIDsByMessageID[messageID] else {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
let allowedIDs = Set(peerIDs.map { ConversationID.directPeer($0) })
|
||||||
|
return applyDeliveryStatus(
|
||||||
|
status,
|
||||||
|
forMessageID: messageID,
|
||||||
|
among: indexedIDs.intersection(allowedIDs)
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
private func applyDeliveryStatus(
|
||||||
|
_ status: DeliveryStatus,
|
||||||
|
forMessageID messageID: String,
|
||||||
|
among ids: Set<ConversationID>
|
||||||
|
) -> Bool {
|
||||||
|
guard !ids.isEmpty else { return false }
|
||||||
var applied = false
|
var applied = false
|
||||||
var skipped: [ConversationID] = []
|
var skipped: [ConversationID] = []
|
||||||
for id in ids {
|
for id in ids {
|
||||||
|
|||||||
@@ -106,13 +106,14 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
dedupKey: String?,
|
dedupKey: String?,
|
||||||
operationID: UUID?
|
operationID: UUID?
|
||||||
)] = []
|
)] = []
|
||||||
/// Message IDs already published as drops (sender-side dedup) and drop
|
/// Opaque recipient/message keys already published as drops (sender-side
|
||||||
/// event IDs already handled (multi-relay dedup). Both persist across
|
/// dedup) and drop event IDs already handled (multi-relay dedup). Both
|
||||||
/// relaunches: relays hold drops for the full 24h NIP-40 window and the
|
/// persist across relaunches: relays hold drops for the full 24h NIP-40
|
||||||
/// persisted outbox keeps re-depositing, so in-memory-only dedup meant
|
/// window and the persisted outbox keeps re-depositing, so in-memory-only
|
||||||
/// every relaunch republished the same message as a fresh drop and every
|
/// dedup meant every relaunch republished the same message as a fresh drop
|
||||||
/// gateway relaunch re-delivered the whole backlog (field-verified
|
/// and every gateway relaunch re-delivered the whole backlog
|
||||||
/// amplification storm). Entries age out with the 24h drop window.
|
/// (field-verified amplification storm). Entries age out with the 24h drop
|
||||||
|
/// window.
|
||||||
private var publishedDropKeys: ExpiringIDSet
|
private var publishedDropKeys: ExpiringIDSet
|
||||||
private var seenDropEventIDs: ExpiringIDSet
|
private var seenDropEventIDs: ExpiringIDSet
|
||||||
private var subscriptionOpen = false
|
private var subscriptionOpen = false
|
||||||
@@ -126,7 +127,7 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
}
|
}
|
||||||
/// Sender operations queued locally or awaiting relay confirmation.
|
/// Sender operations queued locally or awaiting relay confirmation.
|
||||||
/// The per-attempt ID prevents a stale pre-wipe callback from completing
|
/// The per-attempt ID prevents a stale pre-wipe callback from completing
|
||||||
/// a newer attempt for the same message.
|
/// a newer attempt for the same recipient/message pair.
|
||||||
private var activeDropOperations: [String: ActiveDropOperation] = [:]
|
private var activeDropOperations: [String: ActiveDropOperation] = [:]
|
||||||
/// Held-envelope publishes have no sender message ID, but still need an
|
/// Held-envelope publishes have no sender message ID, but still need an
|
||||||
/// in-flight identity: repeated refreshes inside the relay-OK wait window
|
/// in-flight identity: repeated refreshes inside the relay-OK wait window
|
||||||
@@ -214,11 +215,46 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
|
|
||||||
// MARK: - Sender role
|
// MARK: - Sender role
|
||||||
|
|
||||||
|
/// Stable, opaque sender-side dedup key. Recipient scope prevents one
|
||||||
|
/// conversation's colliding message ID from suppressing another's drop,
|
||||||
|
/// while hashing keeps recipient keys out of the persisted snapshot.
|
||||||
|
private static func senderDropKey(
|
||||||
|
messageID: String,
|
||||||
|
recipientNoiseKey: Data
|
||||||
|
) -> String {
|
||||||
|
var material = Data("bitchat-bridge-drop-dedup-v2".utf8)
|
||||||
|
appendLengthPrefixed(Data(messageID.utf8), to: &material)
|
||||||
|
appendLengthPrefixed(recipientNoiseKey, to: &material)
|
||||||
|
return "v2:\(material.sha256Hex())"
|
||||||
|
}
|
||||||
|
|
||||||
|
private static func appendLengthPrefixed(_ value: Data, to output: inout Data) {
|
||||||
|
let length = UInt32(value.count)
|
||||||
|
output.append(UInt8((length >> 24) & 0xFF))
|
||||||
|
output.append(UInt8((length >> 16) & 0xFF))
|
||||||
|
output.append(UInt8((length >> 8) & 0xFF))
|
||||||
|
output.append(UInt8(length & 0xFF))
|
||||||
|
output.append(value)
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Previous releases persisted raw message IDs without recipient scope.
|
||||||
|
/// They remain conservative wildcards for their original 24-hour
|
||||||
|
/// lifetime: assigning one to a recipient would be guesswork and could
|
||||||
|
/// republish the original drop. New acceptances persist only v2 keys.
|
||||||
|
private func wasPublished(
|
||||||
|
legacyMessageID: String,
|
||||||
|
dedupKey: String,
|
||||||
|
now date: Date
|
||||||
|
) -> Bool {
|
||||||
|
publishedDropKeys.contains(dedupKey, now: date)
|
||||||
|
|| publishedDropKeys.contains(legacyMessageID, now: date)
|
||||||
|
}
|
||||||
|
|
||||||
/// Parallel-deposit a sealed copy of an outbound private message as a
|
/// Parallel-deposit a sealed copy of an outbound private message as a
|
||||||
/// relay drop. Called by the message router alongside physical courier
|
/// relay drop. Called by the message router alongside physical courier
|
||||||
/// deposits; idempotent per message ID. Completion becomes true only
|
/// deposits; idempotent per recipient/message pair. Completion becomes
|
||||||
/// after a real relay acceptance arrives, which is when the router may
|
/// true only after a real relay acceptance arrives, which is when the
|
||||||
/// show "carried".
|
/// router may show "carried".
|
||||||
func depositDrop(
|
func depositDrop(
|
||||||
content: String,
|
content: String,
|
||||||
messageID: String,
|
messageID: String,
|
||||||
@@ -229,9 +265,18 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
completion(false)
|
completion(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
guard !publishedDropKeys.contains(messageID, now: now()),
|
let date = now()
|
||||||
activeDropOperations[messageID] == nil,
|
let dedupKey = Self.senderDropKey(
|
||||||
!rejectedDropKeys.contains(messageID, now: now()) else {
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: recipientNoiseKey
|
||||||
|
)
|
||||||
|
guard !wasPublished(
|
||||||
|
legacyMessageID: messageID,
|
||||||
|
dedupKey: dedupKey,
|
||||||
|
now: date
|
||||||
|
),
|
||||||
|
activeDropOperations[dedupKey] == nil,
|
||||||
|
!rejectedDropKeys.contains(dedupKey, now: date) else {
|
||||||
completion(false)
|
completion(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -244,13 +289,13 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
// of the sealing); suppress it in-memory so the retry sweep does not
|
// of the sealing); suppress it in-memory so the retry sweep does not
|
||||||
// churn, but never persist it as a published drop.
|
// churn, but never persist it as a published drop.
|
||||||
guard let encoded = envelope.encode(), encoded.count <= Limits.maxDropEnvelopeBytes else {
|
guard let encoded = envelope.encode(), encoded.count <= Limits.maxDropEnvelopeBytes else {
|
||||||
rejectedDropKeys.insert(messageID, now: now())
|
rejectedDropKeys.insert(dedupKey, now: date)
|
||||||
completion(false)
|
completion(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
let operationID = UUID()
|
let operationID = UUID()
|
||||||
activeDropOperations[messageID] = ActiveDropOperation(id: operationID, completion: completion)
|
activeDropOperations[dedupKey] = ActiveDropOperation(id: operationID, completion: completion)
|
||||||
publishDrop(envelope, messageID: messageID, operationID: operationID)
|
publishDrop(envelope, dedupKey: dedupKey, operationID: operationID)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Publishes held envelopes (mail we carry for others) as drops,
|
/// Publishes held envelopes (mail we carry for others) as drops,
|
||||||
@@ -272,13 +317,14 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Publishes a drop, or queues it when relays are down. `messageID` is the
|
/// Publishes a drop, or queues it when relays are down. `dedupKey` is the
|
||||||
/// sender-side dedup key (nil for held/relayed envelopes we don't track);
|
/// opaque sender-side recipient/message key (nil for held/relayed
|
||||||
/// it rides the pending queue so an evicted or failed drop can release its
|
/// envelopes we don't track); it rides the pending queue so an evicted or
|
||||||
/// in-flight slot. Completion reports actual NIP-01 relay acceptance.
|
/// failed drop can release its in-flight slot. Completion reports actual
|
||||||
|
/// NIP-01 relay acceptance.
|
||||||
private func publishDrop(
|
private func publishDrop(
|
||||||
_ envelope: CourierEnvelope,
|
_ envelope: CourierEnvelope,
|
||||||
messageID: String? = nil,
|
dedupKey: String? = nil,
|
||||||
operationID: UUID? = nil,
|
operationID: UUID? = nil,
|
||||||
untrackedCompletion: (@MainActor (Bool) -> Void)? = nil
|
untrackedCompletion: (@MainActor (Bool) -> Void)? = nil
|
||||||
) {
|
) {
|
||||||
@@ -286,7 +332,7 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
encoded.count <= Limits.maxDropEnvelopeBytes,
|
encoded.count <= Limits.maxDropEnvelopeBytes,
|
||||||
!envelope.isExpired else {
|
!envelope.isExpired else {
|
||||||
finishPublish(
|
finishPublish(
|
||||||
messageID: messageID,
|
dedupKey: dedupKey,
|
||||||
operationID: operationID,
|
operationID: operationID,
|
||||||
succeeded: false,
|
succeeded: false,
|
||||||
untrackedCompletion: untrackedCompletion
|
untrackedCompletion: untrackedCompletion
|
||||||
@@ -297,15 +343,15 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
// Held mail remains in CourierStore and has no sender operation to
|
// Held mail remains in CourierStore and has no sender operation to
|
||||||
// recover after an in-memory queue loss. Leave its cooldown unset
|
// recover after an in-memory queue loss. Leave its cooldown unset
|
||||||
// and let the next connected refresh offer it again.
|
// and let the next connected refresh offer it again.
|
||||||
guard messageID != nil else {
|
guard dedupKey != nil else {
|
||||||
untrackedCompletion?(false)
|
untrackedCompletion?(false)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
pendingDrops.append((envelope, messageID, operationID))
|
pendingDrops.append((envelope, dedupKey, operationID))
|
||||||
while pendingDrops.count > Limits.maxPendingDrops {
|
while pendingDrops.count > Limits.maxPendingDrops {
|
||||||
let evicted = pendingDrops.removeFirst()
|
let evicted = pendingDrops.removeFirst()
|
||||||
finishPublish(
|
finishPublish(
|
||||||
messageID: evicted.dedupKey,
|
dedupKey: evicted.dedupKey,
|
||||||
operationID: evicted.operationID,
|
operationID: evicted.operationID,
|
||||||
succeeded: false
|
succeeded: false
|
||||||
)
|
)
|
||||||
@@ -321,7 +367,7 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
) else {
|
) else {
|
||||||
SecureLogger.error("📦🌉 Failed to compose courier drop", category: .encryption)
|
SecureLogger.error("📦🌉 Failed to compose courier drop", category: .encryption)
|
||||||
finishPublish(
|
finishPublish(
|
||||||
messageID: messageID,
|
dedupKey: dedupKey,
|
||||||
operationID: operationID,
|
operationID: operationID,
|
||||||
succeeded: false,
|
succeeded: false,
|
||||||
untrackedCompletion: untrackedCompletion
|
untrackedCompletion: untrackedCompletion
|
||||||
@@ -331,7 +377,7 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
guard let publishEvent else {
|
guard let publishEvent else {
|
||||||
SecureLogger.error("📦🌉 Courier drop publisher is not configured", category: .session)
|
SecureLogger.error("📦🌉 Courier drop publisher is not configured", category: .session)
|
||||||
finishPublish(
|
finishPublish(
|
||||||
messageID: messageID,
|
dedupKey: dedupKey,
|
||||||
operationID: operationID,
|
operationID: operationID,
|
||||||
succeeded: false,
|
succeeded: false,
|
||||||
untrackedCompletion: untrackedCompletion
|
untrackedCompletion: untrackedCompletion
|
||||||
@@ -341,7 +387,7 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
publishEvent(event) { [weak self] succeeded in
|
publishEvent(event) { [weak self] succeeded in
|
||||||
guard let self else { return }
|
guard let self else { return }
|
||||||
guard self.finishPublish(
|
guard self.finishPublish(
|
||||||
messageID: messageID,
|
dedupKey: dedupKey,
|
||||||
operationID: operationID,
|
operationID: operationID,
|
||||||
succeeded: succeeded,
|
succeeded: succeeded,
|
||||||
untrackedCompletion: untrackedCompletion
|
untrackedCompletion: untrackedCompletion
|
||||||
@@ -356,23 +402,23 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
|
|
||||||
@discardableResult
|
@discardableResult
|
||||||
private func finishPublish(
|
private func finishPublish(
|
||||||
messageID: String?,
|
dedupKey: String?,
|
||||||
operationID: UUID?,
|
operationID: UUID?,
|
||||||
succeeded: Bool,
|
succeeded: Bool,
|
||||||
untrackedCompletion: (@MainActor (Bool) -> Void)? = nil
|
untrackedCompletion: (@MainActor (Bool) -> Void)? = nil
|
||||||
) -> Bool {
|
) -> Bool {
|
||||||
guard let messageID else {
|
guard let dedupKey else {
|
||||||
untrackedCompletion?(succeeded)
|
untrackedCompletion?(succeeded)
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
// Missing/mismatched means this callback was duplicated, invalidated
|
// Missing/mismatched means this callback was duplicated, invalidated
|
||||||
// by panic wipe, or belongs to an older attempt for the same key.
|
// by panic wipe, or belongs to an older attempt for the same key.
|
||||||
guard let operationID,
|
guard let operationID,
|
||||||
let operation = activeDropOperations[messageID],
|
let operation = activeDropOperations[dedupKey],
|
||||||
operation.id == operationID else { return false }
|
operation.id == operationID else { return false }
|
||||||
activeDropOperations.removeValue(forKey: messageID)
|
activeDropOperations.removeValue(forKey: dedupKey)
|
||||||
if succeeded {
|
if succeeded {
|
||||||
publishedDropKeys.insert(messageID, now: now())
|
publishedDropKeys.insert(dedupKey, now: now())
|
||||||
persistDedup()
|
persistDedup()
|
||||||
}
|
}
|
||||||
operation.completion(succeeded)
|
operation.completion(succeeded)
|
||||||
@@ -387,7 +433,7 @@ final class BridgeCourierService: ObservableObject {
|
|||||||
for item in queued {
|
for item in queued {
|
||||||
publishDrop(
|
publishDrop(
|
||||||
item.envelope,
|
item.envelope,
|
||||||
messageID: item.dedupKey,
|
dedupKey: item.dedupKey,
|
||||||
operationID: item.operationID
|
operationID: item.operationID
|
||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -64,11 +64,11 @@ struct ExpiringIDSet {
|
|||||||
/// fresh drop (fresh throwaway seal, undeduplicatable downstream) and every
|
/// fresh drop (fresh throwaway seal, undeduplicatable downstream) and every
|
||||||
/// gateway relaunch re-delivered the whole backlog. Field-verified: ~20
|
/// gateway relaunch re-delivered the whole backlog. Field-verified: ~20
|
||||||
/// copies of one DM delivered in 40ms fed the storm behind a permanent
|
/// copies of one DM delivered in 40ms fed the storm behind a permanent
|
||||||
/// device freeze. Persisting both sides caps this at one drop per message
|
/// device freeze. Persisting both sides caps this at one drop per
|
||||||
/// ID per 24h regardless of relaunch count.
|
/// recipient/message pair per 24h regardless of relaunch count.
|
||||||
///
|
///
|
||||||
/// Contents are opaque IDs (message UUIDs, relay event IDs) — no plaintext,
|
/// Contents are opaque hashes and relay event IDs — no plaintext or peer
|
||||||
/// no peer identities — so until-first-unlock protection matches
|
/// identities — so until-first-unlock protection matches
|
||||||
/// `NostrProcessedEventStore`, and the file must load during a
|
/// `NostrProcessedEventStore`, and the file must load during a
|
||||||
/// locked-background restoration relaunch. Wiped on panic with the rest of
|
/// locked-background restoration relaunch. Wiped on panic with the rest of
|
||||||
/// the courier state.
|
/// the courier state.
|
||||||
|
|||||||
@@ -34,6 +34,11 @@ struct CourierDirectory {
|
|||||||
final class MessageRouter {
|
final class MessageRouter {
|
||||||
typealias QueuedMessage = MessageOutboxStore.QueuedMessage
|
typealias QueuedMessage = MessageOutboxStore.QueuedMessage
|
||||||
|
|
||||||
|
private struct PeerMessageKey: Hashable {
|
||||||
|
let peerID: PeerID
|
||||||
|
let messageID: String
|
||||||
|
}
|
||||||
|
|
||||||
private let transports: [Transport]
|
private let transports: [Transport]
|
||||||
private let now: () -> Date
|
private let now: () -> Date
|
||||||
private let courierDirectory: CourierDirectory
|
private let courierDirectory: CourierDirectory
|
||||||
@@ -104,15 +109,17 @@ final class MessageRouter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private var bridgeSweepTask: Task<Void, Never>?
|
private var bridgeSweepTask: Task<Void, Never>?
|
||||||
private var bridgeDepositsInFlight = Set<String>()
|
private var bridgeDepositsInFlight = Set<PeerMessageKey>()
|
||||||
|
|
||||||
private var outbox: [PeerID: [QueuedMessage]] = [:]
|
private var outbox: [PeerID: [QueuedMessage]] = [:]
|
||||||
/// IDs whose latest router-owned transmission used an already-established
|
/// Peer/message pairs whose latest router-owned transmission used an
|
||||||
/// secure session and still awaits an ack. This deliberately excludes
|
/// already-established secure session and still await an ack. Peer scope
|
||||||
/// messages handed to BLE while a handshake is pending: BLE owns those
|
/// is required because message IDs are not globally unique across direct
|
||||||
/// sends and drains its queue after authentication, so retrying them here
|
/// conversations. This deliberately excludes messages handed to BLE while
|
||||||
/// would duplicate every normal first-handshake DM.
|
/// a handshake is pending: BLE owns those sends and drains its queue after
|
||||||
private var securelyTransmittedMessageIDs = Set<String>()
|
/// authentication, so retrying them here would duplicate every normal
|
||||||
|
/// first-handshake DM.
|
||||||
|
private var secureTransmissions = Set<PeerMessageKey>()
|
||||||
|
|
||||||
// Outbox limits to prevent unbounded memory growth
|
// Outbox limits to prevent unbounded memory growth
|
||||||
private static let maxMessagesPerPeer = 100
|
private static let maxMessagesPerPeer = 100
|
||||||
@@ -205,7 +212,7 @@ final class MessageRouter {
|
|||||||
// replacement handshake will retry this same message ID, which
|
// replacement handshake will retry this same message ID, which
|
||||||
// receivers deduplicate.
|
// receivers deduplicate.
|
||||||
enqueue(message, for: peerID)
|
enqueue(message, for: peerID)
|
||||||
securelyTransmittedMessageIDs.insert(messageID)
|
secureTransmissions.insert(PeerMessageKey(peerID: peerID, messageID: 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
|
||||||
@@ -229,7 +236,7 @@ final class MessageRouter {
|
|||||||
// 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)
|
||||||
enqueue(message, for: peerID)
|
enqueue(message, for: peerID)
|
||||||
securelyTransmittedMessageIDs.remove(messageID)
|
secureTransmissions.remove(PeerMessageKey(peerID: peerID, messageID: messageID))
|
||||||
transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID)
|
transport.sendPrivateMessage(content, to: peerID, recipientNickname: recipientNickname, messageID: messageID)
|
||||||
attemptCourierDeposit(messageID: messageID, for: peerID)
|
attemptCourierDeposit(messageID: messageID, for: peerID)
|
||||||
return
|
return
|
||||||
@@ -365,11 +372,12 @@ final class MessageRouter {
|
|||||||
for peerID: PeerID,
|
for peerID: PeerID,
|
||||||
recipientKey: Data
|
recipientKey: Data
|
||||||
) {
|
) {
|
||||||
|
let inFlightKey = PeerMessageKey(peerID: peerID, messageID: message.messageID)
|
||||||
guard let bridgeCourierDeposit,
|
guard let bridgeCourierDeposit,
|
||||||
bridgeDepositsInFlight.insert(message.messageID).inserted else { return }
|
bridgeDepositsInFlight.insert(inFlightKey).inserted else { return }
|
||||||
bridgeCourierDeposit(message.content, message.messageID, recipientKey) { [weak self] succeeded in
|
bridgeCourierDeposit(message.content, message.messageID, recipientKey) { [weak self] succeeded in
|
||||||
guard let self else { return }
|
guard let self else { return }
|
||||||
self.bridgeDepositsInFlight.remove(message.messageID)
|
self.bridgeDepositsInFlight.remove(inFlightKey)
|
||||||
// A direct delivery may have cleared the outbox while the relay
|
// A direct delivery may have cleared the outbox while the relay
|
||||||
// relay confirmation was in flight; do not regress its UI state.
|
// relay confirmation was in flight; do not regress its UI state.
|
||||||
guard succeeded, self.queuedMessage(message.messageID, for: peerID) != nil else { return }
|
guard succeeded, self.queuedMessage(message.messageID, for: peerID) != nil else { return }
|
||||||
@@ -401,7 +409,7 @@ final class MessageRouter {
|
|||||||
/// retry state.
|
/// retry state.
|
||||||
func markDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {
|
func markDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {
|
||||||
guard !peerIDs.isEmpty else { return }
|
guard !peerIDs.isEmpty else { return }
|
||||||
clearRetainedMessage(messageID, allowedPeerIDs: peerIDs)
|
_ = markDelivered(messageID, for: Array(peerIDs))
|
||||||
}
|
}
|
||||||
|
|
||||||
private func clearRetainedMessage(
|
private func clearRetainedMessage(
|
||||||
@@ -418,11 +426,10 @@ 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
|
let matchingSecureTransmissions = secureTransmissions.filter {
|
||||||
queue.contains { $0.messageID == messageID }
|
$0.messageID == messageID
|
||||||
}) {
|
|
||||||
securelyTransmittedMessageIDs.remove(messageID)
|
|
||||||
}
|
}
|
||||||
|
secureTransmissions.subtract(matchingSecureTransmissions)
|
||||||
// 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.
|
||||||
@@ -449,10 +456,8 @@ 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
|
for peerID in peerIDs {
|
||||||
queue.contains { $0.messageID == messageID }
|
secureTransmissions.remove(PeerMessageKey(peerID: peerID, 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.
|
||||||
@@ -514,7 +519,7 @@ final class MessageRouter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
private func dropMessage(_ messageID: String, for peerID: PeerID) {
|
private func dropMessage(_ messageID: String, for peerID: PeerID) {
|
||||||
securelyTransmittedMessageIDs.remove(messageID)
|
secureTransmissions.remove(PeerMessageKey(peerID: peerID, messageID: messageID))
|
||||||
metrics?.record(.outboxDropped)
|
metrics?.record(.outboxDropped)
|
||||||
onMessageDropped?(messageID, peerID)
|
onMessageDropped?(messageID, peerID)
|
||||||
}
|
}
|
||||||
@@ -558,7 +563,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()
|
secureTransmissions.removeAll()
|
||||||
outboxStore?.wipe()
|
outboxStore?.wipe()
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -592,7 +597,7 @@ final class MessageRouter {
|
|||||||
/// A peer restart can leave that local session looking usable until the
|
/// A peer restart can leave that local session looking usable until the
|
||||||
/// replacement handshake arrives; the first ciphertext is then
|
/// replacement handshake arrives; the first ciphertext is then
|
||||||
/// undecryptable remotely. Normal pre-handshake sends are intentionally
|
/// undecryptable remotely. Normal pre-handshake sends are intentionally
|
||||||
/// absent from `securelyTransmittedMessageIDs` because BLE already queues
|
/// absent from `secureTransmissions` because BLE already queues
|
||||||
/// and drains them when authentication completes.
|
/// and drains them when authentication completes.
|
||||||
func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) {
|
func retrySecurePrivateMessagesAfterAuthentication(for peerIDAliases: [PeerID]) {
|
||||||
typealias Candidate = (
|
typealias Candidate = (
|
||||||
@@ -617,7 +622,8 @@ final class MessageRouter {
|
|||||||
}
|
}
|
||||||
|
|
||||||
for (queueOrder, message) in queued.enumerated() {
|
for (queueOrder, message) in queued.enumerated() {
|
||||||
guard securelyTransmittedMessageIDs.contains(message.messageID) else { continue }
|
let key = PeerMessageKey(peerID: peerID, messageID: message.messageID)
|
||||||
|
guard secureTransmissions.contains(key) else { continue }
|
||||||
candidates.append((
|
candidates.append((
|
||||||
peerID: peerID,
|
peerID: peerID,
|
||||||
message: message,
|
message: message,
|
||||||
@@ -647,8 +653,9 @@ final class MessageRouter {
|
|||||||
for candidate in candidates {
|
for candidate in candidates {
|
||||||
let peerID = candidate.peerID
|
let peerID = candidate.peerID
|
||||||
let message = candidate.message
|
let message = candidate.message
|
||||||
|
let key = PeerMessageKey(peerID: peerID, messageID: message.messageID)
|
||||||
guard retriedMessageIDs.insert(message.messageID).inserted,
|
guard retriedMessageIDs.insert(message.messageID).inserted,
|
||||||
securelyTransmittedMessageIDs.contains(message.messageID),
|
secureTransmissions.contains(key),
|
||||||
queuedMessage(message.messageID, for: peerID) != nil,
|
queuedMessage(message.messageID, for: peerID) != nil,
|
||||||
let transport = connectedTransport(for: peerID),
|
let transport = connectedTransport(for: peerID),
|
||||||
transport.canDeliverSecurely(to: peerID) else {
|
transport.canDeliverSecurely(to: peerID) else {
|
||||||
@@ -732,7 +739,9 @@ final class MessageRouter {
|
|||||||
continue
|
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)
|
secureTransmissions.insert(
|
||||||
|
PeerMessageKey(peerID: peerID, messageID: 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
|
outboxChanged = incrementSendAttemptsIfQueued(message.messageID, for: peerID) || outboxChanged
|
||||||
@@ -748,7 +757,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)
|
secureTransmissions.remove(
|
||||||
|
PeerMessageKey(peerID: peerID, messageID: 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)
|
||||||
} else if let transport = reachableTransport(for: peerID) {
|
} else if let transport = reachableTransport(for: peerID) {
|
||||||
|
|||||||
@@ -97,7 +97,7 @@ enum EncryptionStatus: Equatable {
|
|||||||
case noiseHandshaking // Currently establishing
|
case noiseHandshaking // Currently establishing
|
||||||
case noiseSecured // Established but not verified
|
case noiseSecured // Established but not verified
|
||||||
case noiseVerified // Established and verified
|
case noiseVerified // Established and verified
|
||||||
|
|
||||||
var icon: String? { // Made optional to hide icon when no handshake
|
var icon: String? { // Made optional to hide icon when no handshake
|
||||||
switch self {
|
switch self {
|
||||||
case .none:
|
case .none:
|
||||||
@@ -660,7 +660,7 @@ final class NoiseEncryptionService {
|
|||||||
guard let packetData = packet.toBinaryDataForSigning() else {
|
guard let packetData = packet.toBinaryDataForSigning() else {
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// Sign with the noise private key (converted to Ed25519 for signing)
|
// Sign with the noise private key (converted to Ed25519 for signing)
|
||||||
guard let signature = signData(packetData) else {
|
guard let signature = signData(packetData) else {
|
||||||
return nil
|
return nil
|
||||||
@@ -919,7 +919,7 @@ final class NoiseEncryptionService {
|
|||||||
}
|
}
|
||||||
return try sessionManager.encrypt(data, for: peerID)
|
return try sessionManager.encrypt(data, for: peerID)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Decrypt data from a specific peer
|
/// Decrypt data from a specific peer
|
||||||
func decrypt(_ data: Data, from peerID: PeerID) throws -> Data {
|
func decrypt(_ data: Data, from peerID: PeerID) throws -> Data {
|
||||||
try decryptWithSessionGeneration(data, from: peerID).plaintext
|
try decryptWithSessionGeneration(data, from: peerID).plaintext
|
||||||
@@ -936,7 +936,7 @@ final class NoiseEncryptionService {
|
|||||||
let isStandardCiphertext = NoiseSecurityValidator.validateCiphertextSize(data)
|
let isStandardCiphertext = NoiseSecurityValidator.validateCiphertextSize(data)
|
||||||
let isAdmittedCiphertext = isStandardCiphertext
|
let isAdmittedCiphertext = isStandardCiphertext
|
||||||
|| NoiseSecurityValidator.validatePrivateFileCiphertextSize(data)
|
|| NoiseSecurityValidator.validatePrivateFileCiphertextSize(data)
|
||||||
|
|
||||||
// A quarantined transport is deliberately unavailable for outbound
|
// A quarantined transport is deliberately unavailable for outbound
|
||||||
// state, but remains receive-only until the responder proves identity
|
// state, but remains receive-only until the responder proves identity
|
||||||
// or the bounded rollback restores it.
|
// or the bounded rollback restores it.
|
||||||
|
|||||||
@@ -21,6 +21,14 @@ protocol ChatDeliveryContext: AnyObject {
|
|||||||
/// message is unknown or no copy changed.
|
/// message is unknown or no copy changed.
|
||||||
@discardableResult
|
@discardableResult
|
||||||
func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool
|
func setDeliveryStatus(_ status: DeliveryStatus, forMessageID messageID: String) -> Bool
|
||||||
|
/// Applies an authenticated receipt only to the direct conversations
|
||||||
|
/// represented by the supplied peer aliases.
|
||||||
|
@discardableResult
|
||||||
|
func setDeliveryStatus(
|
||||||
|
_ status: DeliveryStatus,
|
||||||
|
forMessageID messageID: String,
|
||||||
|
inDirectPeerAliases peerIDs: Set<PeerID>
|
||||||
|
) -> Bool
|
||||||
/// Current delivery status of the message in whichever conversation holds it.
|
/// Current delivery status of the message in whichever conversation holds it.
|
||||||
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus?
|
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus?
|
||||||
/// Message IDs across all direct conversations (read-receipt pruning).
|
/// Message IDs across all direct conversations (read-receipt pruning).
|
||||||
@@ -47,6 +55,19 @@ extension ChatViewModel: ChatDeliveryContext {
|
|||||||
conversations.setDeliveryStatus(status, forMessageID: messageID)
|
conversations.setDeliveryStatus(status, forMessageID: messageID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@discardableResult
|
||||||
|
func setDeliveryStatus(
|
||||||
|
_ status: DeliveryStatus,
|
||||||
|
forMessageID messageID: String,
|
||||||
|
inDirectPeerAliases peerIDs: Set<PeerID>
|
||||||
|
) -> Bool {
|
||||||
|
conversations.setDeliveryStatus(
|
||||||
|
status,
|
||||||
|
forMessageID: messageID,
|
||||||
|
inDirectPeerAliases: peerIDs
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus? {
|
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus? {
|
||||||
conversations.deliveryStatus(forMessageID: messageID)
|
conversations.deliveryStatus(forMessageID: messageID)
|
||||||
}
|
}
|
||||||
@@ -68,6 +89,12 @@ extension ChatViewModel: ChatDeliveryContext {
|
|||||||
|
|
||||||
func markMessageDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {
|
func markMessageDelivered(_ messageID: String, from peerIDs: Set<PeerID>) {
|
||||||
messageRouter.markDelivered(messageID, from: peerIDs)
|
messageRouter.markDelivered(messageID, from: peerIDs)
|
||||||
|
// The caller has already bound this receipt to one of our outgoing
|
||||||
|
// conversations. Release the matching media retry only after that
|
||||||
|
// peer-scoped validation and accepted delivery-status transition.
|
||||||
|
mediaTransferCoordinator.confirmPrivateMediaDelivery(
|
||||||
|
messageID: messageID
|
||||||
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set<PeerID>) -> Bool {
|
func isOutgoingPrivateMessage(_ messageID: String, toAny peerIDs: Set<PeerID>) -> Bool {
|
||||||
@@ -157,7 +184,11 @@ final class ChatDeliveryCoordinator {
|
|||||||
}
|
}
|
||||||
guard !peerIDAliases.isEmpty,
|
guard !peerIDAliases.isEmpty,
|
||||||
context.isOutgoingPrivateMessage(messageID, toAny: peerIDAliases),
|
context.isOutgoingPrivateMessage(messageID, toAny: peerIDAliases),
|
||||||
context.setDeliveryStatus(status, forMessageID: messageID) else {
|
context.setDeliveryStatus(
|
||||||
|
status,
|
||||||
|
forMessageID: messageID,
|
||||||
|
inDirectPeerAliases: peerIDAliases
|
||||||
|
) else {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
context.markMessageDelivered(messageID, from: peerIDAliases)
|
context.markMessageDelivered(messageID, from: peerIDAliases)
|
||||||
|
|||||||
@@ -357,6 +357,78 @@ struct ChatViewModelDeliveryStatusTests {
|
|||||||
#expect(transport.sentPrivateMessages.count == 2)
|
#expect(transport.sentPrivateMessages.count == 2)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test @MainActor
|
||||||
|
func authenticatedNoiseAckClearsOnlyIntendedPeersPrivateMediaRetry() async throws {
|
||||||
|
let (viewModel, transport) = makeTestableViewModel()
|
||||||
|
let intendedPeer = PeerID(str: "0102030405060708")
|
||||||
|
let otherPeer = PeerID(str: "1112131415161718")
|
||||||
|
let fileName = "voice_0011223344556677.m4a"
|
||||||
|
let content = Data("voice".utf8)
|
||||||
|
|
||||||
|
transport.privateMediaPolicies[intendedPeer] = .encrypted
|
||||||
|
transport.privateMediaReceiptSessionGenerations[intendedPeer] = UUID()
|
||||||
|
viewModel.selectedPrivateChatPeer = intendedPeer
|
||||||
|
let coordinator = ChatMediaTransferCoordinator(
|
||||||
|
context: viewModel,
|
||||||
|
prepareVoiceNotePacket: { _ in
|
||||||
|
BitchatFilePacket(
|
||||||
|
fileName: fileName,
|
||||||
|
fileSize: UInt64(content.count),
|
||||||
|
mimeType: "audio/mp4",
|
||||||
|
content: content
|
||||||
|
)
|
||||||
|
},
|
||||||
|
transferIDFactory: { "\($0)-receipt-ack" }
|
||||||
|
)
|
||||||
|
viewModel.mediaTransferCoordinator = coordinator
|
||||||
|
|
||||||
|
let directory = FileManager.default.temporaryDirectory
|
||||||
|
.appendingPathComponent(
|
||||||
|
"scoped-media-ack-\(UUID().uuidString)",
|
||||||
|
isDirectory: true
|
||||||
|
)
|
||||||
|
try FileManager.default.createDirectory(
|
||||||
|
at: directory,
|
||||||
|
withIntermediateDirectories: true
|
||||||
|
)
|
||||||
|
let sourceURL = directory.appendingPathComponent(fileName)
|
||||||
|
try content.write(to: sourceURL)
|
||||||
|
defer { try? FileManager.default.removeItem(at: directory) }
|
||||||
|
|
||||||
|
coordinator.sendVoiceNote(at: sourceURL)
|
||||||
|
#expect(await TestHelpers.waitUntil(
|
||||||
|
{
|
||||||
|
transport.sentPrivateFiles.count == 1
|
||||||
|
&& coordinator.retainedReconnectRetryCount == 1
|
||||||
|
},
|
||||||
|
timeout: TestConstants.longTimeout
|
||||||
|
))
|
||||||
|
let messageID = try #require(
|
||||||
|
viewModel.privateChats[intendedPeer]?.first?.id
|
||||||
|
)
|
||||||
|
let transferID = try #require(
|
||||||
|
transport.sentPrivateFiles.first?.transferID
|
||||||
|
)
|
||||||
|
|
||||||
|
#expect(!viewModel.deliveryCoordinator
|
||||||
|
.updateAcknowledgedMessageDeliveryStatus(
|
||||||
|
messageID,
|
||||||
|
status: .delivered(to: "Other", at: Date()),
|
||||||
|
from: [otherPeer]
|
||||||
|
))
|
||||||
|
#expect(coordinator.retainedReconnectRetryCount == 1)
|
||||||
|
#expect(transport.cancelledTransfers.isEmpty)
|
||||||
|
|
||||||
|
#expect(viewModel.deliveryCoordinator
|
||||||
|
.updateAcknowledgedMessageDeliveryStatus(
|
||||||
|
messageID,
|
||||||
|
status: .delivered(to: "Intended", at: Date()),
|
||||||
|
from: [intendedPeer]
|
||||||
|
))
|
||||||
|
#expect(coordinator.retainedReconnectRetryCount == 0)
|
||||||
|
#expect(transport.cancelledTransfers == [transferID])
|
||||||
|
}
|
||||||
|
|
||||||
@Test @MainActor
|
@Test @MainActor
|
||||||
func cleanupOldReadReceipts_removesReceiptIDsWithoutMessages() async {
|
func cleanupOldReadReceipts_removesReceiptIDsWithoutMessages() async {
|
||||||
let (viewModel, transport) = makeTestableViewModel()
|
let (viewModel, transport) = makeTestableViewModel()
|
||||||
@@ -643,6 +715,19 @@ private final class MockChatDeliveryContext: ChatDeliveryContext {
|
|||||||
store.setDeliveryStatus(status, forMessageID: messageID)
|
store.setDeliveryStatus(status, forMessageID: messageID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@discardableResult
|
||||||
|
func setDeliveryStatus(
|
||||||
|
_ status: DeliveryStatus,
|
||||||
|
forMessageID messageID: String,
|
||||||
|
inDirectPeerAliases peerIDs: Set<PeerID>
|
||||||
|
) -> Bool {
|
||||||
|
store.setDeliveryStatus(
|
||||||
|
status,
|
||||||
|
forMessageID: messageID,
|
||||||
|
inDirectPeerAliases: peerIDs
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus? {
|
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus? {
|
||||||
store.deliveryStatus(forMessageID: messageID)
|
store.deliveryStatus(forMessageID: messageID)
|
||||||
}
|
}
|
||||||
@@ -769,6 +854,94 @@ struct ChatDeliveryCoordinatorContextTests {
|
|||||||
#expect(context.peerBoundDeliveredMessages[0].peerIDs == [peerID])
|
#expect(context.peerBoundDeliveredMessages[0].peerIDs == [peerID])
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test @MainActor
|
||||||
|
func authenticatedReceiptWithCollidingIDUpdatesOnlyAuthenticatedAliases() async {
|
||||||
|
let context = MockChatDeliveryContext()
|
||||||
|
let coordinator = ChatDeliveryCoordinator(context: context)
|
||||||
|
let ephemeralPeerID = PeerID(str: "0102030405060708")
|
||||||
|
let stablePeerID = PeerID(hexData: Data(repeating: 0x08, count: 32))
|
||||||
|
let otherPeerID = PeerID(str: "1112131415161718")
|
||||||
|
let messageID = "authenticated-receipt-collision"
|
||||||
|
let mirroredMessage = makePrivateMessage(id: messageID, status: .sent)
|
||||||
|
|
||||||
|
context.store.append(mirroredMessage, to: .directPeer(ephemeralPeerID))
|
||||||
|
context.store.append(mirroredMessage, to: .directPeer(stablePeerID))
|
||||||
|
context.store.append(
|
||||||
|
makePrivateMessage(id: messageID, status: .sent),
|
||||||
|
to: .directPeer(otherPeerID)
|
||||||
|
)
|
||||||
|
context.store.append(makePublicMessage(id: messageID, status: .sent), to: .mesh)
|
||||||
|
|
||||||
|
let aliases: Set<PeerID> = [ephemeralPeerID, stablePeerID]
|
||||||
|
let didUpdate = coordinator.updateAcknowledgedMessageDeliveryStatus(
|
||||||
|
messageID,
|
||||||
|
status: .delivered(to: "Peer", at: Date()),
|
||||||
|
from: aliases
|
||||||
|
)
|
||||||
|
|
||||||
|
#expect(didUpdate)
|
||||||
|
#expect(isDelivered(
|
||||||
|
context.store.conversation(for: .directPeer(ephemeralPeerID))
|
||||||
|
.message(withID: messageID)?.deliveryStatus
|
||||||
|
))
|
||||||
|
#expect(isDelivered(
|
||||||
|
context.store.conversation(for: .directPeer(stablePeerID))
|
||||||
|
.message(withID: messageID)?.deliveryStatus
|
||||||
|
))
|
||||||
|
#expect(isSent(
|
||||||
|
context.store.conversation(for: .directPeer(otherPeerID))
|
||||||
|
.message(withID: messageID)?.deliveryStatus
|
||||||
|
))
|
||||||
|
#expect(isSent(
|
||||||
|
context.store.conversation(for: .mesh)
|
||||||
|
.message(withID: messageID)?.deliveryStatus
|
||||||
|
))
|
||||||
|
#expect(context.peerBoundDeliveredMessages.count == 1)
|
||||||
|
#expect(context.peerBoundDeliveredMessages[0].messageID == messageID)
|
||||||
|
#expect(context.peerBoundDeliveredMessages[0].peerIDs == aliases)
|
||||||
|
#expect(context.notifyUIChangedCount == 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test @MainActor
|
||||||
|
func authenticatedReceiptRemainsScopedAfterPeerAliasMigration() async {
|
||||||
|
let context = MockChatDeliveryContext()
|
||||||
|
let coordinator = ChatDeliveryCoordinator(context: context)
|
||||||
|
let ephemeralPeerID = PeerID(str: "2122232425262728")
|
||||||
|
let stablePeerID = PeerID(hexData: Data(repeating: 0x28, count: 32))
|
||||||
|
let otherPeerID = PeerID(str: "3132333435363738")
|
||||||
|
let messageID = "authenticated-receipt-after-migration"
|
||||||
|
|
||||||
|
context.store.append(
|
||||||
|
makePrivateMessage(id: messageID, status: .sent),
|
||||||
|
to: .directPeer(ephemeralPeerID)
|
||||||
|
)
|
||||||
|
context.store.append(
|
||||||
|
makePrivateMessage(id: messageID, status: .sent),
|
||||||
|
to: .directPeer(otherPeerID)
|
||||||
|
)
|
||||||
|
context.store.migrateConversation(
|
||||||
|
from: .directPeer(ephemeralPeerID),
|
||||||
|
to: .directPeer(stablePeerID)
|
||||||
|
)
|
||||||
|
|
||||||
|
let didUpdate = coordinator.updateAcknowledgedMessageDeliveryStatus(
|
||||||
|
messageID,
|
||||||
|
status: .read(by: "Peer", at: Date()),
|
||||||
|
from: [ephemeralPeerID, stablePeerID]
|
||||||
|
)
|
||||||
|
|
||||||
|
#expect(didUpdate)
|
||||||
|
#expect(context.store.conversationsByID[.directPeer(ephemeralPeerID)] == nil)
|
||||||
|
#expect(isRead(
|
||||||
|
context.store.conversation(for: .directPeer(stablePeerID))
|
||||||
|
.message(withID: messageID)?.deliveryStatus
|
||||||
|
))
|
||||||
|
#expect(isSent(
|
||||||
|
context.store.conversation(for: .directPeer(otherPeerID))
|
||||||
|
.message(withID: messageID)?.deliveryStatus
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
@Test @MainActor
|
@Test @MainActor
|
||||||
func receiptFromWrongPeerDoesNotUpdateOrTerminalizeOutgoingMessage() async {
|
func receiptFromWrongPeerDoesNotUpdateOrTerminalizeOutgoingMessage() async {
|
||||||
let context = MockChatDeliveryContext()
|
let context = MockChatDeliveryContext()
|
||||||
|
|||||||
@@ -755,6 +755,71 @@ struct ConversationStoreTests {
|
|||||||
#expect(statusChangedIDs.isEmpty)
|
#expect(statusChangedIDs.isEmpty)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test("peer-scoped receipt updates only authenticated direct aliases")
|
||||||
|
@MainActor
|
||||||
|
func peerScopedReceiptUpdatesOnlyAuthenticatedDirectAliases() {
|
||||||
|
let store = ConversationStore()
|
||||||
|
let ephemeralPeer = PeerID(str: "0102030405060708")
|
||||||
|
let stablePeer = PeerID(hexData: Data(repeating: 0x08, count: 32))
|
||||||
|
let otherPeer = PeerID(str: "1112131415161718")
|
||||||
|
let ephemeral = ConversationID.directPeer(ephemeralPeer)
|
||||||
|
let stable = ConversationID.directPeer(stablePeer)
|
||||||
|
let other = ConversationID.directPeer(otherPeer)
|
||||||
|
let messageID = "scoped-receipt"
|
||||||
|
let mirrored = makeMessage(
|
||||||
|
id: messageID,
|
||||||
|
timestamp: 1,
|
||||||
|
isPrivate: true,
|
||||||
|
deliveryStatus: .sent
|
||||||
|
)
|
||||||
|
store.upsertByID(mirrored, in: ephemeral)
|
||||||
|
store.upsertByID(mirrored, in: stable)
|
||||||
|
store.upsertByID(
|
||||||
|
makeMessage(
|
||||||
|
id: messageID,
|
||||||
|
timestamp: 1,
|
||||||
|
isPrivate: true,
|
||||||
|
deliveryStatus: .sent
|
||||||
|
),
|
||||||
|
in: other
|
||||||
|
)
|
||||||
|
store.upsertByID(
|
||||||
|
makeMessage(id: messageID, timestamp: 1, deliveryStatus: .sent),
|
||||||
|
in: .mesh
|
||||||
|
)
|
||||||
|
|
||||||
|
var cancellables = Set<AnyCancellable>()
|
||||||
|
var publishedIDs: [ConversationID] = []
|
||||||
|
for id in [ephemeral, stable, other, .mesh] {
|
||||||
|
store.conversation(for: id).objectWillChange
|
||||||
|
.sink { publishedIDs.append(id) }
|
||||||
|
.store(in: &cancellables)
|
||||||
|
}
|
||||||
|
var statusChangedIDs: [ConversationID] = []
|
||||||
|
store.changes
|
||||||
|
.sink { change in
|
||||||
|
if case .statusChanged(let id, messageID, _) = change,
|
||||||
|
messageID == "scoped-receipt" {
|
||||||
|
statusChangedIDs.append(id)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
.store(in: &cancellables)
|
||||||
|
|
||||||
|
let delivered = DeliveryStatus.delivered(to: "bob", at: Date())
|
||||||
|
#expect(store.setDeliveryStatus(
|
||||||
|
delivered,
|
||||||
|
forMessageID: messageID,
|
||||||
|
inDirectPeerAliases: [ephemeralPeer, stablePeer]
|
||||||
|
))
|
||||||
|
|
||||||
|
#expect(Set(publishedIDs) == Set([ephemeral, stable]))
|
||||||
|
#expect(Set(statusChangedIDs) == Set([ephemeral, stable]))
|
||||||
|
#expect(store.conversation(for: ephemeral).message(withID: messageID)?.deliveryStatus == delivered)
|
||||||
|
#expect(store.conversation(for: stable).message(withID: messageID)?.deliveryStatus == delivered)
|
||||||
|
#expect(store.conversation(for: other).message(withID: messageID)?.deliveryStatus == .sent)
|
||||||
|
#expect(store.conversation(for: .mesh).message(withID: messageID)?.deliveryStatus == .sent)
|
||||||
|
}
|
||||||
|
|
||||||
// MARK: - Invariant audit (field observability)
|
// MARK: - Invariant audit (field observability)
|
||||||
|
|
||||||
/// A store exercised through every intent family: public + geohash +
|
/// A store exercised through every intent family: public + geohash +
|
||||||
|
|||||||
@@ -65,6 +65,7 @@ final class MockTransport: Transport, PrivateMediaDeletionPersisting {
|
|||||||
var peerFingerprints: [PeerID: String] = [:]
|
var peerFingerprints: [PeerID: String] = [:]
|
||||||
var peerNoiseStates: [PeerID: LazyHandshakeState] = [:]
|
var peerNoiseStates: [PeerID: LazyHandshakeState] = [:]
|
||||||
var privateMediaPolicies: [PeerID: PrivateMediaSendPolicy] = [:]
|
var privateMediaPolicies: [PeerID: PrivateMediaSendPolicy] = [:]
|
||||||
|
var privateMediaReceiptSessionGenerations: [PeerID: UUID] = [:]
|
||||||
var persistDeletedPrivateMediaResult = true
|
var persistDeletedPrivateMediaResult = true
|
||||||
var deferDeletedPrivateMediaPersistence = false
|
var deferDeletedPrivateMediaPersistence = false
|
||||||
private var pendingDeletedPrivateMediaCompletions: [
|
private var pendingDeletedPrivateMediaCompletions: [
|
||||||
@@ -223,6 +224,12 @@ final class MockTransport: Transport, PrivateMediaDeletionPersisting {
|
|||||||
privateMediaPolicies[peerID] ?? .encrypted
|
privateMediaPolicies[peerID] ?? .encrypted
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func authenticatedPrivateMediaReceiptSessionGeneration(
|
||||||
|
to peerID: PeerID
|
||||||
|
) -> UUID? {
|
||||||
|
privateMediaReceiptSessionGenerations[peerID]
|
||||||
|
}
|
||||||
|
|
||||||
func resolvePrivateMediaSendPolicy(
|
func resolvePrivateMediaSendPolicy(
|
||||||
to peerID: PeerID,
|
to peerID: PeerID,
|
||||||
completion: @escaping @MainActor (PrivateMediaSendPolicy) -> Void
|
completion: @escaping @MainActor (PrivateMediaSendPolicy) -> Void
|
||||||
|
|||||||
@@ -700,6 +700,19 @@ private final class PerfDeliveryContext: ChatDeliveryContext {
|
|||||||
store.setDeliveryStatus(status, forMessageID: messageID)
|
store.setDeliveryStatus(status, forMessageID: messageID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@discardableResult
|
||||||
|
func setDeliveryStatus(
|
||||||
|
_ status: DeliveryStatus,
|
||||||
|
forMessageID messageID: String,
|
||||||
|
inDirectPeerAliases peerIDs: Set<PeerID>
|
||||||
|
) -> Bool {
|
||||||
|
store.setDeliveryStatus(
|
||||||
|
status,
|
||||||
|
forMessageID: messageID,
|
||||||
|
inDirectPeerAliases: peerIDs
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus? {
|
func deliveryStatus(forMessageID messageID: String) -> DeliveryStatus? {
|
||||||
store.deliveryStatus(forMessageID: messageID)
|
store.deliveryStatus(forMessageID: messageID)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -203,6 +203,134 @@ struct BridgeCourierServiceTests {
|
|||||||
#expect(confirmed.sealRequests.isEmpty)
|
#expect(confirmed.sealRequests.isEmpty)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test func sameMessageIDIsScopedByRecipientAcrossRejectedActiveAndPersistedState() {
|
||||||
|
let fileURL = FileManager.default.temporaryDirectory
|
||||||
|
.appendingPathComponent("bridge-dedup-\(UUID().uuidString).json")
|
||||||
|
defer { try? FileManager.default.removeItem(at: fileURL) }
|
||||||
|
let rejectedKey = Fixture.randomKey()
|
||||||
|
let firstKey = Fixture.randomKey()
|
||||||
|
let secondKey = Fixture.randomKey()
|
||||||
|
let thirdKey = Fixture.randomKey()
|
||||||
|
let messageID = "recipient-scoped-collision"
|
||||||
|
|
||||||
|
let fixture = Fixture(dedupStore: BridgeDropDedupStore(fileURL: fileURL))
|
||||||
|
fixture.sealResult = makeEnvelope(
|
||||||
|
recipientKey: rejectedKey,
|
||||||
|
ciphertext: Data(
|
||||||
|
repeating: 7,
|
||||||
|
count: BridgeCourierService.Limits.maxDropEnvelopeBytes + 1
|
||||||
|
)
|
||||||
|
)
|
||||||
|
var rejectedResults: [Bool] = []
|
||||||
|
fixture.service.depositDrop(
|
||||||
|
content: "rejected",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: rejectedKey
|
||||||
|
) { rejectedResults.append($0) }
|
||||||
|
#expect(rejectedResults == [false])
|
||||||
|
|
||||||
|
fixture.sealResult = makeEnvelope(recipientKey: firstKey)
|
||||||
|
fixture.automaticPublishResult = nil
|
||||||
|
var firstResults: [Bool] = []
|
||||||
|
var secondResults: [Bool] = []
|
||||||
|
var duplicateFirstResults: [Bool] = []
|
||||||
|
fixture.service.depositDrop(
|
||||||
|
content: "first",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: firstKey
|
||||||
|
) { firstResults.append($0) }
|
||||||
|
fixture.service.depositDrop(
|
||||||
|
content: "second",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: secondKey
|
||||||
|
) { secondResults.append($0) }
|
||||||
|
fixture.service.depositDrop(
|
||||||
|
content: "first duplicate",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: firstKey
|
||||||
|
) { duplicateFirstResults.append($0) }
|
||||||
|
|
||||||
|
#expect(fixture.publishedEvents.count == 2)
|
||||||
|
#expect(fixture.pendingPublishCompletions.count == 2)
|
||||||
|
#expect(duplicateFirstResults == [false])
|
||||||
|
#expect(firstResults.isEmpty)
|
||||||
|
#expect(secondResults.isEmpty)
|
||||||
|
|
||||||
|
fixture.resolveNextPublish(true)
|
||||||
|
fixture.resolveNextPublish(true)
|
||||||
|
#expect(firstResults == [true])
|
||||||
|
#expect(secondResults == [true])
|
||||||
|
fixture.service.flushDedupSnapshot()
|
||||||
|
|
||||||
|
let relaunched = Fixture(
|
||||||
|
dedupStore: BridgeDropDedupStore(fileURL: fileURL)
|
||||||
|
)
|
||||||
|
relaunched.sealResult = makeEnvelope(recipientKey: thirdKey)
|
||||||
|
var relaunchResults: [Bool] = []
|
||||||
|
relaunched.service.depositDrop(
|
||||||
|
content: "first",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: firstKey
|
||||||
|
) { relaunchResults.append($0) }
|
||||||
|
relaunched.service.depositDrop(
|
||||||
|
content: "second",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: secondKey
|
||||||
|
) { relaunchResults.append($0) }
|
||||||
|
relaunched.service.depositDrop(
|
||||||
|
content: "third",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: thirdKey
|
||||||
|
) { relaunchResults.append($0) }
|
||||||
|
|
||||||
|
#expect(relaunchResults == [false, false, true])
|
||||||
|
#expect(relaunched.sealRequests.count == 1)
|
||||||
|
#expect(relaunched.sealRequests.first?.key == thirdKey)
|
||||||
|
#expect(relaunched.publishedEvents.count == 1)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test func legacyPublishedMessageIDIsWildcardUntilItsOriginalExpiry() {
|
||||||
|
let fileURL = FileManager.default.temporaryDirectory
|
||||||
|
.appendingPathComponent("bridge-dedup-\(UUID().uuidString).json")
|
||||||
|
defer { try? FileManager.default.removeItem(at: fileURL) }
|
||||||
|
var date = Date()
|
||||||
|
let messageID = "legacy-wildcard"
|
||||||
|
let recipientKey = Fixture.randomKey()
|
||||||
|
let store = BridgeDropDedupStore(fileURL: fileURL)
|
||||||
|
store.save(BridgeDropDedupStore.Snapshot(
|
||||||
|
publishedDropKeys: [messageID: date],
|
||||||
|
seenDropEventIDs: [:]
|
||||||
|
))
|
||||||
|
|
||||||
|
let fixture = Fixture(
|
||||||
|
now: { date },
|
||||||
|
dedupStore: BridgeDropDedupStore(fileURL: fileURL)
|
||||||
|
)
|
||||||
|
fixture.sealResult = makeEnvelope(recipientKey: recipientKey)
|
||||||
|
var results: [Bool] = []
|
||||||
|
fixture.service.depositDrop(
|
||||||
|
content: "legacy",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: recipientKey
|
||||||
|
) { results.append($0) }
|
||||||
|
#expect(results == [false])
|
||||||
|
#expect(fixture.sealRequests.isEmpty)
|
||||||
|
|
||||||
|
date = date.addingTimeInterval(CourierEnvelope.maxLifetimeSeconds + 1)
|
||||||
|
fixture.service.depositDrop(
|
||||||
|
content: "after expiry",
|
||||||
|
messageID: messageID,
|
||||||
|
recipientNoiseKey: recipientKey
|
||||||
|
) { results.append($0) }
|
||||||
|
#expect(results == [false, true])
|
||||||
|
#expect(fixture.publishedEvents.count == 1)
|
||||||
|
fixture.service.flushDedupSnapshot()
|
||||||
|
|
||||||
|
let snapshot = BridgeDropDedupStore(fileURL: fileURL).load()
|
||||||
|
#expect(snapshot.publishedDropKeys[messageID] == nil)
|
||||||
|
#expect(snapshot.publishedDropKeys.count == 1)
|
||||||
|
}
|
||||||
|
|
||||||
@Test func panicWipeInvalidatesInFlightPublishCompletion() throws {
|
@Test func panicWipeInvalidatesInFlightPublishCompletion() throws {
|
||||||
let fileURL = FileManager.default.temporaryDirectory
|
let fileURL = FileManager.default.temporaryDirectory
|
||||||
.appendingPathComponent("bridge-dedup-\(UUID().uuidString).json")
|
.appendingPathComponent("bridge-dedup-\(UUID().uuidString).json")
|
||||||
@@ -297,8 +425,10 @@ struct BridgeCourierServiceTests {
|
|||||||
#expect(firstResults == [false])
|
#expect(firstResults == [false])
|
||||||
|
|
||||||
// The evicted first drop is deposit-able again (slot released).
|
// The evicted first drop is deposit-able again (slot released).
|
||||||
|
let sealCountBeforeRetry = fixture.sealRequests.count
|
||||||
fixture.service.depositDrop(content: "0-retry", messageID: firstID, recipientNoiseKey: key)
|
fixture.service.depositDrop(content: "0-retry", messageID: firstID, recipientNoiseKey: key)
|
||||||
#expect(fixture.service.pendingDrops.last?.dedupKey == firstID)
|
#expect(fixture.sealRequests.count == sealCountBeforeRetry + 1)
|
||||||
|
#expect(fixture.service.pendingDrops.count == BridgeCourierService.Limits.maxPendingDrops)
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test func oversizeDropConsumesSlotInsteadOfChurning() {
|
@Test func oversizeDropConsumesSlotInsteadOfChurning() {
|
||||||
|
|||||||
@@ -188,6 +188,58 @@ struct MessageRouterTests {
|
|||||||
router.markDelivered("normal-handshake")
|
router.markDelivered("normal-handshake")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test @MainActor
|
||||||
|
func authenticationRetry_scopesCollidingMessageIDsByPeer() async {
|
||||||
|
let securePeer = PeerID(str: "0000000000000025")
|
||||||
|
let pendingPeer = PeerID(str: "0000000000000026")
|
||||||
|
let transport = MockTransport()
|
||||||
|
transport.connectedPeers = [securePeer, pendingPeer]
|
||||||
|
transport.securePeers = [securePeer]
|
||||||
|
|
||||||
|
let router = MessageRouter(transports: [transport])
|
||||||
|
let promotedID = "collision-promoted"
|
||||||
|
let clearedID = "collision-cleared"
|
||||||
|
|
||||||
|
// Pending B then secure A: an ID-global marker falsely promotes B.
|
||||||
|
router.sendPrivate(
|
||||||
|
"pending promoted",
|
||||||
|
to: pendingPeer,
|
||||||
|
recipientNickname: "Pending",
|
||||||
|
messageID: promotedID
|
||||||
|
)
|
||||||
|
router.sendPrivate(
|
||||||
|
"secure promoted",
|
||||||
|
to: securePeer,
|
||||||
|
recipientNickname: "Secure",
|
||||||
|
messageID: promotedID
|
||||||
|
)
|
||||||
|
|
||||||
|
// Secure A then pending B: an ID-global removal falsely clears A.
|
||||||
|
router.sendPrivate(
|
||||||
|
"secure cleared",
|
||||||
|
to: securePeer,
|
||||||
|
recipientNickname: "Secure",
|
||||||
|
messageID: clearedID
|
||||||
|
)
|
||||||
|
router.sendPrivate(
|
||||||
|
"pending cleared",
|
||||||
|
to: pendingPeer,
|
||||||
|
recipientNickname: "Pending",
|
||||||
|
messageID: clearedID
|
||||||
|
)
|
||||||
|
|
||||||
|
transport.resetRecordings()
|
||||||
|
transport.securePeers = [securePeer, pendingPeer]
|
||||||
|
|
||||||
|
router.retrySecurePrivateMessagesAfterAuthentication(for: [pendingPeer])
|
||||||
|
#expect(transport.sentPrivateMessages.isEmpty)
|
||||||
|
|
||||||
|
router.retrySecurePrivateMessagesAfterAuthentication(for: [securePeer])
|
||||||
|
#expect(transport.sentPrivateMessages.count == 2)
|
||||||
|
#expect(Set(transport.sentPrivateMessages.map(\.messageID)) == [promotedID, clearedID])
|
||||||
|
#expect(transport.sentPrivateMessages.allSatisfy { $0.peerID == securePeer })
|
||||||
|
}
|
||||||
|
|
||||||
@Test @MainActor
|
@Test @MainActor
|
||||||
func authenticationRetry_doesNotDuplicateMessageRequeuedByBLEForHandshake() async {
|
func authenticationRetry_doesNotDuplicateMessageRequeuedByBLEForHandshake() async {
|
||||||
let peerID = PeerID(str: "0000000000000021")
|
let peerID = PeerID(str: "0000000000000021")
|
||||||
@@ -713,6 +765,52 @@ struct MessageRouterTests {
|
|||||||
#expect(carried == ["bridge-ack"])
|
#expect(carried == ["bridge-ack"])
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test @MainActor
|
||||||
|
func bridgeDepositsScopeCollidingMessageIDsByRecipient() async {
|
||||||
|
let firstRecipient = PeerID(str: "00000000000000b1")
|
||||||
|
let secondRecipient = PeerID(str: "00000000000000b2")
|
||||||
|
let firstKey = Data(repeating: 0xB1, count: 32)
|
||||||
|
let secondKey = Data(repeating: 0xB2, count: 32)
|
||||||
|
let recipientKeys = [
|
||||||
|
firstRecipient: firstKey,
|
||||||
|
secondRecipient: secondKey
|
||||||
|
]
|
||||||
|
let router = MessageRouter(
|
||||||
|
transports: [MockTransport()],
|
||||||
|
courierDirectory: CourierDirectory(
|
||||||
|
noiseKey: { recipientKeys[$0] },
|
||||||
|
isTrustedCourier: { _ in false }
|
||||||
|
)
|
||||||
|
)
|
||||||
|
var requestedKeys: [Data] = []
|
||||||
|
var completions: [@MainActor (Bool) -> Void] = []
|
||||||
|
router.bridgeCourierDeposit = { _, _, recipientKey, completion in
|
||||||
|
requestedKeys.append(recipientKey)
|
||||||
|
completions.append(completion)
|
||||||
|
}
|
||||||
|
var carriedPeers: [PeerID] = []
|
||||||
|
router.onMessageCarried = { _, peerID in carriedPeers.append(peerID) }
|
||||||
|
|
||||||
|
router.sendPrivate(
|
||||||
|
"First",
|
||||||
|
to: firstRecipient,
|
||||||
|
recipientNickname: "First",
|
||||||
|
messageID: "bridge-collision"
|
||||||
|
)
|
||||||
|
router.sendPrivate(
|
||||||
|
"Second",
|
||||||
|
to: secondRecipient,
|
||||||
|
recipientNickname: "Second",
|
||||||
|
messageID: "bridge-collision"
|
||||||
|
)
|
||||||
|
|
||||||
|
#expect(completions.count == 2)
|
||||||
|
#expect(Set(requestedKeys) == [firstKey, secondKey])
|
||||||
|
|
||||||
|
completions.forEach { $0(true) }
|
||||||
|
#expect(Set(carriedPeers) == [firstRecipient, secondRecipient])
|
||||||
|
}
|
||||||
|
|
||||||
// MARK: - Outbox persistence
|
// MARK: - Outbox persistence
|
||||||
|
|
||||||
@Test @MainActor
|
@Test @MainActor
|
||||||
@@ -829,7 +927,10 @@ struct MessageRouterTests {
|
|||||||
transport.reachablePeers.formUnion([acknowledgedPeer, otherPeer])
|
transport.reachablePeers.formUnion([acknowledgedPeer, otherPeer])
|
||||||
let router = MessageRouter(transports: [transport], outboxStore: restoredStore)
|
let router = MessageRouter(transports: [transport], outboxStore: restoredStore)
|
||||||
|
|
||||||
#expect(!router.markDelivered("shared-locked-id", for: [acknowledgedPeer]))
|
router.markDelivered(
|
||||||
|
"shared-locked-id",
|
||||||
|
from: [acknowledgedPeer]
|
||||||
|
)
|
||||||
protectedDataUnavailable = false
|
protectedDataUnavailable = false
|
||||||
restoredStore.retryDeferredLoad()
|
restoredStore.retryDeferredLoad()
|
||||||
await Task.yield()
|
await Task.yield()
|
||||||
|
|||||||
Reference in New Issue
Block a user