mirror of
https://github.com/permissionlesstech/bitchat.git
synced 2026-07-25 02:05:19 +00:00
Surface router drops as failed status; bound relay pending subscriptions; jitter backoff
MessageRouter outbox drops (attempt cap, TTL expiry in flush and cleanup, per-peer overflow eviction) now invoke onMessageDropped, wired to mark the message .failed in the ConversationStore - guarded so a late failure never downgrades an already delivered/read status. NostrRelayManager pending subscriptions gain a per-relay cap (64, oldest-by-sequence eviction; durable intent still replays from subscriptionRequestState) and a 10-minute age sweep on the existing connect path. Reconnect backoff gets injectable +/-20% jitter so recovering relays don't thundering-herd. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -349,6 +349,80 @@ struct ChatViewModelDeliveryStatusTests {
|
||||
#expect(isRead(viewModel.privateChats[peerID]?.last?.deliveryStatus))
|
||||
}
|
||||
|
||||
// MARK: - MessageRouter Drop Wiring Tests
|
||||
|
||||
/// Drives a real outbox drop (per-peer overflow eviction with no
|
||||
/// reachable transport) and proves the bootstrapper wiring marks the
|
||||
/// dropped message `.failed` in the conversation store.
|
||||
@Test @MainActor
|
||||
func messageRouterDrop_marksMessageFailedInStore() async {
|
||||
let (viewModel, transport) = makeTestableViewModel()
|
||||
let peerID = PeerID(str: "0102030405060708")
|
||||
let droppedID = "router-drop-0"
|
||||
|
||||
let message = BitchatMessage(
|
||||
id: droppedID,
|
||||
sender: viewModel.nickname,
|
||||
content: "Will be dropped",
|
||||
timestamp: Date(),
|
||||
isRelay: false,
|
||||
isPrivate: true,
|
||||
recipientNickname: "Peer",
|
||||
senderPeerID: transport.myPeerID,
|
||||
deliveryStatus: .sent
|
||||
)
|
||||
viewModel.seedPrivateChat([message], for: peerID)
|
||||
|
||||
// No transport is reachable, so every send is queued; the 101st
|
||||
// enqueue for this peer evicts the oldest queued message.
|
||||
viewModel.messageRouter.sendPrivate("Will be dropped", to: peerID, recipientNickname: "Peer", messageID: droppedID)
|
||||
for i in 1...100 {
|
||||
viewModel.messageRouter.sendPrivate("Filler \(i)", to: peerID, recipientNickname: "Peer", messageID: "router-drop-\(i)")
|
||||
}
|
||||
|
||||
let status = viewModel.conversations.deliveryStatus(forMessageID: droppedID)
|
||||
#expect({
|
||||
if case .failed = status { return true }
|
||||
return false
|
||||
}())
|
||||
}
|
||||
|
||||
/// The store's no-downgrade rule does not cover `.failed` over confirmed
|
||||
/// receipts, so the wiring guards it: a drop of an already-delivered
|
||||
/// message must not downgrade its status.
|
||||
@Test @MainActor
|
||||
func messageRouterDrop_doesNotDowngradeDeliveredStatus() async {
|
||||
let (viewModel, transport) = makeTestableViewModel()
|
||||
let peerID = PeerID(str: "0102030405060708")
|
||||
let droppedID = "router-drop-delivered"
|
||||
|
||||
let message = BitchatMessage(
|
||||
id: droppedID,
|
||||
sender: viewModel.nickname,
|
||||
content: "Already delivered",
|
||||
timestamp: Date(),
|
||||
isRelay: false,
|
||||
isPrivate: true,
|
||||
recipientNickname: "Peer",
|
||||
senderPeerID: transport.myPeerID,
|
||||
deliveryStatus: .delivered(to: "Peer", at: Date())
|
||||
)
|
||||
viewModel.seedPrivateChat([message], for: peerID)
|
||||
|
||||
// Same eviction-driven drop as above, but the store already recorded
|
||||
// a delivery confirmation for the message.
|
||||
viewModel.messageRouter.sendPrivate("Already delivered", to: peerID, recipientNickname: "Peer", messageID: droppedID)
|
||||
for i in 1...100 {
|
||||
viewModel.messageRouter.sendPrivate("Filler \(i)", to: peerID, recipientNickname: "Peer", messageID: "router-keep-\(i)")
|
||||
}
|
||||
|
||||
let status = viewModel.conversations.deliveryStatus(forMessageID: droppedID)
|
||||
#expect({
|
||||
if case .delivered = status { return true }
|
||||
return false
|
||||
}())
|
||||
}
|
||||
|
||||
// MARK: - Status Rank Tests (for deduplication)
|
||||
|
||||
@Test @MainActor
|
||||
|
||||
@@ -110,6 +110,98 @@ struct MessageRouterTests {
|
||||
#expect(transport.sentPrivateMessages.count == 8)
|
||||
}
|
||||
|
||||
// MARK: - Drop visibility (onMessageDropped)
|
||||
|
||||
@Test @MainActor
|
||||
func flushOutbox_attemptCapDropInvokesOnMessageDropped() async {
|
||||
let peerID = PeerID(str: "0000000000000009")
|
||||
let transport = MockTransport()
|
||||
transport.reachablePeers.insert(peerID)
|
||||
|
||||
let router = MessageRouter(transports: [transport])
|
||||
var dropped: [(messageID: String, peerID: PeerID)] = []
|
||||
router.onMessageDropped = { dropped.append(($0, $1)) }
|
||||
|
||||
router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "m9")
|
||||
for _ in 0..<10 {
|
||||
router.flushOutbox(for: peerID)
|
||||
}
|
||||
|
||||
#expect(dropped.count == 1)
|
||||
#expect(dropped.first?.messageID == "m9")
|
||||
#expect(dropped.first?.peerID == peerID)
|
||||
}
|
||||
|
||||
@Test @MainActor
|
||||
func flushOutbox_ttlExpiryInvokesOnMessageDroppedAndDoesNotResend() async {
|
||||
let peerID = PeerID(str: "000000000000000a")
|
||||
let transport = MockTransport()
|
||||
let clock = MutableTestClock()
|
||||
|
||||
let router = MessageRouter(transports: [transport], now: { clock.now })
|
||||
var dropped: [String] = []
|
||||
router.onMessageDropped = { messageID, _ in dropped.append(messageID) }
|
||||
|
||||
// No reachable transport: the message is queued, never sent.
|
||||
router.sendPrivate("Hello", to: peerID, recipientNickname: "Peer", messageID: "m10")
|
||||
#expect(transport.sentPrivateMessages.isEmpty)
|
||||
|
||||
// Past the 24h TTL the flush must drop it (visibly), not send it.
|
||||
clock.now = clock.now.addingTimeInterval(24 * 60 * 60 + 1)
|
||||
transport.reachablePeers.insert(peerID)
|
||||
router.flushOutbox(for: peerID)
|
||||
|
||||
#expect(dropped == ["m10"])
|
||||
#expect(transport.sentPrivateMessages.isEmpty)
|
||||
|
||||
// The drop is final: nothing is retained for later flushes.
|
||||
router.flushOutbox(for: peerID)
|
||||
#expect(dropped == ["m10"])
|
||||
}
|
||||
|
||||
@Test @MainActor
|
||||
func cleanupExpiredMessages_invokesOnMessageDroppedForExpiredOnly() async {
|
||||
let peerID = PeerID(str: "000000000000000b")
|
||||
let transport = MockTransport()
|
||||
let clock = MutableTestClock()
|
||||
|
||||
let router = MessageRouter(transports: [transport], now: { clock.now })
|
||||
var dropped: [String] = []
|
||||
router.onMessageDropped = { messageID, _ in dropped.append(messageID) }
|
||||
|
||||
router.sendPrivate("Old", to: peerID, recipientNickname: "Peer", messageID: "m11-old")
|
||||
clock.now = clock.now.addingTimeInterval(24 * 60 * 60 - 60)
|
||||
router.sendPrivate("Fresh", to: peerID, recipientNickname: "Peer", messageID: "m11-fresh")
|
||||
clock.now = clock.now.addingTimeInterval(120)
|
||||
|
||||
router.cleanupExpiredMessages()
|
||||
|
||||
#expect(dropped == ["m11-old"])
|
||||
|
||||
// The fresh message survived and still flushes once reachable.
|
||||
transport.reachablePeers.insert(peerID)
|
||||
router.flushOutbox(for: peerID)
|
||||
#expect(transport.sentPrivateMessages.map(\.messageID) == ["m11-fresh"])
|
||||
}
|
||||
|
||||
@Test @MainActor
|
||||
func enqueue_perPeerOverflowEvictionInvokesOnMessageDropped() async {
|
||||
let peerID = PeerID(str: "000000000000000c")
|
||||
let transport = MockTransport()
|
||||
|
||||
let router = MessageRouter(transports: [transport])
|
||||
var dropped: [String] = []
|
||||
router.onMessageDropped = { messageID, _ in dropped.append(messageID) }
|
||||
|
||||
// No reachable transport: everything queues. The cap is 100 per peer,
|
||||
// so the 101st enqueue evicts the oldest.
|
||||
for i in 0...100 {
|
||||
router.sendPrivate("Hello \(i)", to: peerID, recipientNickname: "Peer", messageID: "q\(i)")
|
||||
}
|
||||
|
||||
#expect(dropped == ["q0"])
|
||||
}
|
||||
|
||||
@Test @MainActor
|
||||
func sendReadReceipt_usesReachableTransport() async {
|
||||
let peerID = PeerID(str: "0000000000000003")
|
||||
@@ -135,3 +227,9 @@ struct MessageRouterTests {
|
||||
#expect(transport.sentFavoriteNotifications.count == 1)
|
||||
}
|
||||
}
|
||||
|
||||
/// Mutable wall clock injected into `MessageRouter` so TTL expiry is testable
|
||||
/// without real waiting.
|
||||
private final class MutableTestClock {
|
||||
var now = Date(timeIntervalSince1970: 1_700_000_000)
|
||||
}
|
||||
|
||||
@@ -1195,6 +1195,98 @@ final class NostrRelayManagerTests: XCTestCase {
|
||||
XCTAssertTrue(reconnected)
|
||||
}
|
||||
|
||||
func test_pendingSubscriptions_perRelayCapEvictsOldestByInsertionOrder() async {
|
||||
let relayURL = "wss://pending-cap.example"
|
||||
// Tor stalled: nothing flushes, so every REQ stays pending.
|
||||
let context = makeContext(permission: .denied, userTorEnabled: true, torEnforced: true, torIsReady: false)
|
||||
let cap = TransportConfig.nostrPendingSubscriptionsPerRelayCap
|
||||
|
||||
for i in 0..<(cap + 3) {
|
||||
context.manager.subscribe(filter: makeFilter(), id: "cap-sub-\(i)", relayUrls: [relayURL], handler: { _ in })
|
||||
}
|
||||
|
||||
XCTAssertEqual(context.manager.debugPendingSubscriptionCount(for: relayURL), cap)
|
||||
let pendingIDs = context.manager.debugPendingSubscriptionIDs(for: relayURL)
|
||||
// The three oldest entries were evicted; the newest survive.
|
||||
for i in 0..<3 {
|
||||
XCTAssertFalse(pendingIDs.contains("cap-sub-\(i)"), "expected cap-sub-\(i) to be evicted")
|
||||
}
|
||||
for i in 3..<(cap + 3) {
|
||||
XCTAssertTrue(pendingIDs.contains("cap-sub-\(i)"), "expected cap-sub-\(i) to be retained")
|
||||
}
|
||||
}
|
||||
|
||||
func test_pendingSubscriptions_staleEntriesSweptOnConnectAttempt() async {
|
||||
let relayURL = "wss://pending-sweep.example"
|
||||
// Tor stalled: the REQ stays pending and no socket ever opens.
|
||||
let context = makeContext(permission: .denied, userTorEnabled: true, torEnforced: true, torIsReady: false)
|
||||
|
||||
context.manager.subscribe(filter: makeFilter(), id: "stale-pending-sub", relayUrls: [relayURL], handler: { _ in })
|
||||
XCTAssertEqual(context.manager.debugPendingSubscriptionCount(for: relayURL), 1)
|
||||
|
||||
// Just under the TTL the entry survives a connect attempt.
|
||||
context.clock.now = context.clock.now.addingTimeInterval(TransportConfig.nostrPendingSubscriptionTTLSeconds - 1)
|
||||
context.manager.ensureConnections(to: [relayURL])
|
||||
XCTAssertEqual(context.manager.debugPendingSubscriptionCount(for: relayURL), 1)
|
||||
|
||||
// Past the TTL the next connect attempt sweeps it.
|
||||
context.clock.now = context.clock.now.addingTimeInterval(2)
|
||||
context.manager.ensureConnections(to: [relayURL])
|
||||
XCTAssertEqual(context.manager.debugPendingSubscriptionCount(for: relayURL), 0)
|
||||
}
|
||||
|
||||
func test_reconnectBackoff_appliesJitterWithinConfiguredBounds() async {
|
||||
let relayURL = "wss://jitter-bounds.example"
|
||||
// Pin the jitter source to the extremes and the midpoint of [0, 1).
|
||||
let jitter = JitterSequence([0.0, 1.0.nextDown, 0.25])
|
||||
let context = makeContext(permission: .denied, jitterUnit: { jitter.next() })
|
||||
// Persistent ping failure: every connect attempt fails and schedules
|
||||
// the next reconnect with an increasing attempt count.
|
||||
context.sessionFactory.pingErrorByURL[relayURL] = NSError(domain: NSURLErrorDomain, code: NSURLErrorTimedOut)
|
||||
|
||||
context.manager.ensureConnections(to: [relayURL])
|
||||
|
||||
var delays: [TimeInterval] = []
|
||||
for attempt in 1...3 {
|
||||
let scheduled = await waitUntil { context.scheduler.scheduled.count == 1 }
|
||||
XCTAssertTrue(scheduled, "reconnect for attempt \(attempt) was not scheduled")
|
||||
delays.append(context.scheduler.scheduled[0].delay)
|
||||
context.scheduler.runNext()
|
||||
}
|
||||
|
||||
// Bases: 1s, 2s, 4s. Jitter factors: 0.8, ~1.2, 0.9.
|
||||
XCTAssertEqual(delays[0], 0.8 * TransportConfig.nostrRelayInitialBackoffSeconds, accuracy: 1e-9)
|
||||
XCTAssertEqual(delays[1], 1.2 * TransportConfig.nostrRelayInitialBackoffSeconds * TransportConfig.nostrRelayBackoffMultiplier, accuracy: 1e-6)
|
||||
XCTAssertEqual(delays[2], 0.9 * TransportConfig.nostrRelayInitialBackoffSeconds * pow(TransportConfig.nostrRelayBackoffMultiplier, 2), accuracy: 1e-9)
|
||||
}
|
||||
|
||||
func test_reconnectBackoff_realRandomJitterStaysInBoundsAndVaries() async {
|
||||
let relayURL = "wss://jitter-random.example"
|
||||
let context = makeContext(permission: .denied, jitterUnit: { Double.random(in: 0..<1) })
|
||||
context.sessionFactory.pingErrorByURL[relayURL] = NSError(domain: NSURLErrorDomain, code: NSURLErrorTimedOut)
|
||||
|
||||
context.manager.ensureConnections(to: [relayURL])
|
||||
|
||||
var factors: [Double] = []
|
||||
for attempt in 1...5 {
|
||||
let scheduled = await waitUntil { context.scheduler.scheduled.count == 1 }
|
||||
XCTAssertTrue(scheduled, "reconnect for attempt \(attempt) was not scheduled")
|
||||
let base = min(
|
||||
TransportConfig.nostrRelayInitialBackoffSeconds * pow(TransportConfig.nostrRelayBackoffMultiplier, Double(attempt - 1)),
|
||||
TransportConfig.nostrRelayMaxBackoffSeconds
|
||||
)
|
||||
let factor = context.scheduler.scheduled[0].delay / base
|
||||
XCTAssertGreaterThanOrEqual(factor, 1.0 - TransportConfig.nostrRelayBackoffJitterRatio)
|
||||
XCTAssertLessThan(factor, 1.0 + TransportConfig.nostrRelayBackoffJitterRatio)
|
||||
factors.append(factor)
|
||||
context.scheduler.runNext()
|
||||
}
|
||||
|
||||
// A real RNG must not produce a constant delay across attempts
|
||||
// (5 identical uniform doubles is probability ~0).
|
||||
XCTAssertGreaterThan(Set(factors).count, 1)
|
||||
}
|
||||
|
||||
private func makeContext(
|
||||
permission: LocationChannelManager.PermissionState,
|
||||
favorites: Set<Data> = [],
|
||||
@@ -1202,7 +1294,8 @@ final class NostrRelayManagerTests: XCTestCase {
|
||||
userTorEnabled: Bool = false,
|
||||
torEnforced: Bool = false,
|
||||
torIsReady: Bool = true,
|
||||
torIsForeground: Bool = true
|
||||
torIsForeground: Bool = true,
|
||||
jitterUnit: @escaping () -> Double = { 0.5 } // 0.5 -> jitter factor 1.0 (no jitter)
|
||||
) -> RelayManagerTestContext {
|
||||
let permissionSubject = CurrentValueSubject<LocationChannelManager.PermissionState, Never>(permission)
|
||||
let favoritesSubject = CurrentValueSubject<Set<Data>, Never>(favorites)
|
||||
@@ -1228,7 +1321,8 @@ final class NostrRelayManagerTests: XCTestCase {
|
||||
scheduleAfter: { delay, action in
|
||||
scheduler.schedule(delay: delay, action: action)
|
||||
},
|
||||
now: { clock.now }
|
||||
now: { clock.now },
|
||||
jitterUnit: jitterUnit
|
||||
)
|
||||
)
|
||||
return RelayManagerTestContext(
|
||||
@@ -1305,6 +1399,20 @@ private final class MutableClock {
|
||||
}
|
||||
}
|
||||
|
||||
/// Deterministic jitter source: returns the queued values in order, then a
|
||||
/// neutral 0.5 (jitter factor 1.0) once exhausted.
|
||||
private final class JitterSequence {
|
||||
private var values: [Double]
|
||||
|
||||
init(_ values: [Double]) {
|
||||
self.values = values
|
||||
}
|
||||
|
||||
func next() -> Double {
|
||||
values.isEmpty ? 0.5 : values.removeFirst()
|
||||
}
|
||||
}
|
||||
|
||||
private final class MutableBool {
|
||||
var value: Bool
|
||||
|
||||
|
||||
Reference in New Issue
Block a user