diff --git a/apps/swift-ios/App/NativeFeatureClient.swift b/apps/swift-ios/App/NativeFeatureClient.swift index b37257fc6600..c45cf05e6cc7 100644 --- a/apps/swift-ios/App/NativeFeatureClient.swift +++ b/apps/swift-ios/App/NativeFeatureClient.swift @@ -123,6 +123,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, private var activeThreadSequence: Int? private var activeThreadPage: FeatureThreadPage? private var threadHistoryEpoch = 0 + private var detailSnapshotRequiredAfterEpoch: Int? private var pendingOlderThreadPage: PendingOlderThreadPage? init( @@ -3524,6 +3525,14 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, guard activeThreadID == threadID, activeThreadEnvironmentID == client.environment.id else { return } guard force || detailStreamTask == nil else { return } + if force { + // This required read owns recovery now. An older fallback must not + // replace its loading state with an error from a stale snapshot. + detailCatchUpTask?.cancel() + detailCatchUpTask = nil + detailCatchUpID = nil + continuation.yield(.threadSync(id: threadID, state: .catchingUp)) + } guard detailRefreshTask == nil else { detailRefreshPending = true return @@ -3534,9 +3543,11 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, let sessionGeneration = environmentGeneration detailRefreshTask = Task { [weak self] in do { - // Four updates per second keeps streaming text responsive while - // coalescing bursty shell events into one detail snapshot. - try await Task.sleep(for: .milliseconds(250)) + // Shell updates can be coalesced. A required replacement cannot + // apply more thread events until its snapshot arrives. + if !force { + try await Task.sleep(for: .milliseconds(250)) + } } catch { self?.finishDetailRefresh(generation: generation, client: client) return @@ -3549,7 +3560,20 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, environmentID: client.environment.id, generation: sessionGeneration ) { - try? await self.refreshThread(id: threadID, client: client) + do { + try await self.refreshThread(id: threadID, client: client) + } catch is CancellationError { + // Closing a thread cancels its read without changing its status. + } catch { + if !Task.isCancelled, + self.detailRefreshGeneration == generation, + self.activeThreadID == threadID, + self.activeRawThread == nil || self.detailStreamTask == nil { + self.continuation.yield(.threadSync( + id: threadID, state: .failed(error.localizedDescription) + )) + } + } } self.finishDetailRefresh(generation: generation, client: client) } @@ -3572,6 +3596,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, let supportsPagination = self?.serverConfigsByEnvironmentID[ route.environmentID ]?.threadSnapshotPagination == true + let subscriptionEpoch = self?.threadHistoryEpoch ?? 0 do { for try await item in await route.client.threadEvents( threadID: route.wireID, @@ -3582,7 +3607,9 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, guard !Task.isCancelled, let self, self.isCurrentDetail(route, generation: streamGeneration), self.environmentGeneration == sessionGeneration else { return } - self.consumeDetailStreamItem(item, route: route) + self.consumeDetailStreamItem( + item, route: route, subscriptionEpoch: subscriptionEpoch + ) } } catch is CancellationError { return @@ -3607,7 +3634,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, } private func ensureDetailCatchUpFallback(_ route: NativeThreadRoute, generation: Int) { - guard detailCatchUpTask == nil else { return } + guard detailCatchUpTask == nil, detailRefreshTask == nil else { return } let id = UUID() detailCatchUpID = id let delay = catchUpDelay @@ -3628,12 +3655,14 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, expectedStreamGeneration: generation ) guard !Task.isCancelled, let self, - self.isCurrentDetail(route, generation: generation) else { return } + self.isCurrentDetail(route, generation: generation), + self.activeRawThread != nil, + !self.detailRefreshPending else { return } if self.serverConfigsByEnvironmentID[route.environmentID]? .threadResumeCompletionMarker == true { self.continuation.yield(.threadSync(id: route.uiID, state: .reconnecting)) } else { - self.continuation.yield(.threadSync(id: route.uiID, state: .live)) + self.markDetailSynchronized(route) } } catch is CancellationError { return @@ -3646,7 +3675,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, } private func markDetailSynchronized(_ route: NativeThreadRoute) { - guard activeRawThread != nil else { return } + guard activeRawThread != nil, !detailRefreshPending else { return } // Flush the final message before publishing the completion state. // Otherwise the loading label can vanish one render before the text. flushDetailPublish(route) @@ -3658,7 +3687,8 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, private func consumeDetailStreamItem( _ item: ThreadStreamItem, - route: NativeThreadRoute + route: NativeThreadRoute, + subscriptionEpoch: Int ) { switch item { case .synchronized: @@ -3667,21 +3697,41 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, markDetailSynchronized(route) return case let .snapshot(snapshot): - guard activeRawThread == nil || snapshot.snapshotSequence > (activeThreadSequence ?? 0) else { return } + // A cursor-less event cannot prove that an already-requested + // snapshot includes it. Keep the post-event read until it does. + if let requiredEpoch = detailSnapshotRequiredAfterEpoch, + subscriptionEpoch < requiredEpoch { return } + guard snapshot.snapshotSequence >= (activeThreadSequence ?? 0), + activeRawThread == nil || snapshot.snapshotSequence > (activeThreadSequence ?? 0) else { return } + resetDetailRefresh() + detailSnapshotRequiredAfterEpoch = nil threadHistoryEpoch &+= 1 pendingOlderThreadPage = nil activeThreadSequence = snapshot.snapshotSequence activeRawThread = snapshot.thread activeThreadPage = featurePage(snapshot.page) scheduleRawDetailPublish(route: route, mutation: .full) + if detailCompletionReceived { markDetailSynchronized(route) } case let .event(event): guard let current = activeRawThread else { + // Do not apply later events to a snapshot that missed earlier + // ones. It must cover every event skipped while replacing it. + if case let .number(value) = event["sequence"], + let sequence = Int(exactly: value), sequence >= 0 { + activeThreadSequence = max(activeThreadSequence ?? 0, sequence) + } else { + // An event without a cursor needs a read started after it. + threadHistoryEpoch &+= 1 + detailSnapshotRequiredAfterEpoch = threadHistoryEpoch + pendingOlderThreadPage = nil + } scheduleDetailRefresh(threadID: route.uiID, client: route.client, force: true) return } let reduction = NativeThreadDetailReducer.apply(event, to: current) if reduction.sequence < 0 { threadHistoryEpoch &+= 1 + detailSnapshotRequiredAfterEpoch = threadHistoryEpoch pendingOlderThreadPage = nil activeRawThread = nil discardPendingDetailPublish() @@ -3722,13 +3772,13 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, detailPublishTask = Task { [weak self] in try? await Task.sleep(for: .milliseconds(80)) guard let self else { return } - self.detailPublishTask = nil guard !Task.isCancelled, self.detailStreamGeneration == streamGeneration, self.activeThreadID == route.uiID, self.activeRawThread != nil else { return } + self.detailPublishTask = nil self.flushDetailPublish(route) } } @@ -3774,7 +3824,9 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, let needsTrailingRefresh = detailRefreshPending detailRefreshPending = false if needsTrailingRefresh, let threadID = activeThreadID { - scheduleDetailRefresh(threadID: threadID, client: client) + // Events received without a base snapshot cannot be reduced. Read + // again even when the stream is open so those events are included. + scheduleDetailRefresh(threadID: threadID, client: client, force: true) } } @@ -3787,6 +3839,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, private func resetDetailStream() { detailStreamGeneration &+= 1 + detailSnapshotRequiredAfterEpoch = nil detailStreamTask?.cancel() detailStreamTask = nil detailCatchUpTask?.cancel() @@ -4090,6 +4143,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, } let environment = route.client.environment let generation = environmentGeneration + let historyEpoch = threadHistoryEpoch let supportsPagination = serverConfigsByEnvironmentID[ environment.id ]?.threadSnapshotPagination == true @@ -4104,8 +4158,25 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, throw CancellationError() } if activeThreadID == route.uiID { - guard snapshot.snapshotSequence >= (activeThreadSequence ?? 0) else { return } + if activeRawThread == nil, historyEpoch != threadHistoryEpoch { + return + } + guard snapshot.snapshotSequence >= (activeThreadSequence ?? 0) else { + if activeRawThread == nil { + if !detailRefreshPending { + throw NativeFeatureClientError.threadSnapshotOutdated + } + } else if detailCompletionReceived + || serverConfigsByEnvironmentID[environment.id]?.threadResumeCompletionMarker != true { + markDetailSynchronized(route) + } + return + } discardPendingDetailPublish() + // This snapshot includes the skipped events, so their pending + // request is satisfied without another HTTP read. + detailRefreshPending = false + detailSnapshotRequiredAfterEpoch = nil threadHistoryEpoch &+= 1 pendingOlderThreadPage = nil activeRawThread = snapshot.thread @@ -4126,22 +4197,17 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, client: client, thread: snapshot.thread, sequence: snapshot.snapshotSequence, page: featurePage(snapshot.page) ) - if activeThreadID == route.uiID, detailCompletionReceived { + if activeThreadID == route.uiID, + detailCompletionReceived + || serverConfigsByEnvironmentID[environment.id]?.threadResumeCompletionMarker != true { markDetailSynchronized(route) } - let hydrationBase = latestDetails[route.uiID] ?? detail - let hydrated = await hydratedAttachmentURLs( - in: hydrationBase, + scheduleAttachmentHydration( + in: detail, + threadID: route.uiID, client: client, - environmentID: environment.id, - generation: generation + environmentID: environment.id ) - guard isKnownClient(client, environmentID: environment.id, generation: generation), - latestDetails[route.uiID] == hydrationBase, - hydrated != hydrationBase else { - return - } - publish(hydrated, threadID: route.uiID, synchronizeRenderedMessages: true) } private func emitSnapshot( @@ -5956,9 +6022,7 @@ final class NativeFeatureClient: FeatureClient, FeatureDeviceManaging, environmentID: String ) { guard detail.messages.contains(where: { message in - message.attachments.contains { - $0.mimeType.hasPrefix("image/") && $0.url == nil - } + message.attachments.contains { $0.url == nil } }) else { return } @@ -7054,6 +7118,7 @@ private enum NativeFeatureClientError: LocalizedError { case environmentNotFound case projectNotFound case threadNotFound + case threadSnapshotOutdated case workspaceNotFound case approvalNotFound case inputRequestNotFound @@ -7071,6 +7136,7 @@ private enum NativeFeatureClientError: LocalizedError { case .environmentNotFound: "That T3 environment is no longer available." case .projectNotFound: "The selected project is no longer available." case .threadNotFound: "The selected thread is no longer available." + case .threadSnapshotOutdated: "The computer has not finished updating this thread. Try again." case .workspaceNotFound: "The thread workspace is no longer available." case .approvalNotFound: "The approval request is no longer active." case .inputRequestNotFound: "The input request is no longer active." diff --git a/apps/swift-ios/Features/Chat/ThreadDetailView.swift b/apps/swift-ios/Features/Chat/ThreadDetailView.swift index f2468cbaf8d0..39d192ed0dc8 100644 --- a/apps/swift-ios/Features/Chat/ThreadDetailView.swift +++ b/apps/swift-ios/Features/Chat/ThreadDetailView.swift @@ -308,11 +308,12 @@ public struct ThreadDetailView: View { Spacer(minLength: 6) - // The per-second timeline only exists for the live working - // duration; idle threads render a static status instead of - // waking every second forever. + // Cached work state is not proof that the agent is still + // running. Only current, working threads need a live timer. Group { - if currentThread.homeStatus == .working { + if refreshPresentation != nil { + EmptyView() + } else if currentThread.homeStatus == .working { TimelineView(.periodic(from: .now, by: 1)) { context in headerStatus(at: context.date) } @@ -331,7 +332,7 @@ public struct ThreadDetailView: View { .accessibilityElement(children: .combine) .accessibilityAddTraits(.isHeader) .accessibilityAddTraits( - currentThread.hasLiveWorkingDuration ? .updatesFrequently : [] + refreshPresentation == nil && currentThread.hasLiveWorkingDuration ? .updatesFrequently : [] ) .transaction { transaction in transaction.animation = nil @@ -631,17 +632,23 @@ public struct ThreadDetailView: View { } private func timeline(_ detail: FeatureThreadDetail) -> some View { - let isWorking = detail.thread.state == .working + let hasActiveWork = detail.thread.state == .working || detail.thread.state == .queued || detail.thread.state == .monitoring + let isWorking = hasActiveWork && refreshPresentation == nil return Group { - if detail.messages.isEmpty, !isWorking { - ContentUnavailableView( - "Ready for a task", - systemImage: "sparkles", - description: Text("Tell the agent what you want to build.") - ) - .frame(maxWidth: .infinity, maxHeight: .infinity) + if detail.messages.isEmpty, !hasActiveWork { + if refreshPresentation == nil { + ContentUnavailableView( + "Ready for a task", + systemImage: "sparkles", + description: Text("Tell the agent what you want to build.") + ) + .frame(maxWidth: .infinity, maxHeight: .infinity) + } else { + Color.clear + .frame(maxWidth: .infinity, maxHeight: .infinity) + } } else { FeatureTranscriptCollectionView( threadID: thread.id, @@ -1126,6 +1133,9 @@ enum ThreadRefreshPresentation: Equatable { } if isOpening || loadState == .loading { return .loading } if case .failed = loadState { return .failed } + // A synchronized thread subscription is newer evidence than an + // environment's periodic shell probe. Socket loss has its own state. + if syncState == .live { return nil } switch connectionState { case .connecting, .reconnecting: return .reconnecting case .disconnected: return .offline diff --git a/apps/swift-ios/Features/Root/FeatureRootModel.swift b/apps/swift-ios/Features/Root/FeatureRootModel.swift index af3dd65ddcc1..4dd9e2917772 100644 --- a/apps/swift-ios/Features/Root/FeatureRootModel.swift +++ b/apps/swift-ios/Features/Root/FeatureRootModel.swift @@ -965,7 +965,9 @@ public final class FeatureRootModel { store(value, delta: delta) upsert(value.thread) case let .threadSync(id, state): - threadSyncStates[id] = state + if threadSyncStates[id] != state { + threadSyncStates[id] = state + } if state == .live, case .failed = detailLoadStates[id] { detailLoadStates[id] = nil } diff --git a/apps/swift-ios/Tests/FeatureTests/FeatureRootModelTests.swift b/apps/swift-ios/Tests/FeatureTests/FeatureRootModelTests.swift index a130cea033ad..455dbc5774e4 100644 --- a/apps/swift-ios/Tests/FeatureTests/FeatureRootModelTests.swift +++ b/apps/swift-ios/Tests/FeatureTests/FeatureRootModelTests.swift @@ -36,6 +36,49 @@ struct FeatureRootModelTests { ) == nil) } + @Test + func liveThreadSyncOutranksStaleEnvironmentReachability() { + for connectionState in [FeatureConnection.State.disconnected, .reconnecting] { + #expect(ThreadRefreshPresentation.resolve( + loadState: nil, connectionState: connectionState, isOpening: false, syncState: .live + ) == nil) + #expect(ThreadRefreshPresentation.resolve( + loadState: .loading, connectionState: connectionState, isOpening: false, syncState: .live + ) == .loading) + #expect(ThreadRefreshPresentation.resolve( + loadState: .failed("Timeout"), connectionState: connectionState, isOpening: false, syncState: .live + ) == .failed) + } + } + + @Test + func repeatedCatchUpEventsDoNotInvalidateViewState() async { + let client = FeatureClientStub() + let model = testRootModel(client: client) + let run = Task { await model.start() } + await withCheckedContinuation { continuation in + withObservationTracking { + _ = model.threadSyncStates["thread"] + } onChange: { + continuation.resume() + } + client.emit(.threadSync(id: "thread", state: .catchingUp)) + } + let changes = AsyncStream.makeStream() + withObservationTracking { + _ = model.threadSyncStates["thread"] + } onChange: { + changes.continuation.yield() + } + client.emit(.threadSync(id: "thread", state: .catchingUp)) + client.finishEvents() + await run.value + changes.continuation.finish() + let didChange = await changes.stream.contains { _ in true } + #expect(!didChange) + #expect(model.threadSyncStates["thread"] == .catchingUp) + } + @Test func appearanceAppliesImmediatelyAndPersistsWithoutSavingTheDraft() async { let client = FeatureClientStub() diff --git a/apps/swift-ios/Tests/FeatureTests/NativeThreadCatchUpTests.swift b/apps/swift-ios/Tests/FeatureTests/NativeThreadCatchUpTests.swift index f35609dbb14d..eb507f584a01 100644 --- a/apps/swift-ios/Tests/FeatureTests/NativeThreadCatchUpTests.swift +++ b/apps/swift-ios/Tests/FeatureTests/NativeThreadCatchUpTests.swift @@ -153,6 +153,366 @@ final class NativeThreadCatchUpTests: XCTestCase { await fixture.client.disconnect() } + func testFailedRequiredSnapshotKeepsCachedTextAndOffersFreshRetry() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + await fixture.http.setResponse(text: "Cached answer", sequence: 2) + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + + await fixture.http.holdThreadReads(true) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + let read = try await nextHeldRead(&reads) + read.fail() + while let event = await events.next(isolation: #isolation) { + if case .threadSync(fixture.firstID, .failed) = event { break } + if case .threadSync(fixture.firstID, .live) = event { + XCTFail("A failed snapshot must not mark cached content current.") + } + } + + await fixture.http.holdThreadReads(false) + await fixture.http.setResponse(text: "Fresh answer", sequence: 11) + let detail = try await fixture.client.loadThread(id: fixture.firstID, fresh: true) + XCTAssertEqual(detail.messages.map(\.text), ["Fresh answer"]) + await fixture.client.disconnect() + } + + func testRequiredSnapshotsCoverEveryEventReceivedDuringReplacement() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + + await fixture.http.setResponse(text: "Before new messages", sequence: 10) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + let first = try await nextHeldRead(&reads) + try await stream.sendMessage(text: "Eleven", sequence: 11) + await nextCatchUp(&events, threadID: fixture.firstID) + await fixture.http.setResponse(text: "Eleven", sequence: 11) + first.succeed() + + let second = try await nextHeldRead(&reads) + await nextCatchUp(&events, threadID: fixture.firstID) + try await stream.sendMessage(text: "Twelve", sequence: 12) + await nextCatchUp(&events, threadID: fixture.firstID) + await fixture.http.setResponse(texts: ["Eleven", "Twelve"], sequence: 12) + second.succeed() + + let third = try await nextHeldRead(&reads) + await nextCatchUp(&events, threadID: fixture.firstID) + third.succeed() + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["Eleven", "Twelve"]) + let count = await fixture.http.threadRequests.count + XCTAssertEqual(count, 4, "Each stale response needs one coalesced follow-up, not one read per event.") + await fixture.client.disconnect() + } + + func testEventWithoutCursorNeedsSnapshotStartedAfterIt() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + await fixture.http.setResponse(text: "Before unknown event", sequence: 10) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + let first = try await nextHeldRead(&reads) + try await stream.socket.chunk(id: stream.id, values: [.object(["kind": .string("unknown")])]) + await nextCatchUp(&events, threadID: fixture.firstID) + await fixture.http.setResponse(text: "After unknown event", sequence: 11) + first.succeed() + let second = try await nextHeldRead(&reads) + await nextCatchUp(&events, threadID: fixture.firstID) + second.succeed() + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["After unknown event"]) + await fixture.client.disconnect() + } + + func testSocketSnapshotCancelsFailedHTTPReplacement() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + let read = try await nextHeldRead(&reads) + try await stream.snapshot(texts: ["Recovered over the socket"], sequence: 11) + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["Recovered over the socket"]) + read.fail() + let cancelled = await read.finished.first { _ in true } + XCTAssertEqual(cancelled, true, "The obsolete HTTP read must not replace live state with an error.") + await fixture.client.disconnect() + } + + func testOldSocketSnapshotDoesNotCancelReadAfterCursorlessEvent() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + await fixture.http.setResponse(text: "Fresh after unknown event", sequence: 20) + try await stream.socket.chunk(id: stream.id, values: [.object(["kind": .string("unknown")])]) + await nextCatchUp(&events, threadID: fixture.firstID) + let read = try await nextHeldRead(&reads) + + try await stream.snapshot(texts: ["Requested before unknown event"], sequence: 3) + try await stream.synchronize() + // This next event is a receipt that the old snapshot was handled first. + try await stream.sendMessage(text: "After old snapshot", sequence: 19) + await nextCatchUp(&events, threadID: fixture.firstID) + read.succeed() + let cancelled = await read.finished.first { _ in true } + XCTAssertEqual(cancelled, false, "The required post-event read must still run.") + guard cancelled == false else { + await fixture.client.disconnect() + return + } + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["Fresh after unknown event"]) + await fixture.client.disconnect() + } + + func testNewSubscriptionSnapshotCanRecoverAfterCursorlessEvent() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + try await stream.socket.chunk(id: stream.id, values: [.object(["kind": .string("unknown")])]) + await nextCatchUp(&events, threadID: fixture.firstID) + let read = try await nextHeldRead(&reads) + + await stream.socket.close() + let resumed = try await nextThreadRequest(&requests) + try await resumed.snapshot(texts: ["New socket snapshot"], sequence: 20) + try await resumed.synchronize() + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["New socket snapshot"]) + read.fail() + let cancelled = await read.finished.first { _ in true } + XCTAssertEqual(cancelled, true, "The new subscription can replace the required HTTP read.") + await fixture.client.disconnect() + } + + func testLeavingThreadCancelsHeldReplacementAndItsPendingFollowUp() async throws { + for disconnect in [false, true] { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + let read = try await nextHeldRead(&reads) + try await stream.sendMessage(text: "Do not publish after leaving", sequence: 11) + await nextCatchUp(&events, threadID: fixture.firstID) + + if disconnect { + await fixture.client.disconnect() + } else { + fixture.client.releaseThread(id: fixture.firstID) + await fixture.http.holdThreadReads(false) + await fixture.http.setResponse(text: "Second thread", sequence: 20) + let detail = try await fixture.client.loadThread(id: fixture.secondID) + XCTAssertEqual(detail.thread.id, fixture.secondID) + XCTAssertEqual(detail.messages.map(\.text), ["Second thread"]) + } + read.succeed() + let cancelled = await read.finished.first { _ in true } + XCTAssertEqual(cancelled, true) + let count = await fixture.http.threadRequests.count + XCTAssertEqual(count, disconnect ? 2 : 3, "A closed thread must not start its pending read.") + await fixture.client.disconnect() + } + } + + func testAttachmentLookupDoesNotBlockRequiredTextRefresh() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + await fixture.http.holdThreadReads(true) + await fixture.http.setResponse(texts: ["Text ready"], sequence: 11, withImage: true) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + let first = try await nextHeldRead(&reads) + try await stream.sendMessage(text: "Text ready", sequence: 11) + await nextCatchUp(&events, threadID: fixture.firstID) + first.succeed() + while let request = await requests.next(isolation: #isolation) { + if request.tag == RPCMethod.assetsCreateURL.rawValue { break } + } + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["Text ready"]) + let firstReadCount = await fixture.http.threadRequests.count + XCTAssertEqual(firstReadCount, 2, "The first snapshot already includes the skipped message.") + + // Leave the asset RPC unanswered. A later text refresh must start anyway. + await fixture.http.setResponse(texts: ["New text ready"], sequence: 12, withImage: true) + try await stream.invalidate(sequence: 12) + await nextCatchUp(&events, threadID: fixture.firstID) + let second = try await nextHeldRead(&reads) + second.succeed() + let updated = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(updated, ["New text ready"]) + await fixture.client.disconnect() + } + + func testNonImageAttachmentsResolveAfterTextCatchUpCompletes() async throws { + for (name, mimeType) in [("document.pdf", "application/pdf"), ("clip.mp4", "video/mp4")] { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + _ = try await fixture.client.loadThread(id: fixture.firstID) + let stream = try await nextThreadRequest(&requests) + try await stream.synchronize() + _ = await messagesBeforeLive(&events, threadID: fixture.firstID) + + await fixture.http.setResponse(text: "File ready", sequence: 10, attachment: .init( + type: "file", id: "file", name: name, mimeType: mimeType, sizeBytes: 20 + )) + try await stream.invalidate(sequence: 10) + await nextCatchUp(&events, threadID: fixture.firstID) + // Text must become current before the asset URL request completes. + let messages = await messagesBeforeLive(&events, threadID: fixture.firstID) + XCTAssertEqual(messages, ["File ready"]) + + while let request = await requests.next(isolation: #isolation) { + guard request.tag == RPCMethod.assetsCreateURL.rawValue else { continue } + XCTAssertEqual(request.payload["resource"]?["mimeType"], .string(mimeType)) + try await request.socket.succeed(id: request.id, value: .object([ + "relativeUrl": .string("/assets/\(name)"), + "expiresAt": .number(Date.now.addingTimeInterval(3_600).timeIntervalSince1970 * 1_000), + ])) + break + } + var resolved: FeatureMessageAttachment? + while let event = await events.next(isolation: #isolation) { + guard case let .detail(detail) = event, detail.thread.id == fixture.firstID, + let attachment = detail.messages.first?.attachments.first else { continue } + resolved = attachment + break + } + XCTAssertEqual(resolved?.mimeType, mimeType) + XCTAssertEqual(resolved?.url, URL(string: "https://one.example/assets/\(name)")) + await fixture.client.disconnect() + } + } + + func testRequiredReadReplacesColdFallbackWithoutHidingItsOwnFailure() async throws { + let fixture = try await CatchUpFixture.make() + defer { fixture.cleanUp() } + var requests = fixture.requests.makeAsyncIterator() + var events = fixture.client.events().makeAsyncIterator() + var reads = fixture.http.heldRequests.makeAsyncIterator() + await fixture.http.holdThreadReads(true) + + // Fail the cold open so catch-up starts without a base snapshot. + let opening = Task { try await fixture.client.loadThread(id: fixture.firstID) } + let initial = try await nextHeldRead(&reads) + initial.fail() + do { + _ = try await opening.value + XCTFail("The initial snapshot should fail.") + } catch {} + let stream = try await nextThreadRequest(&requests) + while let event = await events.next(isolation: #isolation) { + if case .threadSync(fixture.firstID, .failed) = event { break } + } + await nextCatchUp(&events, threadID: fixture.firstID) + await fixture.delay.release() + let fallback = try await nextHeldRead(&reads) + + // The event needs a newer snapshot than the fallback captured. + await fixture.http.setResponse(text: "New message", sequence: 3) + try await stream.sendMessage(text: "New message", sequence: 3) + await nextCatchUp(&events, threadID: fixture.firstID) + let replacement = try await nextHeldRead(&reads) + fallback.succeed() + let wasCancelled = await fallback.finished.first { _ in true } + XCTAssertEqual(wasCancelled, true, "The older fallback must stop when the required read takes over.") + + replacement.fail() + var failure: String? + while let event = await events.next(isolation: #isolation) { + guard case let .threadSync(id, .failed(message)) = event, + id == fixture.firstID else { continue } + failure = message + break + } + XCTAssertEqual(failure, URLError(.notConnectedToInternet).localizedDescription) + await fixture.client.disconnect() + } + + private func nextHeldRead( + _ iterator: inout AsyncStream.Iterator + ) async throws -> CatchUpHTTPRead { + let read = await iterator.next(isolation: #isolation) + return try XCTUnwrap(read) + } + + private func nextCatchUp( + _ iterator: inout AsyncStream.Iterator, threadID: String + ) async { + while let event = await iterator.next(isolation: #isolation) { + if case .threadSync(threadID, .catchingUp) = event { return } + if case .threadSync(threadID, .live) = event { + XCTFail("The thread became live before its required snapshot was complete.") + return + } + } + XCTFail("The thread did not report its pending refresh.") + } + private func nextThreadRequest( _ iterator: inout AsyncStream.Iterator ) async throws -> CatchUpRequest { @@ -171,6 +531,9 @@ final class NativeThreadCatchUpTests: XCTestCase { case let .detail(detail), let .detailDelta(detail, _): if detail.thread.id == threadID { messages = detail.messages.map(\.text) } case .threadSync(threadID, .live): return messages + case let .threadSync(id, .failed(error)) where id == threadID: + XCTFail("The thread failed synchronization: \(error)") + return messages default: break } } @@ -223,15 +586,33 @@ private struct CatchUpFixture { private actor CatchUpHTTPTransport: HTTPTransport { private(set) var threadRequests: [URLRequest] = [] - private var message: String? + private var messages: [OrchestrationMessage] = [] private var sequence = 2 + private var holdsThreadReads = false + private let heldReadContinuation: AsyncStream.Continuation + nonisolated let heldRequests: AsyncStream - func setResponse(text: String, sequence: Int) { - message = text + init() { + let reads = AsyncStream.makeStream() + heldRequests = reads.stream + heldReadContinuation = reads.continuation + } + + func setResponse(text: String, sequence: Int, attachment: ChatAttachment? = nil) { + messages = [catchUpMessage(text, index: 0, attachment: attachment)] self.sequence = sequence } - func data(for request: URLRequest) throws -> (Data, HTTPURLResponse) { + func setResponse(texts: [String], sequence: Int, withImage: Bool = false) { + messages = texts.enumerated().map { index, text in + catchUpMessage(text, index: index, withImage: withImage) + } + self.sequence = sequence + } + + func holdThreadReads(_ hold: Bool) { holdsThreadReads = hold } + + func data(for request: URLRequest) async throws -> (Data, HTTPURLResponse) { let value: JSONValue switch request.url!.path { case "/api/auth/websocket-ticket": @@ -248,24 +629,51 @@ private actor CatchUpHTTPTransport: HTTPTransport { throw URLError(.unsupportedURL) } threadRequests.append(request) - let messages = message.map { text in - [OrchestrationMessage( - id: "answer", role: "assistant", text: text, attachments: [], - turnId: nil, streaming: false, createdAt: "2026-09-02T12:00:00Z", - updatedAt: "2026-09-02T12:00:00Z" - )] - } ?? [] value = try .encode(multiEnvironmentDetail( projectID: "project", threadID: request.url!.lastPathComponent, snapshotSequence: sequence, messages: messages )) } - return (try JSONEncoder.t3.encode(value), HTTPURLResponse( + let response = (try JSONEncoder.t3.encode(value), HTTPURLResponse( url: request.url!, statusCode: 200, httpVersion: "HTTP/1.1", headerFields: nil )!) + if holdsThreadReads, request.url!.path.hasPrefix("/api/orchestration/threads/") { + let finished = AsyncStream.makeStream() + defer { + finished.continuation.yield(Task.isCancelled) + finished.continuation.finish() + } + return try await withCheckedThrowingContinuation { continuation in + heldReadContinuation.yield(CatchUpHTTPRead( + response: response, continuation: continuation, finished: finished.stream + )) + } + } + return response } } +private struct CatchUpHTTPRead: Sendable { + let response: (Data, HTTPURLResponse) + let continuation: CheckedContinuation<(Data, HTTPURLResponse), any Error> + let finished: AsyncStream + func succeed() { continuation.resume(returning: response) } + func fail() { continuation.resume(throwing: URLError(.notConnectedToInternet)) } +} + +private func catchUpMessage( + _ text: String, index: Int, withImage: Bool = false, attachment: ChatAttachment? = nil +) -> OrchestrationMessage { + OrchestrationMessage( + id: "answer-\(index)", role: "assistant", text: text, + attachments: attachment.map { [$0] } ?? (withImage ? [.init( + type: "image", id: "image-\(index)", name: "test.png", mimeType: "image/png", sizeBytes: 20 + )] : []), + turnId: nil, streaming: false, createdAt: "2026-09-02T12:00:00Z", + updatedAt: "2026-09-02T12:00:00Z" + ) +} + private struct CatchUpConnector: WebSocketConnecting { let requests: AsyncStream.Continuation let completionMarker: Bool @@ -284,6 +692,27 @@ private struct CatchUpRequest: Sendable { try await socket.chunk(id: id, values: [.object(["kind": .string("synchronized")])]) } + func invalidate(sequence: Int) async throws { + try await socket.chunk(id: id, values: [.object([ + "kind": .string("event"), "event": .object([ + "type": .string("thread.reverted"), "sequence": .number(Double(sequence)), + "occurredAt": .string("2026-09-02T12:00:00Z"), + "payload": .object(["threadId": payload["threadId"]!]), + ]), + ])]) + } + + func snapshot(texts: [String], sequence: Int) async throws { + let snapshot = multiEnvironmentDetail( + projectID: "project", threadID: payload["threadId"]!.stringValue!, + snapshotSequence: sequence, + messages: texts.enumerated().map { catchUpMessage($0.element, index: $0.offset) } + ) + try await socket.chunk(id: id, values: [.object([ + "kind": .string("snapshot"), "snapshot": try .encode(snapshot), + ])]) + } + func sendMessage(text: String, sequence: Int) async throws { try await socket.chunk(id: id, values: [.object([ "kind": .string("event"), "event": .object([ @@ -347,6 +776,13 @@ private actor CatchUpSocket: WebSocketConnection { ])) } + func succeed(id: Int, value: JSONValue) throws { + try enqueue(.object([ + "_tag": .string("Exit"), "requestId": .number(Double(id)), + "exit": .object(["_tag": .string("Success"), "value": value]), + ])) + } + private func enqueue(_ value: JSONValue) throws { let data = try JSONEncoder.t3.encode(value) if let receiver { diff --git a/docs/user/swiftui-mobile.md b/docs/user/swiftui-mobile.md index 3d0dd17179de..a6804f4028d7 100644 --- a/docs/user/swiftui-mobile.md +++ b/docs/user/swiftui-mobile.md @@ -28,6 +28,13 @@ settings. Supported environments offer an icon picker in their connection prefer Use **Refresh prices** in Usage to fetch new model rates without waiting for the daily refresh. A failed environment does not remove the usage data from other environments. +## Thread connection state + +Cached messages stay readable while a thread catches up. A status above the composer shows when +the content is incomplete or the computer cannot be reached. Use **Retry** if an update fails. +Working indicators return after the thread is current. File previews can finish loading after +the text is ready. + ## Attachments and sharing One message can contain up to eight photos, videos, or files. Images can be up to 10 MB. Other