From ddfa08b7f55e5697995e417f4b76ce1bdeed396e Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 16:06:22 +0900 Subject: [PATCH 1/8] Update Codex runtime while retaining queued reviews and MCP sessions --- .../CodexReview/ReviewRuntimeLifecycle.swift | 1 + .../CodexReview/Store/CodexReviewStore.swift | 178 ++++++++-- .../Store/CodexReviewStoreBackend.swift | 2 + .../Store/CodexReviewStoreCancellation.swift | 10 +- .../Store/CodexReviewStoreUpdate.swift | 159 +++++++++ .../Store/ReviewRuntimeTeardownIntent.swift | 9 +- .../LiveCodexReviewStoreBackend.swift | 12 +- Sources/CodexReviewTesting/TestSupport.swift | 19 +- .../CodexReviewHostTests.swift | 41 ++- .../CodexReviewMCPHTTPServerTests.swift | 62 +++- .../CodexReviewStoreHistoryTests.swift | 31 ++ .../CodexReviewStoreUpdateTests.swift | 326 ++++++++++++++++++ 12 files changed, 797 insertions(+), 53 deletions(-) create mode 100644 Sources/CodexReview/Store/CodexReviewStoreUpdate.swift create mode 100644 Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift diff --git a/Sources/CodexReview/ReviewRuntimeLifecycle.swift b/Sources/CodexReview/ReviewRuntimeLifecycle.swift index b22a654e..0fc30fd0 100644 --- a/Sources/CodexReview/ReviewRuntimeLifecycle.swift +++ b/Sources/CodexReview/ReviewRuntimeLifecycle.swift @@ -14,6 +14,7 @@ package struct ReviewRuntimeGeneration: Hashable, Sendable { package enum ReviewRuntimeCleanupRecoverySuppression: Equatable, Sendable { case explicitStop + case codexUpdate case staleSource case successorAlreadyFinished } diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index c8e90988..41c1fda0 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -77,7 +77,7 @@ package enum StoreReviewAttemptOwnership { @MainActor @Observable public final class CodexReviewStore { - private struct RuntimeStartOperation { + package struct RuntimeStartOperation { let task: Task let sourceCloseReceiptOwner: ReviewRuntimeRecoveryReplacement? } @@ -88,6 +88,13 @@ public final class CodexReviewStore { package var timeoutTask: Task? } + package var codexUpdate: CodexUpdateState = .idle + @_spi(ApplicationHostSupport) public var codexUpdateState: CodexUpdateState { codexUpdate } + @ObservationIgnored package var codexUpdateRuntimeFailure: String? + @ObservationIgnored package var codexUpdateTask: Task? + @ObservationIgnored package var runtimeAccountOperations: [UUID: Task] = [:] + @ObservationIgnored package var unclosedCodexUpdateRuntime: PreparedRuntime? + public package(set) var serverState: CodexReviewServerState = .stopped public let auth: CodexReviewAuthModel package let settings: SettingsStore @@ -206,6 +213,7 @@ public final class CodexReviewStore { } isolated deinit { + codexUpdateTask?.cancel() accountRateLimitAutoRefreshDriver?.cancel() historyLoadTask?.cancel() applicationShutdownTask?.cancel() @@ -272,6 +280,10 @@ public final class CodexReviewStore { } public func start(forceRestartIfNeeded: Bool = false) async { + if let update = codexUpdateTask { + _ = try? await update.value + if case .running = runtimeState { return } + } guard applicationShutdownRequested == false else { return } @@ -281,6 +293,12 @@ public final class CodexReviewStore { else { return } + switch runtimeState { + case .stopped, .failed(_, nil, _): + do { try await closeUnclosedCodexUpdateRuntimeIfNeeded() } + catch { transitionToFailed(error.localizedDescription); return } + default: break + } guard let operation = admitRuntimeStart( forceRestartIfNeeded: forceRestartIfNeeded ) else { @@ -301,7 +319,7 @@ public final class CodexReviewStore { } } - private func admitRuntimeStart( + package func admitRuntimeStart( forceRestartIfNeeded: Bool ) -> RuntimeStartOperation? { let previousState = runtimeState @@ -327,8 +345,9 @@ public final class CodexReviewStore { case .failed(let sourceGeneration, let retainedMCP?, _): return admitRuntimeReplacement( sourceGeneration: sourceGeneration, - retiringRuntime: nil, - retainedMCP: retainedMCP + retiringRuntime: unclosedCodexUpdateRuntime, + retainedMCP: retainedMCP, + preservingQueuedReviews: hasQueuedCodexUpdateRecovery ) case .acquiring(let sourceGeneration, let context, let task): task.cancel() @@ -381,8 +400,16 @@ public final class CodexReviewStore { } package func stop(intent: ReviewRuntimeTeardownIntent) async { + let update = codexUpdateTask + if codexUpdate == .waitingForReviews { + update?.cancel() + } else { + _ = try? await update?.value + } let task = admitRuntimeTeardown(intent: intent) await task.value + _ = try? await update?.value + if intent == .explicitStop { reviewStartsAreSuspended = false } } package func requestRuntimeTeardown( @@ -396,6 +423,21 @@ public final class CodexReviewStore { handle: any RuntimeLifecycleHandle, cause: String ) -> Bool { + if codexUpdateTask != nil { + switch runtimeState { + case .running(_, let runtime, _) where runtime.handle === handle: + for job in jobs where job.isTerminal == false && queuedReviewStarts[job.id] == nil { + markReviewFailed(job, message: cause, terminal: .interrupted(.transport(message: cause))) + } + serverState = .failed(cause) + writeDiagnosticsIfNeeded() + return true + case .replacing(let replacement, _) where replacement.ownsPublishedRuntime(handle: handle): + codexUpdateRuntimeFailure = cause + return true + default: break + } + } let failureIncident: ReviewRuntimeFailureIncident switch runtimeState { case .running(let generation, let runtime, _) @@ -427,6 +469,11 @@ public final class CodexReviewStore { sourceGeneration: ReviewRuntimeGeneration, cause: String ) -> ReviewRuntimeCleanupRecoveryAdmission { + if codexUpdateTask != nil, + case .running(let generation, let runtime, _) = runtimeState, + runtime.handle === sourceHandle, generation == sourceGeneration { + return .suppressed(.codexUpdate) + } switch runtimeState { case .running(let generation, let runtime, _) where runtime.handle === sourceHandle && generation == sourceGeneration: @@ -640,6 +687,12 @@ public final class CodexReviewStore { await task.value case .failed(_, let retainedMCP, _): + await stopPublishedRuntimeSemantics(intent: intent) + if let runtime = unclosedCodexUpdateRuntime { + if await closeRuntime(runtime, purpose: .stop) == nil { + unclosedCodexUpdateRuntime = nil + } + } if retainedMCP != nil { await stopMCPServer() } @@ -1080,13 +1133,16 @@ public final class CodexReviewStore { } } - private func admitRuntimeReplacement( + package func admitRuntimeReplacement( sourceGeneration: ReviewRuntimeGeneration, retiringRuntime: PreparedRuntime?, - retainedMCP: RetainedMCPServer + retainedMCP: RetainedMCPServer, + preservingQueuedReviews: Bool = false, + install: (@MainActor @Sendable () async -> Void)? = nil ) -> RuntimeStartOperation { - storeWorkRegistry.closeReviewAdmission() - let pendingHistoryStarts = requestHistoryStartCancellations( + if preservingQueuedReviews { codexUpdateRuntimeFailure = nil } + else { storeWorkRegistry.closeReviewAdmission() } + let pendingHistoryStarts = preservingQueuedReviews ? [] : requestHistoryStartCancellations( cancellation: .system(message: "Review runtime restarted.") ) let replacement = ReviewRuntimeRecoveryReplacement( @@ -1102,7 +1158,7 @@ public final class CodexReviewStore { return } await self.waitForHistoryStarts(pendingHistoryStarts) - await self.performRuntimeReplacement(replacement) + await self.performRuntimeReplacement(replacement, preservingQueuedReviews: preservingQueuedReviews, install: install) } runtimeState = .replacing( replacement: replacement, @@ -1112,7 +1168,9 @@ public final class CodexReviewStore { } private func performRuntimeReplacement( - _ replacement: ReviewRuntimeRecoveryReplacement + _ replacement: ReviewRuntimeRecoveryReplacement, + preservingQueuedReviews: Bool, + install: (@MainActor @Sendable () async -> Void)? ) async { guard isCurrentReplacement(replacement) else { return @@ -1125,12 +1183,17 @@ public final class CodexReviewStore { let token = try await settingsService.beginRuntimeCutover() cutoverToken = token - await closeRetiringRuntime(for: replacement) + if preservingQueuedReviews { + try await closeRetiringRuntimeForCodexUpdate(replacement) + } else { + await closeRetiringRuntime(for: replacement) + } guard isCurrentReplacement(replacement) else { cancelRuntimeCutover(token) return } + await install?() let runtime = try await backend.prepareRuntime( generation: replacement.replacementGeneration, purpose: .restartSameAccount @@ -1175,6 +1238,9 @@ public final class CodexReviewStore { guard isCurrentReplacement(replacement) else { return } + if preservingQueuedReviews, let failure = codexUpdateRuntimeFailure { + throw CodexReviewAPI.Error.io(failure) + } guard let publishedRuntime = replacement.takePublishedRuntime() else { preconditionFailure( "ReviewRuntimeRecoveryReplacement must own its published runtime until exact-identity transfer." @@ -1189,9 +1255,18 @@ public final class CodexReviewStore { publishRuntime(serverURL: replacement.retainedMCP.serverURL) replacement.finish(.running(replacement.replacementGeneration)) } catch { - await closeRetiringRuntime(for: replacement) - if let preparedRuntime { - await closeRuntime(preparedRuntime, purpose: .restartSameAccount) + var failureMessage = error.localizedDescription + if preservingQueuedReviews { + do { try await closeRetiringRuntimeForCodexUpdate(replacement) } + catch { failureMessage += "; Runtime cleanup failed: \(error.localizedDescription)" } + } else { + await closeRetiringRuntime(for: replacement) + } + let failedRuntime = preparedRuntime ?? replacement.takePublishedRuntime() + if let failedRuntime, + let cleanupFailure = await closeRuntime(failedRuntime, purpose: .restartSameAccount) { + failureMessage += "; Runtime cleanup failed: \(cleanupFailure.localizedDescription)" + if preservingQueuedReviews { unclosedCodexUpdateRuntime = failedRuntime } } let isCurrentGeneration = isCurrentReplacement(replacement) let wasIntentionallyCancelled = Task.isCancelled || isCurrentGeneration == false @@ -1199,7 +1274,7 @@ public final class CodexReviewStore { if wasIntentionallyCancelled { cancelRuntimeCutover(cutoverToken) } else { - abortRuntimeCutover(cutoverToken, message: error.localizedDescription) + abortRuntimeCutover(cutoverToken, message: failureMessage) } } guard wasIntentionallyCancelled == false else { @@ -1211,12 +1286,35 @@ public final class CodexReviewStore { failureIncident: nil ) serverURL = replacement.retainedMCP.serverURL - serverState = .failed(error.localizedDescription) + serverState = .failed(failureMessage) writeDiagnosticsIfNeeded() - replacement.finish(.failed(error.localizedDescription)) + replacement.finish(.failed(failureMessage)) } } + private func closeRetiringRuntimeForCodexUpdate( + _ replacement: ReviewRuntimeRecoveryReplacement + ) async throws { + guard let retiring = replacement.takeRetiringRuntime() else { + replacement.finishSourceClose(.closed) + return + } + unclosedCodexUpdateRuntime = retiring + await stopPublishedRuntimeSemantics(intent: .codexUpdate) + if let failure = await closeRuntime(retiring, purpose: .restartSameAccount) { + replacement.finishSourceClose(.failed(failure)) + throw failure + } + unclosedCodexUpdateRuntime = nil + replacement.finishSourceClose(.closed) + } + + package func closeUnclosedCodexUpdateRuntimeIfNeeded() async throws { + guard let runtime = unclosedCodexUpdateRuntime else { return } + if let failure = await closeRuntime(runtime, purpose: .restartSameAccount) { throw failure } + unclosedCodexUpdateRuntime = nil + } + private func closePublishedRuntimeForReplacement( _ runtime: PreparedRuntime ) async -> ReviewRuntimeCloseFailure? { @@ -1251,7 +1349,8 @@ public final class CodexReviewStore { locallyCancelledJobIDs = [] } else { let cancellationOutcome = await requestActiveReviewCancellationsForRuntimeStop( - reason: intent.reviewCancellation + reason: intent.reviewCancellation, + includingQueued: intent != .codexUpdate ) locallyCancelledJobIDs = cancellationOutcome.jobIDs if let failure = cancellationOutcome.firstFailure { @@ -1262,10 +1361,12 @@ public final class CodexReviewStore { } await backend.stop(store: self, intent: intent) let remainingLocallyCancelledJobIDs = cancelActiveReviewsLocallyForRuntimeStop( - reason: intent.reviewCancellation + reason: intent.reviewCancellation, + includingQueued: intent != .codexUpdate ) await cancelAndDetachReviewWorkersForRuntimeStop( - jobIDs: Array(Set(locallyCancelledJobIDs + remainingLocallyCancelledJobIDs)), + jobIDs: Array(Set(locallyCancelledJobIDs + remainingLocallyCancelledJobIDs + + (intent == .codexUpdate ? reviewWorkerJobIDsForRuntimeStop : []))), reason: intent.reviewCancellation ) } @@ -1400,7 +1501,7 @@ public final class CodexReviewStore { for intent: ReviewRuntimeTeardownIntent ) -> ReviewRuntimeTransitionPurpose { switch intent { - case .explicitStop: + case .explicitStop, .codexUpdate: .stop case .unexpectedFailure: .runtimeFailure @@ -1410,27 +1511,30 @@ public final class CodexReviewStore { private func publishRuntime(serverURL: URL?) { storeWorkRegistry.openReviewAdmission() transitionToRunning(serverURL: serverURL) + if hasQueuedCodexUpdateRecovery, codexUpdateTask == nil { + resumeReviewsAfterCodexUpdateIfPossible() + } startAccountRateLimitAutoRefresh() } public func refreshAuthentication() async { - await backend.refreshAuth(auth: auth) + await performRuntimeAuthentication { store in await store.backend.refreshAuth(auth: store.auth) } } public func signIn() async { - await backend.signIn(auth: auth, using: .chatGPT) + await performRuntimeAuthentication { store in await store.backend.signIn(auth: store.auth, using: .chatGPT) } } package func signIn(apiKey: CodexReviewAPIKey) async { - await backend.signIn(auth: auth, using: .apiKey(apiKey)) + await performRuntimeAuthentication { store in await store.backend.signIn(auth: store.auth, using: .apiKey(apiKey)) } } public func addAccount() async { - await backend.addAccount(auth: auth, using: .chatGPT) + await performRuntimeAuthentication { store in await store.backend.addAccount(auth: store.auth, using: .chatGPT) } } package func addAccount(apiKey: CodexReviewAPIKey) async { - await backend.addAccount(auth: auth, using: .apiKey(apiKey)) + await performRuntimeAuthentication { store in await store.backend.addAccount(auth: store.auth, using: .apiKey(apiKey)) } } public func cancelAuthentication() async { @@ -1463,7 +1567,7 @@ public final class CodexReviewStore { else { return } - await backend.signIn(auth: auth, using: method) + await performRuntimeAuthentication { store in await store.backend.signIn(auth: store.auth, using: method) } } public func logout() async { @@ -1481,7 +1585,9 @@ public final class CodexReviewStore { } public func signOutActiveAccount() async throws { - try await backend.signOutActiveAccount(auth: auth) + try await performRuntimeAccountChange { store in + try await store.backend.signOutActiveAccount(auth: store.auth) + } } package func switchAccount(_ account: CodexAccount) async throws { @@ -1496,7 +1602,9 @@ public final class CodexReviewStore { defer { targetAccount?.updateIsSwitching(false) } - try await backend.switchAccount(auth: auth, accountKey: account.accountKey) + try await performRuntimeAccountChange { store in + try await store.backend.switchAccount(auth: store.auth, accountKey: account.accountKey) + } } package func requestSwitchAccount(_ account: CodexAccount, requiresConfirmation: Bool) { @@ -1583,15 +1691,21 @@ public final class CodexReviewStore { } package func removeAccount(accountKey: String) async throws { - try await backend.removeAccount(auth: auth, accountKey: accountKey) + try await performRuntimeAccountChange { store in + try await store.backend.removeAccount(auth: store.auth, accountKey: accountKey) + } } package func reorderPersistedAccount(accountKey: String, toIndex: Int) async throws { - try await backend.reorderPersistedAccount(auth: auth, accountKey: accountKey, toIndex: toIndex) + try await performRuntimeAccountChange { store in + try await store.backend.reorderPersistedAccount(auth: store.auth, accountKey: accountKey, toIndex: toIndex) + } } package func refreshAccountRateLimits(accountKey: String) async { - await backend.refreshAccountRateLimits(auth: auth, accountKey: accountKey) + await performRuntimeAuthentication { store in + await store.backend.refreshAccountRateLimits(auth: store.auth, accountKey: accountKey) + } } package func startStartupAuthRefresh() { diff --git a/Sources/CodexReview/Store/CodexReviewStoreBackend.swift b/Sources/CodexReview/Store/CodexReviewStoreBackend.swift index 3396d1a7..1975ac87 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreBackend.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreBackend.swift @@ -27,6 +27,7 @@ package struct CodexReviewStoreSeed { package protocol CodexReviewStoreBackend: CodexReviewSettingsBackend, Sendable { var seed: CodexReviewStoreSeed { get } var isActive: Bool { get } + var shutdownCleanupTimeout: Duration { get } var handlesActiveReviewStopCleanup: Bool { get } var mcpServerLifecycle: any MCPServerLifecycleOwner { get } @@ -85,6 +86,7 @@ package protocol CodexReviewStoreBackend: CodexReviewSettingsBackend, Sendable { } extension CodexReviewStoreBackend { + package var shutdownCleanupTimeout: Duration { .seconds(2) } package var handlesActiveReviewStopCleanup: Bool { false } diff --git a/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift b/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift index ee64f43d..f782378b 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift @@ -192,10 +192,11 @@ extension CodexReviewStore { } package func requestActiveReviewCancellationsForRuntimeStop( - reason: ReviewCancellation = .system(message: "Review runtime stopped.") + reason: ReviewCancellation = .system(message: "Review runtime stopped."), + includingQueued: Bool = true ) async -> ReviewRuntimeCancellationRequestOutcome { let activeJobIDs = orderedJobs - .filter { $0.isTerminal == false } + .filter { $0.isTerminal == false && (includingQueued || queuedReviewStarts[$0.id] == nil) } .map(\.id) var firstFailure: ReviewRuntimeCloseFailure? for jobID in activeJobIDs { @@ -221,10 +222,11 @@ extension CodexReviewStore { @discardableResult package func cancelActiveReviewsLocallyForRuntimeStop( - reason: ReviewCancellation = .system(message: "Review runtime stopped.") + reason: ReviewCancellation = .system(message: "Review runtime stopped."), + includingQueued: Bool = true ) -> [String] { let activeJobIDs = orderedJobs - .filter { $0.isTerminal == false } + .filter { $0.isTerminal == false && (includingQueued || queuedReviewStarts[$0.id] == nil) } .map(\.id) guard activeJobIDs.isEmpty == false else { return [] diff --git a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift new file mode 100644 index 00000000..d806ab9e --- /dev/null +++ b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift @@ -0,0 +1,159 @@ +import Foundation + +extension CodexReviewStore { + @_spi(ApplicationHostSupport) public enum CodexUpdateTiming: Sendable { + case afterCurrentReviews + case immediately + } + + @_spi(ApplicationHostSupport) public enum CodexUpdateState: Equatable, Sendable { + case idle + case waitingForReviews + case stoppingRuntime + case installing + case restarting + case failed(String) + } + + /// Keeps MCP sessions and accepted jobs alive while replacing the Codex runtime. + /// Concurrent callers join the existing operation; only its installation closure runs. + /// An installation failure is thrown even when the runtime recovers and jobs resume. + @_spi(ApplicationHostSupport) public func updateCodex( + when timing: CodexUpdateTiming, + install: @escaping @MainActor @Sendable () async throws -> Void + ) async throws { + if let task = codexUpdateTask { + try await task.value + return + } + guard applicationShutdownRequested == false else { throw CancellationError() } + let previousAccountOperations = Array(runtimeAccountOperations.values) + let previousRuntimeTask: Task? = switch runtimeState { + case .acquiring(_, _, let task), .replacing(_, let task), .tearingDown(_, _, _, _, let task): task + case .running, .stopped, .failed: nil + } + suspendReviewStarts() + codexUpdate = .waitingForReviews + let task = Task { @MainActor [self] in + defer { + codexUpdateTask = nil + if Task.isCancelled == false { resumeReviewsAfterCodexUpdateIfPossible() } + } + do { + for operation in previousAccountOperations { _ = try? await operation.value } + await previousRuntimeTask?.value + try await performCodexUpdate(when: timing, install: install) + codexUpdate = .idle + } catch { + if error is CancellationError { + codexUpdate = .idle + } else { + codexUpdate = .failed(error.localizedDescription) + } + throw error + } + } + codexUpdateTask = task + try await task.value + } + + private func performCodexUpdate( + when timing: CodexUpdateTiming, + install: @escaping @MainActor @Sendable () async throws -> Void + ) async throws { + // A dispatcher already saving an execution start belongs to the current batch. + await queuedReviewDispatchTask?.value + try Task.checkCancellation() + if timing == .afterCurrentReviews { + let executingIDs = jobs.filter { + $0.isTerminal == false && queuedReviewStarts[$0.id] == nil + }.map(\.id) + for id in executingIDs { + _ = try await awaitReview(sessionID: nil, jobID: id) + try Task.checkCancellation() + } + _ = await drainReviewWorkersForRuntimeStop(timeout: backend.shutdownCleanupTimeout) + try Task.checkCancellation() + } + guard applicationShutdownRequested == false else { throw CancellationError() } + codexUpdate = .stoppingRuntime + var installationFailure: String? + let installation: @MainActor @Sendable () async -> Void = { [self] in + codexUpdate = .installing + do { try await install() } catch { installationFailure = error.localizedDescription } + codexUpdate = .restarting + } + let operation: RuntimeStartOperation + switch runtimeState { + case .running(let generation, let runtime, let mcp): + operation = admitRuntimeReplacement( + sourceGeneration: generation, retiringRuntime: runtime, retainedMCP: mcp, + preservingQueuedReviews: true, + install: installation + ) + case .failed(let generation, let mcp?, _): + operation = admitRuntimeReplacement( + sourceGeneration: generation, retiringRuntime: unclosedCodexUpdateRuntime, retainedMCP: mcp, + preservingQueuedReviews: true, + install: installation + ) + case .stopped, .failed: + try await closeUnclosedCodexUpdateRuntimeIfNeeded() + await installation() + guard let start = admitRuntimeStart(forceRestartIfNeeded: false) else { throw CancellationError() } + operation = start + case .acquiring, .replacing, .tearingDown: + throw CodexReviewAPI.Error.io("The Codex runtime changed while waiting to update.") + } + await operation.task.value + let failures = [installationFailure, serverState.failureMessage].compactMap { $0 } + if failures.isEmpty == false { + throw CodexReviewAPI.Error.io(failures.joined(separator: "; ")) + } + guard case .running = runtimeState else { + throw CodexReviewAPI.Error.io("Codex update finished without an available runtime.") + } + } + + package func resumeReviewsAfterCodexUpdateIfPossible() { + guard runtimeAccountOperations.isEmpty, applicationShutdownRequested == false, + case .running = runtimeState else { return } + resumeReviewStarts() + } + + package func performRuntimeAuthentication( + _ operation: @escaping @MainActor @Sendable (CodexReviewStore) async -> Void + ) async { + do { try await performRuntimeAccountChange { store in await operation(store) } } + catch is CancellationError { } + catch { auth.updatePhase(.failed(message: error.localizedDescription)) } + } + + package func performRuntimeAccountChange( + _ operation: @escaping @MainActor @Sendable (CodexReviewStore) async throws -> Void + ) async throws { + let update = codexUpdateTask + let id = UUID() + let task = Task { @MainActor [self] in + _ = try? await update?.value + defer { + runtimeAccountOperations.removeValue(forKey: id) + if codexUpdateTask == nil { resumeReviewsAfterCodexUpdateIfPossible() } + } + try Task.checkCancellation() + guard applicationShutdownRequested == false else { throw CancellationError() } + try await operation(self) + } + runtimeAccountOperations[id] = task + try await withTaskCancellationHandler { + try await task.value + } onCancel: { + task.cancel() + } + } + + package var hasQueuedCodexUpdateRecovery: Bool { + if case .failed = codexUpdateState { return reviewStartsAreSuspended } + return false + } +} diff --git a/Sources/CodexReview/Store/ReviewRuntimeTeardownIntent.swift b/Sources/CodexReview/Store/ReviewRuntimeTeardownIntent.swift index 0c27f466..cc2e2eac 100644 --- a/Sources/CodexReview/Store/ReviewRuntimeTeardownIntent.swift +++ b/Sources/CodexReview/Store/ReviewRuntimeTeardownIntent.swift @@ -6,6 +6,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable { case failed(String) } + case codexUpdate case explicitStop case unexpectedFailure(String) @@ -15,7 +16,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable { package var finalState: FinalState { switch self { - case .explicitStop: + case .explicitStop, .codexUpdate: .stopped case .unexpectedFailure: .failed(message) @@ -24,7 +25,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable { package var diagnosticContext: String { switch self { - case .explicitStop: + case .explicitStop, .codexUpdate: "runtime stop" case .unexpectedFailure: "runtime failure" @@ -33,7 +34,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable { package var cleanupTimeoutWarning: String { switch self { - case .explicitStop: + case .explicitStop, .codexUpdate: "Timed out cleaning active reviews before stopping runtime" case .unexpectedFailure: "Timed out cleaning active reviews after runtime failure" @@ -46,7 +47,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable { private var message: String { switch self { - case .explicitStop: + case .explicitStop, .codexUpdate: "Review runtime stopped." case .unexpectedFailure(let errorDescription): "Review runtime stopped unexpectedly: \(errorDescription)" diff --git a/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift b/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift index 6b0146d8..6894140d 100644 --- a/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift +++ b/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift @@ -726,7 +726,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend, MCPSer private let mcpPortOwnerResolver: CodexReviewMCPPortOwnerResolver private let mcpHTTPServerBindChecker: CodexReviewMCPHTTPServerBindChecker private let appServerRuntimeFactory: AppServerRuntimeFactory - private let shutdownCleanupTimeout: Duration + let shutdownCleanupTimeout: Duration private weak var attachedStore: CodexReviewStore? private var preparingMCPServer: PreparedMCPServer? private var preparedMCPServer: PreparedMCPServer? @@ -1169,13 +1169,14 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend, MCPSer store: CodexReviewStore, appServerBackend: AppServerCodexReviewBackend, reason: ReviewCancellation, - timeoutWarning: String + timeoutWarning: String, + includingQueued: Bool = true ) async { let workerJobIDs = store.reviewWorkerJobIDsForRuntimeStop let cancellationCleanup = await runRuntimeShutdownCleanup( timeout: shutdownCleanupTimeout ) { - await store.requestActiveReviewCancellationsForRuntimeStop(reason: reason) + await store.requestActiveReviewCancellationsForRuntimeStop(reason: reason, includingQueued: includingQueued) } let cancellationJobIDs: [String] let didRequestCancellation: Bool @@ -1207,7 +1208,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend, MCPSer } } let locallyCancelledJobIDs = store.cancelActiveReviewsLocallyForRuntimeStop( - reason: reason + reason: reason, includingQueued: includingQueued ) let currentWorkerJobIDs = store.reviewWorkerJobIDsForRuntimeStop await store.cancelAndDetachReviewWorkersForRuntimeStop( @@ -1252,7 +1253,8 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend, MCPSer store: store, appServerBackend: appServerBackend, reason: intent.reviewCancellation, - timeoutWarning: intent.cleanupTimeoutWarning + timeoutWarning: intent.cleanupTimeoutWarning, + includingQueued: intent != .codexUpdate ) } await cleanupLoginRuntime(loginCleanup) diff --git a/Sources/CodexReviewTesting/TestSupport.swift b/Sources/CodexReviewTesting/TestSupport.swift index a9365639..47bfe89b 100644 --- a/Sources/CodexReviewTesting/TestSupport.swift +++ b/Sources/CodexReviewTesting/TestSupport.swift @@ -223,6 +223,7 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { private var auth: CodexReviewBackendModel.Auth.Snapshot private var commands: [Command] = [] private var startAdmissionIdentities: [ObjectIdentifier] = [] + private var scriptedRuns: [CodexReviewBackendModel.Review.Run] = [] private var nextRun: CodexReviewBackendModel.Review.Run private var nextRecoveredRun: CodexReviewBackendModel.Review.Run? private var interruptFailureMessage: String? @@ -331,6 +332,10 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { interruptReviewGate = gate } + package func scriptReviewRuns(_ runs: [CodexReviewBackendModel.Review.Run]) { + scriptedRuns = runs + } + package func setNextRun(_ run: CodexReviewBackendModel.Review.Run) { nextRun = run } @@ -576,6 +581,8 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { _ request: CodexReviewBackendModel.Review.Start, admission: ReviewStartAdmission ) async throws -> BackendReviewAttempt { + if scriptedRuns.isEmpty == false { nextRun = scriptedRuns.removeFirst() } + let run = nextRun startAdmissionIdentities.append(ObjectIdentifier(admission)) try await admission.admitThreadStartDispatch() commands.append(.startReview(request)) @@ -585,10 +592,10 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { waiter.resume() } let provisionalRun = CodexReviewBackendModel.Review.Run( - attemptID: nextRun.attemptID, - threadID: nextRun.threadID, - reviewThreadID: nextRun.threadID, - model: nextRun.model + attemptID: run.attemptID, + threadID: run.threadID, + reviewThreadID: run.threadID, + model: run.model ) try await admission.recordPreparedThread(provisionalRun) do { @@ -604,8 +611,8 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { await startReviewGate.wait() } } - try await admission.recordActiveRun(nextRun) - return .init(run: nextRun, events: eventMailbox(for: nextRun)) + try await admission.recordActiveRun(run) + return .init(run: run, events: eventMailbox(for: run)) } package func receivedStartAdmission(_ admission: ReviewStartAdmission) -> Bool { diff --git a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift index ce9d41f8..b4d687c1 100644 --- a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift +++ b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift @@ -2,7 +2,7 @@ import Foundation import AppKit import AuthenticationServices import Testing -import CodexReview +@_spi(ApplicationHostSupport) import CodexReview import CodexReviewAppServer import CodexReviewHost import CodexReviewMCPServer @@ -46,6 +46,45 @@ private actor HostCloseFailureTransport: JSONRPC.Transport { @Suite("host composition") @MainActor struct CodexReviewHostTests { + @Test func liveRuntimeUpdatePreservesMCPAndDispatchesQueueOnNewTransport() async throws { + let homeURL = try temporaryHome() + let old = FakeJSONRPCTransport() + let new = FakeJSONRPCTransport() + try await enqueueRuntimeStartResponses(old) + try await enqueueRuntimeStartResponses(new) + try await enqueueLiveRouteReviewStartResponses(new, threadID: "new-thread", turnID: "new-turn") + try await new.enqueue(EmptyResponse(), for: "turn/interrupt") + try await enqueueReviewCleanupResponses(new) + let server = ControlledMCPHTTPServer(endpoint: try #require(URL(string: "http://127.0.0.1:19438/mcp"))) + var transports = [old, new] + let store = CodexReviewStore.makeLiveStoreForTesting( + environment: ["HOME": homeURL.path], + webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, + mcpHTTPServerFactory: { _, _ in server }, + mcpHTTPServerBindChecker: { _ in }, + transportFactory: { _ in transports.removeFirst() } + ) + await store.start() + store.suspendReviewStarts() + let queued = try await store.startReview( + sessionID: "client", request: .init(cwd: "/tmp/project", target: .uncommittedChanges), waitTimeout: .zero + ) + try await store.updateCodex(when: .immediately) { + #expect(await old.isClosedForTesting()) + #expect(store.job(id: queued.jobID)?.core.lifecycle.status == .queued) + #expect(server.stopCallCount == 0) + } + try #require(await StoreSnapshotProbe(store: store).waitUntil { $0.job()?.activeRun?.turnID == "new-turn" } != nil) + #expect(await old.recordedRequests().contains { $0.method == "thread/start" } == false) + #expect(server.startCallCount == 1) + #expect(store.serverURL == server.endpoint) + let cancel = Task { try await store.cancelReview(jobID: queued.jobID, sessionID: "client") } + try #require(await waitUntil(timeout: .seconds(2)) { await new.recordedRequests().contains { $0.method == "turn/interrupt" } }) + try await emitInterruptedTurn(new, threadID: "new-thread", turnID: "new-turn", message: "Cancelled.") + _ = try await cancel.value + await store.stop() + } + @Test func hostStartsAndStopsRuntimeWithFakeBackend() async throws { let backend = FakeCodexReviewBackend() let host = CodexReviewHost( diff --git a/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift b/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift index 473e2858..046cf0e0 100644 --- a/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift +++ b/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift @@ -3,13 +3,73 @@ import Foundation import MCP @preconcurrency import NIOCore import Testing -@_spi(Testing) @testable import CodexReview +@_spi(Testing) @_spi(ApplicationHostSupport) @testable import CodexReview import CodexReviewMCPServer import CodexReviewTesting @Suite("MCP Streamable HTTP server") @MainActor struct CodexReviewMCPHTTPServerTests { + @Test func sameMCPSessionWaitsThroughDeferredUpdateAndCompletesTwoQueuedCalls() async throws { + let reviews = FakeCodexReviewBackend() + let runs = (0..<3).map { CodexReviewBackendModel.Review.Run( + threadID: "thread-\($0)", turnID: "turn-\($0)", reviewThreadID: "review-\($0)" + ) } + await reviews.scriptReviewRuns(runs) + let owner = TestingMCPServerLifecycleOwner() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews, mcpServerLifecycle: owner) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + try await withHTTPServer(store: store) { server in + let endpoint = await server.url + let session = try await initializeSession(endpoint: endpoint) + let first = Task { try await postJSONRPCData( + endpoint: endpoint, sessionID: session, bodyData: makeReviewStartBody(id: 2) + ) } + try await reviews.waitForStartReview(timeout: .seconds(2)) + var installs = 0 + let entered = AsyncGate() + let release = AsyncGate() + let update = Task { try await store.updateCodex(when: .afterCurrentReviews) { + installs += 1 + await entered.open() + await release.wait() + } } + try #require(await waitUntil(timeout: .seconds(2)) { store.codexUpdateState == .waitingForReviews }) + var queuedResponses = 0 + let second = Task { + defer { queuedResponses += 1 } + return try await postJSONRPCData(endpoint: endpoint, sessionID: session, bodyData: makeReviewStartBody(id: 3)) + } + let third = Task { + defer { queuedResponses += 1 } + return try await postJSONRPCData(endpoint: endpoint, sessionID: session, bodyData: makeReviewStartBody(id: 4)) + } + try #require(await waitUntil(timeout: .seconds(2)) { store.jobs.filter { $0.core.lifecycle.status == .queued }.count == 2 }) + let queuedIDs = Set(store.jobs.filter { $0.core.lifecycle.status == .queued }.map(\.id)) + await reviews.yield(.completed(summary: "Done", result: "No findings."), for: runs[0]) + let firstResult = try decodeSSEJSON(from: try await first.value) + #expect(firstResult.value(for: ["result", "structuredContent", "lifecycle", "status"]) as? String == "succeeded") + await entered.wait() + #expect(queuedResponses == 0) + #expect(owner.stopCallCount == 0) + await release.open() + try await update.value + try #require(await waitUntil(timeout: .seconds(2)) { + queuedIDs.allSatisfy { store.job(id: $0)?.core.run.threadID != nil } + }) + await reviews.yield(.completed(summary: "Done", result: "No findings."), for: runs[1]) + await reviews.yield(.completed(summary: "Done", result: "No findings."), for: runs[2]) + let results = try [decodeSSEJSON(from: await second.value), decodeSSEJSON(from: await third.value)] + #expect(Set(results.compactMap { $0.value(for: ["result", "structuredContent", "jobId"]) as? String }) == queuedIDs) + #expect(results.allSatisfy { $0.value(for: ["result", "structuredContent", "lifecycle", "status"]) as? String == "succeeded" }) + #expect(installs == 1) + #expect(owner.preparedServers.count == 1) + #expect(owner.activatedServers.count == 1) + } + await store.stop() + } + @Test func reviewStartRemainsPendingAcrossKitQueueWithoutClientResubmission() async throws { let backend = FakeCodexReviewBackend() let store = CodexReviewStore.makeTestingStore( diff --git a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift index 0a59dd55..c8862578 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift @@ -7,6 +7,37 @@ import CodexReviewTesting @Suite("review history store", .serialized) @MainActor struct CodexReviewStoreHistoryTests { + @Test func updateKeepsRequestAcceptedDuringHistoryWriteInTheQueue() async throws { + let acceptedEntered = AsyncGate() + let acceptedRelease = AsyncGate() + let history = ReviewHistoryPersistenceProbe(startedWriteEntered: acceptedEntered, startedWriteRelease: acceptedRelease) + let backend = FakeCodexReviewBackend() + let store = makeStore(history: history, backend: backend) + await store.start() + let request = Task { try await store.startReview( + sessionID: "client", request: .init(cwd: "/tmp/queued", target: .uncommittedChanges), waitTimeout: .zero + ) } + await acceptedEntered.wait() + let installEntered = AsyncGate() + let installRelease = AsyncGate() + let update = Task { try await store.updateCodex(when: .afterCurrentReviews) { + await installEntered.open() + await installRelease.wait() + } } + await installEntered.wait() + await acceptedRelease.open() + let queued = try await request.value + #expect(queued.core.lifecycle.status == .queued) + #expect(queued.core.lifecycle.startedAt == nil) + #expect(await backend.recordedCommands().contains { if case .startReview = $0 { true } else { false } } == false) + await installRelease.open() + try await update.value + try await backend.waitForStartReview(timeout: .seconds(2)) + await backend.yield(.completed(summary: "Done", result: "No findings.")) + #expect(try await store.awaitReview(sessionID: "client", jobID: queued.jobID).core.lifecycle.status == .succeeded) + await store.stop() + } + @Test func executionPersistenceFailureReturnsTheAcceptedJobWithoutDispatch() async throws { let history = ReviewHistoryPersistenceProbe(executionWriteFailure: "Execution start could not be saved.") let backend = FakeCodexReviewBackend() diff --git a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift new file mode 100644 index 00000000..4ce2e3f4 --- /dev/null +++ b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift @@ -0,0 +1,326 @@ +import Foundation +import Testing +@_spi(ApplicationHostSupport) @testable import CodexReview +import CodexReviewTesting + +@Suite("Codex runtime updates", .serialized) +@MainActor +struct CodexReviewStoreUpdateTests { + @Test func deferredUpdateWaitsForCleanupAndRetainsQueuedCalls() async throws { + let reviews = FakeCodexReviewBackend() + let mcp = TestingMCPServerLifecycleOwner() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews, mcpServerLifecycle: mcp) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let originalRuntime = try #require(backend.lastPreparedRuntimeHandle) + let cleanup = AsyncGate() + await reviews.holdCleanupReview(with: cleanup) + let first = Task { try await store.startReview(sessionID: "client", request: request("first")) } + try await reviews.waitForStartReview(timeout: .seconds(2)) + var installations = 0 + let installing = AsyncGate() + let releaseInstall = AsyncGate() + let update = Task { + try await store.updateCodex(when: .afterCurrentReviews) { + installations += 1 + #expect(originalRuntime.waitUntilClosedCallCount == 1) + await installing.open() + await releaseInstall.wait() + } + } + try #require(await waitUntil { store.codexUpdateState == .waitingForReviews }) + var queuedReturned = false + let queued = Task { + defer { queuedReturned = true } + return try await store.startReview(sessionID: "client", request: request("queued")) + } + try #require(await waitUntil { store.jobs.contains { $0.cwd == "/tmp/queued" } }) + let queuedID = try #require(store.jobs.first { $0.cwd == "/tmp/queued" }?.id) + await reviews.yield(.completed(summary: "Done", result: "No findings.")) + await reviews.waitForCleanupReview() + #expect(installations == 0) + #expect(originalRuntime.closePurposes.isEmpty) + await cleanup.open() + #expect(try await first.value.core.lifecycle.status == .succeeded) + await installing.wait() + #expect(queuedReturned == false) + #expect(try store.readReview(sessionID: "client", jobID: queuedID).core.lifecycle.status == .queued) + #expect(mcp.stopCallCount == 0) + let duplicate = Task { + try await store.updateCodex(when: .immediately) { installations += 100 } + } + await releaseInstall.open() + try await update.value + try await duplicate.value + #expect(installations == 1) + #expect(backend.startRequests == [false, true]) + #expect(mcp.preparedServers.count == 1) + #expect(mcp.activatedServers.count == 1) + try #require(await waitUntil { store.job(id: queuedID)?.core.run.threadID != nil }) + await reviews.yield(.completed(summary: "Done", result: "No findings.")) + #expect(try await queued.value.jobID == queuedID) + #expect(store.codexUpdateState == .idle) + await store.stop() + } + + @Test func failedInstallRestartsRuntimeAndResumesQueueButKeepsError() async throws { + let reviews = FakeCodexReviewBackend() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + var queuedID: String? + do { + try await store.updateCodex(when: .afterCurrentReviews) { + queuedID = try await store.startReview(sessionID: "client", request: request("queued"), waitTimeout: .zero).jobID + throw UpdateTestError.install + } + Issue.record("Expected installation failure") + } catch { #expect(error.localizedDescription.contains("install failed")) } + #expect(store.serverState == .running) + #expect(store.codexUpdateState == .failed("install failed")) + let id = try #require(queuedID) + try #require(await waitUntil { store.job(id: id)?.core.run.threadID != nil }) + await reviews.yield(.completed(summary: "Done", result: "No findings.")) + #expect(try await store.awaitReview(sessionID: "client", jobID: id).core.lifecycle.status == .succeeded) + await store.stop() + } + + @Test func failedRestartRetainsQueueAndRetryDoesNotInstallAgain() async throws { + let reviews = FakeCodexReviewBackend() + let mcp = TestingMCPServerLifecycleOwner() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews, mcpServerLifecycle: mcp) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + backend.failNextRuntimePreparation(message: "restart failed") + var queuedID: String? + var installs = 0 + do { + try await store.updateCodex(when: .afterCurrentReviews) { + installs += 1 + queuedID = try await store.startReview(sessionID: "client", request: request("queued"), waitTimeout: .zero).jobID + } + Issue.record("Expected restart failure") + } catch { #expect(error.localizedDescription.contains("restart failed")) } + let id = try #require(queuedID) + #expect(try store.readReview(sessionID: "client", jobID: id).core.lifecycle.status == .queued) + #expect(mcp.stopCallCount == 0) + await store.restart() + #expect(store.serverState == .running) + #expect(installs == 1) + try #require(await waitUntil { store.job(id: id)?.core.run.threadID != nil }) + await reviews.yield(.completed(summary: "Done", result: "No findings.")) + #expect(try await store.awaitReview(sessionID: "client", jobID: id).core.lifecycle.status == .succeeded) + await store.stop() + } + + @Test func failedSourceClosePreventsInstallationAndKeepsRecoverableRuntime() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let runtime = try #require(backend.lastPreparedRuntimeHandle) + runtime.failClose(with: .process("close failed")) + var installs = 0 + do { + try await store.updateCodex(when: .afterCurrentReviews) { installs += 1 } + Issue.record("Expected close failure") + } catch { #expect(error.localizedDescription.contains("close failed")) } + #expect(installs == 0) + #expect(store.unclosedCodexUpdateRuntime?.handle === runtime) + #expect(backend.startRequests == [false]) + await store.restart() + #expect(store.serverState == .running) + #expect(store.unclosedCodexUpdateRuntime == nil) + #expect(installs == 0) + await store.stop() + } + + @Test func shutdownWaitsForInstallationThenCancelsQueuedJobs() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let entered = AsyncGate() + let release = AsyncGate() + let update = Task { + try await store.updateCodex(when: .immediately) { + await entered.open() + await release.wait() + } + } + await entered.wait() + let queued = try await store.startReview(sessionID: "client", request: request("queued"), waitTimeout: .zero) + var stopped = false + let shutdown = Task { await store.shutdown(); stopped = true } + try #require(await waitUntil { store.applicationShutdownRequested }) + #expect(stopped == false) + await release.open() + try await update.value + await shutdown.value + #expect(store.serverState == .stopped) + #expect(store.job(id: queued.jobID)?.core.lifecycle.status == .cancelled) + } + + @Test func immediateUpdateCancelsExecutingReviewButKeepsQueuedReview() async throws { + let reviews = FakeCodexReviewBackend() + await reviews.holdStartReview(with: AsyncGate()) + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let active = Task { try await store.startReview(sessionID: "owner", request: request("active")) } + try await reviews.waitForStartReview(timeout: .seconds(2)) + var queuedID: String? + try await store.updateCodex(when: .immediately) { + queuedID = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero).jobID + #expect(try await active.value.core.lifecycle.status == .cancelled) + } + let id = try #require(queuedID) + #expect(store.job(id: id)?.isTerminal == false) + await store.stop() + } + + @Test func shutdownCancelsDeferredUpdateWithoutWaitingForReviewCompletion() async throws { + let reviews = FakeCodexReviewBackend() + await reviews.holdStartReview(with: AsyncGate()) + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let active = Task { try await store.startReview(sessionID: "owner", request: request("active")) } + try await reviews.waitForStartReview(timeout: .seconds(2)) + var installs = 0 + let update = Task { try await store.updateCodex(when: .afterCurrentReviews) { installs += 1 } } + try #require(await waitUntil { store.codexUpdateState == .waitingForReviews }) + await store.shutdown() + do { try await update.value; Issue.record("Expected update cancellation") } + catch { #expect(error is CancellationError) } + #expect(installs == 0) + #expect(try await active.value.core.lifecycle.status == .cancelled) + } + + @Test func manualRestartJoinsInstallationInsteadOfReplacingAgain() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let entered = AsyncGate() + let release = AsyncGate() + let update = Task { try await store.updateCodex(when: .immediately) { + await entered.open() + await release.wait() + } } + await entered.wait() + var restartReturned = false + let restart = Task { await store.restart(); restartReturned = true } + await Task.yield() + #expect(restartReturned == false) + await release.open() + try await update.value + await restart.value + #expect(backend.startRequests == [false, true]) + await store.stop() + } + + @Test func accountChangesAndUpdateExecuteInOrderWithoutDispatchingQueueBetweenThem() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let firstEntered = AsyncGate() + let firstRelease = AsyncGate() + var events: [String] = [] + let first = Task { try await store.performRuntimeAccountChange { _ in + events.append("first") + await firstEntered.open() + await firstRelease.wait() + } } + await firstEntered.wait() + let installEntered = AsyncGate() + let installRelease = AsyncGate() + let update = Task { try await store.updateCodex(when: .immediately) { + events.append("install") + await installEntered.open() + await installRelease.wait() + } } + try #require(await waitUntil { store.codexUpdateTask != nil }) + #expect(events == ["first"]) + await firstRelease.open() + try await first.value + await installEntered.wait() + let queued = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero) + let second = Task { try await store.performRuntimeAccountChange { store in + events.append("second") + #expect(store.job(id: queued.jobID)?.core.lifecycle.status == .queued) + await store.closeActiveReviewSessions(reason: .system(message: "Account switched.")) + } } + try #require(await waitUntil { store.runtimeAccountOperations.isEmpty == false }) + await installRelease.open() + try await update.value + try await second.value + #expect(events == ["first", "install", "second"]) + #expect(store.job(id: queued.jobID)?.core.lifecycle.status == .cancelled) + await store.stop() + } + + @Test func failedRuntimeDuringDeferredUpdateDoesNotDiscardQueuedJobs() async throws { + let reviews = FakeCodexReviewBackend() + let mcp = TestingMCPServerLifecycleOwner() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews, mcpServerLifecycle: mcp) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let runtime = try #require(backend.lastPreparedRuntimeHandle) + let active = Task { try await store.startReview(sessionID: "owner", request: request("active")) } + try await reviews.waitForStartReview(timeout: .seconds(2)) + var installs = 0 + let update = Task { try await store.updateCodex(when: .afterCurrentReviews) { installs += 1 } } + try #require(await waitUntil { store.codexUpdateState == .waitingForReviews }) + let queued = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero) + #expect(store.requestRuntimeFailure(handle: runtime, cause: "process exited")) + await reviews.finishEvents(throwing: UpdateTestError.install) + try await update.value + #expect(installs == 1) + #expect(try await active.value.core.lifecycle.terminal == .interrupted(.transport(message: "process exited"))) + #expect(store.job(id: queued.jobID)?.isTerminal == false) + #expect(mcp.stopCallCount == 0) + await store.stop() + } + + @Test func publicationFailureRetainsQueueAndCanBeRetried() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + backend.runOnNextRuntimePublication { + if let handle = backend.lastPreparedRuntimeHandle { + #expect(store.requestRuntimeFailure(handle: handle, cause: "publication failed")) + } + } + var queuedID: String? + do { + try await store.updateCodex(when: .immediately) { + queuedID = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero).jobID + } + Issue.record("Expected publication failure") + } catch { #expect(error.localizedDescription.contains("publication failed")) } + let id = try #require(queuedID) + #expect(store.job(id: id)?.core.lifecycle.status == .queued) + #expect(store.serverState == .failed("publication failed")) + await store.restart() + #expect(store.serverState == .running) + await store.stop() + } + + private func request(_ name: String) -> CodexReviewAPI.Start.Request { + .init(cwd: "/tmp/\(name)", target: .uncommittedChanges) + } +} + +private enum UpdateTestError: LocalizedError { + case install + var errorDescription: String? { "install failed" } +} + +@MainActor +private func waitUntil(condition: () -> Bool) async -> Bool { + let clock = ContinuousClock() + let deadline = clock.now + .seconds(2) + while condition() == false { + if clock.now >= deadline { return false } + await Task.yield() + } + return true +} From 6941e402300c568b46ee250b57ed7f9d8c1e03d6 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 16:12:57 +0900 Subject: [PATCH 2/8] Preserve update recovery across retries and history deletion --- .../CodexReview/Store/CodexReviewStore.swift | 4 +- .../Store/CodexReviewStoreUpdate.swift | 2 + .../CodexReviewStoreHistoryTests.swift | 39 +++++++++++++++++++ .../CodexReviewStoreUpdateTests.swift | 9 +++++ 4 files changed, 52 insertions(+), 2 deletions(-) diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index 41c1fda0..bdd8552c 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -423,7 +423,7 @@ public final class CodexReviewStore { handle: any RuntimeLifecycleHandle, cause: String ) -> Bool { - if codexUpdateTask != nil { + if codexUpdateTask != nil || hasQueuedCodexUpdateRecovery { switch runtimeState { case .running(_, let runtime, _) where runtime.handle === handle: for job in jobs where job.isTerminal == false && queuedReviewStarts[job.id] == nil { @@ -469,7 +469,7 @@ public final class CodexReviewStore { sourceGeneration: ReviewRuntimeGeneration, cause: String ) -> ReviewRuntimeCleanupRecoveryAdmission { - if codexUpdateTask != nil, + if codexUpdateTask != nil || hasQueuedCodexUpdateRecovery, case .running(let generation, let runtime, _) = runtimeState, runtime.handle === sourceHandle, generation == sourceGeneration { return .suppressed(.codexUpdate) diff --git a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift index d806ab9e..b98567f6 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift @@ -69,6 +69,8 @@ extension CodexReviewStore { $0.isTerminal == false && queuedReviewStarts[$0.id] == nil }.map(\.id) for id in executingIDs { + // Another completed review can be removed while we await this batch. + guard job(id: id)?.isTerminal == false else { continue } _ = try await awaitReview(sessionID: nil, jobID: id) try Task.checkCancellation() } diff --git a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift index c8862578..25349d02 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift @@ -7,6 +7,34 @@ import CodexReviewTesting @Suite("review history store", .serialized) @MainActor struct CodexReviewStoreHistoryTests { + @Test func deferredUpdateContinuesWhenAnotherCompletedReviewIsDeleted() async throws { + let history = ReviewHistoryPersistenceProbe() + let backend = FakeCodexReviewBackend() + let runs = (0..<2).map { CodexReviewBackendModel.Review.Run(threadID: "thread-\($0)", turnID: "turn-\($0)") } + await backend.scriptReviewRuns(runs) + let store = makeStore(history: history, backend: backend) + await store.start() + for index in 0..<2 { + _ = try await store.startReview(sessionID: "client", request: .init(cwd: "/tmp/\(index)", target: .uncommittedChanges), waitTimeout: .zero) + } + try #require(await waitUntil { store.jobs.count == 2 && store.jobs.allSatisfy { $0.core.run.threadID != nil } }) + var installs = 0 + let update = Task { try await store.updateCodex(when: .afterCurrentReviews) { installs += 1 } } + try #require(await waitUntil { store.reviewTerminalWaiters.isEmpty == false }) + let awaitedID = try #require(store.reviewTerminalWaiters.keys.first) + let other = try #require(store.jobs.first { $0.id != awaitedID }) + let otherRun = try #require(runs.first { $0.threadID == other.core.run.threadID }) + let awaitedRun = try #require(runs.first { $0.threadID == store.job(id: awaitedID)?.core.run.threadID }) + await backend.yield(.completed(summary: "Done", result: "No findings."), for: otherRun) + try #require(await waitUntil { other.isTerminal && store.reviewWorkerTasks[other.id] == nil }) + await store.deleteReviewHistory(withIDs: [other.id]) + #expect(store.job(id: other.id) == nil) + await backend.yield(.completed(summary: "Done", result: "No findings."), for: awaitedRun) + try await update.value + #expect(installs == 1) + await store.stop() + } + @Test func updateKeepsRequestAcceptedDuringHistoryWriteInTheQueue() async throws { let acceptedEntered = AsyncGate() let acceptedRelease = AsyncGate() @@ -2168,3 +2196,14 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { Set(records.map { $0.started.id } + terminals.map(\.id)) } } + +@MainActor +private func waitUntil(condition: () -> Bool) async -> Bool { + let clock = ContinuousClock() + let deadline = clock.now + .seconds(2) + while condition() == false { + if clock.now >= deadline { return false } + await Task.yield() + } + return true +} diff --git a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift index 4ce2e3f4..2c9af38e 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift @@ -299,6 +299,15 @@ struct CodexReviewStoreUpdateTests { let id = try #require(queuedID) #expect(store.job(id: id)?.core.lifecycle.status == .queued) #expect(store.serverState == .failed("publication failed")) + backend.runOnNextRuntimePublication { + if let handle = backend.lastPreparedRuntimeHandle { + #expect(store.requestRuntimeFailure(handle: handle, cause: "retry publication failed")) + } + } + await store.restart() + await store.waitUntilStopped() + #expect(store.serverState == .failed("retry publication failed")) + #expect(store.job(id: id)?.core.lifecycle.status == .queued) await store.restart() #expect(store.serverState == .running) await store.stop() From 7854079eaa12d2fd4e2dc08500a43c38bb82bc35 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 16:22:17 +0900 Subject: [PATCH 3/8] Use one publication path for stopped-runtime updates --- .../CodexReview/Store/CodexReviewStore.swift | 13 +++++++---- .../Store/CodexReviewStoreUpdate.swift | 19 ++++++++++++--- .../CodexReviewStoreUpdateTests.swift | 23 +++++++++++++++++++ 3 files changed, 47 insertions(+), 8 deletions(-) diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index bdd8552c..f7d0de0b 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -306,6 +306,7 @@ public final class CodexReviewStore { } let sourceCloseJoin = operation.sourceCloseReceiptOwner?.sourceCloseJoin() await operation.task.value + if hasQueuedCodexUpdateRecovery { resumeReviewsAfterCodexUpdateIfPossible() } guard let replacement = operation.sourceCloseReceiptOwner, let sourceCloseJoin else { @@ -319,7 +320,7 @@ public final class CodexReviewStore { } } - package func admitRuntimeStart( + private func admitRuntimeStart( forceRestartIfNeeded: Bool ) -> RuntimeStartOperation? { let previousState = runtimeState @@ -425,7 +426,12 @@ public final class CodexReviewStore { ) -> Bool { if codexUpdateTask != nil || hasQueuedCodexUpdateRecovery { switch runtimeState { - case .running(_, let runtime, _) where runtime.handle === handle: + case .running(let generation, let runtime, let mcp) where runtime.handle === handle: + if codexUpdate != .waitingForReviews { + runtime.handle.closeAdmission() + unclosedCodexUpdateRuntime = runtime + runtimeState = .failed(generation: generation, retainedMCP: mcp, failureIncident: nil) + } for job in jobs where job.isTerminal == false && queuedReviewStarts[job.id] == nil { markReviewFailed(job, message: cause, terminal: .interrupted(.transport(message: cause))) } @@ -1511,9 +1517,6 @@ public final class CodexReviewStore { private func publishRuntime(serverURL: URL?) { storeWorkRegistry.openReviewAdmission() transitionToRunning(serverURL: serverURL) - if hasQueuedCodexUpdateRecovery, codexUpdateTask == nil { - resumeReviewsAfterCodexUpdateIfPossible() - } startAccountRateLimitAutoRefresh() } diff --git a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift index b98567f6..3d7cdeb0 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift @@ -99,11 +99,24 @@ extension CodexReviewStore { preservingQueuedReviews: true, install: installation ) - case .stopped, .failed: + case .stopped(let generation), .failed(let generation, nil, _): try await closeUnclosedCodexUpdateRuntimeIfNeeded() await installation() - guard let start = admitRuntimeStart(forceRestartIfNeeded: false) else { throw CancellationError() } - operation = start + let preparation = try await backend.mcpServerLifecycle.prepare() + let snapshot: MCPServerPublicationSnapshot + do { + snapshot = try await backend.mcpServerLifecycle.activate(preparation) + } catch { + var message = error.localizedDescription + do { try await backend.mcpServerLifecycle.stop() } + catch { message += "; MCP cleanup failed: \(error.localizedDescription)" } + throw CodexReviewAPI.Error.io(message) + } + operation = admitRuntimeReplacement( + sourceGeneration: generation, retiringRuntime: nil, + retainedMCP: RetainedMCPServer(serverURL: snapshot.serverURL), + preservingQueuedReviews: true + ) case .acquiring, .replacing, .tearingDown: throw CodexReviewAPI.Error.io("The Codex runtime changed while waiting to update.") } diff --git a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift index 2c9af38e..1529c59b 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift @@ -313,6 +313,29 @@ struct CodexReviewStoreUpdateTests { await store.stop() } + @Test func updateFromStoppedRuntimeRetainsQueueWhenPublicationFails() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + backend.runOnNextRuntimePublication { + if let handle = backend.lastPreparedRuntimeHandle { + #expect(store.requestRuntimeFailure(handle: handle, cause: "initial publication failed")) + } + } + var queuedID: String? + do { + try await store.updateCodex(when: .immediately) { + queuedID = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero).jobID + } + Issue.record("Expected initial publication failure") + } catch { #expect(error.localizedDescription.contains("initial publication failed")) } + let id = try #require(queuedID) + #expect(store.job(id: id)?.core.lifecycle.status == .queued) + #expect(store.serverState == .failed("initial publication failed")) + await store.start() + #expect(store.serverState == .running) + await store.stop() + } + private func request(_ name: String) -> CodexReviewAPI.Start.Request { .init(cwd: "/tmp/\(name)", target: .uncommittedChanges) } From 558a0f9d962ecc62be0fb55d3a66d6d1616cc46e Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 16:29:24 +0900 Subject: [PATCH 4/8] Cancel queued reviews independently of runtime teardown state --- .../CodexReview/Store/CodexReviewStore.swift | 7 +++ .../CodexReviewStoreUpdateTests.swift | 59 +++++++++++++++++++ 2 files changed, 66 insertions(+) diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index f7d0de0b..eeea2b21 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -620,6 +620,12 @@ public final class CodexReviewStore { let pendingHistoryStarts = requestHistoryStartCancellations( cancellation: intent.reviewCancellation ) + let queuedJobIDs = orderedJobs.filter { queuedReviewStarts[$0.id] != nil }.map(\.id) + for id in queuedJobIDs { + if let job = job(id: id) { + try? completeCancellationLocally(jobID: id, sessionID: job.sessionID, cancellation: intent.reviewCancellation) + } + } let generation = previousState.generation.successor() if case .replacing(let replacement, _) = previousState { replacement.finish(.superseded(runtimeTransitionPurpose(for: intent))) @@ -632,6 +638,7 @@ public final class CodexReviewStore { return } await self.waitForHistoryStarts(pendingHistoryStarts) + for id in queuedJobIDs { await self.waitForHistoryTerminalCommitIfNeeded(jobID: id) } await self.performRuntimeTeardown( previousState: previousState, generation: generation, diff --git a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift index 1529c59b..a8aae87d 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift @@ -336,6 +336,46 @@ struct CodexReviewStoreUpdateTests { await store.stop() } + @Test(arguments: [false, true]) + func stopCancelsQueuedCallsAfterUpdateMCPStartupFails(failActivation: Bool) async throws { + let mcp = FailingUpdateMCPServer(failActivation: failActivation) + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend(), mcpServerLifecycle: mcp) + let store = CodexReviewStore.makeTestingStore(backend: backend) + var requestTask: Task? + do { + try await store.updateCodex(when: .immediately) { + requestTask = Task { try await store.startReview(sessionID: "owner", request: request("queued")) } + try #require(await waitUntil { store.jobs.isEmpty == false }) + } + Issue.record("Expected MCP startup failure") + } catch { #expect(error.localizedDescription.contains("MCP startup failed")) } + let accepted = try #require(store.jobs.first) + #expect(accepted.core.lifecycle.status == .queued) + await store.stop() + #expect(accepted.core.lifecycle.status == .cancelled) + #expect(try await requestTask?.value.core.lifecycle.status == .cancelled) + #expect(store.queuedReviewStarts.isEmpty) + } + + @Test(arguments: [false, true]) + func stopOwnsQueueBeforeRuntimePublication(starting: Bool) async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + let release = AsyncGate() + backend.holdRuntimePreparation(with: release) + let start: Task? = starting ? Task { await store.start() } : nil + if starting { await backend.waitForRuntimePreparation() } + store.suspendReviewStarts() + let queued = Task { try await store.startReview(sessionID: "owner", request: request("queued")) } + try #require(await waitUntil { store.jobs.isEmpty == false }) + let stop = Task { await store.stop() } + try #require(await waitUntil { store.queuedReviewStarts.isEmpty }) + await release.open() + await start?.value + await stop.value + #expect(try await queued.value.core.lifecycle.status == .cancelled) + } + private func request(_ name: String) -> CodexReviewAPI.Start.Request { .init(cwd: "/tmp/\(name)", target: .uncommittedChanges) } @@ -356,3 +396,22 @@ private func waitUntil(condition: () -> Bool) async -> Bool { } return true } + +@MainActor +private final class FailingUpdateMCPServer: MCPServerLifecycleOwner { + private let failActivation: Bool + private let owner = TestingMCPServerLifecycleOwner() + + init(failActivation: Bool) { self.failActivation = failActivation } + + func prepare() async throws -> PreparedMCPServer { + if failActivation == false { throw CodexReviewAPI.Error.io("MCP startup failed") } + return try await owner.prepare() + } + + func activate(_ preparation: PreparedMCPServer) async throws -> MCPServerPublicationSnapshot { + throw CodexReviewAPI.Error.io("MCP startup failed") + } + + func stop() async throws { try await owner.stop() } +} From 701885799c7135ef961851392f16746eeecca061 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 16:35:38 +0900 Subject: [PATCH 5/8] Join Codex updates after startup history loading --- .../CodexReview/Store/CodexReviewStore.swift | 9 ++++--- .../CodexReviewStoreHistoryTests.swift | 26 +++++++++++++++++++ 2 files changed, 31 insertions(+), 4 deletions(-) diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index eeea2b21..6ee7d9c6 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -280,10 +280,6 @@ public final class CodexReviewStore { } public func start(forceRestartIfNeeded: Bool = false) async { - if let update = codexUpdateTask { - _ = try? await update.value - if case .running = runtimeState { return } - } guard applicationShutdownRequested == false else { return } @@ -299,6 +295,11 @@ public final class CodexReviewStore { catch { transitionToFailed(error.localizedDescription); return } default: break } + while let update = codexUpdateTask { + _ = try? await update.value + if case .running = runtimeState { return } + } + guard Task.isCancelled == false, applicationShutdownRequested == false else { return } guard let operation = admitRuntimeStart( forceRestartIfNeeded: forceRestartIfNeeded ) else { diff --git a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift index 25349d02..494b47d4 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift @@ -7,6 +7,32 @@ import CodexReviewTesting @Suite("review history store", .serialized) @MainActor struct CodexReviewStoreHistoryTests { + @Test func startupJoinsUpdateThatBeganWhileHistoryWasLoading() async throws { + let loadEntered = AsyncGate() + let loadRelease = AsyncGate() + let history = ReviewHistoryPersistenceProbe(loadEntered: loadEntered, loadRelease: loadRelease) + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend, historyPersistence: history) + let startup = Task { await store.start() } + await loadEntered.wait() + let installEntered = AsyncGate() + let installRelease = AsyncGate() + let update = Task { try await store.updateCodex(when: .immediately) { + await installEntered.open() + await installRelease.wait() + } } + await installEntered.wait() + await loadRelease.open() + try #require(await waitUntil { store.historyLoadWasApplied }) + #expect(backend.startRequests.isEmpty) + await installRelease.open() + try await update.value + await startup.value + #expect(backend.startRequests.count == 1) + #expect(store.serverState == .running) + await store.stop() + } + @Test func deferredUpdateContinuesWhenAnotherCompletedReviewIsDeleted() async throws { let history = ReviewHistoryPersistenceProbe() let backend = FakeCodexReviewBackend() From b90e17c7c3d8f87e5d8c9ee7e46da3219bc61742 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 16:56:57 +0900 Subject: [PATCH 6/8] Cancel update dependency waits and honor pending runtime stops --- .../CodexReview/Store/CodexReviewStore.swift | 3 + .../Store/CodexReviewStoreUpdate.swift | 32 +++++++- .../CodexReviewStoreUpdateTests.swift | 73 +++++++++++++++++++ 3 files changed, 105 insertions(+), 3 deletions(-) diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index 6ee7d9c6..20934405 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -91,6 +91,7 @@ public final class CodexReviewStore { package var codexUpdate: CodexUpdateState = .idle @_spi(ApplicationHostSupport) public var codexUpdateState: CodexUpdateState { codexUpdate } @ObservationIgnored package var codexUpdateRuntimeFailure: String? + @ObservationIgnored package var pendingRuntimeStopCount = 0 @ObservationIgnored package var codexUpdateTask: Task? @ObservationIgnored package var runtimeAccountOperations: [UUID: Task] = [:] @ObservationIgnored package var unclosedCodexUpdateRuntime: PreparedRuntime? @@ -402,6 +403,8 @@ public final class CodexReviewStore { } package func stop(intent: ReviewRuntimeTeardownIntent) async { + pendingRuntimeStopCount += 1 + defer { pendingRuntimeStopCount -= 1 } let update = codexUpdateTask if codexUpdate == .waitingForReviews { update?.cancel() diff --git a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift index 3d7cdeb0..a94c2ced 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreUpdate.swift @@ -26,7 +26,7 @@ extension CodexReviewStore { try await task.value return } - guard applicationShutdownRequested == false else { throw CancellationError() } + guard applicationShutdownRequested == false, pendingRuntimeStopCount == 0 else { throw CancellationError() } let previousAccountOperations = Array(runtimeAccountOperations.values) let previousRuntimeTask: Task? = switch runtimeState { case .acquiring(_, _, let task), .replacing(_, let task), .tearingDown(_, _, _, _, let task): task @@ -40,8 +40,9 @@ extension CodexReviewStore { if Task.isCancelled == false { resumeReviewsAfterCodexUpdateIfPossible() } } do { - for operation in previousAccountOperations { _ = try? await operation.value } - await previousRuntimeTask?.value + try await Self.waitForPrecedingRuntimeWork( + accountOperations: previousAccountOperations, runtimeTask: previousRuntimeTask + ) try await performCodexUpdate(when: timing, install: install) codexUpdate = .idle } catch { @@ -57,6 +58,30 @@ extension CodexReviewStore { try await task.value } + private static func waitForPrecedingRuntimeWork( + accountOperations: [Task], + runtimeTask: Task? + ) async throws { + // Cancelling the update abandons its join, not the predecessor's lifecycle ownership. + let (completion, continuation) = AsyncStream.makeStream() + let waiter = Task { + for operation in accountOperations { + _ = try? await operation.value + guard Task.isCancelled == false else { return } + } + await runtimeTask?.value + continuation.yield(()) + continuation.finish() + } + defer { + waiter.cancel() + continuation.finish() + } + var iterator = completion.makeAsyncIterator() + _ = await iterator.next() + try Task.checkCancellation() + } + private func performCodexUpdate( when timing: CodexUpdateTiming, install: @escaping @MainActor @Sendable () async throws -> Void @@ -132,6 +157,7 @@ extension CodexReviewStore { package func resumeReviewsAfterCodexUpdateIfPossible() { guard runtimeAccountOperations.isEmpty, applicationShutdownRequested == false, + pendingRuntimeStopCount == 0, case .running = runtimeState else { return } resumeReviewStarts() } diff --git a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift index a8aae87d..d0afc799 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift @@ -376,6 +376,79 @@ struct CodexReviewStoreUpdateTests { #expect(try await queued.value.core.lifecycle.status == .cancelled) } + @Test func stopCancelsUpdateJoinWhileEarlierAccountOperationRemainsPending() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let accountEntered = AsyncGate() + let accountRelease = AsyncGate() + let account = Task { try await store.performRuntimeAccountChange { _ in + await accountEntered.open() + await accountRelease.wait() + } } + await accountEntered.wait() + var installs = 0 + let update = Task { try await store.updateCodex(when: .afterCurrentReviews) { installs += 1 } } + #expect(await waitUntil { store.codexUpdateState == .waitingForReviews }) + var stopped = false + let stop = Task { await store.stop(); stopped = true } + let stoppedBeforeAccountFinished = await waitUntil { stopped } + await accountRelease.open() + try await account.value + await stop.value + do { try await update.value; Issue.record("Expected cancellation") } + catch { #expect(error is CancellationError) } + #expect(stoppedBeforeAccountFinished) + #expect(installs == 0) + } + + @Test func cancelledUpdateCanLeaveItsPrecedingRuntimeTaskWithItsOwner() async throws { + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let release = AsyncGate() + backend.holdRuntimePreparation(with: release) + let store = CodexReviewStore.makeTestingStore(backend: backend) + let startup = Task { await store.start() } + await backend.waitForRuntimePreparation() + var updateReturned = false + let update = Task { + defer { updateReturned = true } + try await store.updateCodex(when: .afterCurrentReviews) { Issue.record("Unexpected install") } + } + #expect(await waitUntil { store.codexUpdateTask != nil }) + store.codexUpdateTask?.cancel() + let returnedBeforeRuntimeStarted = await waitUntil { updateReturned } + await release.open() + await startup.value + do { try await update.value; Issue.record("Expected cancellation") } + catch { #expect(error is CancellationError) } + #expect(returnedBeforeRuntimeStarted) + await store.stop() + } + + @Test func pendingStopPreventsQueueDispatchAfterInstallation() async throws { + let reviews = FakeCodexReviewBackend() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let entered = AsyncGate() + let release = AsyncGate() + let update = Task { try await store.updateCodex(when: .immediately) { + await entered.open() + await release.wait() + } } + await entered.wait() + let queued = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero) + var stopRequested = false + let stop = Task { stopRequested = true; await store.stop() } + #expect(await waitUntil { stopRequested }) + await release.open() + try await update.value + await stop.value + #expect(store.job(id: queued.jobID)?.core.lifecycle.status == .cancelled) + #expect(store.job(id: queued.jobID)?.core.lifecycle.startedAt == nil) + #expect(await reviews.recordedCommands().contains { if case .startReview = $0 { true } else { false } } == false) + } + private func request(_ name: String) -> CodexReviewAPI.Start.Request { .init(cwd: "/tmp/\(name)", target: .uncommittedChanges) } From bb0c3daa5f5e5a34302c290f7ec09818954c1779 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 17:08:17 +0900 Subject: [PATCH 7/8] Confirm current process closure when retrying failed updates --- .../CodexReview/ReviewRuntimeLifecycle.swift | 5 ++ .../CodexReview/Store/CodexReviewStore.swift | 16 ++++- .../AppServerClient.swift | 4 ++ .../AppServerProcessTransport.swift | 30 +++++++++- Sources/CodexReviewAppServer/JSONRPC.swift | 5 ++ .../LiveCodexReviewStoreBackend.swift | 8 +++ Sources/CodexReviewTesting/TestSupport.swift | 1 + .../AppServerClientTests.swift | 4 ++ .../CodexReviewHostTests.swift | 60 +++++++++++++++++++ 9 files changed, 130 insertions(+), 3 deletions(-) diff --git a/Sources/CodexReview/ReviewRuntimeLifecycle.swift b/Sources/CodexReview/ReviewRuntimeLifecycle.swift index 0fc30fd0..6fa3cf0d 100644 --- a/Sources/CodexReview/ReviewRuntimeLifecycle.swift +++ b/Sources/CodexReview/ReviewRuntimeLifecycle.swift @@ -90,6 +90,11 @@ package protocol RuntimeLifecycleHandle: AnyObject, Sendable { func closeAdmission() func close(purpose: ReviewRuntimeTransitionPurpose) async throws func waitUntilClosed() async throws + func confirmClosed() async throws +} + +extension RuntimeLifecycleHandle { + package func confirmClosed() async throws { try await waitUntilClosed() } } @MainActor diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index 20934405..52ea3722 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -1316,6 +1316,16 @@ public final class CodexReviewStore { replacement.finishSourceClose(.closed) return } + if unclosedCodexUpdateRuntime?.handle === retiring.handle { + do { try await retiring.handle.confirmClosed() } + catch { + replacement.finishSourceClose(.failed(runtimeCloseFailure(from: error))) + throw error + } + unclosedCodexUpdateRuntime = nil + replacement.finishSourceClose(.closed) + return + } unclosedCodexUpdateRuntime = retiring await stopPublishedRuntimeSemantics(intent: .codexUpdate) if let failure = await closeRuntime(retiring, purpose: .restartSameAccount) { @@ -1328,8 +1338,10 @@ public final class CodexReviewStore { package func closeUnclosedCodexUpdateRuntimeIfNeeded() async throws { guard let runtime = unclosedCodexUpdateRuntime else { return } - if let failure = await closeRuntime(runtime, purpose: .restartSameAccount) { throw failure } - unclosedCodexUpdateRuntime = nil + try await runtime.handle.confirmClosed() + if unclosedCodexUpdateRuntime?.handle === runtime.handle { + unclosedCodexUpdateRuntime = nil + } } private func closePublishedRuntimeForReplacement( diff --git a/Sources/CodexReviewAppServer/AppServerClient.swift b/Sources/CodexReviewAppServer/AppServerClient.swift index dadf4b7d..618e95b4 100644 --- a/Sources/CodexReviewAppServer/AppServerClient.swift +++ b/Sources/CodexReviewAppServer/AppServerClient.swift @@ -290,6 +290,10 @@ package actor AppServerClient { try await transport.close() } + package func confirmClosed() async throws { + try await transport.confirmClosed() + } + private func allocateRequestID() -> Int { defer { nextRequestID += 1 } return nextRequestID diff --git a/Sources/CodexReviewAppServer/AppServerProcessTransport.swift b/Sources/CodexReviewAppServer/AppServerProcessTransport.swift index 4383d61f..d0352c7e 100644 --- a/Sources/CodexReviewAppServer/AppServerProcessTransport.swift +++ b/Sources/CodexReviewAppServer/AppServerProcessTransport.swift @@ -176,6 +176,15 @@ package actor AppServerProcessTransport: JSONRPC.Transport { try await closeTransport(terminateProcess: true, readerTask: nil) } + package func confirmClosed() async throws { + guard let closeTask else { + throw JSONRPC.Error.invalidMessage("Transport closure has not been requested.") + } + _ = await closeTask.result + await waitForReaderTasks(excluding: nil) + try process.confirmTerminated() + } + private func closeTransport( terminateProcess: Bool, readerTask: ReaderTask? @@ -745,6 +754,7 @@ private final class AppServerSpawnedProcess: @unchecked Sendable { private let processGroupID: pid_t private let stateLock = NSLock() private var didReap = false + private var terminationProcesses: Set? private init(processIdentifier: pid_t) { self.processIdentifier = processIdentifier @@ -850,7 +860,7 @@ private final class AppServerSpawnedProcess: @unchecked Sendable { killDuration: Duration = .seconds(1) ) async throws { let processIdentity = Self.processSnapshot(processIdentifier)?.identity - let trackedProcesses = descendantProcesses() + let trackedProcesses = beginTerminationTracking() guard isFullyTerminated(trackedProcesses: trackedProcesses) == false else { return } @@ -875,6 +885,24 @@ private final class AppServerSpawnedProcess: @unchecked Sendable { } } + private func beginTerminationTracking() -> Set { + let descendants = descendantProcesses() + stateLock.withLock { terminationProcesses = descendants } + return descendants + } + + func confirmTerminated() throws { + let (reaped, tracked) = stateLock.withLock { (didReap, terminationProcesses) } + // Reaping follows complete tree termination; never inspect a reused PID/group afterward. + if reaped { return } + guard let tracked, isFullyTerminated(trackedProcesses: tracked) else { + throw AppServerProcessTransportError.processDidNotTerminate( + processIdentifier, + liveProcessIdentifiers: liveProcessIdentifiers(trackedProcesses: tracked ?? []) + ) + } + } + private func signalProcessTree( _ signal: Int32, processIdentity: ProcessIdentity?, diff --git a/Sources/CodexReviewAppServer/JSONRPC.swift b/Sources/CodexReviewAppServer/JSONRPC.swift index f6990773..68f630e3 100644 --- a/Sources/CodexReviewAppServer/JSONRPC.swift +++ b/Sources/CodexReviewAppServer/JSONRPC.swift @@ -56,6 +56,7 @@ package enum JSONRPC { func notificationStream() async -> AsyncThrowingStream func notificationHighWatermark() async -> NotificationReceipt func close() async throws + func confirmClosed() async throws } package enum Error: Swift.Error, Equatable, Sendable, LocalizedError { @@ -142,3 +143,7 @@ package struct AnyEncodable: Encodable { package struct EmptyResponse: Codable, Equatable, Sendable { package init() {} } + +extension JSONRPC.Transport { + package func confirmClosed() async throws { try await close() } +} diff --git a/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift b/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift index 6894140d..c815bc77 100644 --- a/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift +++ b/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift @@ -4184,6 +4184,14 @@ private final class LiveRuntimeLifecycleHandle: RuntimeLifecycleHandle { } try await closeTask.value.get() } + func confirmClosed() async throws { + guard let closeTask else { throw CancellationError() } + switch await closeTask.value { + case .success: return + case .failure: try await client.confirmClosed() + } + } + } @MainActor diff --git a/Sources/CodexReviewTesting/TestSupport.swift b/Sources/CodexReviewTesting/TestSupport.swift index 47bfe89b..b52ffc93 100644 --- a/Sources/CodexReviewTesting/TestSupport.swift +++ b/Sources/CodexReviewTesting/TestSupport.swift @@ -932,6 +932,7 @@ package final class TestingRuntimeLifecycleHandle: RuntimeLifecycleHandle { await closeGate?.waitIgnoringCancellation() closeGate = nil guard didClose == false else { + if let closeFailure { throw closeFailure } return } didClose = true diff --git a/Tests/CodexReviewAppServerTests/AppServerClientTests.swift b/Tests/CodexReviewAppServerTests/AppServerClientTests.swift index c1e34e6d..e463ce2c 100644 --- a/Tests/CodexReviewAppServerTests/AppServerClientTests.swift +++ b/Tests/CodexReviewAppServerTests/AppServerClientTests.swift @@ -685,9 +685,13 @@ struct AppServerClientTests { } ) + await #expect(throws: JSONRPC.Error.self) { + try await transport.confirmClosed() + } await #expect(throws: failure) { try await transport.close() } + try await transport.confirmClosed() await #expect(throws: failure) { try await transport.close() } diff --git a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift index b4d687c1..e66ae78f 100644 --- a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift +++ b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift @@ -46,6 +46,45 @@ private actor HostCloseFailureTransport: JSONRPC.Transport { @Suite("host composition") @MainActor struct CodexReviewHostTests { + @Test func updateRecoveryConfirmsPhysicalClosureWithoutReplayingCachedFailure() async throws { + let homeURL = try temporaryHome() + let original = FakeJSONRPCTransport() + let replacement = FakeJSONRPCTransport() + try await enqueueRuntimeStartResponses(original) + try await enqueueRuntimeStartResponses(replacement) + let delayedClose = DelayedCloseConfirmationTransport(base: original) + let server = ControlledMCPHTTPServer(endpoint: try #require(URL(string: "http://127.0.0.1:19439/mcp"))) + var transports: [any JSONRPC.Transport] = [delayedClose, replacement] + let store = CodexReviewStore.makeLiveStoreForTesting( + environment: ["HOME": homeURL.path], + webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, + mcpHTTPServerFactory: { _, _ in server }, + mcpHTTPServerBindChecker: { _ in }, + transportFactory: { _ in transports.removeFirst() } + ) + await store.start() + store.suspendReviewStarts() + let queued = try await store.startReview(sessionID: "owner", request: .init(cwd: "/tmp/project", target: .uncommittedChanges), waitTimeout: .zero) + var installations = 0 + do { + try await store.updateCodex(when: .immediately) { installations += 1 } + Issue.record("Expected close failure") + } catch { #expect(error.localizedDescription.contains("Injected host close failure")) } + await store.restart() + #expect(transports.count == 1) + #expect(store.job(id: queued.jobID)?.core.lifecycle.status == .queued) + #expect(server.stopCallCount == 0) + _ = try await store.cancelReview(jobID: queued.jobID, sessionID: "owner") + await delayedClose.finishProcessTermination() + await store.restart() + #expect(store.serverState == .running) + #expect(transports.isEmpty) + #expect(installations == 0) + #expect(await delayedClose.closeCalls == 1) + #expect(server.startCallCount == 1) + await store.stop() + } + @Test func liveRuntimeUpdatePreservesMCPAndDispatchesQueueOnNewTransport() async throws { let homeURL = try temporaryHome() let old = FakeJSONRPCTransport() @@ -7193,3 +7232,24 @@ private actor CompletionFlag { completed } } + +private actor DelayedCloseConfirmationTransport: JSONRPC.Transport { + let base: FakeJSONRPCTransport + private var processTerminated = false + private(set) var closeCalls = 0 + + init(base: FakeJSONRPCTransport) { self.base = base } + func send(_ request: JSONRPC.Request) async throws -> Data { try await base.send(request) } + func notify(_ notification: JSONRPC.Notification) async throws { try await base.notify(notification) } + func notificationStream() async -> AsyncThrowingStream { await base.notificationStream() } + func notificationHighWatermark() async -> JSONRPC.NotificationReceipt { await base.notificationHighWatermark() } + func close() async throws { + closeCalls += 1 + try await base.close() + throw HostCloseFailure.injected + } + func confirmClosed() async throws { + guard processTerminated else { throw HostCloseFailure.injected } + } + func finishProcessTermination() { processTerminated = true } +} From 8ff02dbd04c65b90b016d96c2d47c8625e8df6d1 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 26 Sep 2026 17:17:17 +0900 Subject: [PATCH 8/8] Request runtime closure before confirming recovery --- .../CodexReview/Store/CodexReviewStore.swift | 2 + .../CodexReviewStoreUpdateTests.swift | 38 +++++++++++++++++++ 2 files changed, 40 insertions(+) diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index 52ea3722..75497d3e 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -1317,6 +1317,7 @@ public final class CodexReviewStore { return } if unclosedCodexUpdateRuntime?.handle === retiring.handle { + _ = await closeRuntime(retiring, purpose: .restartSameAccount) do { try await retiring.handle.confirmClosed() } catch { replacement.finishSourceClose(.failed(runtimeCloseFailure(from: error))) @@ -1338,6 +1339,7 @@ public final class CodexReviewStore { package func closeUnclosedCodexUpdateRuntimeIfNeeded() async throws { guard let runtime = unclosedCodexUpdateRuntime else { return } + _ = await closeRuntime(runtime, purpose: .restartSameAccount) try await runtime.handle.confirmClosed() if unclosedCodexUpdateRuntime?.handle === runtime.handle { unclosedCodexUpdateRuntime = nil diff --git a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift index d0afc799..0b4e344a 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreUpdateTests.swift @@ -449,6 +449,44 @@ struct CodexReviewStoreUpdateTests { #expect(await reviews.recordedCommands().contains { if case .startReview = $0 { true } else { false } } == false) } + @Test func recoveryRequestsCloseForFailureAfterRuntimePublication() async throws { + let reviews = FakeCodexReviewBackend() + let backend = TestingCodexReviewStoreBackend(reviewBackend: reviews) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let accountEntered = AsyncGate() + let accountRelease = AsyncGate() + var account: Task? + var queuedID: String? + do { + try await store.updateCodex(when: .immediately) { + queuedID = try await store.startReview(sessionID: "owner", request: request("queued"), waitTimeout: .zero).jobID + account = Task { try await store.performRuntimeAccountChange { _ in + await accountEntered.open() + await accountRelease.wait() + } } + #expect(await waitUntil { store.runtimeAccountOperations.isEmpty == false }) + throw UpdateTestError.install + } + Issue.record("Expected install failure") + } catch { #expect(error.localizedDescription.contains("install failed")) } + await accountEntered.wait() + let runtime = try #require(backend.lastPreparedRuntimeHandle) + #expect(runtime.closePurposes.isEmpty) + #expect(store.requestRuntimeFailure(handle: runtime, cause: "late runtime failure")) + await accountRelease.open() + try await account?.value + let id = try #require(queuedID) + #expect(store.job(id: id)?.core.lifecycle.status == .queued) + await store.restart() + #expect(store.serverState == .running) + #expect(runtime.closePurposes == [.restartSameAccount]) + try await reviews.waitForStartReview(timeout: .seconds(2)) + await reviews.yield(.completed(summary: "Done", result: "No findings.")) + #expect(try await store.awaitReview(sessionID: "owner", jobID: id).core.lifecycle.status == .succeeded) + await store.stop() + } + private func request(_ name: String) -> CodexReviewAPI.Start.Request { .init(cwd: "/tmp/\(name)", target: .uncommittedChanges) }