Skip to content
6 changes: 6 additions & 0 deletions Sources/CodexReview/ReviewRuntimeLifecycle.swift
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ package struct ReviewRuntimeGeneration: Hashable, Sendable {

package enum ReviewRuntimeCleanupRecoverySuppression: Equatable, Sendable {
case explicitStop
case codexUpdate
case staleSource
case successorAlreadyFinished
}
Expand Down Expand Up @@ -89,6 +90,11 @@ package protocol RuntimeLifecycleHandle: AnyObject, Sendable {
func closeAdmission()
func close(purpose: ReviewRuntimeTransitionPurpose) async throws
func waitUntilClosed() async throws
func confirmClosed() async throws
}

extension RuntimeLifecycleHandle {
package func confirmClosed() async throws { try await waitUntilClosed() }
}

@MainActor
Expand Down
204 changes: 173 additions & 31 deletions Sources/CodexReview/Store/CodexReviewStore.swift

Large diffs are not rendered by default.

2 changes: 2 additions & 0 deletions Sources/CodexReview/Store/CodexReviewStoreBackend.swift
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,7 @@ package struct CodexReviewStoreSeed {
package protocol CodexReviewStoreBackend: CodexReviewSettingsBackend, Sendable {
var seed: CodexReviewStoreSeed { get }
var isActive: Bool { get }
var shutdownCleanupTimeout: Duration { get }
var handlesActiveReviewStopCleanup: Bool { get }
var mcpServerLifecycle: any MCPServerLifecycleOwner { get }

Expand Down Expand Up @@ -85,6 +86,7 @@ package protocol CodexReviewStoreBackend: CodexReviewSettingsBackend, Sendable {
}

extension CodexReviewStoreBackend {
package var shutdownCleanupTimeout: Duration { .seconds(2) }
package var handlesActiveReviewStopCleanup: Bool {
false
}
Expand Down
10 changes: 6 additions & 4 deletions Sources/CodexReview/Store/CodexReviewStoreCancellation.swift
Original file line number Diff line number Diff line change
Expand Up @@ -192,10 +192,11 @@ extension CodexReviewStore {
}

package func requestActiveReviewCancellationsForRuntimeStop(
reason: ReviewCancellation = .system(message: "Review runtime stopped.")
reason: ReviewCancellation = .system(message: "Review runtime stopped."),
includingQueued: Bool = true
) async -> ReviewRuntimeCancellationRequestOutcome {
let activeJobIDs = orderedJobs
.filter { $0.isTerminal == false }
.filter { $0.isTerminal == false && (includingQueued || queuedReviewStarts[$0.id] == nil) }
.map(\.id)
var firstFailure: ReviewRuntimeCloseFailure?
for jobID in activeJobIDs {
Expand All @@ -221,10 +222,11 @@ extension CodexReviewStore {

@discardableResult
package func cancelActiveReviewsLocallyForRuntimeStop(
reason: ReviewCancellation = .system(message: "Review runtime stopped.")
reason: ReviewCancellation = .system(message: "Review runtime stopped."),
includingQueued: Bool = true
) -> [String] {
let activeJobIDs = orderedJobs
.filter { $0.isTerminal == false }
.filter { $0.isTerminal == false && (includingQueued || queuedReviewStarts[$0.id] == nil) }
.map(\.id)
guard activeJobIDs.isEmpty == false else {
return []
Expand Down
200 changes: 200 additions & 0 deletions Sources/CodexReview/Store/CodexReviewStoreUpdate.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,200 @@
import Foundation

extension CodexReviewStore {
@_spi(ApplicationHostSupport) public enum CodexUpdateTiming: Sendable {
case afterCurrentReviews
case immediately
}

@_spi(ApplicationHostSupport) public enum CodexUpdateState: Equatable, Sendable {
case idle
case waitingForReviews
case stoppingRuntime
case installing
case restarting
case failed(String)
}

/// Keeps MCP sessions and accepted jobs alive while replacing the Codex runtime.
/// Concurrent callers join the existing operation; only its installation closure runs.
/// An installation failure is thrown even when the runtime recovers and jobs resume.
@_spi(ApplicationHostSupport) public func updateCodex(
when timing: CodexUpdateTiming,
install: @escaping @MainActor @Sendable () async throws -> Void
) async throws {
if let task = codexUpdateTask {
try await task.value
return
}
guard applicationShutdownRequested == false, pendingRuntimeStopCount == 0 else { throw CancellationError() }
let previousAccountOperations = Array(runtimeAccountOperations.values)
let previousRuntimeTask: Task<Void, Never>? = switch runtimeState {
case .acquiring(_, _, let task), .replacing(_, let task), .tearingDown(_, _, _, _, let task): task
case .running, .stopped, .failed: nil
}
suspendReviewStarts()
codexUpdate = .waitingForReviews
let task = Task { @MainActor [self] in
defer {
codexUpdateTask = nil
if Task.isCancelled == false { resumeReviewsAfterCodexUpdateIfPossible() }
Comment thread
lynnswap marked this conversation as resolved.
}
do {
try await Self.waitForPrecedingRuntimeWork(
accountOperations: previousAccountOperations, runtimeTask: previousRuntimeTask
)
try await performCodexUpdate(when: timing, install: install)
codexUpdate = .idle
} catch {
if error is CancellationError {
codexUpdate = .idle
} else {
codexUpdate = .failed(error.localizedDescription)
}
throw error
}
}
codexUpdateTask = task
try await task.value
}

private static func waitForPrecedingRuntimeWork(
accountOperations: [Task<Void, any Error>],
runtimeTask: Task<Void, Never>?
) async throws {
// Cancelling the update abandons its join, not the predecessor's lifecycle ownership.
let (completion, continuation) = AsyncStream<Void>.makeStream()
let waiter = Task {
for operation in accountOperations {
_ = try? await operation.value
guard Task.isCancelled == false else { return }
}
await runtimeTask?.value
continuation.yield(())
continuation.finish()
}
defer {
waiter.cancel()
continuation.finish()
}
var iterator = completion.makeAsyncIterator()
_ = await iterator.next()
try Task.checkCancellation()
}

private func performCodexUpdate(
when timing: CodexUpdateTiming,
install: @escaping @MainActor @Sendable () async throws -> Void
) async throws {
// A dispatcher already saving an execution start belongs to the current batch.
await queuedReviewDispatchTask?.value
try Task.checkCancellation()
if timing == .afterCurrentReviews {
let executingIDs = jobs.filter {
$0.isTerminal == false && queuedReviewStarts[$0.id] == nil
}.map(\.id)
for id in executingIDs {
// Another completed review can be removed while we await this batch.
guard job(id: id)?.isTerminal == false else { continue }
_ = try await awaitReview(sessionID: nil, jobID: id)
try Task.checkCancellation()
}
_ = await drainReviewWorkersForRuntimeStop(timeout: backend.shutdownCleanupTimeout)
try Task.checkCancellation()
}
guard applicationShutdownRequested == false else { throw CancellationError() }
codexUpdate = .stoppingRuntime
var installationFailure: String?
let installation: @MainActor @Sendable () async -> Void = { [self] in
codexUpdate = .installing
do { try await install() } catch { installationFailure = error.localizedDescription }
codexUpdate = .restarting
}
let operation: RuntimeStartOperation
switch runtimeState {
case .running(let generation, let runtime, let mcp):
operation = admitRuntimeReplacement(
sourceGeneration: generation, retiringRuntime: runtime, retainedMCP: mcp,
preservingQueuedReviews: true,
install: installation
)
case .failed(let generation, let mcp?, _):
operation = admitRuntimeReplacement(
sourceGeneration: generation, retiringRuntime: unclosedCodexUpdateRuntime, retainedMCP: mcp,
preservingQueuedReviews: true,
install: installation
)
case .stopped(let generation), .failed(let generation, nil, _):
try await closeUnclosedCodexUpdateRuntimeIfNeeded()
await installation()
let preparation = try await backend.mcpServerLifecycle.prepare()
let snapshot: MCPServerPublicationSnapshot
do {
snapshot = try await backend.mcpServerLifecycle.activate(preparation)
} catch {
var message = error.localizedDescription
do { try await backend.mcpServerLifecycle.stop() }
catch { message += "; MCP cleanup failed: \(error.localizedDescription)" }
throw CodexReviewAPI.Error.io(message)
}
operation = admitRuntimeReplacement(
sourceGeneration: generation, retiringRuntime: nil,
retainedMCP: RetainedMCPServer(serverURL: snapshot.serverURL),
preservingQueuedReviews: true
)
case .acquiring, .replacing, .tearingDown:
throw CodexReviewAPI.Error.io("The Codex runtime changed while waiting to update.")
}
await operation.task.value
let failures = [installationFailure, serverState.failureMessage].compactMap { $0 }
if failures.isEmpty == false {
throw CodexReviewAPI.Error.io(failures.joined(separator: "; "))
}
guard case .running = runtimeState else {
throw CodexReviewAPI.Error.io("Codex update finished without an available runtime.")
}
}

package func resumeReviewsAfterCodexUpdateIfPossible() {
guard runtimeAccountOperations.isEmpty, applicationShutdownRequested == false,
pendingRuntimeStopCount == 0,
case .running = runtimeState else { return }
resumeReviewStarts()
}

package func performRuntimeAuthentication(
_ operation: @escaping @MainActor @Sendable (CodexReviewStore) async -> Void
) async {
do { try await performRuntimeAccountChange { store in await operation(store) } }
catch is CancellationError { }
catch { auth.updatePhase(.failed(message: error.localizedDescription)) }
}

package func performRuntimeAccountChange(
_ operation: @escaping @MainActor @Sendable (CodexReviewStore) async throws -> Void
) async throws {
let update = codexUpdateTask
let id = UUID()
let task = Task { @MainActor [self] in
_ = try? await update?.value
defer {
runtimeAccountOperations.removeValue(forKey: id)
if codexUpdateTask == nil { resumeReviewsAfterCodexUpdateIfPossible() }
}
try Task.checkCancellation()
guard applicationShutdownRequested == false else { throw CancellationError() }
try await operation(self)
}
runtimeAccountOperations[id] = task
try await withTaskCancellationHandler {
try await task.value
} onCancel: {
task.cancel()
}
}

package var hasQueuedCodexUpdateRecovery: Bool {
if case .failed = codexUpdateState { return reviewStartsAreSuspended }
return false
}
}
9 changes: 5 additions & 4 deletions Sources/CodexReview/Store/ReviewRuntimeTeardownIntent.swift
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable {
case failed(String)
}

case codexUpdate
case explicitStop
case unexpectedFailure(String)

Expand All @@ -15,7 +16,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable {

package var finalState: FinalState {
switch self {
case .explicitStop:
case .explicitStop, .codexUpdate:
.stopped
case .unexpectedFailure:
.failed(message)
Expand All @@ -24,7 +25,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable {

package var diagnosticContext: String {
switch self {
case .explicitStop:
case .explicitStop, .codexUpdate:
"runtime stop"
case .unexpectedFailure:
"runtime failure"
Expand All @@ -33,7 +34,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable {

package var cleanupTimeoutWarning: String {
switch self {
case .explicitStop:
case .explicitStop, .codexUpdate:
"Timed out cleaning active reviews before stopping runtime"
case .unexpectedFailure:
"Timed out cleaning active reviews after runtime failure"
Expand All @@ -46,7 +47,7 @@ package enum ReviewRuntimeTeardownIntent: Equatable, Sendable {

private var message: String {
switch self {
case .explicitStop:
case .explicitStop, .codexUpdate:
"Review runtime stopped."
case .unexpectedFailure(let errorDescription):
"Review runtime stopped unexpectedly: \(errorDescription)"
Expand Down
4 changes: 4 additions & 0 deletions Sources/CodexReviewAppServer/AppServerClient.swift
Original file line number Diff line number Diff line change
Expand Up @@ -290,6 +290,10 @@ package actor AppServerClient {
try await transport.close()
}

package func confirmClosed() async throws {
try await transport.confirmClosed()
}

private func allocateRequestID() -> Int {
defer { nextRequestID += 1 }
return nextRequestID
Expand Down
30 changes: 29 additions & 1 deletion Sources/CodexReviewAppServer/AppServerProcessTransport.swift
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,15 @@ package actor AppServerProcessTransport: JSONRPC.Transport {
try await closeTransport(terminateProcess: true, readerTask: nil)
}

package func confirmClosed() async throws {
guard let closeTask else {
throw JSONRPC.Error.invalidMessage("Transport closure has not been requested.")
}
_ = await closeTask.result
await waitForReaderTasks(excluding: nil)
try process.confirmTerminated()
}

private func closeTransport(
terminateProcess: Bool,
readerTask: ReaderTask?
Expand Down Expand Up @@ -745,6 +754,7 @@ private final class AppServerSpawnedProcess: @unchecked Sendable {
private let processGroupID: pid_t
private let stateLock = NSLock()
private var didReap = false
private var terminationProcesses: Set<ProcessIdentity>?

private init(processIdentifier: pid_t) {
self.processIdentifier = processIdentifier
Expand Down Expand Up @@ -850,7 +860,7 @@ private final class AppServerSpawnedProcess: @unchecked Sendable {
killDuration: Duration = .seconds(1)
) async throws {
let processIdentity = Self.processSnapshot(processIdentifier)?.identity
let trackedProcesses = descendantProcesses()
let trackedProcesses = beginTerminationTracking()
guard isFullyTerminated(trackedProcesses: trackedProcesses) == false else {
return
}
Expand All @@ -875,6 +885,24 @@ private final class AppServerSpawnedProcess: @unchecked Sendable {
}
}

private func beginTerminationTracking() -> Set<ProcessIdentity> {
let descendants = descendantProcesses()
stateLock.withLock { terminationProcesses = descendants }
return descendants
}

func confirmTerminated() throws {
let (reaped, tracked) = stateLock.withLock { (didReap, terminationProcesses) }
// Reaping follows complete tree termination; never inspect a reused PID/group afterward.
if reaped { return }
guard let tracked, isFullyTerminated(trackedProcesses: tracked) else {
throw AppServerProcessTransportError.processDidNotTerminate(
processIdentifier,
liveProcessIdentifiers: liveProcessIdentifiers(trackedProcesses: tracked ?? [])
)
}
}

private func signalProcessTree(
_ signal: Int32,
processIdentity: ProcessIdentity?,
Expand Down
5 changes: 5 additions & 0 deletions Sources/CodexReviewAppServer/JSONRPC.swift
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ package enum JSONRPC {
func notificationStream() async -> AsyncThrowingStream<ReceivedNotification, Swift.Error>
func notificationHighWatermark() async -> NotificationReceipt
func close() async throws
func confirmClosed() async throws
}

package enum Error: Swift.Error, Equatable, Sendable, LocalizedError {
Expand Down Expand Up @@ -142,3 +143,7 @@ package struct AnyEncodable: Encodable {
package struct EmptyResponse: Codable, Equatable, Sendable {
package init() {}
}

extension JSONRPC.Transport {
package func confirmClosed() async throws { try await close() }
}
Loading