diff --git a/Flipcash/Core/Controllers/BackfillQueue.swift b/Flipcash/Core/Controllers/BackfillQueue.swift new file mode 100644 index 000000000..e3b9b084e --- /dev/null +++ b/Flipcash/Core/Controllers/BackfillQueue.swift @@ -0,0 +1,68 @@ +// +// BackfillQueue.swift +// Flipcash +// +// Copyright © 2026 Code Inc. All rights reserved. +// + +import Foundation +import FlipcashCore + +/// The transcript fetches still waiting to run, in the order workers take them. A chat the user +/// opens fetches for itself and drops out of the queue, which puts it ahead of everything waiting. +struct BackfillQueue: Sendable { + + /// What a conversation needs to bring its local transcript level with the server. + enum Kind: Equatable, Sendable { + /// The local cursor lags the server's head: stream the missed window from it. + case delta + /// Nothing is cached at all: fetch the newest page, which also seats the cursor to head. + case newestPage + } + + struct Item: Equatable, Sendable { + let conversationID: ConversationID + let kind: Kind + } + + private var items: [Item] = [] + private var workers = 0 + + var count: Int { items.count } + + func contains(_ conversationID: ConversationID) -> Bool { + items.contains { $0.conversationID == conversationID } + } + + /// Appends `newItems`, skipping any conversation already queued. + mutating func enqueue(_ newItems: [Item]) { + for item in newItems where !contains(item.conversationID) { + items.append(item) + } + } + + /// Drops a queued conversation, for a caller that is fetching it itself. Returns false when it isn't queued. + @discardableResult + mutating func remove(_ conversationID: ConversationID) -> Bool { + guard let index = items.firstIndex(where: { $0.conversationID == conversationID }) else { return false } + items.remove(at: index) + return true + } + + /// Removes and returns the next item to run. + mutating func popNext() -> Item? { + items.isEmpty ? nil : items.removeFirst() + } + + /// Takes the worker slots a caller may start to drain the queue, up to `limit` running at once. + /// Each started worker calls ``releaseWorker()`` when it finishes. + mutating func claimWorkers(limit: Int) -> Int { + let claimed = max(0, min(limit - workers, items.count)) + workers += claimed + return claimed + } + + mutating func releaseWorker() { + workers -= 1 + } +} diff --git a/Flipcash/Core/Controllers/ConversationController.swift b/Flipcash/Core/Controllers/ConversationController.swift index 14c344536..afccddbbc 100644 --- a/Flipcash/Core/Controllers/ConversationController.swift +++ b/Flipcash/Core/Controllers/ConversationController.swift @@ -210,6 +210,12 @@ final class ConversationController { @ObservationIgnored let database: Database @ObservationIgnored let owner: KeyPair @ObservationIgnored private var startTask: Task? + /// The feed load in flight, joined by every caller that arrives before it finishes. + @ObservationIgnored private var feedLoadTask: Task? + /// How long a reconcile waits for the slowest feed, and for unread counts, before applying what landed. + @ObservationIgnored private let reconcileTiming: FeedReconcileTiming + /// The transcript fetches waiting to run, shared by every feed load. + @ObservationIgnored private var backfillQueue = BackfillQueue() /// The cache read's result, held for `hydrateIfReady()`; set off the main actor. @ObservationIgnored private let finishedCache = OSAllocatedUnfairLock(initialState: nil) @ObservationIgnored private var hasAppliedCache = false @@ -276,8 +282,10 @@ final class ConversationController { incomingTypingExpiry: Duration = .seconds(10), typingStoppedLinger: Duration = ConversationTyping.defaultStoppedLinger, typingExpiryClock: TypingExpiryClock = .continuous, + reconcileTiming: FeedReconcileTiming = .launch, receipts: ConversationReceiptReporter? = nil ) { + self.reconcileTiming = reconcileTiming self.fetching = fetching self.membership = membership self.viewerSettings = viewerSettings @@ -709,6 +717,8 @@ final class ConversationController { func stop() { startTask?.cancel() startTask = nil + feedLoadTask?.cancel() + feedLoadTask = nil if let extensionStoreWriteToken { ChatStoreWriteNotification.stopObserving(extensionStoreWriteToken) self.extensionStoreWriteToken = nil @@ -790,11 +800,24 @@ final class ConversationController { // MARK: - Feed + /// Loads every feed, then backfills transcripts behind it. A caller that arrives while a load is + /// in flight joins it rather than starting a second: launch's `start()`, the foreground hook that + /// fires with it, and a reconnect all converge here, and a second pass would refetch the same + /// feeds and backfill the same chats again. func loadFeed() async { - let loaded = await loadFeeds() - // Outside the loading flag: the feed itself is on screen the moment the - // conversations apply, and the transcripts fill behind it. - await backfillMessages(for: loaded) + if let feedLoadTask { + await feedLoadTask.value + return + } + let task = Task { [self] in + let loaded = await loadFeeds() + // Outside the loading flag, and after the list has been reconciled: backfill only + // changes a row whose newest message is newer than the one it holds. + await backfillMessages(for: loaded) + } + feedLoadTask = task + await task.value + feedLoadTask = nil } private func loadFeeds() async -> [Conversation] { @@ -803,12 +826,103 @@ final class ConversationController { isLoadingFeed = false hasResolvedFeed = true } - // Both DM feeds load concurrently and apply independently, so one - // type's failure doesn't drop the other's conversations. - async let contact = loadFeed(type: .contactDm) - async let tip = loadFeed(type: .tipDm) - async let groups = loadGroupFeed() - return await contact + tip + groups + // With nothing listed there is no list to keep steady, so each feed shows as it lands. Both + // DM feeds load concurrently and apply independently, so one type's failure doesn't drop + // the other's conversations. + guard !store.conversations.isEmpty else { + async let contact = loadFeed(type: .contactDm) + async let tip = loadFeed(type: .tipDm) + async let groups = loadGroupFeed() + return await contact + tip + groups + } + return await reconcileFeeds() + } + + /// Fetches the three feeds together and applies them to the listed chats as one update, so the + /// list moves once instead of once per feed, per preview and per unread count. + /// + /// Waits at most ``FeedReconcileTiming/feedCap`` for the slowest feed and then + /// ``FeedReconcileTiming/unreadCap`` for unread counts; past either it applies what has landed. + /// A feed that lands after the cap applies on its own. Backfill is not waited on. + private func reconcileFeeds() async -> [Conversation] { + let sources: [FeedArrivals.Source] = [.dm(.contactDm), .dm(.tipDm), .groups] + let arrivals = FeedArrivals(expecting: sources.count) + let fetches = sources.map { source in + Task { [self] in + let feed = await fetchFeed(source) + guard !Task.isCancelled else { + arrivals.arrive(source, nil) + return + } + if arrivals.isClosed { + guard let feed else { return } + withTransaction(Self.unanimated) { applyFeed(feed, from: source) } + await backfillMessages(for: feed) + } else { + arrivals.arrive(source, feed) + } + } + } + let landed = await withTaskCancellationHandler { + await arrivals.wait(upTo: reconcileTiming.feedCap) + } onCancel: { + fetches.forEach { $0.cancel() } + } + guard !Task.isCancelled else { return [] } + + await resolveUnreadCounts(in: landed.values.flatMap { $0 }) + + withTransaction(Self.unanimated) { + for (source, feed) in landed { + applyFeed(feed, from: source) + } + watermarkStamps.release() + } + return landed.values.flatMap { $0 } + } + + private static var unanimated: Transaction { + var transaction = Transaction() + transaction.disablesAnimations = true + return transaction + } + + /// Resolves the unread count of every chat in `conversations` that can't show one yet, so the + /// numbers arrive with the list rather than popping in row by row. The feed carries the READ + /// pointer as a message id, not that message's `unreadSeq`, so a count is only local when the + /// pointer's message is stored; the rest cost a `GetMessage` each. Gives up waiting after + /// ``FeedReconcileTiming/unreadCap``; the lookups keep going and fill their rows in as they land. + private func resolveUnreadCounts(in conversations: [Conversation]) async { + unresolvedUnreadQueue = conversations + .filter { $0.hasUnread(for: selfUserID) && unreadCount(for: $0) == nil } + guard !unresolvedUnreadQueue.isEmpty else { return } + watermarkStamps.hold() + let wait = BoundedWait(expecting: unresolvedUnreadQueue.count) + for _ in 0.. [Conversation]? { + switch source { + case .dm(let type): await fetchFeed(type: type) + case .groups: await fetchGroupFeed() + } + } + + private func applyFeed(_ feed: [Conversation], from source: FeedArrivals.Source) { + switch source { + case .dm(let type): applyFeed(feed, type: type) + case .groups: applyGroupFeed(feed) + } } /// Loads the group feed, returning the groups the server reported — empty if the load failed or @@ -820,24 +934,32 @@ final class ConversationController { /// this feed, so a type-scoped replace would delete it out from under its own open screen. @discardableResult func loadGroupFeed() async -> [Conversation] { + guard let groups = await fetchGroupFeed() else { return [] } + applyGroupFeed(groups) + return groups + } + + private func fetchGroupFeed() async -> [Conversation]? { do { - let groups = try await fetching.getGroupChatFeed(owner: owner) - let departed = store.setGroupFeed(groups) - reconcileHidden() - persist(operation: "replace-group-feed") { - try database.replaceGroupFeed(groups.map(withStoreMembers), departed: departed) - } - resendUnsyncedReadPointers() - // Same repair the DM feeds need: the store refuses a tombstone as a preview, so a chat whose - // newest message is deleted seats blank without this. - for group in groups where group.lastMessage?.isDeleted == true { - refreshFeedPreview(for: group.id) - } - return groups + return try await fetching.getGroupChatFeed(owner: owner) } catch { logger.error("Failed to load group chat feed", metadata: ["error": "\(error)"]) ErrorReporting.captureError(error, reason: "Failed to load group chat feed") - return [] + return nil + } + } + + private func applyGroupFeed(_ groups: [Conversation]) { + let departed = store.setGroupFeed(groups) + reconcileHidden() + persist(operation: "replace-group-feed") { + try database.replaceGroupFeed(groups.map(withStoreMembers), departed: departed) + } + resendUnsyncedReadPointers() + // Same repair the DM feeds need: the store refuses a tombstone as a preview, so a chat whose + // newest message is deleted seats blank without this. + for group in groups where group.lastMessage?.isDeleted == true { + refreshFeedPreview(for: group.id) } } @@ -845,27 +967,35 @@ final class ConversationController { /// empty if the load failed. @discardableResult func loadFeed(type: ConversationType) async -> [Conversation] { + guard let conversations = await fetchFeed(type: type) else { return [] } + applyFeed(conversations, type: type) + return conversations + } + + private func fetchFeed(type: ConversationType) async -> [Conversation]? { do { - let conversations = try await fetching.getDmChatFeed(owner: owner, type: type) - store.setFeed(conversations, type: type) - reconcileHidden() - persist(operation: "replace-feed") { try database.replaceConversationFeed(conversations, type: type) } - resendUnsyncedReadPointers() - // The store refuses a tombstone as a preview, so a chat whose newest message is deleted - // seats blank here. Fill it from the newest visible message already cached — the feed - // reloads on every launch and foreground, so without this the row stays blank until the - // transcript is opened. - for conversation in conversations where conversation.lastMessage?.isDeleted == true { - refreshFeedPreview(for: conversation.id) - } - return conversations + return try await fetching.getDmChatFeed(owner: owner, type: type) } catch { logger.error("Failed to load conversation feed", metadata: [ "type": "\(type)", "error": "\(error)", ]) ErrorReporting.captureError(error, reason: "Failed to load conversation feed") - return [] + return nil + } + } + + private func applyFeed(_ conversations: [Conversation], type: ConversationType) { + store.setFeed(conversations, type: type) + reconcileHidden() + persist(operation: "replace-feed") { try database.replaceConversationFeed(conversations, type: type) } + resendUnsyncedReadPointers() + // The store refuses a tombstone as a preview, so a chat whose newest message is deleted + // seats blank here. Fill it from the newest visible message already cached — the feed + // reloads on every launch and foreground, so without this the row stays blank until the + // transcript is opened. + for conversation in conversations where conversation.lastMessage?.isDeleted == true { + refreshFeedPreview(for: conversation.id) } } @@ -1027,20 +1157,6 @@ final class ConversationController { /// rates, and history syncs it shares the launch with. private static let backfillConcurrency = 4 - /// What a conversation needs to bring its local transcript level with the server. - private enum Backfill: Equatable { - /// The local cursor lags the server's head: stream the missed window from it. - case delta - /// Nothing is cached at all: fetch the newest page, which also seats the cursor to head. - case newestPage - } - - /// One conversation's backfill, queued. - private struct BackfillTask: Sendable { - let conversationID: ConversationID - let kind: Backfill - } - /// Brings every conversation in a freshly-loaded feed up to the server's head, so a transcript is /// there when the chat is opened rather than fetched on arrival. /// @@ -1053,9 +1169,9 @@ final class ConversationController { // Deduped by id: the feed loads by type, and a conversation reported under more than one type // must not queue its transcript twice. var seen: Set = [] - let work = conversations.compactMap { conversation -> BackfillTask? in + let work = conversations.compactMap { conversation -> BackfillQueue.Item? in guard seen.insert(conversation.id).inserted else { return nil } - return backfill(for: conversation).map { BackfillTask(conversationID: conversation.id, kind: $0) } + return backfill(for: conversation).map { BackfillQueue.Item(conversationID: conversation.id, kind: $0) } } guard !work.isEmpty else { return } logger.info("Backfilling chat history", metadata: [ @@ -1063,31 +1179,36 @@ final class ConversationController { "delta": "\(work.filter { $0.kind == .delta }.count)", ]) + backfillQueue.enqueue(work) + + // A fixed number of workers draw from the shared queue, so it drains at a steady width instead + // of in lock-stepped batches, and a chat opened meanwhile can leave the queue early. A call + // that arrives while workers are running only tops them up to the width. + let workers = backfillQueue.claimWorkers(limit: Self.backfillConcurrency) await withTaskGroup(of: Void.self) { group in - var next = 0 - // Start a window of tasks, then replace each as it finishes, so the queue drains at a - // steady width instead of in lock-stepped batches. - while next < min(Self.backfillConcurrency, work.count) { - group.addTask { [task = work[next]] in await self.perform(task) } - next += 1 - } - while await group.next() != nil, next < work.count { - group.addTask { [task = work[next]] in await self.perform(task) } - next += 1 + for _ in 0.. Backfill? { + private func backfill(for conversation: Conversation) -> BackfillQueue.Kind? { let conversationID = conversation.id guard !messageLoadsInFlight.contains(conversationID) else { return nil } @@ -1619,6 +1740,8 @@ final class ConversationController { return } messageLoadsInFlight.insert(conversationID) + // This load is the backfill: a queued one would fetch the same page again. + backfillQueue.remove(conversationID) defer { messageLoadsInFlight.remove(conversationID) } do { let messages = try await messaging.getMessages(owner: owner, conversationID: conversationID, before: nil) diff --git a/Flipcash/Core/Controllers/FeedReconcile.swift b/Flipcash/Core/Controllers/FeedReconcile.swift new file mode 100644 index 000000000..79b2ad77b --- /dev/null +++ b/Flipcash/Core/Controllers/FeedReconcile.swift @@ -0,0 +1,100 @@ +// +// FeedReconcile.swift +// Flipcash +// +// Copyright © 2026 Code Inc. All rights reserved. +// + +import Foundation +import FlipcashCore + +/// How long the launch reconcile waits before it applies what has landed. +struct FeedReconcileTiming: Sendable { + /// The longest the reconcile waits for the slowest of the three feeds. + var feedCap: Duration + /// The longest it then waits for unread counts to resolve. + var unreadCap: Duration + + /// Matches the Android launch sync. + static let launch = FeedReconcileTiming(feedCap: .seconds(2), unreadCap: .milliseconds(1500)) +} + +/// Waits for a fixed number of units of work to finish, or for a deadline, whichever comes first. +/// The work itself is never cancelled by the deadline: it keeps running and reports late through +/// ``isClosed``. +@MainActor +final class BoundedWait { + + private var remaining: Int + private var isDone = false + private var waiter: CheckedContinuation? + private var deadline: Task? + + /// Whether ``wait(upTo:)`` has already returned, so a unit finishing now is late. + private(set) var isClosed = false + + init(expecting count: Int) { + remaining = count + isDone = count <= 0 + } + + /// Records one unit finished. + func signal() { + remaining -= 1 + if remaining <= 0 { finish() } + } + + /// Suspends until every unit has signalled or `cap` has passed, then closes the wait. + func wait(upTo cap: Duration) async { + if !isDone { + deadline = Task { [weak self] in + try? await Task.sleep(for: cap) + guard !Task.isCancelled else { return } + self?.finish() + } + await withCheckedContinuation { waiter = $0 } + } + deadline?.cancel() + deadline = nil + isClosed = true + } + + private func finish() { + isDone = true + waiter?.resume() + waiter = nil + } +} + +/// The three feeds' results as they land, with the wait for them. A feed that failed settles +/// without a result. +@MainActor +final class FeedArrivals { + + enum Source: Hashable, Sendable { + case dm(ConversationType) + case groups + } + + private var feeds: [Source: [Conversation]] = [:] + private let wait: BoundedWait + + init(expecting count: Int) { + wait = BoundedWait(expecting: count) + } + + /// Whether ``wait(upTo:)`` has returned, so a feed landing now is late. + var isClosed: Bool { wait.isClosed } + + /// Records a feed's result, `nil` when it failed. + func arrive(_ source: Source, _ feed: [Conversation]?) { + if let feed { feeds[source] = feed } + wait.signal() + } + + /// The feeds that landed once all have settled or `cap` has passed. + func wait(upTo cap: Duration) async -> [Source: [Conversation]] { + await wait.wait(upTo: cap) + return feeds + } +} diff --git a/Flipcash/Core/Controllers/KnownAuthorDirectory.swift b/Flipcash/Core/Controllers/KnownAuthorDirectory.swift index e079b5015..4b185ae0f 100644 --- a/Flipcash/Core/Controllers/KnownAuthorDirectory.swift +++ b/Flipcash/Core/Controllers/KnownAuthorDirectory.swift @@ -6,6 +6,7 @@ // import Foundation +import os import FlipcashCore import FlipcashStore @@ -48,6 +49,9 @@ final class KnownAuthorDirectory { /// would stop a later launch asking again once they have picked a name. @ObservationIgnored private var nameless: [UserID: ConversationMember] = [:] + /// The launch-time read's result, held for ``hydrateIfReady()``; set off the main actor. + @ObservationIgnored private let preloaded = OSAllocatedUnfairLock<[UserID: ConversationMember]?>(initialState: nil) + /// The reload in flight, so concurrent requests collapse onto one read of the same store. @ObservationIgnored private var reloadTask: Task? @@ -69,6 +73,27 @@ final class KnownAuthorDirectory { ) } + /// Starts the first read off the main actor and lands it in ``snapshot`` as soon as it finishes, + /// so the Chats list takes the table in the same frame it paints the cached feed. Without it the + /// snapshot is empty until a `.task` runs after that paint, and every group preview redraws with + /// its sender prefix a moment later. + /// + /// The landing is pushed from here rather than left to the screen's `onAppear`, which can fire + /// before the read has finished and then never asks again. + func preload() { + Task.detached { [weak self, read, preloaded] in + guard let members = try? read() else { return } + preloaded.withLock { $0 = members } + await self?.hydrateIfReady() + } + } + + /// Lands the ``preload()`` result if it has finished and nothing newer has landed. Never waits. + func hydrateIfReady() { + guard snapshot === Snapshot.empty, let members = preloaded.withLock({ $0 }), !members.isEmpty else { return } + snapshot = Snapshot(membersByUserID: members.merging(nameless) { cached, _ in cached }) + } + /// Re-reads the local cache off the main thread, landing a new ``snapshot`` only if the table /// changed. Safe to call whenever a roster might have moved; a repeat while one is running is /// dropped rather than queued. diff --git a/Flipcash/Core/Controllers/ReadWatermarkStamps.swift b/Flipcash/Core/Controllers/ReadWatermarkStamps.swift index b30a80fd2..8559cd6a8 100644 --- a/Flipcash/Core/Controllers/ReadWatermarkStamps.swift +++ b/Flipcash/Core/Controllers/ReadWatermarkStamps.swift @@ -26,6 +26,22 @@ final class ReadWatermarkStamps { private var stamps: [Key: UInt64] = [:] @ObservationIgnored private var inFlight: Set = [] @ObservationIgnored private var missing: Set = [] + @ObservationIgnored private var isHolding = false + @ObservationIgnored private var held: [Key: UInt64] = [:] + + /// Starts keeping fetched stamps out of view, so a batch of lookups shows up as one change at + /// ``release()`` instead of one per row. + func hold() { + isHolding = true + } + + /// Publishes every stamp fetched since ``hold()`` at once; later fetches publish as they land. + func release() { + isHolding = false + guard !held.isEmpty else { return } + stamps.merge(held) { _, new in new } + held = [:] + } /// The fetched stamp for `key`, or nil when it hasn't been fetched. func stamp(for key: Key) -> UInt64? { @@ -35,11 +51,15 @@ final class ReadWatermarkStamps { /// Fetches `key`'s message unless it's already stamped, in flight, or known missing. A failed /// fetch leaves the key unresolved, so the next call tries again. func resolve(_ key: Key, fetch: () async throws -> ConversationMessage?) async throws { - guard stamps[key] == nil, !inFlight.contains(key), !missing.contains(key) else { return } + guard stamps[key] == nil, held[key] == nil, !inFlight.contains(key), !missing.contains(key) else { return } inFlight.insert(key) defer { inFlight.remove(key) } if let message = try await fetch() { - stamps[key] = message.unreadSeq + if isHolding { + held[key] = message.unreadSeq + } else { + stamps[key] = message.unreadSeq + } } else { missing.insert(key) } diff --git a/Flipcash/Core/Screens/Conversation/ConversationScreen.swift b/Flipcash/Core/Screens/Conversation/ConversationScreen.swift index 16787f3a0..9b4f69fd6 100644 --- a/Flipcash/Core/Screens/Conversation/ConversationScreen.swift +++ b/Flipcash/Core/Screens/Conversation/ConversationScreen.swift @@ -526,7 +526,9 @@ struct ConversationScreen: View { // allowed to see — // which is why this reads `withholdsTranscript` and not `obscuresTranscript`: a chat // whose rules haven't landed yet is blurred without yet refusing anything. - showsGatePlaceholder: gate.withholdsTranscript && (coordinator?.items.isEmpty ?? true), + // Also stands in while an empty transcript's first load is out, so a chat with nothing + // cached opens on the placeholder and paints once with its history. + showsGatePlaceholder: (gate.withholdsTranscript || (chatExists && !didInitialRead)) && (coordinator?.items.isEmpty ?? true), gateMintName: gateMintName, onGateAddFunds: addFunds, onGateJoin: joinChat, diff --git a/Flipcash/Core/Screens/Main/Tips/TipConversationsScreen.swift b/Flipcash/Core/Screens/Main/Tips/TipConversationsScreen.swift index e3df919de..18fde2b71 100644 --- a/Flipcash/Core/Screens/Main/Tips/TipConversationsScreen.swift +++ b/Flipcash/Core/Screens/Main/Tips/TipConversationsScreen.swift @@ -63,6 +63,9 @@ struct TipConversationsScreen: View { // Takes a finished cache read before the first frame, instead of drawing empty until launch // work lets the controller's own hydration run. .onAppear { + // Names first: the cache hydration below invalidates the rows, and they must not + // render once without them. + sessionContainer.knownAuthors.hydrateIfReady() conversationController.hydrateIfReady() } // Every counterpart, not just the rows on screen. A row's own `.task` fires when the row is diff --git a/Flipcash/Core/Session/SessionAuthenticator.swift b/Flipcash/Core/Session/SessionAuthenticator.swift index 7ef123ac0..f89983c43 100644 --- a/Flipcash/Core/Session/SessionAuthenticator.swift +++ b/Flipcash/Core/Session/SessionAuthenticator.swift @@ -727,6 +727,7 @@ final class SessionContainer { self.profileAvatars = ProfileAvatarStore(flipClient: flipClient, owner: session.ownerKeyPair) self.knownAuthors = KnownAuthorDirectory(database: database, flipClient: flipClient, owner: session.ownerKeyPair) + knownAuthors.preload() } fileprivate func injectingEnvironment(into view: SomeView) -> some View where SomeView: View { diff --git a/FlipcashCore/Sources/FlipcashCore/Models/Conversation/ConversationStore.swift b/FlipcashCore/Sources/FlipcashCore/Models/Conversation/ConversationStore.swift index 84f34c8d4..864630cff 100644 --- a/FlipcashCore/Sources/FlipcashCore/Models/Conversation/ConversationStore.swift +++ b/FlipcashCore/Sources/FlipcashCore/Models/Conversation/ConversationStore.swift @@ -59,7 +59,7 @@ public struct ConversationStore: Sendable { /// Replace the feed from a paged load, sorted most-recent-activity first. public mutating func setFeed(_ conversations: [Conversation]) { - let merged = conversations.map { keepingSelfReadPointer(seated($0)) } + let merged = conversations.map { keepingNewerActivity(keepingSelfReadPointer(seated($0))) } self.conversations = merged.sorted { $0.lastActivity > $1.lastActivity } } @@ -69,7 +69,7 @@ public struct ConversationStore: Sendable { public mutating func setFeed(_ conversations: [Conversation], type: ConversationType) { // Only the incoming rows are merged: re-merging the other types against themselves would // read the store's own copy as the server acknowledging an unsynced READ pointer. - let incoming = conversations.filter { $0.type == type }.map { keepingSelfReadPointer(seated($0)) } + let incoming = conversations.filter { $0.type == type }.map { keepingNewerActivity(keepingSelfReadPointer(seated($0))) } self.conversations = (self.conversations.filter { $0.type != type } + incoming) .sorted { $0.lastActivity > $1.lastActivity } } @@ -306,9 +306,10 @@ public struct ConversationStore: Sendable { } /// Bump a conversation's last activity and re-sort the feed. No-ops for a conversation not in the - /// feed. + /// feed, and for a date that is not later than the one it holds. public mutating func advanceLastActivity(to date: Date, in conversationID: ConversationID) { - guard let index = conversations.firstIndex(where: { $0.id == conversationID }) else { return } + guard let index = conversations.firstIndex(where: { $0.id == conversationID }), + date > conversations[index].lastActivity else { return } conversations[index].lastActivity = date sort() } @@ -559,7 +560,7 @@ public struct ConversationStore: Sendable { } private mutating func upsert(_ conversation: Conversation) { - var conversation = keepingSelfReadPointer(seated(conversation)) + var conversation = keepingNewerActivity(keepingSelfReadPointer(seated(conversation))) if let index = conversations.firstIndex(where: { $0.id == conversation.id }) { if conversation.type == .group { conversation.members = mergedMembers(conversation.members, over: conversations[index].members) @@ -606,6 +607,22 @@ public struct ConversationStore: Sendable { return conversation } + /// A server copy with the stored row's `lastActivity` and `lastMessage` kept when they are newer. + /// A feed or metadata response can be older than a stream event that landed while it was in + /// flight; taking it as-is would roll the row back, and a launch would show the list move twice. + private func keepingNewerActivity(_ conversation: Conversation) -> Conversation { + guard let stored = conversations.first(where: { $0.id == conversation.id }) else { return conversation } + var conversation = conversation + if stored.lastActivity > conversation.lastActivity { + conversation.lastActivity = stored.lastActivity + } + if let held = stored.lastMessage, let incoming = conversation.lastMessage, + (held.id, held.eventSequence) > (incoming.id, incoming.eventSequence) { + conversation.lastMessage = held + } + return conversation + } + /// `conversation` with the self READ pointer raised to the one the store already holds, when the /// server's copy is behind it. The server never lowers a pointer, so a lower copy means an advance /// it hasn't taken yet; that advance is kept and recorded as unsynced for the caller to re-send. diff --git a/FlipcashCore/Sources/FlipcashStore/Database+Conversations.swift b/FlipcashCore/Sources/FlipcashStore/Database+Conversations.swift index 6cae74324..9435ac3eb 100644 --- a/FlipcashCore/Sources/FlipcashStore/Database+Conversations.swift +++ b/FlipcashCore/Sources/FlipcashStore/Database+Conversations.swift @@ -85,11 +85,17 @@ nonisolated extension Database { let rows = try reader.prepareRowIterator(c.table.order(c.lastActivity.desc)) return try rows.map { row in let id = row[c.id] + let lastMessage = try latestMessage(conversationId: id) + // The row's activity only moves when the app writes the conversation. The notification + // extension writes message rows alone while the app is closed, so a chat that received a + // push would otherwise launch with its new preview at its old position. The live stream + // advances activity to a message's date the same way. + let storedActivity = Date(timeIntervalSinceReferenceDate: row[c.lastActivity]) return Conversation( id: ConversationID(data: id), members: membersByConversation[id] ?? [], - lastMessage: try latestMessage(conversationId: id), - lastActivity: Date(timeIntervalSinceReferenceDate: row[c.lastActivity]), + lastMessage: lastMessage, + lastActivity: max(storedActivity, lastMessage?.date ?? storedActivity), type: ConversationType(rawValue: row[c.type]) ?? .contactDm, isHidden: row[c.isHidden], title: row[c.title], diff --git a/FlipcashCore/Tests/FlipcashCoreTests/ConversationStoreTests.swift b/FlipcashCore/Tests/FlipcashCoreTests/ConversationStoreTests.swift index 21e69ee33..0e52c546b 100644 --- a/FlipcashCore/Tests/FlipcashCoreTests/ConversationStoreTests.swift +++ b/FlipcashCore/Tests/FlipcashCoreTests/ConversationStoreTests.swift @@ -51,6 +51,15 @@ struct ConversationStoreTests { // MARK: - Feed + @Test("A feed older than a live activity bump doesn't roll the row back") + func feedKeepsNewerActivity() { + var store = ConversationStore() + store.setFeed([conversation(1, lastActivity: 100), conversation(2, lastActivity: 50)]) + store.advanceLastActivity(to: Date(timeIntervalSince1970: 900), in: conversationID(2)) + store.setFeed([conversation(1, lastActivity: 100), conversation(2, lastActivity: 50)]) + #expect(store.conversations.map(\.id) == [conversationID(2), conversationID(1)]) + } + @Test("setFeed sorts by last activity, most recent first") func feedSortsByActivity() { var store = ConversationStore() @@ -98,6 +107,46 @@ struct ConversationStoreTests { ) } + @Test("A feed preview older than the stored one doesn't replace it, and a newer one does") + func feedKeepsNewerPreview() { + var store = ConversationStore() + var row = conversation(1, lastActivity: 100) + row.lastMessage = message(5, "live", eventSequence: 5) + store.setFeed([row]) + + var stale = conversation(1, lastActivity: 100) + stale.lastMessage = message(3, "stale", eventSequence: 3) + store.setFeed([stale], type: .contactDm) + #expect(store.conversations.first?.lastMessage?.id.value == 5) + store.setFeed([stale]) + #expect(store.conversations.first?.lastMessage?.id.value == 5) + + var fresh = conversation(1, lastActivity: 100) + fresh.lastMessage = message(8, "fresh", eventSequence: 8) + store.setFeed([fresh], type: .contactDm) + #expect(store.conversations.first?.lastMessage?.id.value == 8) + } + + @Test("advanceLastActivity ignores an equal or older date, so the order never moves back") + func advanceActivityNeverRegresses() { + var store = ConversationStore() + store.setFeed([conversation(1, lastActivity: 100), conversation(2, lastActivity: 200)]) + store.advanceLastActivity(to: Date(timeIntervalSince1970: 50), in: conversationID(2)) + store.advanceLastActivity(to: Date(timeIntervalSince1970: 200), in: conversationID(2)) + #expect(store.conversations.map(\.id) == [conversationID(2), conversationID(1)]) + #expect(store.conversations.first?.lastActivity == Date(timeIntervalSince1970: 200)) + } + + @Test("setFeedPreview with an equal message leaves the row untouched") + func equalPreviewIsANoOp() { + var store = ConversationStore() + var row = conversation(1, lastActivity: 100) + row.lastMessage = message(5, "same", eventSequence: 5) + store.setFeed([row]) + store.setFeedPreview(message(5, "same", eventSequence: 5), in: conversationID(1)) + #expect(store.conversations == [row]) + } + @Test("advanceLastActivity moves the conversation to the front") func advanceActivityResorts() { var store = ConversationStore() diff --git a/FlipcashTests/Chat/KnownAuthorDirectoryTests.swift b/FlipcashTests/Chat/KnownAuthorDirectoryTests.swift index 04b0c0f9a..c253b5b50 100644 --- a/FlipcashTests/Chat/KnownAuthorDirectoryTests.swift +++ b/FlipcashTests/Chat/KnownAuthorDirectoryTests.swift @@ -66,6 +66,32 @@ struct KnownAuthorDirectoryTests { ) } + @Test("A preload that finishes after the screen appeared lands without another hydrateIfReady") + func preloadLandsItselfWhenItFinishesLate() async throws { + let recorder = Recorder() + let sender = UserID() + let gate = DispatchSemaphore(value: 0) + let directory = KnownAuthorDirectory( + read: { + gate.wait() + return recorder.withState { $0.table } + }, + fetch: { _ in throw Unreachable.offline }, + cache: { _, _ in } + ) + recorder.withState { $0.table[sender] = ConversationMember(userID: sender, displayName: "Ada") } + + directory.preload() + directory.hydrateIfReady() // the screen's onAppear, with the read still running + #expect(directory.snapshot.membersByUserID.isEmpty) + + gate.signal() + for _ in 0..<200 where directory.snapshot.membersByUserID.isEmpty { + try await Task.sleep(for: .milliseconds(10)) + } + #expect(directory.snapshot.membersByUserID[sender]?.displayName == "Ada") + } + @Test("A sender the local cache cannot name is fetched, cached and landed in the snapshot") func fetchesAndCachesAnUnknownSender() async { let recorder = Recorder() @@ -79,6 +105,38 @@ struct KnownAuthorDirectoryTests { #expect(directory.snapshot.membersByUserID[sender]?.displayName == "Ada") } + @Test("Names persisted by an earlier session are in the snapshot at the first cached render") + func namesPersistedByAnEarlierSessionLandOnFirstRender() async throws { + let recorder = Recorder() + let sender = UserID() + recorder.withState { $0.table[sender] = ConversationMember(userID: sender, displayName: "Grace") } + let directory = directory(recorder) + + directory.preload() + for _ in 0..<200 where directory.snapshot === KnownAuthorDirectory.Snapshot.empty { + directory.hydrateIfReady() + try await Task.sleep(for: .milliseconds(10)) + } + + #expect(directory.snapshot.membersByUserID[sender]?.displayName == "Grace") + #expect(recorder.requests.isEmpty) + } + + @Test("A reload that finds the same names does not replace the snapshot") + func unchangedReloadKeepsTheSnapshot() async { + let recorder = Recorder() + let sender = UserID() + recorder.withState { $0.table[sender] = ConversationMember(userID: sender, displayName: "Grace") } + let directory = directory(recorder) + await directory.reload() + let before = directory.snapshot + + await directory.reload() + await directory.resolve([sender]) + + #expect(directory.snapshot === before) + } + @Test("A sender the local cache already names costs no round trip") func skipsSendersTheCacheAlreadyNames() async { let recorder = Recorder() diff --git a/FlipcashTests/ChatListLaunchSyncTests.swift b/FlipcashTests/ChatListLaunchSyncTests.swift new file mode 100644 index 000000000..1285faa93 --- /dev/null +++ b/FlipcashTests/ChatListLaunchSyncTests.swift @@ -0,0 +1,322 @@ +// +// ChatListLaunchSyncTests.swift +// FlipcashTests +// + +import Testing +import Foundation +import Observation +import FlipcashCore +import FlipcashStore +@testable import Flipcash + +@MainActor +@Suite("Chat list launch sync") +struct ChatListLaunchSyncTests { + + private let me = UUID() + + private func waitUntil(_ condition: () -> Bool, sourceLocation: SourceLocation = #_sourceLocation) async throws { + for _ in 0..<100 where !condition() { + try? await Task.sleep(for: .milliseconds(20)) + } + try #require(condition(), "Timed out waiting for condition after ~2s", sourceLocation: sourceLocation) + } + + private func makeController( + _ mock: MockConversations, + database: Database? = nil, + timing: FeedReconcileTiming = .launch + ) -> ConversationController { + ConversationController( + fetching: mock, membership: mock, viewerSettings: mock, messaging: mock, streaming: mock, + contactNaming: MockDMContactNaming(), + database: database ?? (try! Database.makeTemp().database), + owner: .generate()!, selfUserID: me, + reconcileTiming: timing + ) + } + + private func message(_ id: UInt64, seq: UInt64? = nil, at: TimeInterval = 0) -> ConversationMessage { + ConversationMessage( + id: MessageID(value: id), senderID: nil, content: .text("m\(id)"), + date: Date(timeIntervalSince1970: at), unreadSeq: seq ?? id, eventSequence: seq ?? id + ) + } + + private func tipDm( + _ n: UInt8, activity: TimeInterval, last: ConversationMessage? = nil, readPointer: UInt64? = nil + ) -> Conversation { + Conversation( + id: .test(n), + members: [ConversationMember(userID: me, displayName: "", readPointer: readPointer.map(MessageID.init(value:)))], + lastMessage: last, lastActivity: Date(timeIntervalSince1970: activity), type: .tipDm + ) + } + + private func group(_ n: UInt8, activity: TimeInterval) -> Conversation { + Conversation(id: .test(n), members: [], lastMessage: nil, lastActivity: Date(timeIntervalSince1970: activity), type: .group) + } + + private func cached(_ rows: [Conversation]) throws -> Database { + let db = try Database.makeTemp().database + try db.replaceConversationFeed(rows, type: .tipDm) + return db + } + + private func ids(_ controller: ConversationController) -> [ConversationID] { + controller.chatListConversations.map(\.id) + } + + // MARK: - One reconcile + + @Test("three feeds land as one update with unread counts already populated") + func feedsApplyAsOneUpdate() async throws { + let unread = tipDm(2, activity: 200, last: message(5), readPointer: 1) + let db = try cached([tipDm(1, activity: 100), tipDm(2, activity: 50, last: message(5), readPointer: 1)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 100), unread] + mock.groupFeed = [group(9, activity: 300)] + mock.feedDelay = .milliseconds(150) + mock.singleMessages = [MessageID(value: 1): message(1, seq: 3)] + let controller = makeController(mock, database: db) + + let observer = ListObserver(controller, unreadFor: unread) + controller.start() + try await waitUntil { ids(controller) == [.test(1), .test(2)] } + observer.start() + + await controller.loadFeed() + try await Task.sleep(for: .milliseconds(50)) + + #expect(ids(controller) == [.test(9), .test(2), .test(1)]) + let first = try #require(observer.shown.first) + #expect(first.ids == [.test(9), .test(2), .test(1)]) + #expect(first.unread != nil) + #expect(observer.shown.allSatisfy { $0.ids == first.ids && $0.unread == first.unread }) + controller.stop() + } + + @Test("a slow feed is capped: the others apply, and the slow one applies when it lands") + func slowFeedIsCapped() async throws { + let db = try cached([tipDm(1, activity: 100)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 100), tipDm(2, activity: 200)] + mock.groupFeed = [group(9, activity: 300)] + mock.groupFeedDelay = .milliseconds(800) + let controller = makeController(mock, database: db, timing: .init(feedCap: .milliseconds(100), unreadCap: .milliseconds(100))) + controller.start() + + try await waitUntil { ids(controller) == [.test(2), .test(1)] } + #expect(mock.groupFeedCalls == 1) + try await waitUntil { ids(controller) == [.test(9), .test(2), .test(1)] } + controller.stop() + } + + @Test("slow unread lookups don't hold the list past their cap") + func unreadIsCapped() async throws { + let unread = tipDm(2, activity: 200, last: message(5), readPointer: 1) + let db = try cached([tipDm(1, activity: 100)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 100), unread] + mock.singleMessages = [MessageID(value: 1): message(1, seq: 3)] + mock.singleMessageDelay = .milliseconds(800) + let controller = makeController(mock, database: db, timing: .init(feedCap: .seconds(2), unreadCap: .milliseconds(100))) + controller.start() + + try await waitUntil { ids(controller) == [.test(2), .test(1)] } + #expect(controller.unreadCount(for: unread) == nil) + try await waitUntil { controller.unreadCount(for: unread) != nil } + controller.stop() + } + + @Test("the reconcile does not wait on backfill") + func reconcileDoesNotWaitOnBackfill() async throws { + let db = try cached([tipDm(1, activity: 100)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 100), tipDm(2, activity: 200)] + mock.messagesDelay = .seconds(30) + let controller = makeController(mock, database: db) + controller.start() + + try await waitUntil { ids(controller) == [.test(2), .test(1)] } + try await waitUntil { !mock.latestPageQueries.isEmpty } + controller.stop() + } + + @Test("the cached list is ready before any feed answers") + func cacheIsInstant() async throws { + let db = try cached([tipDm(1, activity: 100)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 100), tipDm(2, activity: 200)] + mock.feedDelay = .seconds(30) + let controller = makeController(mock, database: db) + controller.start() + + try await waitUntil { ids(controller) == [.test(1)] && mock.dmFeedCalls > 0 } + controller.stop() + } + + // MARK: - Idempotence + + @Test("backfill and a feed with older data leave rows, previews and order alone") + func olderDataDoesNotMutate() async throws { + let db = try cached([ + tipDm(1, activity: 200, last: message(5, at: 200)), + tipDm(2, activity: 100, last: message(4, at: 100)), + ]) + let mock = MockConversations() + // A feed behind what is stored, and a transcript page behind it too. + mock.feed = [ + tipDm(1, activity: 150, last: message(3, at: 150)), + tipDm(2, activity: 100, last: message(4, at: 100)), + ] + mock.messages = [message(2, at: 50)] + let controller = makeController(mock, database: db) + controller.start() + try await waitUntil { ids(controller) == [.test(1), .test(2)] } + + await controller.loadFeed() + + let rows = controller.chatListConversations + #expect(rows.map(\.id) == [.test(1), .test(2)]) + #expect(rows.map(\.lastActivity) == [Date(timeIntervalSince1970: 200), Date(timeIntervalSince1970: 100)]) + #expect(rows.compactMap(\.lastMessage?.id.value) == [5, 4]) + controller.stop() + } + + @Test("a live event older than the stored activity does not reorder the list") + func olderStreamEventDoesNotReorder() async throws { + let db = try cached([tipDm(1, activity: 200), tipDm(2, activity: 100)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 200), tipDm(2, activity: 100)] + let controller = makeController(mock, database: db) + controller.start() + try await waitUntil { mock.streamOpened && ids(controller) == [.test(1), .test(2)] } + await controller.loadFeed() + + mock.emit(.lastActivityChanged(conversationID: .test(2), date: Date(timeIntervalSince1970: 10))) + try await Task.sleep(for: .milliseconds(100)) + #expect(ids(controller) == [.test(1), .test(2)]) + + mock.emit(.lastActivityChanged(conversationID: .test(2), date: Date(timeIntervalSince1970: 900))) + try await waitUntil { ids(controller) == [.test(2), .test(1)] } + controller.stop() + } + + // MARK: - Empty store + + @Test("with nothing cached each feed shows as it lands") + func emptyStoreIsProgressive() async throws { + let mock = MockConversations() + mock.feed = [tipDm(2, activity: 200)] + mock.feedDelay = .milliseconds(500) + mock.groupFeed = [group(9, activity: 300)] + let controller = makeController(mock) + controller.start() + + // The group feed answers at once and shows without waiting on the DM feeds. + try await waitUntil { ids(controller) == [.test(9)] } + #expect(controller.chatListConversations.count == 1) + try await waitUntil { ids(controller) == [.test(9), .test(2)] } + controller.stop() + } + + // MARK: - Dedupe + + @Test("start and a foreground refresh make one fetch per feed") + func dedupesFeedLoads() async throws { + let db = try cached([tipDm(1, activity: 100)]) + let mock = MockConversations() + mock.feed = [tipDm(1, activity: 100)] + mock.feedDelay = .milliseconds(200) + let controller = makeController(mock, database: db) + controller.start() + try await waitUntil { mock.dmFeedCalls > 0 } + await controller.loadFeed() + // contactDm + tipDm + groups, once each. + #expect(mock.dmFeedCalls == 2) + #expect(mock.groupFeedCalls == 1) + controller.stop() + } +} + +/// Records what the Chats list shows each time it changes, read after the change has fully applied. +@MainActor +private final class ListObserver { + struct Shown { let ids: [ConversationID]; let unread: Int? } + + private(set) var shown: [Shown] = [] + private let controller: ConversationController + private let conversation: Conversation + + init(_ controller: ConversationController, unreadFor conversation: Conversation) { + self.controller = controller + self.conversation = conversation + } + + func start() { + withObservationTracking { + _ = controller.chatListConversations + _ = controller.unreadCount(for: conversation) + } onChange: { [weak self] in + Task { @MainActor in + guard let self else { return } + self.shown.append(Shown(ids: self.controller.chatListConversations.map(\.id), unread: self.controller.unreadCount(for: self.conversation))) + self.start() + } + } + } +} + +@MainActor +@Suite("BoundedWait") +struct BoundedWaitTests { + + @Test("returns when every unit has signalled, ahead of the cap") + func returnsWhenDone() async { + let wait = BoundedWait(expecting: 2) + Task { wait.signal(); wait.signal() } + let start = ContinuousClock.now + await wait.wait(upTo: .seconds(10)) + #expect(ContinuousClock.now - start < .seconds(5)) + #expect(wait.isClosed) + } + + @Test("returns at the cap when a unit never signals, then reports later ones as late") + func returnsAtCap() async { + let wait = BoundedWait(expecting: 2) + wait.signal() + await wait.wait(upTo: .milliseconds(50)) + #expect(wait.isClosed) + } +} + +@MainActor +@Suite("BackfillQueue") +struct BackfillQueueTests { + + @Test("enqueue skips duplicates; remove takes a chat out ahead of its turn") + func removeAndDedupe() { + var queue = BackfillQueue() + queue.enqueue([.init(conversationID: .test(1), kind: .delta), .init(conversationID: .test(2), kind: .newestPage)]) + queue.enqueue([.init(conversationID: .test(1), kind: .delta)]) + #expect(queue.count == 2) + let first = queue.remove(.test(1)) + let second = queue.remove(.test(1)) + #expect(first) + #expect(!second) + #expect(queue.popNext()?.conversationID == .test(2)) + #expect(queue.popNext() == nil) + } + + @Test("workers are claimed up to the limit and topped up as they finish") + func claimsWorkersUpToLimit() { + var queue = BackfillQueue() + queue.enqueue((1...6).map { .init(conversationID: .test($0), kind: .delta) }) + #expect(queue.claimWorkers(limit: 4) == 4) + #expect(queue.claimWorkers(limit: 4) == 0) + queue.releaseWorker() + #expect(queue.claimWorkers(limit: 4) == 1) + } +} diff --git a/FlipcashTests/Database/Database+ConversationsTests.swift b/FlipcashTests/Database/Database+ConversationsTests.swift index f8dc16949..ba54ba74d 100644 --- a/FlipcashTests/Database/Database+ConversationsTests.swift +++ b/FlipcashTests/Database/Database+ConversationsTests.swift @@ -317,6 +317,24 @@ struct DatabaseConversationsTests { #expect(previews.contains(UInt64(count))) } + @Test("a message stored after the row's last activity moves the cached activity forward") + func cachedActivityFollowsStoredMessages() throws { + let (database, url) = try Database.makeTemp() + defer { Database.removeTemp(at: url) } + let id = ConversationID.test(1) + try database.upsertConversation(conversation(byte: 1)) + let before = try #require(try database.getConversations().first).lastActivity + // What the notification extension writes: a message row and nothing on the conversation. + let pushed = ConversationMessage( + id: MessageID(value: 7), senderID: nil, content: .text("pushed"), + date: before.addingTimeInterval(3600), unreadSeq: 7, eventSequence: 7 + ) + try database.persistMessages([pushed], cursor: 0, conversationID: id) + let loaded = try #require(try database.getConversations().first) + #expect(loaded.lastActivity == pushed.date) + #expect(loaded.lastMessage?.id.value == 7) + } + // MARK: - Identity + atomicity @Test("clientMessageID round-trips through the cache") diff --git a/FlipcashTests/TestSupport/MockConversations.swift b/FlipcashTests/TestSupport/MockConversations.swift index 2c052adc2..06652f68e 100644 --- a/FlipcashTests/TestSupport/MockConversations.swift +++ b/FlipcashTests/TestSupport/MockConversations.swift @@ -275,9 +275,37 @@ final class MockConversations: ConversationFetching, ConversationMembership, Con // MARK: - ConversationFetching - func getDmChatFeed(owner: KeyPair, type: ConversationType) async throws -> [Conversation] { feed } + /// How many `GetDmChatFeed` calls have been made. + var dmFeedCalls: Int { lock.withLock { _dmFeedCalls } } + private var _dmFeedCalls = 0 + /// Held on every DM feed call before it answers, so a test can overlap loads or outlast a timeout. + var feedDelay: Duration { + get { lock.withLock { _feedDelay } } + set { lock.withLock { _feedDelay = newValue } } + } + private var _feedDelay: Duration = .zero + + func getDmChatFeed(owner: KeyPair, type: ConversationType) async throws -> [Conversation] { + lock.withLock { _dmFeedCalls += 1 } + if feedDelay > .zero { try? await Task.sleep(for: feedDelay) } + return feed + } + + /// How many `GetGroupChatFeed` calls have been made. + var groupFeedCalls: Int { lock.withLock { _groupFeedCalls } } + private var _groupFeedCalls = 0 + /// Held on every group feed call before it answers. + var groupFeedDelay: Duration { + get { lock.withLock { _groupFeedDelay } } + set { lock.withLock { _groupFeedDelay = newValue } } + } + private var _groupFeedDelay: Duration = .zero - func getGroupChatFeed(owner: KeyPair) async throws -> [Conversation] { groupFeed } + func getGroupChatFeed(owner: KeyPair) async throws -> [Conversation] { + lock.withLock { _groupFeedCalls += 1 } + if groupFeedDelay > .zero { try? await Task.sleep(for: groupFeedDelay) } + return groupFeed + } func getChat(owner: KeyPair, conversationID: ConversationID) async throws -> Conversation { guard let conversation = feed.first(where: { $0.id == conversationID }) else { @@ -320,8 +348,23 @@ final class MockConversations: ConversationFetching, ConversationMembership, Con // MARK: - ConversationMessaging + /// Held on every `GetMessage` call before it answers. + var singleMessageDelay: Duration { + get { lock.withLock { _singleMessageDelay } } + set { lock.withLock { _singleMessageDelay = newValue } } + } + private var _singleMessageDelay: Duration = .zero + + /// Held on every newest-page `GetMessages` call before it answers. + var messagesDelay: Duration { + get { lock.withLock { _messagesDelay } } + set { lock.withLock { _messagesDelay = newValue } } + } + private var _messagesDelay: Duration = .zero + func getMessage(owner: KeyPair, conversationID: ConversationID, messageID: MessageID) async throws -> ConversationMessage? { lock.withLock { _singleMessageQueries.append(messageID) } + if singleMessageDelay > .zero { try? await Task.sleep(for: singleMessageDelay) } if let error = singleMessageError { throw error } return singleMessages[messageID] } @@ -329,6 +372,7 @@ final class MockConversations: ConversationFetching, ConversationMembership, Con func getMessages(owner: KeyPair, conversationID: ConversationID, before: MessageID?) async throws -> [ConversationMessage] { guard let before else { lock.withLock { _latestPageQueries.append(conversationID) } + if messagesDelay > .zero { try? await Task.sleep(for: messagesDelay) } return messages } lock.withLock { _olderQueries.append(before) }