From e74c6e3b3c3b5676b9257fc497f6883a396c7775 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 22 Aug 2026 06:09:37 +0900 Subject: [PATCH 1/3] refactor(runtime): own prepared runtime generations --- .../CodexReview/ReviewRuntimeLifecycle.swift | 201 +++++++ .../CodexReview/Store/CodexReviewStore.swift | 369 ++++++++++++- .../Store/CodexReviewStoreBackend.swift | 29 +- .../PreviewCodexReviewStoreBackend.swift | 9 +- Sources/CodexReviewHost/CodexReviewHost.swift | 107 +++- .../LiveCodexReviewStoreBackend.swift | 506 +++++++++++++----- Sources/CodexReviewTesting/TestSupport.swift | 222 +++++++- .../CodexReviewHostTests.swift | 304 ++++++++++- .../CodexReviewStoreLifecycleTests.swift | 212 ++++++++ Tests/ReviewUITests/ReviewUITests.swift | 32 +- 10 files changed, 1794 insertions(+), 197 deletions(-) create mode 100644 Sources/CodexReview/ReviewRuntimeLifecycle.swift create mode 100644 Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift diff --git a/Sources/CodexReview/ReviewRuntimeLifecycle.swift b/Sources/CodexReview/ReviewRuntimeLifecycle.swift new file mode 100644 index 00000000..cc2a49b3 --- /dev/null +++ b/Sources/CodexReview/ReviewRuntimeLifecycle.swift @@ -0,0 +1,201 @@ +import Foundation + +package struct ReviewRuntimeGeneration: Hashable, Sendable { + package let rawValue: UInt64 + + package init(rawValue: UInt64) { + self.rawValue = rawValue + } + + package func successor() -> Self { + .init(rawValue: rawValue &+ 1) + } +} +package enum ReviewRuntimeTransitionPurpose: Equatable, Sendable { + case start + case restartSameAccount + case stop +} +package struct RuntimePublicationSnapshot: Sendable { + package let authentication: CodexReviewBackendModel.Auth.Snapshot + package let settings: CodexReviewSettings.Snapshot + + package init( + authentication: CodexReviewBackendModel.Auth.Snapshot, + settings: CodexReviewSettings.Snapshot + ) { + self.authentication = authentication + self.settings = settings + } +} +@MainActor +package func applyRuntimeAuthenticationSnapshot( + _ snapshot: CodexReviewBackendModel.Auth.Snapshot, + to auth: CodexReviewAuthModel +) { + let observedAccounts = snapshot.accounts.compactMap { account -> CodexAccount? in + let label = account.label.trimmingCharacters(in: .whitespacesAndNewlines) + let accountKey = CodexAccount.normalizedEmail(account.id.rawValue) + guard label.isEmpty == false, accountKey.isEmpty == false else { + return nil + } + return CodexAccount( + accountKey: accountKey, + email: label, + planType: account.planType, + kind: account.kind, + capabilities: account.capabilities + ) + } + let activeAccountKey = snapshot.activeAccountID.map { + CodexAccount.normalizedEmail($0.rawValue) + } + var accounts = auth.persistedAccounts + for observedAccount in observedAccounts { + if let index = accounts.firstIndex(where: { + $0.accountKey == observedAccount.accountKey + }) { + accounts[index].updateEmail(observedAccount.email) + accounts[index].updateKind( + observedAccount.kind, + capabilities: observedAccount.capabilities + ) + accounts[index].updatePlanType(observedAccount.planType) + } else { + accounts.insert(observedAccount, at: 0) + } + } + auth.applyPersistedAccountStates( + accounts.map(savedAccountPayload(from:)), + activeAccountKey: activeAccountKey + ) + auth.selectPersistedAccount(activeAccountKey) + auth.updatePhase(.signedOut) +} +@MainActor +package protocol RuntimeLifecycleHandle: AnyObject, Sendable { + func activate() async throws + func closeAdmission() async + func close(purpose: ReviewRuntimeTransitionPurpose) async throws + func waitUntilClosed() async throws +} +package struct PreparedRuntime: Sendable { + package let snapshot: RuntimePublicationSnapshot + package let handle: any RuntimeLifecycleHandle + + package init( + snapshot: RuntimePublicationSnapshot, + handle: any RuntimeLifecycleHandle + ) { + self.snapshot = snapshot + self.handle = handle + } +} +package struct MCPServerGeneration: Hashable, Sendable { + package let rawValue: UInt64 + + package init(rawValue: UInt64) { + self.rawValue = rawValue + } +} +package struct PreparedMCPServer: Sendable { + package let generation: MCPServerGeneration + + package init(generation: MCPServerGeneration) { + self.generation = generation + } +} +package struct MCPServerPublicationSnapshot: Sendable { + package let serverURL: URL? + + package init(serverURL: URL?) { + self.serverURL = serverURL + } +} +package enum ReviewStoreRuntimeState { + case stopped(ReviewRuntimeGeneration) + case acquiring( + generation: ReviewRuntimeGeneration, + task: Task + ) + case running( + generation: ReviewRuntimeGeneration, + runtime: PreparedRuntime, + mcpGeneration: MCPServerGeneration + ) + case transitioning( + generation: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose, + task: Task + ) + case failed( + generation: ReviewRuntimeGeneration, + retainedMCPGeneration: MCPServerGeneration, + retainedMCPServerURL: URL? + ) + + package var generation: ReviewRuntimeGeneration { + switch self { + case .stopped(let generation), + .acquiring(let generation, _), + .running(let generation, _, _), + .transitioning(let generation, _, _), + .failed(let generation, _, _): + generation + } + } +} + +@MainActor +package protocol MCPServerLifecycleOwner: Sendable { + func prepare() async throws -> PreparedMCPServer + func activate( + _ generation: MCPServerGeneration + ) async throws -> MCPServerPublicationSnapshot + func stop() async throws +} + +@MainActor +package final class NoMCPServerLifecycleOwner: MCPServerLifecycleOwner { + private enum State { + case stopped + case prepared(MCPServerGeneration) + case running(MCPServerGeneration) + } + + private var serverURL: URL? + private var nextGeneration: UInt64 = 0 + private var state: State = .stopped + + package init(serverURL: URL? = nil) { + self.serverURL = serverURL + } + + package func updateServerURL(_ serverURL: URL?) { + self.serverURL = serverURL + } + + package func prepare() async throws -> PreparedMCPServer { + guard case .stopped = state else { + throw CancellationError() + } + nextGeneration &+= 1 + let generation = MCPServerGeneration(rawValue: nextGeneration) + state = .prepared(generation) + return .init(generation: generation) + } + + package func activate( + _ generation: MCPServerGeneration + ) async throws -> MCPServerPublicationSnapshot { + guard case .prepared(generation) = state else { + throw CancellationError() + } + state = .running(generation) + return .init(serverURL: serverURL) + } + + package func stop() async throws { + state = .stopped + } +} diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index d7ab5001..d1fa372a 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -37,6 +37,9 @@ public final class CodexReviewStore { @ObservationIgnored package var reviewTerminalWaiters: [String: [ReviewTerminalWaiter]] = [:] @ObservationIgnored package var closedSessions: Set = [] @ObservationIgnored package var accountRateLimitAutoRefreshDriver: CodexReviewStoreRateLimitAutoRefreshDriver? + @ObservationIgnored package var runtimeState: ReviewStoreRuntimeState = .stopped( + .init(rawValue: 0) + ) package init( backend: any CodexReviewStoreBackend = PreviewCodexReviewStoreBackend(), @@ -78,6 +81,12 @@ public final class CodexReviewStore { isolated deinit { accountRateLimitAutoRefreshDriver?.cancel() + switch runtimeState { + case .acquiring(_, let task), .transitioning(_, _, let task): + task.cancel() + case .stopped, .running, .failed: + break + } for task in reviewWorkerTasks.values { task.cancel() } @@ -126,25 +135,307 @@ public final class CodexReviewStore { } public func start(forceRestartIfNeeded: Bool = false) async { - switch serverState { - case .stopped, .failed: - break - case .starting: + switch runtimeState { + case .acquiring: + return + case .transitioning(_, _, let task): + await task.value return case .running where forceRestartIfNeeded == false: return - case .running: + case .running(let generation, let runtime, let mcpGeneration): + await beginRuntimeReplacement( + previousGeneration: generation, + previousRuntime: runtime, + retainedMCPGeneration: mcpGeneration, + retainedMCPServerURL: serverURL + ) + return + case .failed(let generation, let retainedMCPGeneration, let retainedMCPServerURL): + await beginRuntimeReplacement( + previousGeneration: generation, + previousRuntime: nil, + retainedMCPGeneration: retainedMCPGeneration, + retainedMCPServerURL: retainedMCPServerURL + ) + return + case .stopped: break } + guard case .stopped(let previousGeneration) = runtimeState else { + return + } + let generation = previousGeneration.successor() serverState = .starting serverURL = nil writeDiagnosticsIfNeeded() - await backend.start(store: self, forceRestartIfNeeded: forceRestartIfNeeded) - await settingsService.refreshIfRunning(serverState: serverState) - startAccountRateLimitAutoRefresh() + let task = Task { @MainActor [weak self] in + await self?.performRuntimeAcquisition(generation: generation) + } + runtimeState = .acquiring(generation: generation, task: task) + await task.value } public func stop() async { + let previousState = runtimeState + switch previousState { + case .stopped: + transitionToStopped() + return + case .transitioning(_, .stop, let task): + await task.value + return + case .acquiring, .running, .transitioning, .failed: + break + } + let invalidatedGeneration = previousState.generation.successor() + let task = Task { @MainActor [weak self] in + await self?.performRuntimeStop( + previousState: previousState, + invalidatedGeneration: invalidatedGeneration + ) + } + runtimeState = .transitioning( + generation: invalidatedGeneration, + purpose: .stop, + task: task + ) + await task.value + } + + public func restart() async { + switch runtimeState { + case .acquiring, .transitioning: + await stop() + case .stopped, .running, .failed: + break + } + await start(forceRestartIfNeeded: true) + } + + public func waitUntilStopped() async { + if case .transitioning(_, .stop, let task) = runtimeState { + await task.value + } + await backend.waitUntilStopped() + } + + private func performRuntimeAcquisition( + generation: ReviewRuntimeGeneration + ) async { + guard isCurrentAcquisition(generation) else { + return + } + var preparedRuntime: PreparedRuntime? + do { + let preparedMCPServer = try await backend.mcpServerLifecycle.prepare() + guard isCurrentAcquisition(generation) else { + return + } + + let runtime = try await backend.prepareRuntime( + generation: generation, + purpose: .start + ) + preparedRuntime = runtime + guard isCurrentAcquisition(generation) else { + await closeRuntime(runtime, purpose: .start) + return + } + + try await runtime.handle.activate() + guard isCurrentAcquisition(generation) else { + await closeRuntime(runtime, purpose: .start) + return + } + + let mcpSnapshot = try await backend.mcpServerLifecycle.activate( + preparedMCPServer.generation + ) + guard isCurrentAcquisition(generation) else { + await closeRuntime(runtime, purpose: .start) + return + } + + try backend.commitRuntimePublication( + runtime.snapshot, + handle: runtime.handle, + auth: auth + ) + runtimeState = .running( + generation: generation, + runtime: runtime, + mcpGeneration: preparedMCPServer.generation + ) + publishRuntime(runtime, serverURL: mcpSnapshot.serverURL) + await backend.waitForRuntimePublication(handle: runtime.handle) + } catch { + if let preparedRuntime { + await closeRuntime(preparedRuntime, purpose: .start) + } + guard isCurrentAcquisition(generation) else { + return + } + await stopMCPServer() + guard isCurrentAcquisition(generation) else { + return + } + runtimeState = .stopped(generation) + transitionToFailed(error.localizedDescription) + } + } + + package func failRuntime( + handle: any RuntimeLifecycleHandle, + message: String + ) async { + guard case .running(let generation, let runtime, _) = runtimeState, + runtime.handle === handle + else { + return + } + let stoppedGeneration = generation.successor() + await stop() + guard case .stopped(stoppedGeneration) = runtimeState else { + return + } + transitionToFailed(message) + } + + private func beginRuntimeReplacement( + previousGeneration: ReviewRuntimeGeneration, + previousRuntime: PreparedRuntime?, + retainedMCPGeneration: MCPServerGeneration, + retainedMCPServerURL: URL? + ) async { + let generation = previousGeneration.successor() + serverState = .starting + writeDiagnosticsIfNeeded() + let task = Task { @MainActor [weak self] in + await self?.performRuntimeReplacement( + generation: generation, + previousRuntime: previousRuntime, + retainedMCPGeneration: retainedMCPGeneration, + retainedMCPServerURL: retainedMCPServerURL + ) + } + runtimeState = .transitioning( + generation: generation, + purpose: .restartSameAccount, + task: task + ) + await task.value + } + + private func performRuntimeReplacement( + generation: ReviewRuntimeGeneration, + previousRuntime: PreparedRuntime?, + retainedMCPGeneration: MCPServerGeneration, + retainedMCPServerURL: URL? + ) async { + var preparedRuntime: PreparedRuntime? + if let previousRuntime { + await previousRuntime.handle.closeAdmission() + await stopPublishedRuntimeSemantics() + await closeRuntime( + previousRuntime, + purpose: .restartSameAccount, + admissionAlreadyClosed: true + ) + } + guard isCurrentTransition(generation, purpose: .restartSameAccount) else { + return + } + + do { + let runtime = try await backend.prepareRuntime( + generation: generation, + purpose: .restartSameAccount + ) + preparedRuntime = runtime + guard isCurrentTransition(generation, purpose: .restartSameAccount) else { + await closeRuntime(runtime, purpose: .restartSameAccount) + return + } + + try await runtime.handle.activate() + guard isCurrentTransition(generation, purpose: .restartSameAccount) else { + await closeRuntime(runtime, purpose: .restartSameAccount) + return + } + + try backend.commitRuntimePublication( + runtime.snapshot, + handle: runtime.handle, + auth: auth + ) + runtimeState = .running( + generation: generation, + runtime: runtime, + mcpGeneration: retainedMCPGeneration + ) + publishRuntime(runtime, serverURL: retainedMCPServerURL) + await backend.waitForRuntimePublication(handle: runtime.handle) + } catch { + if let preparedRuntime { + await closeRuntime(preparedRuntime, purpose: .restartSameAccount) + } + guard isCurrentTransition(generation, purpose: .restartSameAccount) else { + return + } + runtimeState = .failed( + generation: generation, + retainedMCPGeneration: retainedMCPGeneration, + retainedMCPServerURL: retainedMCPServerURL + ) + serverURL = retainedMCPServerURL + serverState = .failed(error.localizedDescription) + writeDiagnosticsIfNeeded() + } + } + + private func performRuntimeStop( + previousState: ReviewStoreRuntimeState, + invalidatedGeneration: ReviewRuntimeGeneration + ) async { + switch previousState { + case .acquiring(_, let task): + task.cancel() + await stopMCPServer() + await task.value + + case .running(_, let runtime, _): + await runtime.handle.closeAdmission() + await stopPublishedRuntimeSemantics() + await stopMCPServer() + await closeRuntime(runtime, purpose: .stop, admissionAlreadyClosed: true) + + case .transitioning(_, _, let task): + task.cancel() + await stopMCPServer() + await task.value + + case .failed: + await stopMCPServer() + + case .stopped: + break + } + + guard case .transitioning( + let currentGeneration, + .stop, + _ + ) = runtimeState, + currentGeneration == invalidatedGeneration + else { + return + } + runtimeState = .stopped(invalidatedGeneration) + transitionToStopped() + } + + private func stopPublishedRuntimeSemantics() async { let locallyCancelledJobIDs: [String] if backend.handlesActiveReviewStopCleanup { locallyCancelledJobIDs = [] @@ -156,16 +447,66 @@ public final class CodexReviewStore { cancelAndDetachReviewWorkersForRuntimeStop( jobIDs: Array(Set(locallyCancelledJobIDs + remainingLocallyCancelledJobIDs)) ) - transitionToStopped() } - public func restart() async { - await stop() - await start(forceRestartIfNeeded: true) + private func closeRuntime( + _ runtime: PreparedRuntime, + purpose: ReviewRuntimeTransitionPurpose, + admissionAlreadyClosed: Bool = false + ) async { + if admissionAlreadyClosed == false { + await runtime.handle.closeAdmission() + } + do { + try await runtime.handle.close(purpose: purpose) + } catch { + writeDiagnosticsIfNeeded() + } + do { + try await runtime.handle.waitUntilClosed() + } catch { + writeDiagnosticsIfNeeded() + } } - public func waitUntilStopped() async { - await backend.waitUntilStopped() + private func stopMCPServer() async { + do { + try await backend.mcpServerLifecycle.stop() + } catch { + writeDiagnosticsIfNeeded() + } + } + + private func isCurrentAcquisition( + _ generation: ReviewRuntimeGeneration + ) -> Bool { + guard case .acquiring(let currentGeneration, _) = runtimeState else { + return false + } + return currentGeneration == generation + } + + private func isCurrentTransition( + _ generation: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose + ) -> Bool { + guard case .transitioning( + let currentGeneration, + let currentPurpose, + _ + ) = runtimeState else { + return false + } + return currentGeneration == generation && currentPurpose == purpose + } + + private func publishRuntime( + _ runtime: PreparedRuntime, + serverURL: URL? + ) { + settings.apply(snapshot: runtime.snapshot.settings) + transitionToRunning(serverURL: serverURL) + startAccountRateLimitAutoRefresh() } public func refreshAuthentication() async { diff --git a/Sources/CodexReview/Store/CodexReviewStoreBackend.swift b/Sources/CodexReview/Store/CodexReviewStoreBackend.swift index 7caad062..27b180f2 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreBackend.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreBackend.swift @@ -24,13 +24,25 @@ package struct CodexReviewStoreSeed { } @MainActor -package protocol CodexReviewStoreBackend: CodexReviewSettingsBackend { +package protocol CodexReviewStoreBackend: CodexReviewSettingsBackend, Sendable { var seed: CodexReviewStoreSeed { get } var isActive: Bool { get } var handlesActiveReviewStopCleanup: Bool { get } + var mcpServerLifecycle: any MCPServerLifecycleOwner { get } func attachStore(_ store: CodexReviewStore) - func start(store: CodexReviewStore, forceRestartIfNeeded: Bool) async + func prepareRuntime( + generation: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose + ) async throws -> PreparedRuntime + func commitRuntimePublication( + _ snapshot: RuntimePublicationSnapshot, + handle: any RuntimeLifecycleHandle, + auth: CodexReviewAuthModel + ) throws + func waitForRuntimePublication( + handle: any RuntimeLifecycleHandle + ) async func stop(store: CodexReviewStore) async func waitUntilStopped() async func refreshAuth(auth: CodexReviewAuthModel) async @@ -65,6 +77,19 @@ extension CodexReviewStoreBackend { false } + package func commitRuntimePublication( + _ snapshot: RuntimePublicationSnapshot, + handle _: any RuntimeLifecycleHandle, + auth: CodexReviewAuthModel + ) throws { + applyRuntimeAuthenticationSnapshot(snapshot.authentication, to: auth) + } + + package func waitForRuntimePublication( + handle _: any RuntimeLifecycleHandle + ) async {} + + package func startReview( _ request: CodexReviewBackendModel.Review.Start ) async throws -> BackendReviewAttempt { diff --git a/Sources/CodexReview/Store/PreviewCodexReviewStoreBackend.swift b/Sources/CodexReview/Store/PreviewCodexReviewStoreBackend.swift index 3f059a4b..fd60169f 100644 --- a/Sources/CodexReview/Store/PreviewCodexReviewStoreBackend.swift +++ b/Sources/CodexReview/Store/PreviewCodexReviewStoreBackend.swift @@ -5,6 +5,7 @@ package class PreviewCodexReviewStoreBackend: CodexReviewStoreBackend { package let seed: CodexReviewStoreSeed package var isActive = false package var currentSettingsSnapshot: CodexReviewSettings.Snapshot + package let mcpServerLifecycle: any MCPServerLifecycleOwner = NoMCPServerLifecycleOwner() package init(seed: CodexReviewStoreSeed = .init()) { self.seed = seed @@ -17,9 +18,11 @@ package class PreviewCodexReviewStoreBackend: CodexReviewStoreBackend { package func attachStore(_: CodexReviewStore) {} - package func start(store: CodexReviewStore, forceRestartIfNeeded _: Bool) async { - isActive = true - store.transitionToFailed(Self.previewUnavailableMessage) + package func prepareRuntime( + generation _: ReviewRuntimeGeneration, + purpose _: ReviewRuntimeTransitionPurpose + ) async throws -> PreparedRuntime { + throw CodexReviewAPI.Error.io(Self.previewUnavailableMessage) } package func stop(store _: CodexReviewStore) async { diff --git a/Sources/CodexReviewHost/CodexReviewHost.swift b/Sources/CodexReviewHost/CodexReviewHost.swift index 1af2c8a8..71f8d372 100644 --- a/Sources/CodexReviewHost/CodexReviewHost.swift +++ b/Sources/CodexReviewHost/CodexReviewHost.swift @@ -7,6 +7,7 @@ import CodexReviewMCPServer package final class CodexReviewHost { package let store: CodexReviewStore package let mcpServer: CodexReviewMCPServer + private let directBackend: DirectCodexReviewStoreBackend private let shutdown: @Sendable () async throws -> Void private var endpoint: URL? @@ -17,10 +18,15 @@ package final class CodexReviewHost { endpoint: URL? = nil, shutdown: @escaping @Sendable () async throws -> Void = {} ) { - self.shutdown = shutdown self.endpoint = endpoint + let directBackend = DirectCodexReviewStoreBackend( + backend: backend, + endpoint: endpoint + ) + self.directBackend = directBackend + self.shutdown = shutdown let store = CodexReviewStore( - backend: DirectCodexReviewStoreBackend(backend: backend), + backend: directBackend, clock: clock, idGenerator: idGenerator ) @@ -44,11 +50,15 @@ package final class CodexReviewHost { } package func start(endpoint: URL? = nil) async { + let endpointChanged = endpoint != nil && endpoint != self.endpoint if let endpoint { self.endpoint = endpoint } - store.transitionToRunning(serverURL: self.endpoint) - await store.refreshSettings() + directBackend.updateEndpoint(self.endpoint) + if endpointChanged { + await store.stop() + } + await store.start() } package func stop() async throws { @@ -60,7 +70,9 @@ package final class CodexReviewHost { @MainActor private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { let seed = CodexReviewStoreSeed() + let mcpServerLifecycle: any MCPServerLifecycleOwner private let backend: any CodexReviewBackend + private let mcpLifecycleOwner: NoMCPServerLifecycleOwner private var currentSettingsSnapshot = CodexReviewSettings.Snapshot() private var loginChallenge: CodexReviewBackendModel.Login.Challenge? private var active = false @@ -73,14 +85,39 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { currentSettingsSnapshot } - init(backend: any CodexReviewBackend) { + init( + backend: any CodexReviewBackend, + endpoint: URL? + ) { self.backend = backend + let mcpLifecycleOwner = NoMCPServerLifecycleOwner(serverURL: endpoint) + self.mcpLifecycleOwner = mcpLifecycleOwner + self.mcpServerLifecycle = mcpLifecycleOwner } func attachStore(_: CodexReviewStore) {} - func start(store _: CodexReviewStore, forceRestartIfNeeded _: Bool) async { - active = true + func updateEndpoint(_ endpoint: URL?) { + mcpLifecycleOwner.updateServerURL(endpoint) + } + + func prepareRuntime( + generation _: ReviewRuntimeGeneration, + purpose _: ReviewRuntimeTransitionPurpose + ) async throws -> PreparedRuntime { + let settings = try await Self.monitorSettings(from: backend.readSettings()) + let authentication = try await backend.readAuth() + return PreparedRuntime( + snapshot: .init( + authentication: authentication, + settings: settings + ), + handle: DirectRuntimeLifecycleHandle( + onActivate: { [weak self] in self?.active = true }, + onCloseAdmission: { [weak self] in self?.active = false }, + onClose: { [weak self] in self?.active = false } + ) + ) } func stop(store _: CodexReviewStore) async { @@ -90,6 +127,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { func waitUntilStopped() async {} func refreshSettings() async throws -> CodexReviewSettings.Snapshot { + guard active else { return currentSettingsSnapshot } currentSettingsSnapshot = try await Self.monitorSettings(from: backend.readSettings()) return currentSettingsSnapshot } @@ -101,6 +139,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { serviceTier: CodexReviewSettings.ServiceTier?, persistServiceTier: Bool ) async throws { + guard active else { return } var change = CodexReviewBackendModel.Settings.Change( model: model, updatesModel: true @@ -119,6 +158,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { func updateSettingsReasoningEffort( _ reasoningEffort: CodexReviewSettings.ReasoningEffort? ) async throws { + guard active else { return } currentSettingsSnapshot = try await Self.monitorSettings( from: backend.applySettings(.init( reasoningEffort: reasoningEffort?.rawValue, @@ -130,6 +170,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { func updateSettingsServiceTier( _ serviceTier: CodexReviewSettings.ServiceTier? ) async throws { + guard active else { return } currentSettingsSnapshot = try await Self.monitorSettings( from: backend.applySettings(.init( serviceTier: serviceTier?.rawValue, @@ -139,6 +180,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { } func refreshAuth(auth: CodexReviewAuthModel) async { + guard active else { return } do { Self.applyAuthSnapshot(try await backend.readAuth(), to: auth) } catch { @@ -147,6 +189,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { } func signIn(auth: CodexReviewAuthModel) async { + guard active else { return } do { let challenge = try await backend.startLogin(.init()) loginChallenge = challenge @@ -162,6 +205,7 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { } func addAccount(auth: CodexReviewAuthModel) async { + guard active else { return } await signIn(auth: auth) } @@ -237,7 +281,10 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { _ request: CodexReviewBackendModel.Review.Start, admission: ReviewStartAdmission ) async throws -> BackendReviewAttempt { - try await backend.startReview(request, admission: admission) + guard active else { + throw CodexReviewAPI.Error.io("Review runtime is not running.") + } + return try await backend.startReview(request, admission: admission) } func interruptReview( @@ -313,6 +360,50 @@ private final class DirectCodexReviewStoreBackend: CodexReviewStoreBackend { } } +@MainActor +private final class DirectRuntimeLifecycleHandle: RuntimeLifecycleHandle { + private let onActivate: @MainActor @Sendable () -> Void + private let onCloseAdmission: @MainActor @Sendable () -> Void + private let onClose: @MainActor @Sendable () -> Void + private var isActivated = false + private var didClose = false + + init( + onActivate: @escaping @MainActor @Sendable () -> Void, + onCloseAdmission: @escaping @MainActor @Sendable () -> Void, + onClose: @escaping @MainActor @Sendable () -> Void + ) { + self.onActivate = onActivate + self.onCloseAdmission = onCloseAdmission + self.onClose = onClose + } + + func activate() async throws { + guard isActivated == false, didClose == false else { + throw CancellationError() + } + isActivated = true + onActivate() + } + + func closeAdmission() async { + onCloseAdmission() + } + + func close(purpose _: ReviewRuntimeTransitionPurpose) async throws { + if didClose == false { + didClose = true + onClose() + } + } + + func waitUntilClosed() async throws { + guard didClose else { + throw CancellationError() + } + } +} + extension CodexReviewBackendModel.Login.Challenge { func signInDetail(nativeAuthentication: Bool) -> String { if let userCode = userCode?.trimmingCharacters(in: .whitespacesAndNewlines).nilIfEmpty { diff --git a/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift b/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift index ed0eec4e..63124274 100644 --- a/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift +++ b/Sources/CodexReviewHost/LiveCodexReviewStoreBackend.swift @@ -187,7 +187,7 @@ public extension CodexReviewStore { } @MainActor -private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { +private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend, MCPServerLifecycleOwner { typealias MCPHTTPServerFactory = @MainActor @Sendable ( CodexReviewStore, CodexReviewMCPHTTPServer.Configuration @@ -197,6 +197,8 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { private var client: AppServerClient? private var appServerBackend: AppServerCodexReviewBackend? + private var activeRuntimeHandle: LiveRuntimeLifecycleHandle? + private var acceptsRuntimeRequests = false private var mcpHTTPServer: (any CodexReviewMCPHTTPServing)? private var loginChallenge: CodexReviewBackendModel.Login.Challenge? private var loginBackend: AppServerCodexReviewBackend? @@ -220,6 +222,10 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { private let appServerRuntimeFactory: AppServerRuntimeFactory private let shutdownCleanupTimeout: Duration private weak var attachedStore: CodexReviewStore? + private var nextMCPGeneration: UInt64 = 0 + private var preparingMCPGeneration: MCPServerGeneration? + private var preparedMCPGeneration: MCPServerGeneration? + private var runningMCPGeneration: MCPServerGeneration? init( environment: [String: String] = ProcessInfo.processInfo.environment, @@ -269,7 +275,11 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } var isActive: Bool { - client != nil + acceptsRuntimeRequests + } + + var mcpServerLifecycle: any MCPServerLifecycleOwner { + self } var handlesActiveReviewStopCleanup: Bool { @@ -395,62 +405,209 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { attachedStore = store } - func start(store: CodexReviewStore, forceRestartIfNeeded: Bool) async { - logger.info("Starting review runtime; forceRestartIfNeeded=\(forceRestartIfNeeded, privacy: .public)") - if appServerBackend != nil, forceRestartIfNeeded == false { - logger.info("Review runtime already has an app-server backend") - store.transitionToRunning(serverURL: await mcpHTTPServer?.url) - return - } - if forceRestartIfNeeded { - await stop(store: store) + func prepare() async throws -> PreparedMCPServer { + guard mcpHTTPServer == nil, + preparingMCPGeneration == nil, + preparedMCPGeneration == nil, + runningMCPGeneration == nil + else { + throw CancellationError() } - - var startedClient: AppServerClient? - var startedBackend: AppServerCodexReviewBackend? - var startedHTTPServer: (any CodexReviewMCPHTTPServing)? + nextMCPGeneration &+= 1 + let generation = MCPServerGeneration(rawValue: nextMCPGeneration) + preparingMCPGeneration = generation do { if mcpHTTPServerFactory != nil { try await mcpHTTPServerBindChecker(mcpHTTPServerConfiguration) } - let runtime = try await appServerRuntimeFactory(codexHomeURL) - let client = runtime.client - let backend = runtime.backend - startedClient = client - startedBackend = backend - self.client = client - self.appServerBackend = backend - observeAuthNotifications(client: client, backend: backend, store: store) + guard preparingMCPGeneration == generation else { + throw CancellationError() + } if let mcpHTTPServerFactory { - let mcpHTTPServer = mcpHTTPServerFactory(store, mcpHTTPServerConfiguration) - startedHTTPServer = mcpHTTPServer - try await mcpHTTPServer.start() - self.mcpHTTPServer = mcpHTTPServer + guard let attachedStore else { + throw CancellationError() + } + mcpHTTPServer = mcpHTTPServerFactory( + attachedStore, + mcpHTTPServerConfiguration + ) + } + preparingMCPGeneration = nil + preparedMCPGeneration = generation + return .init(generation: generation) + } catch { + if preparingMCPGeneration == generation { + preparingMCPGeneration = nil } - store.transitionToRunning(serverURL: await self.mcpHTTPServer?.url) - let authSnapshot = try await backend.readAuth() - applyAuthSnapshot(authSnapshot, to: store.auth) - await refreshSelectedAccountRateLimits(auth: store.auth) - logger.info("Review runtime started") + throw CodexReviewAPI.Error.io(await runtimeStartupFailureMessage(for: error)) + } + } + + func activate( + _ generation: MCPServerGeneration + ) async throws -> MCPServerPublicationSnapshot { + guard preparedMCPGeneration == generation else { + throw CancellationError() + } + guard let server = mcpHTTPServer else { + preparedMCPGeneration = nil + runningMCPGeneration = generation + return .init(serverURL: nil) + } + do { + try await server.start() + guard mcpHTTPServer === server, + preparedMCPGeneration == generation + else { + throw CancellationError() + } + let serverURL = await server.url + guard mcpHTTPServer === server, + preparedMCPGeneration == generation + else { + throw CancellationError() + } + preparedMCPGeneration = nil + runningMCPGeneration = generation + return .init(serverURL: serverURL) } catch { - let failureMessage = await runtimeStartupFailureMessage(for: error) - logger.error("Review runtime failed to start: \(failureMessage, privacy: .public)") - self.client = nil - self.appServerBackend = nil - self.mcpHTTPServer = nil - authNotificationTask?.cancel() - authNotificationTask = nil - await stopMCPHTTPServer( - startedHTTPServer, - context: "runtime startup cleanup" + if mcpHTTPServer === server { + mcpHTTPServer = nil + preparedMCPGeneration = nil + try? await server.stop() + } + throw error + } + } + + func stop() async throws { + preparingMCPGeneration = nil + preparedMCPGeneration = nil + runningMCPGeneration = nil + guard let server = mcpHTTPServer else { + return + } + mcpHTTPServer = nil + try await server.stop() + } + + func prepareRuntime( + generation _: ReviewRuntimeGeneration, + purpose _: ReviewRuntimeTransitionPurpose + ) async throws -> PreparedRuntime { + logger.info("Preparing review runtime") + let runtime = try await appServerRuntimeFactory(codexHomeURL) + do { + let authNotificationStream = await runtime.client.notificationStream() + let authentication = try await runtime.backend.readAuth() + let settings = try await Self.monitorSettings(from: runtime.backend.readSettings()) + let handle = LiveRuntimeLifecycleHandle( + owner: self, + client: runtime.client, + backend: runtime.backend, + authNotificationStream: authNotificationStream, + snapshot: .init( + authentication: authentication, + settings: settings + ) ) + logger.info("Review runtime prepared") + return .init(snapshot: handle.snapshot, handle: handle) + } catch { await closeAppServerRuntime( - backend: startedBackend, - fallbackClient: startedClient, - context: "runtime startup cleanup" + backend: runtime.backend, + fallbackClient: runtime.client, + context: "runtime preparation cleanup" + ) + throw error + } + } + + func commitRuntimePublication( + _ snapshot: RuntimePublicationSnapshot, + handle: any RuntimeLifecycleHandle, + auth: CodexReviewAuthModel + ) throws { + guard let handle = handle as? LiveRuntimeLifecycleHandle, + activeRuntimeHandle === handle, + let attachedStore + else { + throw CancellationError() + } + applyRuntimeAuthenticationSnapshot(snapshot.authentication, to: auth) + if let activeAccountID = snapshot.authentication.activeAccountID?.rawValue, + let account = auth.persistedAccounts.first(where: { + $0.accountKey == CodexAccount.normalizedEmail(activeAccountID) + }) + { + try? CodexReviewAccountRegistry.saveAccounts( + auth.persistedAccounts, + activeAccountKey: account.accountKey, + codexHomeURL: codexHomeURL + ) + try? CodexReviewAccountRegistry.saveSharedAuth( + for: account, + codexHomeURL: codexHomeURL ) - store.transitionToFailed(failureMessage) } + acceptsRuntimeRequests = true + observeAuthNotifications( + stream: handle.authNotificationStream, + backend: handle.backend, + store: attachedStore + ) + handle.initialRateLimitTask = Task { @MainActor [weak self, weak auth] in + guard let self, let auth else { + return + } + await self.refreshSelectedAccountRateLimits(auth: auth) + } + } + + func waitForRuntimePublication( + handle: any RuntimeLifecycleHandle + ) async { + guard let handle = handle as? LiveRuntimeLifecycleHandle else { + return + } + await handle.initialRateLimitTask?.value + } + + func activateRuntime(_ handle: LiveRuntimeLifecycleHandle) throws { + guard activeRuntimeHandle == nil else { + throw CancellationError() + } + guard attachedStore != nil else { + throw CancellationError() + } + activeRuntimeHandle = handle + acceptsRuntimeRequests = false + client = handle.client + appServerBackend = handle.backend + settingsSnapshot = handle.snapshot.settings + } + + func closeRuntimeAdmission(_ handle: LiveRuntimeLifecycleHandle) { + guard activeRuntimeHandle === handle else { + return + } + acceptsRuntimeRequests = false + } + + func deactivateRuntime( + _ handle: LiveRuntimeLifecycleHandle + ) -> Task? { + guard activeRuntimeHandle === handle else { + return nil + } + activeRuntimeHandle = nil + acceptsRuntimeRequests = false + client = nil + appServerBackend = nil + let task = authNotificationTask + authNotificationTask = nil + task?.cancel() + return task } private func runtimeStartupFailureMessage(for error: Error) async -> String { @@ -502,15 +659,12 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } func stop(store: CodexReviewStore) async { - let client = client let appServerBackend = appServerBackend - let mcpHTTPServer = mcpHTTPServer - let hasRuntimeState = client != nil || appServerBackend != nil || mcpHTTPServer != nil let loginCleanup = takeLoginRuntimeForCleanup() - guard hasRuntimeState || loginCleanup.isEmpty == false else { + guard appServerBackend != nil || loginCleanup.isEmpty == false else { return } - logger.info("Stopping review runtime") + logger.info("Stopping review runtime semantic work") if let appServerBackend { let reason = ReviewCancellation.system(message: "Review runtime stopped.") await cancelActiveReviewsForRuntimeTeardown( @@ -520,25 +674,14 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { timeoutWarning: "Timed out cleaning active reviews before stopping runtime" ) } - self.client = nil - self.mcpHTTPServer = nil - authNotificationTask?.cancel() - authNotificationTask = nil - await stopMCPHTTPServer(mcpHTTPServer, context: "runtime stop") - self.appServerBackend = nil await cleanupLoginRuntime(loginCleanup) - await closeAppServerRuntime( - backend: appServerBackend, - fallbackClient: client, - context: "runtime stop" - ) - logger.info("Review runtime stopped") + logger.info("Review runtime semantic work stopped") } func waitUntilStopped() async {} func refreshSettings() async throws -> CodexReviewSettings.Snapshot { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { return settingsSnapshot } settingsSnapshot = try await Self.monitorSettings(from: appServerBackend.readSettings()) @@ -552,7 +695,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { serviceTier: CodexReviewSettings.ServiceTier?, persistServiceTier: Bool ) async throws { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { return } var change = CodexReviewBackendModel.Settings.Change( @@ -573,7 +716,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { func updateSettingsReasoningEffort( _ reasoningEffort: CodexReviewSettings.ReasoningEffort? ) async throws { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { return } settingsSnapshot = try await Self.monitorSettings( @@ -587,7 +730,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { func updateSettingsServiceTier( _ serviceTier: CodexReviewSettings.ServiceTier? ) async throws { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { return } settingsSnapshot = try await Self.monitorSettings( @@ -599,14 +742,29 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } func refreshAuth(auth: CodexReviewAuthModel) async { + guard acceptsRuntimeRequests, + let expectedRuntimeHandle = activeRuntimeHandle + else { + return + } do { guard let appServerBackend else { auth.updatePhase(.signedOut) return } let snapshot = try await appServerBackend.readAuth() + guard activeRuntimeHandle === expectedRuntimeHandle, + acceptsRuntimeRequests + else { + return + } applyAuthSnapshot(snapshot, to: auth) } catch { + guard activeRuntimeHandle === expectedRuntimeHandle, + acceptsRuntimeRequests + else { + return + } auth.updatePhase(.failed(message: error.localizedDescription)) } } @@ -678,8 +836,8 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { return } await attachedStore.closeActiveReviewSessions(reason: .system(message: "Account switched.")) - await stop(store: attachedStore) - await start(store: attachedStore, forceRestartIfNeeded: true) + await attachedStore.stop() + await attachedStore.start(forceRestartIfNeeded: true) } func removeAccount(auth: CodexReviewAuthModel, accountKey: String) async throws { @@ -715,8 +873,8 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { return } await attachedStore.closeActiveReviewSessions(reason: .system(message: "Account removed.")) - await stop(store: attachedStore) - await start(store: attachedStore, forceRestartIfNeeded: true) + await attachedStore.stop() + await attachedStore.start(forceRestartIfNeeded: true) } } @@ -771,8 +929,8 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { auth.selectPersistedAccount(nil) auth.applyPersistedAccountStates(remaining.map(savedAccountPayload(from:)), activeAccountKey: nil) if shouldRecycleRuntime, let attachedStore { - await stop(store: attachedStore) - await start(store: attachedStore, forceRestartIfNeeded: true) + await attachedStore.stop() + await attachedStore.start(forceRestartIfNeeded: true) } } @@ -788,6 +946,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } private func startLogin(auth: CodexReviewAuthModel, activation: LoginActivation) async { + let expectedRuntimeHandle = activeRuntimeHandle var isolatedLoginClient: AppServerClient? var isolatedLoginCodexHomeURL: URL? do { @@ -797,6 +956,10 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { let loginClient = runtime.usesPrimaryRuntime ? nil : runtime.client isolatedLoginClient = loginClient isolatedLoginCodexHomeURL = loginCodexHomeURL + guard isCurrentRuntime(expectedRuntimeHandle) else { + await closeIsolatedLoginRuntime(client: loginClient, codexHomeURL: loginCodexHomeURL) + return + } guard runtime.usesPrimaryRuntime || self.appServerBackend != nil else { logger.error("Cannot start login because review runtime is not running") updateAuthenticationFailure( @@ -811,6 +974,14 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { let challenge = try await appServerBackend.startLogin(.init( nativeWebAuthenticationCallbackScheme: nativeAuthenticationConfiguration?.callbackScheme )) + guard let expectedRuntimeHandle, + activeRuntimeHandle === expectedRuntimeHandle, + acceptsRuntimeRequests + else { + try? await appServerBackend.cancelLogin(challenge) + await closeIsolatedLoginRuntime(client: loginClient, codexHomeURL: loginCodexHomeURL) + return + } loginChallenge = challenge loginBackend = appServerBackend self.loginClient = loginClient @@ -857,6 +1028,15 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { nativeAuthenticationConfiguration.browserSessionPolicy, nativeAuthenticationConfiguration.presentationAnchorProvider ) + guard activeRuntimeHandle === expectedRuntimeHandle, + acceptsRuntimeRequests, + loginChallenge?.id == challenge.id + else { + await session.cancel() + try? await appServerBackend.cancelLogin(challenge) + await closeIsolatedLoginRuntime(client: loginClient, codexHomeURL: loginCodexHomeURL) + return + } activeAuthenticationSession = session authenticationTask = Task { @MainActor [weak self, weak auth] in guard let self, let auth else { @@ -897,6 +1077,15 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } } + private func isCurrentRuntime( + _ expectedRuntimeHandle: LiveRuntimeLifecycleHandle? + ) -> Bool { + guard let expectedRuntimeHandle else { + return activeRuntimeHandle == nil + } + return activeRuntimeHandle === expectedRuntimeHandle && acceptsRuntimeRequests + } + private func monitorAuthenticationSession( challenge: CodexReviewBackendModel.Login.Challenge, session: any CodexReviewNativeAuthentication.WebSession, @@ -993,7 +1182,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { private func loginRuntime(for activation: LoginActivation) async throws -> LoginRuntime { switch activation { case .activateAuthenticatedAccount: - guard let client, let appServerBackend else { + guard acceptsRuntimeRequests, let client, let appServerBackend else { throw CodexReviewAPI.Error.io("Review runtime is not running.") } return .init( @@ -1050,14 +1239,14 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { _ request: CodexReviewBackendModel.Review.Start, admission: ReviewStartAdmission ) async throws -> BackendReviewAttempt { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { throw CodexReviewAPI.Error.io("Review runtime is not running.") } return try await appServerBackend.startReview(request, admission: admission) } func interruptReview(_ run: CodexReviewBackendModel.Review.Run, reason: CodexReviewBackendModel.CancellationReason) async throws { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { throw CodexReviewAPI.Error.io("Review runtime is not running.") } try await appServerBackend.interruptReview(run, reason: reason) @@ -1067,7 +1256,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { _ run: CodexReviewBackendModel.Review.Run, reason: CodexReviewBackendModel.CancellationReason ) async throws -> CodexReviewBackendModel.Review.RecoveryToken { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { throw CodexReviewAPI.Error.io("Review runtime is not running.") } return try await appServerBackend.beginReviewRecovery(run, reason: reason) @@ -1077,7 +1266,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { _ token: CodexReviewBackendModel.Review.RecoveryToken, request: CodexReviewBackendModel.Review.Start ) async throws -> BackendReviewAttempt { - guard let appServerBackend else { + guard acceptsRuntimeRequests, let appServerBackend else { throw CodexReviewAPI.Error.io("Review runtime is not running.") } return try await appServerBackend.resumeReviewRecovery(token, request: request) @@ -1154,7 +1343,7 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } private func observeAuthNotifications( - client: AppServerClient, + stream: AsyncThrowingStream, backend: AppServerCodexReviewBackend, store: CodexReviewStore ) { @@ -1163,7 +1352,6 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { guard let self, let store else { return } - let stream = await client.notificationStream() do { for try await notification in stream { await self.handleAuthNotification( @@ -1184,37 +1372,16 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { _ error: any Error, store: CodexReviewStore ) async { - let loginCleanup = takeLoginRuntimeForCleanup() - guard client != nil || appServerBackend != nil || mcpHTTPServer != nil || loginCleanup.isEmpty == false else { + guard let activeRuntimeHandle else { return } let message = "Review runtime stopped unexpectedly: \(error.localizedDescription)" - if let appServerBackend { - let reason = ReviewCancellation.system(message: message) - await cancelActiveReviewsForRuntimeTeardown( - store: store, - appServerBackend: appServerBackend, - reason: reason, - timeoutWarning: "Timed out cleaning active reviews after runtime failure" - ) - } - let failedClient = client - let failedBackend = appServerBackend - let failedMCPHTTPServer = mcpHTTPServer - client = nil - appServerBackend = nil - mcpHTTPServer = nil + // The current notification Task is the caller. Remove it from the + // handle's join set before Store teardown so close cannot await itself. authNotificationTask = nil - store.transitionToFailed(message) - await stopMCPHTTPServer( - failedMCPHTTPServer, - context: "failed runtime cleanup" - ) - await cleanupLoginRuntime(loginCleanup) - await closeAppServerRuntime( - backend: failedBackend, - fallbackClient: failedClient, - context: "notification stream failure" + await store.failRuntime( + handle: activeRuntimeHandle, + message: message ) } @@ -1438,7 +1605,23 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { await refreshSavedAccountRateLimits(for: account) return } - let didRefresh = await refreshRateLimits(for: account, using: appServerBackend, source: "active-runtime") + guard acceptsRuntimeRequests else { + return + } + guard let expectedRuntimeHandle = activeRuntimeHandle else { + return + } + let didRefresh = await refreshRateLimits( + for: account, + using: appServerBackend, + source: "active-runtime", + expectedRuntimeHandle: expectedRuntimeHandle + ) + guard activeRuntimeHandle === expectedRuntimeHandle, + acceptsRuntimeRequests + else { + return + } if didRefresh { persistRefreshedSharedAuth( from: codexHomeURL, @@ -1494,7 +1677,8 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { private func refreshRateLimits( for account: CodexAccount, using backend: AppServerCodexReviewBackend?, - source: String + source: String, + expectedRuntimeHandle: LiveRuntimeLifecycleHandle? = nil ) async -> Bool { do { guard let backend else { @@ -1507,6 +1691,13 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { ) } let response = try await backend.readRateLimits() + if let expectedRuntimeHandle { + guard activeRuntimeHandle === expectedRuntimeHandle, + acceptsRuntimeRequests + else { + return false + } + } applyRateLimits( windows: response.codexRateLimitWindows, planType: response.codexPlanType, @@ -1518,6 +1709,11 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { ) return true } catch { + if let expectedRuntimeHandle, + (activeRuntimeHandle !== expectedRuntimeHandle || acceptsRuntimeRequests == false) + { + return false + } recordRateLimitRefreshFailure(error, account: account) try? CodexReviewAccountRegistry.updateCachedRateLimits( from: account, @@ -1591,22 +1787,6 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } } - private func stopMCPHTTPServer( - _ server: (any CodexReviewMCPHTTPServing)?, - context: String - ) async { - guard let server else { - return - } - do { - try await server.stop() - } catch { - logger.error( - "Failed to stop MCP HTTP server during \(context, privacy: .public): \(error.localizedDescription, privacy: .public)" - ) - } - } - private func closeClientRecordingFailure( _ client: AppServerClient?, context: String @@ -1735,6 +1915,82 @@ private final class LiveCodexReviewStoreBackend: CodexReviewStoreBackend { } } +@MainActor +private final class LiveRuntimeLifecycleHandle: RuntimeLifecycleHandle { + fileprivate let client: AppServerClient + fileprivate let backend: AppServerCodexReviewBackend + fileprivate let authNotificationStream: AsyncThrowingStream + fileprivate let snapshot: RuntimePublicationSnapshot + fileprivate var initialRateLimitTask: Task? + + private weak var owner: LiveCodexReviewStoreBackend? + private var isActivated = false + private var closeTask: Task, Never>? + + init( + owner: LiveCodexReviewStoreBackend, + client: AppServerClient, + backend: AppServerCodexReviewBackend, + authNotificationStream: AsyncThrowingStream, + snapshot: RuntimePublicationSnapshot + ) { + self.owner = owner + self.client = client + self.backend = backend + self.authNotificationStream = authNotificationStream + self.snapshot = snapshot + } + + func activate() async throws { + guard isActivated == false, closeTask == nil, let owner else { + throw CancellationError() + } + try owner.activateRuntime(self) + isActivated = true + } + + func closeAdmission() async { + owner?.closeRuntimeAdmission(self) + } + + func close(purpose _: ReviewRuntimeTransitionPurpose) async throws { + let task: Task, Never> + if let closeTask { + task = closeTask + } else { + let authObservationTask = owner?.deactivateRuntime(self) + let initialRateLimitTask = initialRateLimitTask + self.initialRateLimitTask = nil + let lifecycle = backend.runtimeOwnerLifecycleHandle + let created = Task, Never> { @MainActor in + authObservationTask?.cancel() + initialRateLimitTask?.cancel() + let result: Result + do { + await lifecycle.closeAdmission() + try await lifecycle.closeAndWait() + result = .success(()) + } catch { + result = .failure(error) + } + await authObservationTask?.value + await initialRateLimitTask?.value + return result + } + closeTask = created + task = created + } + try await task.value.get() + } + + func waitUntilClosed() async throws { + guard let closeTask else { + throw CancellationError() + } + try await closeTask.value.get() + } +} + @MainActor private struct AppServerRuntime: Sendable { var client: AppServerClient diff --git a/Sources/CodexReviewTesting/TestSupport.swift b/Sources/CodexReviewTesting/TestSupport.swift index 6b728207..44b35513 100644 --- a/Sources/CodexReviewTesting/TestSupport.swift +++ b/Sources/CodexReviewTesting/TestSupport.swift @@ -805,6 +805,157 @@ package struct StoreJobSnapshot: Sendable { package var cancellationRequested: Bool } +@MainActor +package final class TestingRuntimeLifecycleHandle: RuntimeLifecycleHandle { + package private(set) var activateCallCount = 0 + package private(set) var closeAdmissionCallCount = 0 + package private(set) var closePurposes: [ReviewRuntimeTransitionPurpose] = [] + package private(set) var waitUntilClosedCallCount = 0 + + private let onActivate: @MainActor @Sendable () -> Void + private let onClose: @MainActor @Sendable () -> Void + private var didClose = false + + package init( + onActivate: @escaping @MainActor @Sendable () -> Void = {}, + onClose: @escaping @MainActor @Sendable () -> Void = {} + ) { + self.onActivate = onActivate + self.onClose = onClose + } + + package func activate() async throws { + activateCallCount += 1 + onActivate() + } + + package func closeAdmission() async { + closeAdmissionCallCount += 1 + } + + package func close(purpose: ReviewRuntimeTransitionPurpose) async throws { + closePurposes.append(purpose) + guard didClose == false else { + return + } + didClose = true + onClose() + } + + package func waitUntilClosed() async throws { + waitUntilClosedCallCount += 1 + guard didClose else { + throw CancellationError() + } + } +} + +@MainActor +package final class TestingMCPServerLifecycleOwner: MCPServerLifecycleOwner { + package private(set) var preparedGenerations: [MCPServerGeneration] = [] + package private(set) var activatedGenerations: [MCPServerGeneration] = [] + package private(set) var stopCallCount = 0 + + private let serverURL: URL? + private var nextGeneration: UInt64 = 0 + private var preparedGeneration: MCPServerGeneration? + private var isRunning = false + private var preparationGate: AsyncGate? + private var preparationStartedGate = AsyncGate() + private var preparationCancellationGate = AsyncGate() + private var stopGate: AsyncGate? + private var stopStartedGate = AsyncGate() + private var secondStopStartedGate = AsyncGate() + private var stopCallerCount = 0 + private var isStopping = false + + package init(serverURL: URL? = nil) { + self.serverURL = serverURL + } + + package func holdPreparation(with gate: AsyncGate) { + preparationGate = gate + preparationStartedGate = AsyncGate() + preparationCancellationGate = AsyncGate() + } + + package func waitForPreparation() async { + await preparationStartedGate.wait() + } + + package func waitForPreparationCancellation() async { + await preparationCancellationGate.wait() + } + + package func holdStop(with gate: AsyncGate) { + stopGate = gate + stopStartedGate = AsyncGate() + secondStopStartedGate = AsyncGate() + stopCallerCount = 0 + } + + package func waitForStop() async { + await stopStartedGate.wait() + } + + package func waitForSecondStop() async { + await secondStopStartedGate.wait() + } + + package func prepare() async throws -> PreparedMCPServer { + nextGeneration &+= 1 + let generation = MCPServerGeneration(rawValue: nextGeneration) + preparedGeneration = generation + preparedGenerations.append(generation) + await preparationStartedGate.open() + if let preparationGate { + let cancellationGate = preparationCancellationGate + await withTaskCancellationHandler { + await preparationGate.waitIgnoringCancellation() + } onCancel: { + Task { await cancellationGate.open() } + } + self.preparationGate = nil + } + try Task.checkCancellation() + return .init(generation: generation) + } + + package func activate( + _ generation: MCPServerGeneration + ) async throws -> MCPServerPublicationSnapshot { + guard preparedGeneration == generation else { + throw CancellationError() + } + preparedGeneration = nil + isRunning = true + activatedGenerations.append(generation) + return .init(serverURL: serverURL) + } + + package func stop() async throws { + guard preparedGeneration != nil || isRunning else { + return + } + stopCallerCount += 1 + if stopCallerCount == 2 { + await secondStopStartedGate.open() + } + if isStopping { + await stopGate?.waitIgnoringCancellation() + return + } + isStopping = true + await stopStartedGate.open() + await stopGate?.waitIgnoringCancellation() + stopGate = nil + preparedGeneration = nil + isRunning = false + stopCallCount += 1 + isStopping = false + } +} + @MainActor package final class TestingCodexReviewStoreBackend: CodexReviewStoreBackend { package let reviewBackend: FakeCodexReviewBackend @@ -812,14 +963,22 @@ package final class TestingCodexReviewStoreBackend: CodexReviewStoreBackend { package var currentSettingsSnapshot: CodexReviewSettings.Snapshot package private(set) var isActive = false package private(set) var startRequests: [Bool] = [] + package let mcpServerLifecycle: any MCPServerLifecycleOwner + package private(set) var lastPreparedRuntimeHandle: TestingRuntimeLifecycleHandle? + private var runtimePreparationGate: AsyncGate? + private var runtimePreparationFailureMessage: String? + private var runtimePreparationStartedGate = AsyncGate() + private var runtimePreparationCancellationGate = AsyncGate() package init( reviewBackend: FakeCodexReviewBackend, - seed: CodexReviewStoreSeed = .init() + seed: CodexReviewStoreSeed = .init(), + mcpServerLifecycle: (any MCPServerLifecycleOwner)? = nil ) { self.reviewBackend = reviewBackend self.seed = seed self.currentSettingsSnapshot = seed.initialSettingsSnapshot + self.mcpServerLifecycle = mcpServerLifecycle ?? NoMCPServerLifecycleOwner() } package var initialSettingsSnapshot: CodexReviewSettings.Snapshot { @@ -828,10 +987,55 @@ package final class TestingCodexReviewStoreBackend: CodexReviewStoreBackend { package func attachStore(_: CodexReviewStore) {} - package func start(store: CodexReviewStore, forceRestartIfNeeded: Bool) async { - startRequests.append(forceRestartIfNeeded) - isActive = true - store.transitionToRunning(serverURL: nil) + package func holdRuntimePreparation(with gate: AsyncGate) { + runtimePreparationGate = gate + runtimePreparationStartedGate = AsyncGate() + runtimePreparationCancellationGate = AsyncGate() + } + + package func waitForRuntimePreparation() async { + await runtimePreparationStartedGate.wait() + } + + package func waitForRuntimePreparationCancellation() async { + await runtimePreparationCancellationGate.wait() + } + + package func failNextRuntimePreparation(message: String) { + runtimePreparationFailureMessage = message + } + + package func prepareRuntime( + generation _: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose + ) async throws -> PreparedRuntime { + startRequests.append(purpose == .restartSameAccount) + let handle = TestingRuntimeLifecycleHandle( + onActivate: { [weak self] in self?.isActive = true }, + onClose: { [weak self] in self?.isActive = false } + ) + lastPreparedRuntimeHandle = handle + await runtimePreparationStartedGate.open() + if let runtimePreparationGate { + let cancellationGate = runtimePreparationCancellationGate + await withTaskCancellationHandler { + await runtimePreparationGate.waitIgnoringCancellation() + } onCancel: { + Task { await cancellationGate.open() } + } + self.runtimePreparationGate = nil + } + if let runtimePreparationFailureMessage { + self.runtimePreparationFailureMessage = nil + throw FakeCodexReviewBackendError(message: runtimePreparationFailureMessage) + } + return PreparedRuntime( + snapshot: .init( + authentication: try await reviewBackend.readAuth(), + settings: try await monitoredSettingsSnapshot() + ), + handle: handle + ) } package func stop(store _: CodexReviewStore) async { @@ -970,15 +1174,19 @@ package final class TestingCodexReviewStoreBackend: CodexReviewStoreBackend { } package func refreshSettings() async throws -> CodexReviewSettings.Snapshot { + currentSettingsSnapshot = try await monitoredSettingsSnapshot() + return currentSettingsSnapshot + } + + private func monitoredSettingsSnapshot() async throws -> CodexReviewSettings.Snapshot { let snapshot = try await reviewBackend.readSettings() - currentSettingsSnapshot = .init( + return .init( model: snapshot.model, fallbackModel: snapshot.fallbackModel, reasoningEffort: snapshot.reasoningEffort.flatMap(CodexReviewSettings.ReasoningEffort.init(rawValue:)), serviceTier: snapshot.serviceTier.flatMap(CodexReviewSettings.ServiceTier.init(rawValue:)), models: snapshot.models ) - return currentSettingsSnapshot } package func updateSettingsModel( diff --git a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift index 2289bdc0..9b460b12 100644 --- a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift +++ b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift @@ -44,6 +44,8 @@ private actor HostCloseFailureTransport: JSONRPC.Transport { struct CodexReviewHostTests { @Test func hostStartsAndStopsRuntimeWithFakeBackend() async throws { let backend = FakeCodexReviewBackend() + let interruptGate = AsyncGate() + await backend.holdInterruptReview(with: interruptGate) let host = CodexReviewHost( backend: backend, endpoint: URL(string: "http://localhost:9417/mcp") @@ -53,7 +55,31 @@ struct CodexReviewHostTests { #expect(host.store.serverState == .running) #expect(host.store.serverURL == URL(string: "http://localhost:9417/mcp")) - try await host.stop() + let firstReview = Task { @MainActor in + try await host.store.startReview( + sessionID: "session-1", + request: .init(cwd: "/tmp/project", target: .uncommittedChanges) + ) + } + await backend.waitForStartReview() + let stopTask = Task { @MainActor in try await host.stop() } + try await backend.waitForInterruptReview(timeout: .seconds(2)) + + let rejectedReview = try await host.store.startReview( + sessionID: "session-after-stop", + request: .init(cwd: "/tmp/other-project", target: .uncommittedChanges) + ) + #expect(rejectedReview.core.lifecycle.status == .failed) + await host.store.signIn() + let startCallCount = await backend.recordedCommands().filter { + if case .startReview = $0 { true } else { false } + }.count + #expect(startCallCount == 1) + #expect(await backend.recordedCommands().contains { if case .startLogin = $0 { true } else { false } } == false) + + await interruptGate.open() + try await stopTask.value + _ = try await firstReview.value #expect(host.store.serverState == .stopped) } @@ -110,7 +136,6 @@ struct CodexReviewHostTests { let host = CodexReviewHost(backend: backend) await host.start() - await host.store.refreshAuthentication() #expect(host.store.auth.selectedAccount?.accountKey == "review@example.com") #expect(host.store.auth.selectedAccount?.email == "review@example.com") @@ -238,6 +263,55 @@ struct CodexReviewHostTests { await store.stop() } + @Test func liveStoreCommitsInitialAccountPersistenceAndRateLimitsWithRuntimePublication() async throws { + let homeURL = try temporaryHome() + let codexHomeURL = homeURL.appendingPathComponent(".codex_review", isDirectory: true) + try FileManager.default.createDirectory(at: codexHomeURL, withIntermediateDirectories: true) + let sharedAuth = Data("{\"tokens\":{\"id_token\":\"active-token\"}}".utf8) + try sharedAuth.write(to: codexHomeURL.appendingPathComponent("auth.json")) + let transport = FakeJSONRPCTransport() + try await transport.enqueue(AppServerAPI.Initialize.Response(), for: "initialize") + try await transport.enqueue( + AppServerAPI.Account.Read.Response( + account: .init(email: "active@example.com", planType: "pro") + ), + for: "account/read" + ) + try await transport.enqueue( + AppServerAPI.Config.Read.Response(config: .init(model: "gpt-5")), + for: "config/read" + ) + try await transport.enqueue(AppServerAPI.Model.List.Response(data: []), for: "model/list") + try await transport.enqueue( + AppServerAPI.Account.RateLimits.Response(rateLimits: .init( + limitID: "codex", + primary: .init(usedPercent: 17, windowDurationMins: 300) + )), + for: "account/rateLimits/read" + ) + let store = CodexReviewStore.makeLiveStoreForTesting( + environment: ["HOME": homeURL.path], + webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, + mcpHTTPServerFactory: nil, + transportFactory: { _ in transport } + ) + + await store.start() + try #require(await waitUntil(timeout: .seconds(2)) { + store.auth.selectedAccount?.rateLimits.first?.usedPercent == 17 + }) + + #expect(store.serverState == .running) + #expect(store.auth.selectedAccount?.accountKey == "active@example.com") + #expect(store.auth.selectedAccount?.rateLimits.first?.usedPercent == 17) + #expect(try activeAccountKey(homeURL: homeURL) == "active@example.com") + #expect(try savedAccountAuth(homeURL: homeURL, accountKey: "active@example.com") == sharedAuth) + #expect(await transport.recordedRequests().map(\.method).filter { + $0 == "account/rateLimits/read" + }.count == 1) + await store.stop() + } + @Test func liveStorePassesRuntimePreferenceMCPPortAndPathToHTTPServerFactory() async throws { let homeURL = try temporaryHome() let transport = FakeJSONRPCTransport() @@ -280,6 +354,151 @@ struct CodexReviewHostTests { await store.stop() } + @Test func liveStoreRestartReplacesOnlyAppServerAndRetainsMCPListener() async throws { + let homeURL = try temporaryHome() + let firstTransport = FakeJSONRPCTransport() + let secondTransport = FakeJSONRPCTransport() + for transport in [firstTransport, secondTransport] { + try await enqueueRuntimeStartResponses(transport) + } + let server = ControlledMCPHTTPServer( + endpoint: try #require(URL(string: "http://127.0.0.1:19435/mcp")) + ) + var transports = [firstTransport, secondTransport] + var serverFactoryCallCount = 0 + let store = CodexReviewStore.makeLiveStoreForTesting( + environment: ["HOME": homeURL.path], + webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, + mcpHTTPServerFactory: { _, _ in + serverFactoryCallCount += 1 + return server + }, + mcpHTTPServerBindChecker: { _ in }, + transportFactory: { _ in transports.removeFirst() } + ) + + await store.start() + let firstURL = store.serverURL + await store.restart() + + #expect(store.serverState == .running) + #expect(store.serverURL == firstURL) + #expect(serverFactoryCallCount == 1) + #expect(server.startCallCount == 1) + #expect(server.stopCallCount == 0) + #expect(await firstTransport.isClosedForTesting()) + #expect(transports.isEmpty) + + await store.stop() + #expect(server.stopCallCount == 1) + } + + @Test func liveStoreStopInvalidatesHeldMCPURLPublication() async throws { + let homeURL = try temporaryHome() + let firstTransport = FakeJSONRPCTransport() + let secondTransport = FakeJSONRPCTransport() + for transport in [firstTransport, secondTransport] { + try await enqueueRuntimeStartResponses(transport) + } + let urlRelease = AsyncGate() + let firstServer = ControlledMCPHTTPServer( + endpoint: try #require(URL(string: "http://127.0.0.1:19438/mcp")) + ) + firstServer.holdURLRead(with: urlRelease) + let secondServer = ControlledMCPHTTPServer( + endpoint: try #require(URL(string: "http://127.0.0.1:19439/mcp")) + ) + var transports = [firstTransport, secondTransport] + var servers = [firstServer, secondServer] + let store = CodexReviewStore.makeLiveStoreForTesting( + environment: ["HOME": homeURL.path], + webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, + mcpHTTPServerFactory: { _, _ in servers.removeFirst() }, + mcpHTTPServerBindChecker: { _ in }, + transportFactory: { _ in transports.removeFirst() } + ) + + let startTask = Task { @MainActor in await store.start() } + await firstServer.waitForURLRead() + let stopTask = Task { @MainActor in await store.stop() } + await firstServer.waitForURLReadCancellation() + await urlRelease.open() + await stopTask.value + await startTask.value + + #expect(store.serverState == .stopped) + #expect(store.serverURL == nil) + #expect(firstServer.stopCallCount == 1) + + await store.start() + #expect(store.serverState == .running) + #expect(store.serverURL == secondServer.endpoint) + #expect(secondServer.startCallCount == 1) + #expect(servers.isEmpty) + #expect(transports.isEmpty) + await store.stop() + } + + @Test func liveStorePublishesBeforeHeldInitialRateLimitAndSuppressesStaleResult() async throws { + let homeURL = try temporaryHome() + let rateLimitRelease = AsyncGate() + let transport = FakeJSONRPCTransport() + try await transport.enqueue(AppServerAPI.Initialize.Response(), for: "initialize") + try await transport.enqueue( + AppServerAPI.Account.Read.Response( + account: .init(email: "stale@example.com", planType: "pro") + ), + for: "account/read" + ) + try await transport.enqueue( + AppServerAPI.Config.Read.Response(config: .init(model: "gpt-5")), + for: "config/read" + ) + try await transport.enqueue(AppServerAPI.Model.List.Response(data: []), for: "model/list") + try await transport.enqueue( + AppServerAPI.Account.RateLimits.Response(rateLimits: .init( + limitID: "codex", + primary: .init(usedPercent: 9, windowDurationMins: 300) + )), + for: "account/rateLimits/read" + ) + await transport.holdNextIgnoringCancellation( + method: "account/rateLimits/read", + gate: rateLimitRelease + ) + let stopFinished = CompletionFlag() + let store = CodexReviewStore.makeLiveStoreForTesting( + environment: ["HOME": homeURL.path], + webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, + mcpHTTPServerFactory: nil, + transportFactory: { _ in transport } + ) + + let startTask = Task { @MainActor in await store.start() } + await transport.waitForRequestCount(5) + #expect(store.serverState == .running) + #expect(store.auth.selectedAccount?.accountKey == "stale@example.com") + #expect(store.auth.selectedAccount?.rateLimits.isEmpty == true) + #expect(await transport.notificationStreamCount() == 1) + + let stopTask = Task { @MainActor in + await store.stop() + await stopFinished.complete() + } + await Task.yield() + + #expect(await stopFinished.isCompleted() == false) + #expect(store.auth.selectedAccount?.rateLimits.isEmpty == true) + + await rateLimitRelease.open() + await stopTask.value + await startTask.value + + #expect(store.serverState == .stopped) + #expect(store.auth.selectedAccount?.rateLimits.isEmpty == true) + #expect(await transport.isClosedForTesting()) + } + @Test func liveStoreKeepsThrowingMCPStopAtItsNonthrowingBoundary() async throws { let homeURL = try temporaryHome() let firstTransport = FakeJSONRPCTransport() @@ -329,6 +548,13 @@ struct CodexReviewHostTests { @Test func liveStoreStopsMCPServerAfterItsStartFails() async throws { let homeURL = try temporaryHome() let transport = FakeJSONRPCTransport() + try await transport.enqueue(AppServerAPI.Initialize.Response(), for: "initialize") + try await transport.enqueue(AppServerAPI.Account.Read.Response(), for: "account/read") + try await transport.enqueue( + AppServerAPI.Config.Read.Response(config: .init(model: "gpt-5")), + for: "config/read" + ) + try await transport.enqueue(AppServerAPI.Model.List.Response(data: []), for: "model/list") let server = ControlledMCPHTTPServer( endpoint: try #require(URL(string: "http://127.0.0.1:19434/mcp")), startFailure: .injected @@ -1342,6 +1568,13 @@ struct CodexReviewHostTests { let jobBeforeInterruptCompletes = try #require(store.jobs.first) #expect(jobBeforeInterruptCompletes.cancellationRequested) #expect(jobBeforeInterruptCompletes.core.lifecycle.cancellation?.message == "Review runtime stopped.") + let rejectedReview = try await store.startReview( + sessionID: "session-after-stop", + request: .init(cwd: "/tmp/other-project", target: .uncommittedChanges) + ) + #expect(rejectedReview.core.lifecycle.status == .failed) + #expect(rejectedReview.core.lifecycle.errorMessage == "Review runtime is not running.") + #expect(await transport.recordedRequests().map(\.method).filter { $0 == "thread/start" }.count == 1) await interruptGate.open() await stopTask.value let result = try await reviewRead @@ -1473,23 +1706,21 @@ struct CodexReviewHostTests { @Test func liveStoreMarksRuntimeFailedWhenAppServerNotificationStreamCloses() async throws { let homeURL = try temporaryHome() - let transport = FakeJSONRPCTransport() - try await transport.enqueue(AppServerAPI.Initialize.Response(), for: "initialize") - try await transport.enqueue(AppServerAPI.Account.Read.Response(), for: "account/read") - try await transport.enqueue( - AppServerAPI.Config.Read.Response(config: .init(model: "gpt-5")), - for: "config/read" - ) - try await transport.enqueue(AppServerAPI.Model.List.Response(data: []), for: "model/list") + let firstTransport = FakeJSONRPCTransport() + let secondTransport = FakeJSONRPCTransport() + for transport in [firstTransport, secondTransport] { + try await enqueueRuntimeStartResponses(transport) + } + var transports = [firstTransport, secondTransport] let store = CodexReviewStore.makeLiveStoreForTesting( environment: ["HOME": homeURL.path], webAuthenticationSessionFactory: FakeWebAuthenticationSessions().makeSession, - transport: transport + transportFactory: { _ in transports.removeFirst() } ) await store.start(forceRestartIfNeeded: true) - await transport.waitForNotificationStreamCount(1) - await transport.finishNotificationStreams(throwing: JSONRPC.Error.closed) + await firstTransport.waitForNotificationStreamCount(1) + await firstTransport.finishNotificationStreams(throwing: JSONRPC.Error.closed) await waitUntil { if case .failed = store.serverState { return true @@ -1503,6 +1734,11 @@ struct CodexReviewHostTests { } #expect(message.contains("JSON-RPC transport is closed")) #expect(store.serverURL == nil) + + await store.start() + #expect(store.serverState == .running) + #expect(transports.isEmpty) + await store.stop() } @Test func liveStoreCleansIsolatedLoginRuntimeWhenMainNotificationStreamCloses() async throws { @@ -2042,6 +2278,19 @@ private func writeSavedAccountAuth(homeURL: URL, accountKey: String) throws { try Data("{\"tokens\":{\"id_token\":\"\(accountKey)\"}}".utf8).write(to: authURL) } +private func enqueueRuntimeStartResponses( + _ transport: FakeJSONRPCTransport, + model: String = "gpt-5" +) async throws { + try await transport.enqueue(AppServerAPI.Initialize.Response(), for: "initialize") + try await transport.enqueue(AppServerAPI.Account.Read.Response(), for: "account/read") + try await transport.enqueue( + AppServerAPI.Config.Read.Response(config: .init(model: model)), + for: "config/read" + ) + try await transport.enqueue(AppServerAPI.Model.List.Response(data: []), for: "model/list") +} + private func savedAccountAuth(homeURL: URL, accountKey: String) throws -> Data { try Data(contentsOf: homeURL .appendingPathComponent(".codex_review", isDirectory: true) @@ -2106,6 +2355,9 @@ private final class ControlledMCPHTTPServer: CodexReviewMCPHTTPServing { private let stopFailure: HostCloseFailure? private(set) var startCallCount = 0 private(set) var stopCallCount = 0 + private var urlGate: AsyncGate? + private var urlReadStartedGate = AsyncGate() + private var urlReadCancellationGate = AsyncGate() init( endpoint: URL, @@ -2119,10 +2371,34 @@ private final class ControlledMCPHTTPServer: CodexReviewMCPHTTPServing { var url: URL { get async { - endpoint + await urlReadStartedGate.open() + if let urlGate { + let cancellationGate = urlReadCancellationGate + await withTaskCancellationHandler { + await urlGate.waitIgnoringCancellation() + } onCancel: { + Task { await cancellationGate.open() } + } + self.urlGate = nil + } + return endpoint } } + func holdURLRead(with gate: AsyncGate) { + urlGate = gate + urlReadStartedGate = AsyncGate() + urlReadCancellationGate = AsyncGate() + } + + func waitForURLRead() async { + await urlReadStartedGate.wait() + } + + func waitForURLReadCancellation() async { + await urlReadCancellationGate.wait() + } + func start() async throws { startCallCount += 1 if let startFailure { diff --git a/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift b/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift new file mode 100644 index 00000000..47732307 --- /dev/null +++ b/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift @@ -0,0 +1,212 @@ +import Foundation +import Testing +import CodexReview +import CodexReviewTesting + +@Suite("store runtime lifecycle") +@MainActor +struct CodexReviewStoreLifecycleTests { + @Test func stopInvalidatesHeldRuntimePreparationAndClosesStaleHandleOnce() async throws { + let preparationGate = AsyncGate() + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend() + ) + backend.holdRuntimePreparation(with: preparationGate) + let store = CodexReviewStore.makeTestingStore(backend: backend) + + let startTask = Task { @MainActor in await store.start() } + await backend.waitForRuntimePreparation() + let handle = try #require(backend.lastPreparedRuntimeHandle) + let stopTask = Task { @MainActor in await store.stop() } + await backend.waitForRuntimePreparationCancellation() + + #expect(handle.activateCallCount == 0) + #expect(store.serverState == .starting) + + await preparationGate.open() + await stopTask.value + await startTask.value + + #expect(store.serverState == .stopped) + #expect(handle.activateCallCount == 0) + #expect(handle.closeAdmissionCallCount == 1) + #expect(handle.closePurposes == [.start]) + #expect(handle.waitUntilClosedCallCount == 1) + #expect(backend.isActive == false) + } + @Test func stoppedStoreCanPrepareAndPublishANewRuntimeGeneration() async throws { + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend() + ) + let store = CodexReviewStore.makeTestingStore(backend: backend) + + await store.start() + let firstHandle = try #require(backend.lastPreparedRuntimeHandle) + #expect(store.serverState == .running) + #expect(firstHandle.activateCallCount == 1) + + await store.stop() + #expect(firstHandle.closePurposes == [.stop]) + + await store.start() + let secondHandle = try #require(backend.lastPreparedRuntimeHandle) + #expect(secondHandle !== firstHandle) + #expect(secondHandle.activateCallCount == 1) + #expect(store.serverState == .running) + } + @Test func restartInvalidatesHeldAcquisitionAndStartsAFreshGeneration() async throws { + let preparationGate = AsyncGate() + let mcpOwner = TestingMCPServerLifecycleOwner() + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend(), + mcpServerLifecycle: mcpOwner + ) + backend.holdRuntimePreparation(with: preparationGate) + let store = CodexReviewStore.makeTestingStore(backend: backend) + + let startTask = Task { @MainActor in await store.start() } + await backend.waitForRuntimePreparation() + let staleHandle = try #require(backend.lastPreparedRuntimeHandle) + let restartTask = Task { @MainActor in await store.restart() } + await backend.waitForRuntimePreparationCancellation() + + await preparationGate.open() + await restartTask.value + await startTask.value + + let currentHandle = try #require(backend.lastPreparedRuntimeHandle) + #expect(currentHandle !== staleHandle) + #expect(staleHandle.activateCallCount == 0) + #expect(staleHandle.closePurposes == [.start]) + #expect(currentHandle.activateCallCount == 1) + #expect(store.serverState == .running) + #expect(mcpOwner.preparedGenerations.count == 2) + #expect(mcpOwner.activatedGenerations.count == 1) + #expect(mcpOwner.stopCallCount == 1) + await store.stop() + } + @Test func staleMCPPreparationDoesNotAcquireAnAppServerRuntime() async { + let preparationGate = AsyncGate() + let mcpOwner = TestingMCPServerLifecycleOwner() + mcpOwner.holdPreparation(with: preparationGate) + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend(), + mcpServerLifecycle: mcpOwner + ) + let store = CodexReviewStore.makeTestingStore(backend: backend) + + let startTask = Task { @MainActor in await store.start() } + await mcpOwner.waitForPreparation() + let stopTask = Task { @MainActor in await store.stop() } + await mcpOwner.waitForPreparationCancellation() + await preparationGate.open() + await stopTask.value + await startTask.value + + #expect(store.serverState == .stopped) + #expect(backend.lastPreparedRuntimeHandle == nil) + #expect(mcpOwner.stopCallCount == 1) + #expect(mcpOwner.activatedGenerations.isEmpty) + } + @Test func explicitStopWinsWhileFailedAcquisitionIsDrainingMCP() async { + let stopGate = AsyncGate() + let mcpOwner = TestingMCPServerLifecycleOwner() + mcpOwner.holdStop(with: stopGate) + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend(), + mcpServerLifecycle: mcpOwner + ) + backend.failNextRuntimePreparation(message: "Injected preparation failure.") + let store = CodexReviewStore.makeTestingStore(backend: backend) + + let startTask = Task { @MainActor in await store.start() } + await mcpOwner.waitForStop() + let stopTask = Task { @MainActor in await store.stop() } + await mcpOwner.waitForSecondStop() + await stopGate.open() + await stopTask.value + await startTask.value + + #expect(store.serverState == .stopped) + #expect(store.serverURL == nil) + #expect(mcpOwner.stopCallCount == 1) + } + @Test func sameAccountRestartRetainsOneMCPGenerationAndReplacesOnlyAppServer() async throws { + let endpoint = try #require(URL(string: "http://127.0.0.1:19417/mcp")) + let mcpOwner = TestingMCPServerLifecycleOwner(serverURL: endpoint) + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend(), + mcpServerLifecycle: mcpOwner + ) + let store = CodexReviewStore.makeTestingStore(backend: backend) + + await store.start() + let firstHandle = try #require(backend.lastPreparedRuntimeHandle) + let firstMCPGeneration = try #require(mcpOwner.preparedGenerations.first) + + let replacementGate = AsyncGate() + backend.holdRuntimePreparation(with: replacementGate) + let restartTask = Task { @MainActor in await store.restart() } + await backend.waitForRuntimePreparation() + + #expect(store.serverURL == endpoint) + #expect(firstHandle.closePurposes == [.restartSameAccount]) + #expect(mcpOwner.preparedGenerations == [firstMCPGeneration]) + #expect(mcpOwner.activatedGenerations == [firstMCPGeneration]) + #expect(mcpOwner.stopCallCount == 0) + + await replacementGate.open() + await restartTask.value + + let secondHandle = try #require(backend.lastPreparedRuntimeHandle) + #expect(secondHandle !== firstHandle) + #expect(secondHandle.activateCallCount == 1) + #expect(backend.startRequests == [false, true]) + #expect(store.serverState == .running) + #expect(store.serverURL == endpoint) + #expect(mcpOwner.preparedGenerations == [firstMCPGeneration]) + #expect(mcpOwner.activatedGenerations == [firstMCPGeneration]) + + guard case .running(_, _, let retainedGeneration) = store.runtimeState else { + Issue.record("Replacement must publish the new AppServer runtime.") + return + } + #expect(retainedGeneration == firstMCPGeneration) + await store.stop() + } + @Test func stopInvalidatesHeldRestartBeforeReplacementCanPublish() async throws { + let endpoint = try #require(URL(string: "http://127.0.0.1:19422/mcp")) + let mcpOwner = TestingMCPServerLifecycleOwner(serverURL: endpoint) + let backend = TestingCodexReviewStoreBackend( + reviewBackend: FakeCodexReviewBackend(), + mcpServerLifecycle: mcpOwner + ) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let firstHandle = try #require(backend.lastPreparedRuntimeHandle) + + let replacementGate = AsyncGate() + backend.holdRuntimePreparation(with: replacementGate) + let restartTask = Task { @MainActor in await store.restart() } + await backend.waitForRuntimePreparation() + let staleReplacement = try #require(backend.lastPreparedRuntimeHandle) + let stopTask = Task { @MainActor in await store.stop() } + await backend.waitForRuntimePreparationCancellation() + + #expect(firstHandle.closePurposes == [.restartSameAccount]) + #expect(staleReplacement.activateCallCount == 0) + #expect(store.serverURL == endpoint) + + await replacementGate.open() + await stopTask.value + await restartTask.value + + #expect(store.serverState == .stopped) + #expect(store.serverURL == nil) + #expect(staleReplacement.activateCallCount == 0) + #expect(staleReplacement.closeAdmissionCallCount == 1) + #expect(staleReplacement.closePurposes == [.restartSameAccount]) + #expect(staleReplacement.waitUntilClosedCallCount == 1) + #expect(mcpOwner.stopCallCount == 1) + } +} diff --git a/Tests/ReviewUITests/ReviewUITests.swift b/Tests/ReviewUITests/ReviewUITests.swift index dba17b32..62fa78af 100644 --- a/Tests/ReviewUITests/ReviewUITests.swift +++ b/Tests/ReviewUITests/ReviewUITests.swift @@ -6695,12 +6695,15 @@ func makeStore(backend: AuthActionBackend) -> CodexReviewStore { final class CountingStartBackend: PreviewCodexReviewStoreBackend { private var startCalls = 0 - override func start( - store _: CodexReviewStore, - forceRestartIfNeeded _: Bool - ) async { - isActive = true + override func prepareRuntime( + generation: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose + ) async throws -> PreparedRuntime { startCalls += 1 + return try await super.prepareRuntime( + generation: generation, + purpose: purpose + ) } override func stop(store _: CodexReviewStore) async { @@ -6732,13 +6735,6 @@ final class AuthActionBackend: PreviewCodexReviewStoreBackend { ) } - override func start( - store _: CodexReviewStore, - forceRestartIfNeeded _: Bool - ) async { - isActive = true - } - override func stop(store _: CodexReviewStore) async { isActive = false } @@ -6775,12 +6771,6 @@ final class FailingCancellationBackend: PreviewCodexReviewStoreBackend { ) } - override func start( - store _: CodexReviewStore, - forceRestartIfNeeded _: Bool - ) async { - } - override func stop(store _: CodexReviewStore) async { } @@ -6824,12 +6814,6 @@ final class BlockingSettingsBackend: PreviewCodexReviewStoreBackend { ) } - override func start( - store _: CodexReviewStore, - forceRestartIfNeeded _: Bool - ) async { - } - override func stop(store _: CodexReviewStore) async { } From ad1d873d8f0158a888298689c7e171c47dc8a555 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 22 Aug 2026 06:28:33 +0900 Subject: [PATCH 2/3] fix(runtime): synchronize published settings baseline --- .../Settings/CodexReviewSettingsService.swift | 8 ++++ .../CodexReview/Store/CodexReviewStore.swift | 2 +- Sources/CodexReviewTesting/TestSupport.swift | 9 ++++ .../CodexReviewHostTests.swift | 5 +- .../CodexReviewStoreLifecycleTests.swift | 47 +++++++++++++++++++ 5 files changed, 69 insertions(+), 2 deletions(-) diff --git a/Sources/CodexReview/Settings/CodexReviewSettingsService.swift b/Sources/CodexReview/Settings/CodexReviewSettingsService.swift index ba30fa2d..8d150683 100644 --- a/Sources/CodexReview/Settings/CodexReviewSettingsService.swift +++ b/Sources/CodexReview/Settings/CodexReviewSettingsService.swift @@ -51,6 +51,14 @@ package final class CodexReviewSettingsService { lastPersistedSelection = settings.currentSelection() } + package func applyRuntimeSnapshot(_ snapshot: CodexReviewSettings.Snapshot) { + guard let settingsStore else { + return + } + settingsStore.apply(snapshot: snapshot) + lastPersistedSelection = settingsStore.currentSelection() + } + package func refreshIfRunning(serverState: CodexReviewServerState) async { guard case .running = serverState else { return diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index d1fa372a..508f29a4 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -504,7 +504,7 @@ public final class CodexReviewStore { _ runtime: PreparedRuntime, serverURL: URL? ) { - settings.apply(snapshot: runtime.snapshot.settings) + settingsService.applyRuntimeSnapshot(runtime.snapshot.settings) transitionToRunning(serverURL: serverURL) startAccountRateLimitAutoRefresh() } diff --git a/Sources/CodexReviewTesting/TestSupport.swift b/Sources/CodexReviewTesting/TestSupport.swift index 44b35513..410777b1 100644 --- a/Sources/CodexReviewTesting/TestSupport.swift +++ b/Sources/CodexReviewTesting/TestSupport.swift @@ -184,6 +184,7 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { } private var settings: CodexReviewBackendModel.Settings.Snapshot + private var settingsUpdateFailureMessage: String? private var auth: CodexReviewBackendModel.Auth.Snapshot private var commands: [Command] = [] private var startAdmissionIdentities: [ObjectIdentifier] = [] @@ -255,6 +256,10 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { cleanupFailure = .cleanup(message) } + package func failNextSettingsUpdate(message: String) { + settingsUpdateFailureMessage = message + } + package func holdInterruptReview(with gate: AsyncGate) { interruptReviewGate = gate } @@ -438,6 +443,10 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { package func applySettings(_ change: CodexReviewBackendModel.Settings.Change) async throws -> CodexReviewBackendModel.Settings.Snapshot { commands.append(.applySettings(change)) + if let settingsUpdateFailureMessage { + self.settingsUpdateFailureMessage = nil + throw FakeCodexReviewBackendError(message: settingsUpdateFailureMessage) + } settings = .init( model: change.updatesModel ? change.model : settings.model, fallbackModel: settings.fallbackModel, diff --git a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift index 9b460b12..63229fa0 100644 --- a/Tests/CodexReviewHostTests/CodexReviewHostTests.swift +++ b/Tests/CodexReviewHostTests/CodexReviewHostTests.swift @@ -61,7 +61,10 @@ struct CodexReviewHostTests { request: .init(cwd: "/tmp/project", target: .uncommittedChanges) ) } - await backend.waitForStartReview() + let activeSnapshot = await StoreSnapshotProbe(store: host.store).waitUntil { + $0.job()?.activeRun != nil + } + try #require(activeSnapshot?.job()?.activeRun != nil) let stopTask = Task { @MainActor in try await host.stop() } try await backend.waitForInterruptReview(timeout: .seconds(2)) diff --git a/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift b/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift index 47732307..7bcec873 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift @@ -209,4 +209,51 @@ struct CodexReviewStoreLifecycleTests { #expect(staleReplacement.waitUntilClosedCallCount == 1) #expect(mcpOwner.stopCallCount == 1) } + + @Test func failedSettingsEditRollsBackToPublishedRuntimeSnapshot() async { + let reviewBackend = FakeCodexReviewBackend(settings: .init( + model: "runtime-model", + reasoningEffort: "high", + serviceTier: "fast" + )) + let store = CodexReviewStore.makeTestingStore(backend: TestingCodexReviewStoreBackend( + reviewBackend: reviewBackend, + seed: .init(initialSettingsSnapshot: .init(model: "seed-model")) + )) + await store.start() + await reviewBackend.failNextSettingsUpdate(message: "Injected settings failure.") + + await store.updateSettingsModel("edited-model") + + #expect(store.settings.selectedModel == "runtime-model") + #expect(store.settings.selectedReasoningEffort == .high) + #expect(store.settings.selectedServiceTier == .fast) + #expect(store.settings.lastErrorMessage == "Injected settings failure.") + } + + @Test func modelEditDoesNotRepersistPublishedRuntimeReasoningAndTier() async throws { + let reviewBackend = FakeCodexReviewBackend(settings: .init( + model: "runtime-model", + reasoningEffort: "high", + serviceTier: "fast" + )) + let store = CodexReviewStore.makeTestingStore(backend: TestingCodexReviewStoreBackend( + reviewBackend: reviewBackend, + seed: .init(initialSettingsSnapshot: .init( + model: "seed-model", + reasoningEffort: .low + )) + )) + await store.start() + + await store.updateSettingsModel("edited-model") + + let command = try #require(await reviewBackend.recordedCommands().last) + guard case .applySettings(let change) = command else { + Issue.record("Expected a settings update.") + return + } + #expect(change.updatesReasoningEffort == false) + #expect(change.updatesServiceTier == false) + } } From 3b9f790634fbd1e2f09d889793b582210fc05ff0 Mon Sep 17 00:00:00 2001 From: Kazuki Nakashima <65545348+lynnswap@users.noreply.github.com> Date: Sat, 22 Aug 2026 06:59:56 +0900 Subject: [PATCH 3/3] fix(runtime): serialize publication races --- .../CodexReview/ReviewRuntimeLifecycle.swift | 1 + .../Settings/CodexReviewSettingsService.swift | 38 ++++++--- .../CodexReview/Store/CodexReviewStore.swift | 64 +++++++++++++-- Sources/CodexReviewTesting/TestSupport.swift | 32 ++++++++ .../CodexReviewStoreLifecycleTests.swift | 82 +++++++++++++++++-- 5 files changed, 192 insertions(+), 25 deletions(-) diff --git a/Sources/CodexReview/ReviewRuntimeLifecycle.swift b/Sources/CodexReview/ReviewRuntimeLifecycle.swift index cc2a49b3..5a6e3aae 100644 --- a/Sources/CodexReview/ReviewRuntimeLifecycle.swift +++ b/Sources/CodexReview/ReviewRuntimeLifecycle.swift @@ -15,6 +15,7 @@ package enum ReviewRuntimeTransitionPurpose: Equatable, Sendable { case start case restartSameAccount case stop + case runtimeFailure } package struct RuntimePublicationSnapshot: Sendable { package let authentication: CodexReviewBackendModel.Auth.Snapshot diff --git a/Sources/CodexReview/Settings/CodexReviewSettingsService.swift b/Sources/CodexReview/Settings/CodexReviewSettingsService.swift index 8d150683..03a41af5 100644 --- a/Sources/CodexReview/Settings/CodexReviewSettingsService.swift +++ b/Sources/CodexReview/Settings/CodexReviewSettingsService.swift @@ -32,6 +32,7 @@ package final class CodexReviewSettingsService { private var pendingRefresh = false private var pendingSelection: SettingsStore.Selection? private var lastPersistedSelection: SettingsStore.Selection + private var mutationRevision: UInt64 = 0 package init( initialSnapshot: CodexReviewSettings.Snapshot, @@ -55,6 +56,9 @@ package final class CodexReviewSettingsService { guard let settingsStore else { return } + mutationRevision &+= 1 + pendingRefresh = false + pendingSelection = nil settingsStore.apply(snapshot: snapshot) lastPersistedSelection = settingsStore.currentSelection() } @@ -75,14 +79,19 @@ package final class CodexReviewSettingsService { return } + let revision = mutationRevision settingsStore.beginLoading() do { let snapshot = try await backend.refreshSettings() - settingsStore.apply(snapshot: snapshot) - lastPersistedSelection = settingsStore.currentSelection() + if mutationRevision == revision { + settingsStore.apply(snapshot: snapshot) + lastPersistedSelection = settingsStore.currentSelection() + } settingsStore.finishLoading(errorMessage: nil) } catch { - settingsStore.finishLoading(errorMessage: error.localizedDescription) + settingsStore.finishLoading( + errorMessage: mutationRevision == revision ? error.localizedDescription : nil + ) } await drainPendingWorkIfNeeded() } @@ -172,6 +181,7 @@ package final class CodexReviewSettingsService { return } + let revision = mutationRevision settingsStore.beginLoading() do { try await persistSelection( @@ -179,16 +189,22 @@ package final class CodexReviewSettingsService { previous: previous, candidate: candidate ) - lastPersistedSelection = settingsStore.selectionAfterPersisting( - trigger: trigger, - previous: previous, - candidate: candidate - ) + if mutationRevision == revision { + lastPersistedSelection = settingsStore.selectionAfterPersisting( + trigger: trigger, + previous: previous, + candidate: candidate + ) + } settingsStore.finishLoading(errorMessage: nil) } catch { - settingsStore.apply(snapshot: settingsStore.snapshot(selection: previous)) - lastPersistedSelection = previous - settingsStore.finishLoading(errorMessage: error.localizedDescription) + if mutationRevision == revision { + settingsStore.apply(snapshot: settingsStore.snapshot(selection: previous)) + lastPersistedSelection = previous + } + settingsStore.finishLoading( + errorMessage: mutationRevision == revision ? error.localizedDescription : nil + ) } await drainPendingWorkIfNeeded() } diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index 508f29a4..865005a2 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -185,6 +185,22 @@ public final class CodexReviewStore { case .transitioning(_, .stop, let task): await task.value return + case .transitioning(let generation, .runtimeFailure, let failureTask): + let explicitGeneration = generation.successor() + let task = Task { @MainActor [weak self] in + await failureTask.value + self?.finishRuntimeStop( + generation: explicitGeneration, + purpose: .stop + ) + } + runtimeState = .transitioning( + generation: explicitGeneration, + purpose: .stop, + task: task + ) + await task.value + return case .acquiring, .running, .transitioning, .failed: break } @@ -192,7 +208,8 @@ public final class CodexReviewStore { let task = Task { @MainActor [weak self] in await self?.performRuntimeStop( previousState: previousState, - invalidatedGeneration: invalidatedGeneration + invalidatedGeneration: invalidatedGeneration, + purpose: .stop ) } runtimeState = .transitioning( @@ -294,11 +311,28 @@ public final class CodexReviewStore { else { return } - let stoppedGeneration = generation.successor() - await stop() - guard case .stopped(stoppedGeneration) = runtimeState else { + let invalidatedGeneration = generation.successor() + let previousState = runtimeState + let task = Task { @MainActor [weak self] in + await self?.performRuntimeStop( + previousState: previousState, + invalidatedGeneration: invalidatedGeneration, + purpose: .runtimeFailure + ) + } + runtimeState = .transitioning( + generation: invalidatedGeneration, + purpose: .runtimeFailure, + task: task + ) + await task.value + guard isCurrentTransition( + invalidatedGeneration, + purpose: .runtimeFailure + ) else { return } + runtimeState = .stopped(invalidatedGeneration) transitionToFailed(message) } @@ -396,7 +430,8 @@ public final class CodexReviewStore { private func performRuntimeStop( previousState: ReviewStoreRuntimeState, - invalidatedGeneration: ReviewRuntimeGeneration + invalidatedGeneration: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose ) async { switch previousState { case .acquiring(_, let task): @@ -408,7 +443,7 @@ public final class CodexReviewStore { await runtime.handle.closeAdmission() await stopPublishedRuntimeSemantics() await stopMCPServer() - await closeRuntime(runtime, purpose: .stop, admissionAlreadyClosed: true) + await closeRuntime(runtime, purpose: purpose, admissionAlreadyClosed: true) case .transitioning(_, _, let task): task.cancel() @@ -422,16 +457,27 @@ public final class CodexReviewStore { break } + guard purpose == .stop else { + return + } + finishRuntimeStop(generation: invalidatedGeneration, purpose: purpose) + } + + private func finishRuntimeStop( + generation: ReviewRuntimeGeneration, + purpose: ReviewRuntimeTransitionPurpose + ) { guard case .transitioning( let currentGeneration, - .stop, + let currentPurpose, _ ) = runtimeState, - currentGeneration == invalidatedGeneration + currentGeneration == generation, + currentPurpose == purpose else { return } - runtimeState = .stopped(invalidatedGeneration) + runtimeState = .stopped(generation) transitionToStopped() } diff --git a/Sources/CodexReviewTesting/TestSupport.swift b/Sources/CodexReviewTesting/TestSupport.swift index 410777b1..f2a709e2 100644 --- a/Sources/CodexReviewTesting/TestSupport.swift +++ b/Sources/CodexReviewTesting/TestSupport.swift @@ -185,6 +185,8 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { private var settings: CodexReviewBackendModel.Settings.Snapshot private var settingsUpdateFailureMessage: String? + private var settingsUpdateGate: AsyncGate? + private var settingsUpdateStartedGate = AsyncGate() private var auth: CodexReviewBackendModel.Auth.Snapshot private var commands: [Command] = [] private var startAdmissionIdentities: [ObjectIdentifier] = [] @@ -260,6 +262,19 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { settingsUpdateFailureMessage = message } + package func holdNextSettingsUpdate(with gate: AsyncGate) { + settingsUpdateGate = gate + settingsUpdateStartedGate = AsyncGate() + } + + package func waitForSettingsUpdate() async { + await settingsUpdateStartedGate.wait() + } + + package func setSettingsSnapshot(_ snapshot: CodexReviewBackendModel.Settings.Snapshot) { + settings = snapshot + } + package func holdInterruptReview(with gate: AsyncGate) { interruptReviewGate = gate } @@ -443,6 +458,9 @@ package actor FakeCodexReviewBackend: CodexReviewBackend { package func applySettings(_ change: CodexReviewBackendModel.Settings.Change) async throws -> CodexReviewBackendModel.Settings.Snapshot { commands.append(.applySettings(change)) + await settingsUpdateStartedGate.open() + await settingsUpdateGate?.waitIgnoringCancellation() + settingsUpdateGate = nil if let settingsUpdateFailureMessage { self.settingsUpdateFailureMessage = nil throw FakeCodexReviewBackendError(message: settingsUpdateFailureMessage) @@ -823,6 +841,8 @@ package final class TestingRuntimeLifecycleHandle: RuntimeLifecycleHandle { private let onActivate: @MainActor @Sendable () -> Void private let onClose: @MainActor @Sendable () -> Void + private var closeGate: AsyncGate? + private var closeStartedGate = AsyncGate() private var didClose = false package init( @@ -842,8 +862,20 @@ package final class TestingRuntimeLifecycleHandle: RuntimeLifecycleHandle { closeAdmissionCallCount += 1 } + package func holdClose(with gate: AsyncGate) { + closeGate = gate + closeStartedGate = AsyncGate() + } + + package func waitForClose() async { + await closeStartedGate.wait() + } + package func close(purpose: ReviewRuntimeTransitionPurpose) async throws { closePurposes.append(purpose) + await closeStartedGate.open() + await closeGate?.waitIgnoringCancellation() + closeGate = nil guard didClose == false else { return } diff --git a/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift b/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift index 7bcec873..016852fa 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreLifecycleTests.swift @@ -211,6 +211,7 @@ struct CodexReviewStoreLifecycleTests { } @Test func failedSettingsEditRollsBackToPublishedRuntimeSnapshot() async { + let updateGate = AsyncGate() let reviewBackend = FakeCodexReviewBackend(settings: .init( model: "runtime-model", reasoningEffort: "high", @@ -221,17 +222,30 @@ struct CodexReviewStoreLifecycleTests { seed: .init(initialSettingsSnapshot: .init(model: "seed-model")) )) await store.start() + await reviewBackend.holdNextSettingsUpdate(with: updateGate) await reviewBackend.failNextSettingsUpdate(message: "Injected settings failure.") + let updateTask = Task { @MainActor in + await store.updateSettingsModel("edited-model") + } + await reviewBackend.waitForSettingsUpdate() + await reviewBackend.setSettingsSnapshot(.init( + model: "fresh-runtime-model", + reasoningEffort: "medium", + serviceTier: "flex" + )) + await store.restart() - await store.updateSettingsModel("edited-model") + await updateGate.open() + await updateTask.value - #expect(store.settings.selectedModel == "runtime-model") - #expect(store.settings.selectedReasoningEffort == .high) - #expect(store.settings.selectedServiceTier == .fast) - #expect(store.settings.lastErrorMessage == "Injected settings failure.") + #expect(store.settings.selectedModel == "fresh-runtime-model") + #expect(store.settings.selectedReasoningEffort == .medium) + #expect(store.settings.selectedServiceTier == .flex) + #expect(store.settings.lastErrorMessage == nil) } @Test func modelEditDoesNotRepersistPublishedRuntimeReasoningAndTier() async throws { + let updateGate = AsyncGate() let reviewBackend = FakeCodexReviewBackend(settings: .init( model: "runtime-model", reasoningEffort: "high", @@ -245,7 +259,30 @@ struct CodexReviewStoreLifecycleTests { )) )) await store.start() + await reviewBackend.holdNextSettingsUpdate(with: updateGate) + let staleUpdateTask = Task { @MainActor in + await store.updateSettingsModel("stale-model") + } + await reviewBackend.waitForSettingsUpdate() + await store.updateSettingsModel("queued-stale-model") + await store.refreshSettings() + let freshSnapshot = CodexReviewBackendModel.Settings.Snapshot( + model: "fresh-runtime-model", + reasoningEffort: "medium", + serviceTier: "flex" + ) + await reviewBackend.setSettingsSnapshot(freshSnapshot) + await store.restart() + let commandCountAfterRestart = await reviewBackend.recordedCommands().count + + await updateGate.open() + await staleUpdateTask.value + #expect(store.settings.selectedModel == "fresh-runtime-model") + #expect(store.settings.selectedReasoningEffort == .medium) + #expect(store.settings.selectedServiceTier == .flex) + #expect(await reviewBackend.recordedCommands().count == commandCountAfterRestart) + await reviewBackend.setSettingsSnapshot(freshSnapshot) await store.updateSettingsModel("edited-model") let command = try #require(await reviewBackend.recordedCommands().last) @@ -256,4 +293,39 @@ struct CodexReviewStoreLifecycleTests { #expect(change.updatesReasoningEffort == false) #expect(change.updatesServiceTier == false) } + + @Test func explicitStopSuppressesJoinedRuntimeFailurePublication() async throws { + let closeGate = AsyncGate() + let backend = TestingCodexReviewStoreBackend(reviewBackend: FakeCodexReviewBackend()) + let store = CodexReviewStore.makeTestingStore(backend: backend) + await store.start() + let handle = try #require(backend.lastPreparedRuntimeHandle) + handle.holdClose(with: closeGate) + + let failureTask = Task { @MainActor in + await store.failRuntime(handle: handle, message: "Injected runtime failure.") + } + await handle.waitForClose() + let stopTask = Task { @MainActor in await store.stop() } + for _ in 0..<1_000 { + if case .transitioning(_, .stop, _) = store.runtimeState { + break + } + await Task.yield() + } + guard case .transitioning(_, .stop, _) = store.runtimeState else { + Issue.record("Explicit stop did not supersede runtime failure.") + await closeGate.open() + await stopTask.value + await failureTask.value + return + } + + await closeGate.open() + await stopTask.value + await failureTask.value + + #expect(store.serverState == .stopped) + #expect(handle.closePurposes == [.runtimeFailure]) + } }