Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
158 changes: 158 additions & 0 deletions Docs/deferred-codex-update-design.md

Large diffs are not rendered by default.

6 changes: 3 additions & 3 deletions Sources/CodexReview/History/CodexReviewStoreHistory.swift
Original file line number Diff line number Diff line change
Expand Up @@ -79,15 +79,15 @@ extension CodexReviewStore {
let model = settings.effectiveModel
let workspaceSortOrder = historyWorkspaceSortOrder(cwd: request.cwd)
?? nextHistoryWorkspaceSortOrder()
let record = try StartedReviewRecord(
let record = try AcceptedReviewRecord(
id: id,
cwd: request.cwd,
workspaceMetadata: ReviewWorkspaceMetadata.resolve(cwd: request.cwd),
workspaceSortOrder: workspaceSortOrder,
sortOrder: try nextHistoryJobSortOrder(),
target: request.target,
model: model,
startedAt: clock.now()
acceptedAt: clock.now()
)
let receipt = HistoryStartReceipt(
ordinal: nextHistoryStartOrdinal,
Expand All @@ -107,7 +107,7 @@ extension CodexReviewStore {
intent: receipt.started,
prepare: { $0 },
operation: { record in
try await persistence.recordStarted(record)
try await persistence.recordAccepted(record)
},
apply: { [weak self] record, result in
guard let self else {
Expand Down
4 changes: 2 additions & 2 deletions Sources/CodexReview/History/HistoryStartReceipt.swift
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@
package final class HistoryStartReceipt {
package let ordinal: UInt64
package let sessionID: String
package let started: StartedReviewRecord
package let started: AcceptedReviewRecord
package let workAdmission: ReviewStoreWorkRegistry.Admission

private(set) var cancellation: ReviewCancellation?
Expand All @@ -12,7 +12,7 @@ package final class HistoryStartReceipt {
package init(
ordinal: UInt64,
sessionID: String,
started: StartedReviewRecord,
started: AcceptedReviewRecord,
workAdmission: ReviewStoreWorkRegistry.Admission
) {
self.ordinal = ordinal
Expand Down
7 changes: 5 additions & 2 deletions Sources/CodexReview/History/ReviewHistoryPersistence.swift
Original file line number Diff line number Diff line change
Expand Up @@ -70,7 +70,8 @@ package protocol ReviewHistoryPersistence: Sendable {
retentionPolicy: ReviewHistoryRetentionPolicy
) async throws -> [RestoredReviewRecord]

func recordStarted(_ record: StartedReviewRecord) async throws
func recordAccepted(_ record: AcceptedReviewRecord) async throws
func recordExecutionStarted(id: String, at date: Date) async throws

func recordTerminal(
_ record: TerminalReviewRecord,
Expand All @@ -97,7 +98,9 @@ package struct DisabledReviewHistoryPersistence: ReviewHistoryPersistence {
[]
}

package func recordStarted(_: StartedReviewRecord) async throws {}
package func recordAccepted(_: AcceptedReviewRecord) async throws {}

package func recordExecutionStarted(id: String, at date: Date) async throws {}

package func recordTerminal(
_: TerminalReviewRecord,
Expand Down
17 changes: 10 additions & 7 deletions Sources/CodexReview/History/ReviewHistoryRecord.swift
Original file line number Diff line number Diff line change
Expand Up @@ -12,15 +12,16 @@ package struct ReviewHistoryRecordError: LocalizedError, Sendable, Equatable {
}
}

package struct StartedReviewRecord: Sendable, Hashable {
package struct AcceptedReviewRecord: Sendable, Hashable {
package var id: String
package var cwd: String
package var workspaceMetadata: ReviewWorkspaceMetadata?
package var workspaceSortOrder: Double
package var sortOrder: Double
package var target: CodexReviewAPI.Target
package var model: String?
package var startedAt: Date
package var acceptedAt: Date
package var startedAt: Date?

package init(
id: String,
Expand All @@ -30,7 +31,8 @@ package struct StartedReviewRecord: Sendable, Hashable {
sortOrder: Double,
target: CodexReviewAPI.Target,
model: String?,
startedAt: Date
acceptedAt: Date,
startedAt: Date? = nil
) throws {
guard id.nilIfEmpty != nil else {
throw ReviewHistoryRecordError("A persisted review requires a stable ID.")
Expand All @@ -56,6 +58,7 @@ package struct StartedReviewRecord: Sendable, Hashable {
self.sortOrder = sortOrder
self.target = target
self.model = model?.nilIfEmpty
self.acceptedAt = acceptedAt
self.startedAt = startedAt
}
}
Expand Down Expand Up @@ -209,11 +212,11 @@ package struct TerminalReviewRecord: Sendable, Hashable {
}

package struct RestoredReviewRecord: Sendable, Hashable {
package var started: StartedReviewRecord
package var started: AcceptedReviewRecord
package var terminal: TerminalReviewRecord

package init(
started: StartedReviewRecord,
started: AcceptedReviewRecord,
terminal: TerminalReviewRecord
) throws {
guard started.id == terminal.id else {
Expand Down Expand Up @@ -245,7 +248,7 @@ package struct RestoredReviewRecord: Sendable, Hashable {
target: started.target,
origin: .restoredHistory,
core: core,
logEntries: terminal.compactLogEntries(startedAt: started.startedAt)
logEntries: terminal.compactLogEntries(startedAt: started.startedAt ?? started.acceptedAt)
)
}
}
Expand All @@ -266,7 +269,7 @@ private extension PersistedParsedReviewResult.Finding {
}

private extension TerminalReviewRecord {
func lifecycle(startedAt: Date) -> ReviewJobCore.Lifecycle {
func lifecycle(startedAt: Date?) -> ReviewJobCore.Lifecycle {
let status: ReviewJobState
let cancellation: ReviewCancellation?
let errorMessage: String?
Expand Down
4 changes: 4 additions & 0 deletions Sources/CodexReview/Store/CodexReviewStore.swift
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,10 @@ public final class CodexReviewStore {
@ObservationIgnored package var applicationShutdownTask: Task<Void, Never>?
@ObservationIgnored package var reviewAttemptOwnerships: [String: StoreReviewAttemptOwnership] = [:]
@ObservationIgnored package var reviewWorkerTasks: [String: Task<Void, Never>] = [:]
@ObservationIgnored package var queuedReviewStarts: [String: QueuedReviewStart] = [:]
@ObservationIgnored package var reviewStartsAreSuspended = false
@ObservationIgnored package var queuedReviewDispatchTask: Task<Void, Never>?

@ObservationIgnored package var runtimeStopDetachedReviewWorkerTasks: [String: Task<Void, Never>] = [:]
@ObservationIgnored package var reviewTerminalWaiters: [String: [ReviewTerminalWaiter]] = [:]
@ObservationIgnored package var nextCancellationRequestOrdinal: UInt64 = 0
Expand Down
4 changes: 4 additions & 0 deletions Sources/CodexReview/Store/CodexReviewStoreCancellation.swift
Original file line number Diff line number Diff line change
Expand Up @@ -96,6 +96,10 @@ extension CodexReviewStore {
return
}

if let queued = queuedReviewStarts.removeValue(forKey: jobID) {
queued.finishDispatch()
removeStartingReviewOwnership(for: jobID, ifOwnedBy: queued.admission)
}
let endedAt = clock.now()
job.closeActiveCommandLogEntries(status: "canceled", completedAt: endedAt)
job.pendingCancellationRequest = nil
Expand Down
149 changes: 149 additions & 0 deletions Sources/CodexReview/Store/CodexReviewStoreQueue.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,149 @@
import Foundation

@MainActor
package final class QueuedReviewStart {
package let ordinal: UInt64
package let request: CodexReviewAPI.Start.Request
package let model: String?
package let admission: ReviewStartAdmission
private enum Dispatch {
case pending
case finished(Task<Void, Never>?)
}
private var dispatch = Dispatch.pending
private var dispatchWaiters: [CheckedContinuation<Void, Never>] = []

package init(ordinal: UInt64, request: CodexReviewAPI.Start.Request, model: String?, admission: ReviewStartAdmission) {
self.ordinal = ordinal
self.request = request
self.model = model
self.admission = admission
}

package func finishDispatch(worker: Task<Void, Never>? = nil) {
dispatch = .finished(worker)
let waiters = dispatchWaiters
dispatchWaiters.removeAll()
for waiter in waiters { waiter.resume() }
}

package func waitUntilFinished() async {
if case .pending = dispatch {
await withCheckedContinuation { dispatchWaiters.append($0) }
}
if case .finished(let worker) = dispatch {
await worker?.value
}
}

}

extension CodexReviewStore {
package func suspendReviewStarts() {
reviewStartsAreSuspended = true
}

package func resumeReviewStarts() {
reviewStartsAreSuspended = false
scheduleQueuedReviewStarts()
}

private var nextQueuedReviewStart: (id: String, start: QueuedReviewStart)? {
guard let next = queuedReviewStarts.min(by: { $0.value.ordinal < $1.value.ordinal }),
historyStartReceipts.values.contains(where: { $0.ordinal < next.value.ordinal }) == false
else { return nil }
return (next.key, next.value)
}

package func scheduleQueuedReviewStarts() {
guard reviewStartsAreSuspended == false,
applicationShutdownRequested == false,
queuedReviewDispatchTask == nil,
nextQueuedReviewStart != nil else { return }
queuedReviewDispatchTask = startRegisteredStoreWork(
kind: .reviewMutation("dispatch-queued"),
cancelledBeforeEntry: .runFinalizer { store in
store.queuedReviewDispatchTask = nil
}
) { store in
await store.dispatchQueuedReviews()
}
}

private func dispatchQueuedReviews() async {
defer {
queuedReviewDispatchTask = nil
if Task.isCancelled == false { scheduleQueuedReviewStarts() }
}
while reviewStartsAreSuspended == false, applicationShutdownRequested == false,
Task.isCancelled == false, let next = nextQueuedReviewStart {
queuedReviewStarts.removeValue(forKey: next.id)
var dispatchedWorker: Task<Void, Never>?
defer { next.start.finishDispatch(worker: dispatchedWorker) }
guard let job = job(id: next.id), job.isTerminal == false else { continue }
let startedAt = clock.now()
job.core.lifecycle.status = .running
job.core.lifecycle.startedAt = startedAt
job.core.output.summary = "Review started."
Comment thread
lynnswap marked this conversation as resolved.
Comment thread
lynnswap marked this conversation as resolved.
do {
do {
try await persistReviewExecutionStart(id: next.id, at: startedAt)
} catch {
// The durable job is still queued when its execution-start write fails.
job.core.lifecycle.startedAt = nil
throw error
}
try Task.checkCancellation()
guard job.isTerminal == false else {
removeStartingReviewOwnership(for: next.id, ifOwnedBy: next.start.admission)
continue
}
guard let worker = makeReviewWorker(
jobID: job.id,
sessionID: job.sessionID,
request: next.start.request,
effectiveModel: next.start.model,
admission: next.start.admission
) else { throw CancellationError() }
reviewWorkerTasks[job.id] = worker
dispatchedWorker = worker
writeDiagnosticsIfNeeded()
} catch {
if job.isTerminal == false {
if error is CancellationError || Task.isCancelled {
try? completeCancellationLocally(
jobID: job.id,
sessionID: job.sessionID,
cancellation: job.pendingCancellationRequest?.cancellation ?? .system()
)
} else {
markReviewFailed(job, message: error.localizedDescription)
}
}
removeStartingReviewOwnership(for: next.id, ifOwnedBy: next.start.admission)
await waitForHistoryTerminalCommitIfNeeded(jobID: job.id)
resumeReviewWaiters(for: job.id)
}
}
}

private func persistReviewExecutionStart(id: String, at date: Date) async throws {
let persistence = historyPersistence
guard let receipt = historyMutationCoordinator.enqueue(
intent: id,
prepare: { $0 },
operation: { id in try await persistence.recordExecutionStarted(id: id, at: date) },
apply: { [weak self] _, result in
if case .failure(let error) = result {
self?.publishReviewHistoryFailure(error)
}
}
) else {
throw CodexReviewAPI.Error.io("Review history mutation admission is closed.")
}
switch await receipt.wait() {
case .success: return
case .failure(let error): throw error
}
}
}
Loading