diff --git a/Docs/ARCHITECTURE.md b/Docs/ARCHITECTURE.md index b127d78..91b0eb7 100644 --- a/Docs/ARCHITECTURE.md +++ b/Docs/ARCHITECTURE.md @@ -39,7 +39,9 @@ fallback на plaintext. UI сразу показывает optimistic message Входящая stanza может прийти напрямую, через carbons или MAM. Сервис определяет peer, расшифровывает payload, использует `origin-id`/stanza id для дедупликации и -передаёт value-type envelope в `AppModel`. +передаёт value-type envelope в `AppModel`. Дедупликация выполняется по +хеш-индексам (`origin-id`, `stanza-id`, `id`), поэтому стоимость применения пачки +не растёт квадратично с размером архива. Свежая установка начинает MAM с ограниченной последней страницы RSM, а не с самого старого сообщения. Инкрементальные проходы используют сохранённый @@ -48,14 +50,18 @@ peer, расшифровывает payload, использует `origin-id`/sta Каждая stanza проходит проверку `queryid` и источника личного архива. Martin публикует MAM-результаты на parser queue: ограниченный inbox принимает там всю страницу и делает единственный hand-off на main actor после финального IQ. -Дешифрование выполняется по одной stanza с уступкой event loop, а декодированные -изменения становятся видимыми только одной атомарной пачкой вместе с новой -контрольной точкой. Один foreground-проход имеет конечный бюджет страниц и -времени; ошибка страницы не перезапускает весь архив сама. При уходе приложения -с экрана и на время записи/отправки видеосообщения активный MAM-запрос -закрывается, чтобы дешифрование истории не конкурировало с камерой и upload. -Только после commit `AppModel` обновляет SwiftUI, локальный snapshot и Apple -Watch; незавершённый проход при timeout/disconnect отбрасывается целиком. +Расшифровка OMEMO выполняется на отдельной серийной фоновой очереди, поэтому +Signal-криптография не блокирует главный поток; построение envelope и публикация +остаются на main actor. Видимый индикатор «Синхронизация истории…» показывается +только на первичной странице, а инкрементальная догрузка backlog продолжается в +фоне молча, публикуя декодированные изменения атомарными пачками по мере +готовности страниц. MUC-архив проходит тем же путём и флашит декодированные +групповые сообщения по завершении catch-up комнаты. Ошибка отдельной страницы +автоматически повторяется ограниченное число раз. При уходе приложения с экрана +и на время записи/отправки видеосообщения активный MAM-запрос закрывается, чтобы +дешифрование истории не конкурировало с камерой и upload. Только после commit +`AppModel` обновляет SwiftUI, локальный snapshot и Apple Watch; незавершённый +проход при timeout/disconnect отбрасывается целиком. Редактирование реализовано стандартным XEP-0308: исправление получает новый id и `` с id исходного сообщения. `AppModel` проверяет совпадение bare JID diff --git a/README.md b/README.md index 2c8b780..9f55990 100644 --- a/README.md +++ b/README.md @@ -145,14 +145,18 @@ OMEMO, и для XEP-0084 аватаров. Для медиа и вложени XMPP-сокет. Для production нужны Apple Developer credentials, APNs provider, XEP-0357 push-компонент и регистрация приложения в этом компоненте. - При первом запуске Luma запрашивает ограниченную последнюю страницу - серверного MAM-архива, чтобы огромная история не блокировала интерфейс. + серверного MAM-архива, чтобы огромная история не блокировала интерфейс, и + показывает индикатор «Синхронизация истории…» только на этом первичном этапе. После успешного прохода сохраняются серверный UID последней записи и время; следующие подключения продолжают RSM строго после UID. Временное перекрытие используется только для старых snapshots без UID или один раз, если сервер - удалил сохранённый UID. Сырые stanza собираются вне main queue, расшифровываются - по одной с уступкой UI, а готовый проход публикуется одной пачкой. Один - foreground-проход также ограничен по страницам и времени; после ошибки он не - перезапускает сам себя бесконечно. + удалил сохранённый UID. Сырые stanza собираются вне main queue, а расшифровка + OMEMO выполняется на отдельной фоновой очереди — главный поток освобождается на + время Signal-криптографии, поэтому большая история не лагает интерфейс. + Инкрементальная догрузка backlog продолжается в фоне без вечного спиннера, а + сообщения применяются атомарными пачками с дедупликацией по `origin-id`/stanza + id через хеш-индексы. После сбоя отдельной страницы синхронизация автоматически + повторяется ограниченное число раз, не дожидаясь повторного открытия приложения. - Клиент не вводит собственного ограничения размера вложения: фактический предел сообщает XEP-0363 upload-компонент сервера. Видеосообщение ограничено 60 секундами. diff --git a/Sources/Shared/Models/AppModel.swift b/Sources/Shared/Models/AppModel.swift index a083d39..1597c7e 100644 --- a/Sources/Shared/Models/AppModel.swift +++ b/Sources/Shared/Models/AppModel.swift @@ -78,6 +78,7 @@ final class AppModel: ObservableObject { private var isApplyingArchiveBatch = false private var messageIndexByStorageKey: [String: Int] = [:] private var messageIndexByStanzaKey: [String: Int] = [:] + private var messageIndexByOriginKey: [String: Int] = [:] private var firstMessageIndexByID: [String: Int] = [:] private var selectedMessagesCacheConversationID: String? private var selectedMessagesCache: [ChatMessage] = [] @@ -1949,9 +1950,7 @@ final class AppModel: ObservableObject { return // уже есть это сообщение } if let originID = envelope.originID, - messages.contains(where: { - $0.conversationID == envelope.peerJID.lowercased() && $0.originID == originID - }) { + messageIndex(originID: originID, conversationID: envelope.peerJID) != nil { mergeServerIdentity(from: envelope, matchingOriginID: originID) return } @@ -2049,9 +2048,10 @@ final class AppModel: ObservableObject { from envelope: XMPPService.MessageEnvelope, matchingOriginID originID: String ) { - guard let index = messages.firstIndex(where: { - $0.conversationID == envelope.peerJID.lowercased() && $0.originID == originID - }) else { return } + guard let index = messageIndex( + originID: originID, + conversationID: envelope.peerJID + ) else { return } if messages[index].stanzaID == nil { messages[index].stanzaID = envelope.stanzaID } @@ -2520,6 +2520,13 @@ final class AppModel: ObservableObject { conversationID: merged.conversationID )] = index } + if let originID = merged.originID { + messageIndexByOriginKey[ + Self.localDeletionKey( + messageID: originID, + conversationID: merged.conversationID + )] = index + } updateConversationPreview(for: merged, incrementUnread: false) return false } @@ -2540,6 +2547,13 @@ final class AppModel: ObservableObject { conversationID: message.conversationID )] = insertedIndex } + if let originID = message.originID { + messageIndexByOriginKey[ + Self.localDeletionKey( + messageID: originID, + conversationID: message.conversationID + )] = insertedIndex + } let shouldIncrement = unreadOverride ?? (message.direction == .incoming && selectedConversationID != message.conversationID) @@ -2626,20 +2640,41 @@ final class AppModel: ObservableObject { messageID: referenceID, conversationID: normalizedConversationID ) - if let index = messageIndexByStanzaKey[key], - messages.indices.contains(index), + guard let index = messageIndexByStanzaKey[key] else { return nil } + if messages.indices.contains(index), messages[index].conversationID == normalizedConversationID, messages[index].stanzaID == referenceID { return index } + // Only rebuild when the key exists but points at a stale entry. A + // genuine miss is expected during archive catch-up for every new + // stanza-id and must not trigger an O(n) rebuild per message. rebuildMessageIndex() return messageIndexByStanzaKey[key] } + private func messageIndex(originID: String, conversationID: String) -> Int? { + let normalizedConversationID = conversationID.lowercased() + let key = Self.localDeletionKey( + messageID: originID, + conversationID: normalizedConversationID + ) + guard let index = messageIndexByOriginKey[key] else { return nil } + if messages.indices.contains(index), + messages[index].conversationID == normalizedConversationID, + messages[index].originID == originID + { + return index + } + rebuildMessageIndex() + return messageIndexByOriginKey[key] + } + private func rebuildMessageIndex() { messageIndexByStorageKey.removeAll(keepingCapacity: true) messageIndexByStanzaKey.removeAll(keepingCapacity: true) + messageIndexByOriginKey.removeAll(keepingCapacity: true) firstMessageIndexByID.removeAll(keepingCapacity: true) for (index, message) in messages.enumerated() { let key = Self.localDeletionKey( @@ -2661,6 +2696,15 @@ final class AppModel: ObservableObject { messageIndexByStanzaKey[stanzaKey] = index } } + if let originID = message.originID { + let originKey = Self.localDeletionKey( + messageID: originID, + conversationID: message.conversationID + ) + if messageIndexByOriginKey[originKey] == nil { + messageIndexByOriginKey[originKey] = index + } + } } } diff --git a/Sources/Shared/Models/ArchiveMessageBatchPolicy.swift b/Sources/Shared/Models/ArchiveMessageBatchPolicy.swift index fbef6d5..5aa4a38 100644 --- a/Sources/Shared/Models/ArchiveMessageBatchPolicy.swift +++ b/Sources/Shared/Models/ArchiveMessageBatchPolicy.swift @@ -6,10 +6,14 @@ import Foundation enum ArchiveMessageBatchPolicy { static let pageSize = 24 static let bootstrapMessageLimit = 40 - static let decodeSliceSize = 1 + /// Decrypt this many archived stanzas per main-actor slice before yielding. + /// One-by-one decryption turns a large MAM catch-up into hours of serialized + /// work; a small batch keeps frames mostly intact while cutting the number + /// of run-loop hand-offs by the same factor. + static let decodeSliceSize = 8 static let maximumBufferedStanzas = 48 - /// Leave roughly half of a 60 Hz frame between archived decryptions. This - /// keeps UIScrollView touch handling responsive even when OMEMO work is - /// unusually expensive on older devices. - static let interSliceDelayNanoseconds: UInt64 = 8_000_000 + /// A short yield between archived decryption slices keeps touch handling + /// responsive without the 8 ms per single message that previously dominated + /// catch-up time. + static let interSliceDelayNanoseconds: UInt64 = 1_000_000 } diff --git a/Sources/Shared/Models/ArchiveSyncRecoveryPolicy.swift b/Sources/Shared/Models/ArchiveSyncRecoveryPolicy.swift index 0098e4a..5452327 100644 --- a/Sources/Shared/Models/ArchiveSyncRecoveryPolicy.swift +++ b/Sources/Shared/Models/ArchiveSyncRecoveryPolicy.swift @@ -5,8 +5,15 @@ import Foundation enum ArchiveSyncRecoveryPolicy { static let incrementalOverlap: TimeInterval = 60 static let interPageDelayNanoseconds: UInt64 = 100_000_000 - static let queryTimeoutNanoseconds: UInt64 = 12_000_000_000 - static let pageApplyTimeoutNanoseconds: UInt64 = 8_000_000_000 + static let queryTimeoutNanoseconds: UInt64 = 20_000_000_000 + static let pageApplyTimeoutNanoseconds: UInt64 = 20_000_000_000 static let pageRetryLimit = 1 static let resumeAfterCaptureDelayNanoseconds: UInt64 = 3_000_000_000 + /// Delay before automatically resuming a failed catch-up pass while the + /// app stays connected and foregrounded. Prevents a single slow page from + /// silently leaving history unloaded until the next app activation. + static let retryAfterFailureDelayNanoseconds: UInt64 = 20_000_000_000 + /// Consecutive automatic catch-up retries before falling back to + /// resume-on-activation, so a persistently broken server cannot loop. + static let maximumAutomaticRetries = 3 } diff --git a/Sources/Shared/XMPP/XMPPService.swift b/Sources/Shared/XMPP/XMPPService.swift index 96f3b68..3b919cb 100644 --- a/Sources/Shared/XMPP/XMPPService.swift +++ b/Sources/Shared/XMPP/XMPPService.swift @@ -274,6 +274,8 @@ final class XMPPService { private var archiveRetryTask: Task? private var archiveQueryTimeoutTask: Task? private var archiveQueryCompletionTask: Task? + private var archiveSyncRetryTask: Task? + private var archiveAutoRetryCount = 0 private var archiveActiveQueryID: String? private var archiveResumeAfter: String? private var archiveWorkBudget = ArchiveSyncWorkBudget() @@ -306,6 +308,15 @@ final class XMPPService { private var olderHistoryQueryID: String? private var olderHistoryTimeoutTask: Task? private let olderHistoryQueryTimeoutNanoseconds: UInt64 = 15_000_000_000 + /// Serial background queue for OMEMO decryption. `decode` performs the + /// expensive Signal work (session decrypt + AES-GCM) and previously ran on + /// the main actor for every MAM stanza, stalling the UI during large + /// archive catch-ups. All decodes stay serialized on this single queue so + /// Signal's per-session state is never mutated concurrently. + private let omemoDecodeQueue = DispatchQueue( + label: "app.luma.omemo.decode", + qos: .userInitiated + ) private let archiveStanzaInbox = ArchiveStanzaInbox( maximumCount: ArchiveMessageBatchPolicy.maximumBufferedStanzas ) @@ -373,6 +384,9 @@ final class XMPPService { archiveQueryTimeoutTask = nil archiveQueryCompletionTask?.cancel() archiveQueryCompletionTask = nil + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = nil + archiveAutoRetryCount = 0 archiveActiveQueryID = nil archiveResumeAfter = nil archiveWorkBudget = ArchiveSyncWorkBudget() @@ -438,6 +452,9 @@ final class XMPPService { archiveQueryTimeoutTask = nil archiveQueryCompletionTask?.cancel() archiveQueryCompletionTask = nil + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = nil + archiveAutoRetryCount = 0 archiveActiveQueryID = nil archiveResumeAfter = nil archiveWorkBudget = ArchiveSyncWorkBudget() @@ -495,6 +512,9 @@ final class XMPPService { _ = client.module(.csi).setState(active) if active { archiveRetrySuppressedUntilActivation = false + archiveAutoRetryCount = 0 + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = nil refreshArchiveIfNeeded(client: client) } else { // iOS may freeze network callbacks and timeout tasks while the @@ -927,43 +947,49 @@ final class XMPPService { return } - // History is only allowed while catch-up is idle, so the - // legacy mutation accumulator can safely be used as a local - // page collector here. - let savedMutations = self.archivePassMutations - self.archivePassMutations = [] - for stanza in page.stanzas { - if stanza.message.type == .groupchat, - let roomJID = stanza.message.from?.bareJid { - self.handle( - groupMessage: stanza.message, - archivedRoomJID: roomJID, - archivedTimestamp: stanza.timestamp, - archiveID: stanza.archiveID, - isArchived: true - ) - } else { - self.handle( - message: stanza.message, - timestamp: stanza.timestamp, - archiveID: stanza.archiveID, - isArchived: true, - archivedPeerJID: isGroup ? nil : BareJID(conversationJID.lowercased()) - ) + Task { @MainActor [weak self, weak client] in + guard let self, let client else { + completion(.failure(LumaXMPPError.notConnected)) + return + } + // History is only allowed while catch-up is idle, so the + // legacy mutation accumulator can safely be used as a local + // page collector here. + let savedMutations = self.archivePassMutations + self.archivePassMutations = [] + for stanza in page.stanzas { + if stanza.message.type == .groupchat, + let roomJID = stanza.message.from?.bareJid { + await self.handle( + groupMessage: stanza.message, + archivedRoomJID: roomJID, + archivedTimestamp: stanza.timestamp, + archiveID: stanza.archiveID, + isArchived: true + ) + } else { + await self.handle( + message: stanza.message, + timestamp: stanza.timestamp, + archiveID: stanza.archiveID, + isArchived: true, + archivedPeerJID: isGroup ? nil : BareJID(conversationJID.lowercased()) + ) + } + } + let historyMutations = self.archivePassMutations + self.archivePassMutations = savedMutations + if !historyMutations.isEmpty { + self.eventHandler?(.archiveBatch(historyMutations)) + } + switch result { + case .success(let response): + completion(.success(!response.complete)) + case .failure(let error): + // Do not leave the UI in the loading state. A failed + // interactive query is recoverable by the next scroll. + completion(.failure(error)) } - } - let historyMutations = self.archivePassMutations - self.archivePassMutations = savedMutations - if !historyMutations.isEmpty { - self.eventHandler?(.archiveBatch(historyMutations)) - } - switch result { - case .success(let response): - completion(.success(!response.complete)) - case .failure(let error): - // Do not leave the UI in the loading state. A failed - // interactive query is recoverable by the next scroll. - completion(.failure(error)) } } } @@ -1535,16 +1561,18 @@ final class XMPPService { client.module(.message).messagesPublisher .receive(on: DispatchQueue.main) .sink { [weak self] incoming in -// self?.handle(message: incoming.message, timestamp: Date(), archiveID: nil) - self?.deliverOrDelayDirect(incoming.message, timestamp: Date()) + Task { @MainActor [weak self] in + await self?.deliverOrDelayDirect(incoming.message, timestamp: Date()) + } } .store(in: &cancellables) client.module(.muc).messagesPublisher .receive(on: DispatchQueue.main) .sink { [weak self] incoming in -// self?.handle(groupMessage: incoming.message, room: incoming.room) - self?.deliverOrDelayGroup(incoming.message, room: incoming.room) + Task { @MainActor [weak self] in + await self?.deliverOrDelayGroup(incoming.message, room: incoming.room) + } } .store(in: &cancellables) @@ -1565,8 +1593,9 @@ final class XMPPService { client.module(.messageCarbons).carbonsPublisher .receive(on: DispatchQueue.main) .sink { [weak self] carbon in -// self?.handle(message: carbon.message, timestamp: Date(), archiveID: nil) - self?.deliverOrDelayDirect(carbon.message, timestamp: Date()) + Task { @MainActor [weak self] in + await self?.deliverOrDelayDirect(carbon.message, timestamp: Date()) + } } .store(in: &cancellables) @@ -1622,7 +1651,7 @@ final class XMPPService { .store(in: &cancellables) } - private func deliverOrDelayDirect(_ message: Message, timestamp: Date) { + private func deliverOrDelayDirect(_ message: Message, timestamp: Date) async { guard let client else { return } let archive = MAMArchiveKey.account(client.userBareJid.stringValue) if archiveSyncStarted { @@ -1631,10 +1660,10 @@ final class XMPPService { ) return } - handle(message: message, timestamp: timestamp, archiveID: nil) + await handle(message: message, timestamp: timestamp, archiveID: nil) } - private func deliverOrDelayGroup(_ message: Message, room: RoomProtocol) { + private func deliverOrDelayGroup(_ message: Message, room: RoomProtocol) async { let archive = MAMArchiveKey.muc(room.jid.stringValue) if activeMUCCatchup == archive { delayedLiveByArchive[archive, default: []].append( @@ -1642,17 +1671,35 @@ final class XMPPService { ) return } - handle(groupMessage: message, room: room) + await handle(groupMessage: message, room: room) } - private func replayDelayedLive(for archive: MAMArchiveKey) { + private func replayDelayedLive(for archive: MAMArchiveKey) async { let delayed = delayedLiveByArchive.removeValue(forKey: archive) ?? [] for item in delayed { switch item { case .direct(let message, let timestamp): - handle(message: message, timestamp: timestamp, archiveID: nil) + await handle(message: message, timestamp: timestamp, archiveID: nil) case .group(let message, let room): - handle(groupMessage: message, room: room) + await handle(groupMessage: message, room: room) + } + } + } + + /// Runs OMEMO decryption off the main actor and resumes with the raw + /// result. Kept as a thin wrapper so `handle` keeps its existing switch + /// logic and only the expensive `decode` call leaves the main thread. + private func decodeOmemoOffMain( + _ message: Message, + from sender: BareJID, + serverMsgId: String?, + module: OMEMOModule + ) async -> OMEMOModule.DecryptionResult { + let queue = omemoDecodeQueue + return await withCheckedContinuation { continuation in + queue.async { + let result = module.decode(message: message, from: sender, serverMsgId: serverMsgId) + continuation.resume(returning: result) } } } @@ -1676,7 +1723,7 @@ final class XMPPService { archiveID: String?, isArchived: Bool = false, archivedPeerJID: BareJID? = nil - ) { + ) async { guard message.type != .error, message.type != .groupchat, let client, @@ -1701,8 +1748,12 @@ final class XMPPService { let security: ChatMessage.Security let fingerprint: String? let contentMessage: Message? - switch client.module(.omemo).decode(message: message, from: sender, serverMsgId: archiveID) - { + switch await decodeOmemoOffMain( + message, + from: sender, + serverMsgId: archiveID, + module: client.module(.omemo) + ) { case .successMessage(let decodedMessage, let value): security = .omemo fingerprint = value @@ -1887,7 +1938,7 @@ final class XMPPService { archivedTimestamp: Date? = nil, archiveID: String? = nil, isArchived: Bool = false - ) { + ) async { guard message.type != .error, let client, let from = message.from, @@ -1934,10 +1985,11 @@ final class XMPPService { let fingerprint: String? let contentMessage: Message? if let realSender { - switch client.module(.omemo).decode( - message: message, + switch await decodeOmemoOffMain( + message, from: realSender, - serverMsgId: stanzaID + serverMsgId: stanzaID, + module: client.module(.omemo) ) { case .successMessage(let decodedMessage, let value): security = .omemo @@ -2672,12 +2724,19 @@ final class XMPPService { archiveQueryTimeoutTask = nil archiveQueryCompletionTask?.cancel() archiveQueryCompletionTask = nil + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = nil archiveActiveQueryID = nil archiveStanzaBuffer.removeAll(keepingCapacity: true) archiveBufferOverflowed = false archiveRejectedSource = false archiveStanzaInbox.cancel() - setArchiveSyncIndicator(true) + // Show the visible "Синхронизация истории…" banner only for the + // one-page bootstrap window. Incremental catch-up continues silently so + // a large backlog never leaves the spinner visible for the whole pass. + if archiveIsBootstrapQuery { + setArchiveSyncIndicator(true) + } client.module(.omemo).mamSyncStarted(for: nil) let initialPosition = ArchiveSyncCursorPolicy.requestPosition( checkpoint: archiveSyncCheckpoint, @@ -2756,56 +2815,54 @@ final class XMPPService { self.finishMUCCatchup(archive: archive, succeeded: false) return } - for stanza in page.stanzas { - self.mucCatchupHighWatermark = max( - self.mucCatchupHighWatermark ?? .distantPast, - stanza.timestamp - ) -// guard stanza.message.type == .groupchat, -// let roomJID = stanza.message.from?.bareJid, -// let room = client.module(.muc).roomManager.room( -// for: client, with: roomJID -// ) else { continue } -// self.handle(groupMessage: stanza.message, room: room) - guard stanza.message.type == .groupchat, - let roomJID = stanza.message.from?.bareJid else { continue } - self.handle( - groupMessage: stanza.message, - archivedRoomJID: roomJID, - archivedTimestamp: stanza.timestamp, - archiveID: stanza.archiveID, - isArchived: true - ) - } - - switch result { -// case .failure: -// self.finishMUCCatchup(archive: archive, succeeded: false) - case .failure: - // Same recovery principle as account MAM: a retained local - // UID may have expired on the server. Retry once from the - // durable timestamp overlap, never loop cursor<->time. - if self.mamCheckpoints[archive]?.cursor != nil, - self.mucCursorFallbackArchives.insert(archive).inserted, - let old = self.mamCheckpoints[archive] { - self.mamCheckpoints[archive] = MAMArchiveCheckpoint( - timestamp: old.timestamp, - cursor: nil + Task { @MainActor [weak self, weak client] in + guard let self, let client, + self.activeMUCCatchup == archive else { return } + for stanza in page.stanzas { + self.mucCatchupHighWatermark = max( + self.mucCatchupHighWatermark ?? .distantPast, + stanza.timestamp ) - self.queryMUCArchivePage( - archive: archive, - after: nil, - client: client + guard stanza.message.type == .groupchat, + let roomJID = stanza.message.from?.bareJid else { continue } + await self.handle( + groupMessage: stanza.message, + archivedRoomJID: roomJID, + archivedTimestamp: stanza.timestamp, + archiveID: stanza.archiveID, + isArchived: true ) - } else { - self.finishMUCCatchup(archive: archive, succeeded: false) } - case .success(let response): - self.mucCatchupLastCursor = response.rsm?.last ?? self.mucCatchupLastCursor - if !response.complete, let next = response.rsm?.last, next != after { - self.queryMUCArchivePage(archive: archive, after: next, client: client) - } else { - self.finishMUCCatchup(archive: archive, succeeded: true) + + switch result { + case .failure: + // Same recovery principle as account MAM: a retained + // local UID may have expired on the server. Retry once + // from the durable timestamp overlap, never loop + // cursor<->time. + if self.mamCheckpoints[archive]?.cursor != nil, + self.mucCursorFallbackArchives.insert(archive).inserted, + let old = self.mamCheckpoints[archive] { + self.mamCheckpoints[archive] = MAMArchiveCheckpoint( + timestamp: old.timestamp, + cursor: nil + ) + self.queryMUCArchivePage( + archive: archive, + after: nil, + client: client + ) + } else { + self.finishMUCCatchup(archive: archive, succeeded: false) + } + case .success(let response): + self.mucCatchupLastCursor = + response.rsm?.last ?? self.mucCatchupLastCursor + if !response.complete, let next = response.rsm?.last, next != after { + self.queryMUCArchivePage(archive: archive, after: next, client: client) + } else { + self.finishMUCCatchup(archive: archive, succeeded: true) + } } } } @@ -2817,10 +2874,22 @@ final class XMPPService { activeMUCCatchup = nil mucCatchupHighWatermark = nil mucCatchupLastCursor = nil - replayDelayedLive(for: archive) + Task { @MainActor [weak self] in + await self?.replayDelayedLive(for: archive) + } startNextMUCCatchupIfPossible() startPendingOlderHistoryIfPossible() } + // MUC catch-up decodes into archivePassMutations just like the account + // pass. Publish them once the room catch-up settles, otherwise group + // history fetched via MAM is decoded and then silently dropped. + if succeeded, !archivePassMutations.isEmpty { + let mutations = archivePassMutations + archivePassMutations.removeAll(keepingCapacity: true) + eventHandler?(.archiveBatch(mutations)) + } else { + archivePassMutations.removeAll(keepingCapacity: true) + } guard succeeded else { return } let checkpoint = MAMArchiveCheckpoint( timestamp: max(mucCatchupHighWatermark ?? .distantPast, Date()), @@ -2965,7 +3034,7 @@ final class XMPPService { // Route them through the same group handler instead of // dropping them in the direct-message handler. // handle(groupMessage: stanza.message, room: room) - handle( + await handle( groupMessage: stanza.message, archivedRoomJID: roomJID, archivedTimestamp: stanza.timestamp, @@ -2973,7 +3042,7 @@ final class XMPPService { isArchived: true ) } else { - handle( + await handle( message: stanza.message, timestamp: stanza.timestamp, archiveID: stanza.archiveID, @@ -3063,24 +3132,6 @@ final class XMPPService { caughtUp: false ) archiveLastCompletedCursor = checkpoint.cursor - // if workBudgetReached { - // if checkpoint.advances(over: archiveSyncCheckpoint) { - // finishArchiveSync( - // client: client, - // succeeded: true, - // checkpoint: checkpoint - // ) - // } else { - // archivePassMutations.removeAll(keepingCapacity: true) - // eventHandler?( - // .recoverableError( - // "MAM не продвинул контрольную точку; синхронизация приостановлена." - // )) - // finishArchiveSync(client: client, succeeded: false) - // } - // } else { - // scheduleArchiveNextPage(client: client, after: cursor) - // } if workBudgetReached { guard checkpoint.advances(over: archiveSyncCheckpoint) else { archivePassMutations.removeAll(keepingCapacity: true) @@ -3092,10 +3143,10 @@ final class XMPPService { return } - // The foreground budget is only a UI-yield boundary. It - // must NOT turn response.complete == false into a completed - // MAM pass, otherwise messages after this cursor are missed - // until a later activation. + // Flush the accumulated mutations so the UI advances + // incrementally, then keep catching up in the background. + // The visible sync banner is only shown for the bootstrap + // window, so a large backlog never leaves a spinner hanging. if !archivePassMutations.isEmpty { let mutations = archivePassMutations archivePassMutations.removeAll(keepingCapacity: true) @@ -3313,8 +3364,12 @@ final class XMPPService { archiveSyncCompletedForConnection = true archiveResumeAfter = nil archiveRetrySuppressedUntilActivation = false + archiveAutoRetryCount = 0 eventHandler?(.archiveSyncCompleted(resolvedCheckpoint)) - replayDelayedLive(for: .account(client.userBareJid.stringValue)) + Task { @MainActor [weak self, weak client] in + guard let self, let client else { return } + await self.replayDelayedLive(for: .account(client.userBareJid.stringValue)) + } startNextMUCCatchupIfPossible() return } @@ -3323,11 +3378,39 @@ final class XMPPService { archiveSyncCompletedForConnection = false archiveResumeAfter = archiveSyncCheckpoint?.cursor archiveLastCompletedCursor = archiveSyncCheckpoint?.cursor - archiveRetrySuppressedUntilActivation = true - eventHandler?( - .recoverableError( - "Синхронизация истории приостановлена. Luma продолжит после возврата в приложение." - )) + if archiveAutoRetryCount < ArchiveSyncRecoveryPolicy.maximumAutomaticRetries { + // A single slow/timed-out page must not leave history unloaded + // until the next app activation. Retry automatically a bounded + // number of times while still connected and foregrounded. + archiveAutoRetryCount += 1 + archiveRetrySuppressedUntilActivation = false + scheduleArchiveSyncRetry(client: client) + } else { + archiveAutoRetryCount = 0 + archiveRetrySuppressedUntilActivation = true + eventHandler?( + .recoverableError( + "Синхронизация истории приостановлена. Luma продолжит после возврата в приложение." + )) + } + } + + private func scheduleArchiveSyncRetry(client: XMPPClient) { + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = Task { @MainActor [weak self, weak client] in + try? await Task.sleep( + nanoseconds: ArchiveSyncRecoveryPolicy.retryAfterFailureDelayNanoseconds + ) + guard !Task.isCancelled, + let self, + let client, + self.client === client, + client.state == .connected(), + !self.archiveSyncSuspended + else { return } + self.archiveSyncRetryTask = nil + self.refreshArchiveIfNeeded(client: client) + } } private func refreshArchiveIfNeeded(client: XMPPClient) { @@ -3376,6 +3459,8 @@ final class XMPPService { archiveQueryTimeoutTask = nil archiveQueryCompletionTask?.cancel() archiveQueryCompletionTask = nil + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = nil if archiveSyncStarted { archiveResumeAfter = archiveSyncCheckpoint?.cursor client.module(.omemo).mamSyncFinished(for: nil) @@ -3407,6 +3492,8 @@ final class XMPPService { archiveQueryTimeoutTask = nil archiveQueryCompletionTask?.cancel() archiveQueryCompletionTask = nil + archiveSyncRetryTask?.cancel() + archiveSyncRetryTask = nil archiveStanzaInbox.cancel() olderHistoryTimeoutTask?.cancel() olderHistoryTimeoutTask = nil diff --git a/Tests/ArchiveMessageBatchPolicyTests.swift b/Tests/ArchiveMessageBatchPolicyTests.swift index 8f28a9b..cdccd5e 100644 --- a/Tests/ArchiveMessageBatchPolicyTests.swift +++ b/Tests/ArchiveMessageBatchPolicyTests.swift @@ -7,8 +7,9 @@ final class ArchiveMessageBatchPolicyTests: XCTestCase { XCTAssertGreaterThan(ArchiveMessageBatchPolicy.pageSize, 0) } - func testOnlyOneStanzaIsDecodedPerMainActorSlice() { - XCTAssertEqual(ArchiveMessageBatchPolicy.decodeSliceSize, 1) + func testStanzasAreDecodedInSmallBatchesPerMainActorSlice() { + XCTAssertGreaterThan(ArchiveMessageBatchPolicy.decodeSliceSize, 0) + XCTAssertLessThanOrEqual(ArchiveMessageBatchPolicy.decodeSliceSize, 16) XCTAssertGreaterThan(ArchiveMessageBatchPolicy.interSliceDelayNanoseconds, 0) } diff --git a/Tests/ArchiveSyncRecoveryPolicyTests.swift b/Tests/ArchiveSyncRecoveryPolicyTests.swift index 40f0119..904d609 100644 --- a/Tests/ArchiveSyncRecoveryPolicyTests.swift +++ b/Tests/ArchiveSyncRecoveryPolicyTests.swift @@ -10,7 +10,7 @@ final class ArchiveSyncRecoveryPolicyTests: XCTestCase { XCTAssertGreaterThan(ArchiveSyncRecoveryPolicy.queryTimeoutNanoseconds, 0) XCTAssertLessThanOrEqual( ArchiveSyncRecoveryPolicy.queryTimeoutNanoseconds, - 15_000_000_000 + 20_000_000_000 ) } @@ -26,6 +26,15 @@ final class ArchiveSyncRecoveryPolicyTests: XCTestCase { XCTAssertEqual(ArchiveSyncRecoveryPolicy.pageRetryLimit, 1) } + func testCatchUpRetryDelayAndCapAreBounded() { + XCTAssertGreaterThan( + ArchiveSyncRecoveryPolicy.retryAfterFailureDelayNanoseconds, + 0 + ) + XCTAssertGreaterThan(ArchiveSyncRecoveryPolicy.maximumAutomaticRetries, 0) + XCTAssertLessThanOrEqual(ArchiveSyncRecoveryPolicy.maximumAutomaticRetries, 5) + } + func testCaptureResumeDelayLetsCameraReleaseResources() { XCTAssertGreaterThan( ArchiveSyncRecoveryPolicy.resumeAfterCaptureDelayNanoseconds,