diff --git a/Docs/deferred-codex-update-design.md b/Docs/deferred-codex-update-design.md new file mode 100644 index 00000000..1b498777 --- /dev/null +++ b/Docs/deferred-codex-update-design.md @@ -0,0 +1,158 @@ +# レビュー要求を Kit 内で待たせて Codex を更新する + +状態: 2026-09-26 承認済み。基準は `b19308aa72492797cb4fd7abfbbb28735b4d7cd0`。実装は [親 Issue #382](https://github.com/lynnswap/CodexReviewKit/issues/382) の sub-issue 順に進める。 + +「後でアップデート」を選ぶと、実行中のレビューを最後まで続け、新たな要求を `CodexReviewStore` のキューに受け付ける。実行中のレビューが終わったら Codex を更新し、app-server を再起動してキューを処理する。Monitor と MCP サーバーは動かしたままにする。 + +LLM は従来どおり `review_start` を一度呼ぶ。キュー待ちを理由に即座に応答を返したり、要求の再送・待機判断を LLM に任せたりしない。 + +## 受付から完了までを同じジョブとして扱う + +更新予約後も `review_start` を受け付け、既存のジョブ ID とセッション所有権を割り当てる。要求の妥当性と履歴の保存を確認した後、Monitor では `Queued` と表示する。まだ app-server のスレッドや turn は作成しない。 + +キューに積んだ要求は、実行中レビューの完了、更新、app-server の起動をサーバー内部で待つ。MCP 呼び出しは通常の完了待機につなぎ、同じジョブの最終結果を返す。状態確認とキャンセルは、更新中も通常の MCP API から使える。 + +```mermaid +sequenceDiagram + participant L as MCP クライアント + participant S as CodexReviewStore + participant C as app-server + participant U as Codex 更新処理 + Note over S: 「後でアップデート」を選択 + L->>S: review_start(B) + Note over S: B をキューで保持し、MCP 呼び出しを待機 + C-->>S: 実行中レビュー A の完了 + Note over S,C: A の結果保存と後片付けを完了し、旧プロセスを終了 + S->>U: codex update を一度実行 + U-->>S: 更新結果 + S->>C: 現在の CLI から app-server を起動・初期化 + S->>C: B を開始 + C-->>S: B の最終結果 + S-->>L: review_start(B) の結果 +``` + +「後で」の選択後に届いた要求をキューに保持するため、新しい要求が続いても更新機会は失われない。すでに起動処理に入っていたレビューも、完了を待つ対象に含める。履歴保存中の新規要求は、保存が終わった時点で現在の更新状態を確認してキューへ入れる。 + +再開時は受付順に既存のレビュー実行処理へ渡す。通常のレビュー並行実行は維持し、1件ずつ完了まで待つ新しい直列実行制約は設けない。受付順には履歴保存の完了順ではなく、Store が要求を受け付けた順序を使う。 + +### MCP の既存の待機上限は維持する + +現在の `review_start` と `review_await` は最大540秒で応答する。キュー待ちもこの呼び出し全体の待機時間に含める。HTTP の heartbeat は接続維持に使い、LLM 向けの新しいメッセージにはしない。 + +上限までに完了しなかった場合だけ、既存形式のジョブ ID と `queued` / `running` を返す。クライアントは従来の `review_await` で同じジョブを待てる。キューから削除したり、再度 `review_start` を要求したりしない。接続切断もキューの再投入には結び付けない。 + +この上限を超えても LLM の追加呼び出しを完全になくす保証は、本変更には含めない。固定タイムアウトを持つクライアント側の契約も変える必要があるためである。 + +## 実装済みの処理から変更する点 + +| 現在の処理 | 本変更で必要な動作 | +| --- | --- | +| `beginReview` は履歴を保存すると直ちに worker を作り、`running` にする。 | 要求の受付・保存と、worker の開始を分ける。更新予約中はジョブをキューで保持する。 | +| `hasRunningJobs` はすべての未完了ジョブや待機処理を含む。 | 更新開始の判定には、app-server を使用中のレビューとその後片付けを使う。新たにキューへ入ったジョブを完了待ちの対象にしない。 | +| `admitRuntimeReplacement` はレビュー受付を閉じ、履歴保存中の開始要求をキャンセルする。 | 計画的な更新では受付・参照・キャンセルを維持する。実行開始だけをキューで保留する。 | +| app-server 終了時には未完了レビューをまとめてキャンセルする。 | 更新時の終了処理は、実行前のキューをキャンセルしない。明示的なアプリ終了・サインアウトの契約とは区別する。 | +| 更新承認後に `store.shutdown()`、`codex update`、Monitor 全体の再起動を行う。 | Kit が既存 MCP セッションを保って app-server だけを終了・再起動する。 | +| 自動確認は起動時と確認後20時間の間隔で、設定画面には確認操作がない。 | 起動時と起動から8時間ごとに確認し、同じ checker を設定画面からも呼べるようにする。 | + +参照: + +- [レビューの受付と worker 開始](../Sources/CodexReview/Store/CodexReviewStoreReviews.swift) +- [Store の runtime 差し替えと未完了判定](../Sources/CodexReview/Store/CodexReviewStore.swift) +- [レビュー処理の登録と終了](../Sources/CodexReview/Store/ReviewStoreLifecycle.swift) +- [更新の実行と Monitor の再起動](../Tools/ReviewMonitor/CodexReviewMonitor/ReviewMonitorCodexUpdater.swift) +- [現在の停止確認](../Tools/ReviewMonitor/CodexReviewMonitor/CodexReviewMonitorApp.swift) + +## Store がキューと更新中の状態を所有する + +新しい package や target は作らない。`CodexReviewStore` をレビューと runtime 状態の正本とする現在の構成を維持する。 + +| 責務 | 所有者 | +| --- | --- | +| 更新予約、実行済みレビューの完了待ち、キュー、再開 | `CodexReviewStore` | +| ジョブの受付・実行開始・終端の保存 | 既存の履歴 coordinator と `CodexReviewPersistence` | +| app-server の終了・生成・初期化と設定の引き継ぎ | Store の runtime 差し替え処理と `CodexReviewHost` | +| 使用する CLI の解決とプロセス I/O | 既存の Host / AppServer 層 | +| 更新可否の確認、`codex update` の実行 | 既存の更新処理。Store に渡す依存として扱う | +| 起動時・8時間周期の確認、手動確認と確認結果 | `ReviewMonitorCodexUpdater`。設定画面も同じインスタンスを参照する | +| MCP の呼び出しと完了待機 | 既存の MCP adapter と Store | +| 選択肢、更新待ち・更新中・失敗の表示 | アプリと `ReviewUI`。Store の状態を表示する | + +アプリ側にレビューの一覧や再投入用の配列を持たせない。待機ジョブの対象・cwd・選択済みモデルは、受付時のレビュー要求から保持する。キューは同じジョブを参照し、UI 用のジョブを別途生成しない。 + +更新待ちの内部状態は Store が持つ単一の更新操作にまとめる。`ReviewMonitorCodexUpdater` に同じ段階やレビュー残数を複製しない。 + +### アプリは更新操作を一度呼ぶ + +アプリ向けの利用例は次の形とする。名前は提案で、外部向けの汎用メンテナンス API は追加しない。 + +```swift +try await store.updateCodex(when: .afterCurrentReviews) { + try await installer.install(plan) +} +``` + +`updateCodex` は `CodexReviewStore` の ApplicationHostSupport SPI として提供する。タイミングは `afterCurrentReviews` と `immediately` の2つとし、同じ処理で待機と即時更新を扱う。キューの停止・再開や複数の完了通知を呼び出し側に要求しない。 + +`installer.install(plan)` は `codex update` の完了を返す依存である。テストではこの依存と backend transport を差し替え、Store のキューや runtime 差し替えを本番と同じ経路で動かす。 + +## 更新はレビューの完了後に一度実行する + +更新操作は、受付済みで実行中のレビューについて、最終結果の保存と backend の後片付けまで待つ。既存の cleanup timeout は維持し、追加の無期限待機や「全未完了ジョブがゼロ」という条件は置かない。 + +その後、旧 app-server の終了を確認して `codex update` を実行する。更新中も MCP サーバー、セッション、ジョブ、待機中の呼び出しは存続する。更新後は既存の CLI 探索で実行ファイルを解決し、認証・設定を読み直して新しい runtime を公開する。公開成功後にキューを既存の worker 開始処理へ渡す。 + +本案ではダウンロードだけを先行させない。`codex update` は取得とインストールをまとめて行うため、使用中の補助プログラムまで入れ替わる操作をレビューと並行させない。専用のダウンロード状態や二重の更新キャッシュも追加しない。 + +更新予約中にボタンを再度押しても、更新プロセスやキューのコピーを増やさない。MCP の待機が上限に達しても、同じ更新操作とジョブを維持する。 + +### 失敗しても受け付けたジョブを失わない + +| 状況 | 動作 | +| --- | --- | +| 更新に失敗したが、現在の CLI で app-server を起動できる | 更新エラーを表示し、その runtime でキューを再開する。旧版へ戻せたとは断定しない。 | +| 旧 app-server の終了を確認できない | インストールを始めない。状態と残ったプロセスを報告し、キューを保持する。 | +| 更新後の app-server 起動・認証・設定復元に失敗 | MCP とキューを保持してエラーを表示する。既存の再試行から復旧し、成功後に再開する。 | +| 待機ジョブの `review_cancel`、または所有セッションの明示終了 | app-server に送る前にそのジョブをキャンセルし、履歴と待機中の呼び出しへ結果を返す。 | +| アプリを明示的に終了 | 既存の終了処理に参加する。更新コマンドの実行途中なら、その結果が不明にならないよう完了を待つ。 | + +アカウント切替・サインアウト・手動再起動は更新操作と直列化する。既存の確認とキャンセルの契約を保ち、別アカウントへ保留要求を黙って移さない。通常の状態参照まで更新完了待ちにはしない。 + +## キューを履歴から消さず、実行開始時刻も区別する + +待機中も一覧とキャンセルの対象になるため、受付済みジョブを既存の履歴保存に載せる。受付時刻と実行開始時刻を分け、`queued` の段階では実行時間を計測しない。保存失敗は受付失敗として呼び出し元へ返し、メモリだけに受理済みジョブを残さない。 + +既存の履歴にキュー状態と実行開始の保存を追加し、古いデータは現在と同じ実行済み状態として読めるようにする。新たな保存用ファイルは作らない。 + +今回保証するのは、Monitor を継続稼働させたままの app-server 更新である。Monitor 自体のクラッシュや終了後に、以前の MCP セッションのジョブを自動実行する機能は追加しない。既存の復元契約に沿って中断した履歴を表示し、要求の重複実行を避ける。 + +## 操作画面では「後でアップデート」を選べる + +レビュー中の確認には「後でアップデート」と、既存の「レビューを停止してアップデート」を用意する。「後で」は今回の更新を予約する選択肢で、単にダイアログを閉じる意味にはしない。 + +予約後はサイドバーに更新待ちを示し、新たなレビューは `Queued` と表示する。更新処理に進んだら更新中へ切り替え、完了後に通常表示へ戻す。実行中のレビューがなければ、そのまま更新へ進める。 + +既存のウィンドウ、ログ、アカウント表示は保つ。待機・更新中・失敗・再開を `#Preview` と UI テストから確認できるようにする。 + +## 起動時と8時間ごとに確認し、設定画面からも確認できる + +自動確認は Monitor 起動時に一度、その後は起動から8時間、16時間、24時間という周期で行う。確認処理にかかった時間や、手動で確認した時刻を次回の起点にしない。スリープなどで予定時刻を過ぎた場合は復帰後に一度確認し、過ぎた回数だけ連続実行しない。 + +アプリ全体で一つの確認処理を共有する。自動確認と手動確認が重なった場合は進行中の確認に合流する。Codex の入れ替え中に確認予定が来た場合も、同じ更新対象へ並行してコマンドを起動せず、更新後に一度確認する。 + +設定画面に Codex のアップデート確認用 UI を追加する。手動ボタンは確認だけを行い、インストールやレビューの停止・キュー予約を始めない。表示する結果は「確認中」「更新あり」「最新」「このインストールでは確認できない」「確認失敗」を区別し、最終確認時刻とエラーも示す。CLI の選択失敗や未対応のインストール方法を「最新」と表示しない。 + +確認結果は既存のサイドバーの Update 表示にも反映する。設定画面を閉じても8時間周期は継続し、設定画面を開くたびに監視タスクを作らない。Monitor の確認周期は今回の要求を正本とし、CLI 用の起動時チェック設定を理由に Monitor の手動確認や周期確認を停止しない。 + +## 3つの実装単位で検証する + +1. **Kit のジョブ受付と実行開始を分離する。** キュー、履歴保存、キャンセル、同じ MCP 呼び出しの完了待機を実装する。実行の許可がある通常時には、既存の並行実行を保つ。 +2. **MCP を維持する更新操作を実装する。** 実行中レビューの完了待ち、旧 runtime の終了、更新依存の実行、新 runtime の公開、キュー再開を Store が一度の操作として所有する。 +3. **更新 UI と確認周期を接続して旧再起動経路を削除する。** 新しい選択肢、起動時・8時間ごとの確認、設定画面の手動確認を接続する。Codex 更新専用の Monitor 再起動ヘルパー・失敗起動引数を削除し、Preview と利用文書を更新する。 + +親 Issue で順序と完了条件を管理し、1単位を1つの Ready PR にする。各 PR はローカル codex-review、リモートレビュー、CI を通してから main へ統合する。 + +回帰テストでは、更新予約とレビュー受付が同時に起こる順序、履歴保存中の要求、複数セッションのキュー、待機ジョブのキャンセル、更新・終了・再起動の失敗を確認する。既存の fake backend、transport、時計、gate を使い、時間経過への期待だけで同期しない。 + +最終確認では、一つの MCP セッションでレビュー A を実行中に後回し更新を予約し、B と C を受け付ける。A の結果を返した後に更新が一度だけ実行され、同じセッション・ジョブ ID で B と C が完了することを確かめる。待機上限に届かない条件では、クライアントの再送・追加の待機ツール呼び出しが不要であることも検証する。 + +更新チェックは手動時計で起動直後・8時間・16時間の実行を確認する。途中の手動確認で周期が変わらないこと、確認が重複しないこと、設定画面の表示結果とサイドバーが一致することも対象にする。 diff --git a/Sources/CodexReview/History/CodexReviewStoreHistory.swift b/Sources/CodexReview/History/CodexReviewStoreHistory.swift index 8a605f93..9d71e948 100644 --- a/Sources/CodexReview/History/CodexReviewStoreHistory.swift +++ b/Sources/CodexReview/History/CodexReviewStoreHistory.swift @@ -79,7 +79,7 @@ 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), @@ -87,7 +87,7 @@ extension CodexReviewStore { sortOrder: try nextHistoryJobSortOrder(), target: request.target, model: model, - startedAt: clock.now() + acceptedAt: clock.now() ) let receipt = HistoryStartReceipt( ordinal: nextHistoryStartOrdinal, @@ -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 { diff --git a/Sources/CodexReview/History/HistoryStartReceipt.swift b/Sources/CodexReview/History/HistoryStartReceipt.swift index 2fb3e963..4a5dbf2d 100644 --- a/Sources/CodexReview/History/HistoryStartReceipt.swift +++ b/Sources/CodexReview/History/HistoryStartReceipt.swift @@ -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? @@ -12,7 +12,7 @@ package final class HistoryStartReceipt { package init( ordinal: UInt64, sessionID: String, - started: StartedReviewRecord, + started: AcceptedReviewRecord, workAdmission: ReviewStoreWorkRegistry.Admission ) { self.ordinal = ordinal diff --git a/Sources/CodexReview/History/ReviewHistoryPersistence.swift b/Sources/CodexReview/History/ReviewHistoryPersistence.swift index e0ccb127..76edcf60 100644 --- a/Sources/CodexReview/History/ReviewHistoryPersistence.swift +++ b/Sources/CodexReview/History/ReviewHistoryPersistence.swift @@ -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, @@ -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, diff --git a/Sources/CodexReview/History/ReviewHistoryRecord.swift b/Sources/CodexReview/History/ReviewHistoryRecord.swift index dd589ca3..3fa9457b 100644 --- a/Sources/CodexReview/History/ReviewHistoryRecord.swift +++ b/Sources/CodexReview/History/ReviewHistoryRecord.swift @@ -12,7 +12,7 @@ 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? @@ -20,7 +20,8 @@ package struct StartedReviewRecord: Sendable, Hashable { 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, @@ -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.") @@ -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 } } @@ -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 { @@ -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) ) } } @@ -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? diff --git a/Sources/CodexReview/Store/CodexReviewStore.swift b/Sources/CodexReview/Store/CodexReviewStore.swift index 026f6a61..c8e90988 100644 --- a/Sources/CodexReview/Store/CodexReviewStore.swift +++ b/Sources/CodexReview/Store/CodexReviewStore.swift @@ -146,6 +146,10 @@ public final class CodexReviewStore { @ObservationIgnored package var applicationShutdownTask: Task? @ObservationIgnored package var reviewAttemptOwnerships: [String: StoreReviewAttemptOwnership] = [:] @ObservationIgnored package var reviewWorkerTasks: [String: Task] = [:] + @ObservationIgnored package var queuedReviewStarts: [String: QueuedReviewStart] = [:] + @ObservationIgnored package var reviewStartsAreSuspended = false + @ObservationIgnored package var queuedReviewDispatchTask: Task? + @ObservationIgnored package var runtimeStopDetachedReviewWorkerTasks: [String: Task] = [:] @ObservationIgnored package var reviewTerminalWaiters: [String: [ReviewTerminalWaiter]] = [:] @ObservationIgnored package var nextCancellationRequestOrdinal: UInt64 = 0 diff --git a/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift b/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift index 58cbc4f4..ee64f43d 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreCancellation.swift @@ -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 diff --git a/Sources/CodexReview/Store/CodexReviewStoreQueue.swift b/Sources/CodexReview/Store/CodexReviewStoreQueue.swift new file mode 100644 index 00000000..1a316e77 --- /dev/null +++ b/Sources/CodexReview/Store/CodexReviewStoreQueue.swift @@ -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?) + } + private var dispatch = Dispatch.pending + private var dispatchWaiters: [CheckedContinuation] = [] + + 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? = 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? + 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." + 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 + } + } +} diff --git a/Sources/CodexReview/Store/CodexReviewStoreReviews.swift b/Sources/CodexReview/Store/CodexReviewStoreReviews.swift index 0c3f66a6..ece30d0e 100644 --- a/Sources/CodexReview/Store/CodexReviewStoreReviews.swift +++ b/Sources/CodexReview/Store/CodexReviewStoreReviews.swift @@ -66,11 +66,12 @@ extension CodexReviewStore { waitTimeout: Duration?, workAdmission: ReviewStoreWorkRegistry.Admission ) async throws -> CodexReviewAPI.Read.Result { - let jobID = try await beginReview( + let accepted = try await beginReview( sessionID: sessionID, request: request, workAdmission: workAdmission ) + let jobID = accepted.jobID let historyResultLease = acquireHistoryResultLease(jobID: jobID) defer { releaseHistoryResultLease(historyResultLease) @@ -87,14 +88,13 @@ extension CodexReviewStore { await reviewWorkerTasks[jobID]?.value return try readReview(sessionID: sessionID, jobID: jobID) } - let workerTask = reviewWorkerTasks[jobID] _ = try await awaitReview( sessionID: sessionID, jobID: jobID, timeout: waitTimeout ) if storeWorkRegistry.accepts(workAdmission) == false { - await workerTask?.value + await accepted.start.waitUntilFinished() } return try readReview(sessionID: sessionID, jobID: jobID) } @@ -123,7 +123,7 @@ extension CodexReviewStore { sessionID: String, request: CodexReviewAPI.Start.Request, workAdmission: ReviewStoreWorkRegistry.Admission - ) async throws -> String { + ) async throws -> (jobID: String, start: QueuedReviewStart) { guard closedSessions.contains(sessionID) == false else { throw CodexReviewAPI.Error.invalidArguments("Review session \(sessionID) is closed.") } @@ -136,6 +136,7 @@ extension CodexReviewStore { ) defer { finishHistoryStartReceipt(receipt) + scheduleQueuedReviewStarts() } switch await persistHistoryStart(receipt) { case .success: @@ -164,40 +165,29 @@ extension CodexReviewStore { target: started.target, core: .init( run: .init(model: started.model), - lifecycle: .init(status: .running, startedAt: started.startedAt), - output: .init(summary: "Review started.") + lifecycle: .init(status: .queued), + output: .init(summary: "Review queued.") ), logEntries: [] ) let admission = ReviewStartAdmission() - guard let workerTask = makeReviewWorker( - jobID: jobID, - sessionID: sessionID, - request: validatedRequest, - effectiveModel: started.model, - admission: admission - ) else { - let cancellation = ReviewCancellation.system( - message: "Review start was cancelled before backend dispatch." - ) - await terminalizeStaleHistoryStart( - receipt, - cancellation: cancellation - ) - throw CodexReviewAPI.Error.io("Review Store work admission is closed.") - } insertReviewJob( job, workspaceMetadata: started.workspaceMetadata, workspaceSortOrder: started.workspaceSortOrder ) reviewAttemptOwnerships[jobID] = .starting(admission) - reviewWorkerTasks[jobID]?.cancel() - reviewWorkerTasks[jobID] = workerTask - return jobID + let queued = QueuedReviewStart( + ordinal: receipt.ordinal, + request: validatedRequest, + model: started.model, + admission: admission + ) + queuedReviewStarts[jobID] = queued + return (jobID, queued) } - private func makeReviewWorker( + package func makeReviewWorker( jobID: String, sessionID: String, request: CodexReviewAPI.Start.Request, @@ -387,7 +377,7 @@ extension CodexReviewStore { } @discardableResult - private func removeStartingReviewOwnership( + package func removeStartingReviewOwnership( for jobID: String, ifOwnedBy admission: ReviewStartAdmission ) -> Bool { @@ -1380,7 +1370,7 @@ extension CodexReviewStore { writeDiagnosticsIfNeeded() } - private func markReviewFailed( + package func markReviewFailed( _ job: CodexReviewJob, message: String?, terminal: ReviewTerminalRecord? = nil diff --git a/Sources/CodexReviewHost/ReviewHistoryLocation.swift b/Sources/CodexReviewHost/ReviewHistoryLocation.swift index c0a1c5b1..dd590c72 100644 --- a/Sources/CodexReviewHost/ReviewHistoryLocation.swift +++ b/Sources/CodexReviewHost/ReviewHistoryLocation.swift @@ -120,8 +120,12 @@ package actor OwnedReviewHistoryPersistence: ReviewHistoryPersistence { try await database.load(retentionPolicy: retentionPolicy) } - package func recordStarted(_ record: StartedReviewRecord) async throws { - try await database.recordStarted(record) + package func recordAccepted(_ record: AcceptedReviewRecord) async throws { + try await database.recordAccepted(record) + } + + package func recordExecutionStarted(id: String, at date: Date) async throws { + try await database.recordExecutionStarted(id: id, at: date) } package func recordTerminal( @@ -194,7 +198,11 @@ package struct UnavailableReviewHistoryPersistence: ReviewHistoryPersistence { throw failure } - package func recordStarted(_: StartedReviewRecord) async throws { + package func recordAccepted(_: AcceptedReviewRecord) async throws { + throw failure + } + + package func recordExecutionStarted(id: String, at date: Date) async throws { throw failure } diff --git a/Sources/CodexReviewPersistence/ReviewHistoryDatabase.swift b/Sources/CodexReviewPersistence/ReviewHistoryDatabase.swift index 67b0d9b0..91ccc794 100644 --- a/Sources/CodexReviewPersistence/ReviewHistoryDatabase.swift +++ b/Sources/CodexReviewPersistence/ReviewHistoryDatabase.swift @@ -50,7 +50,7 @@ package actor ReviewHistoryDatabase: ReviewHistoryPersistence { } } - package func recordStarted(_ record: StartedReviewRecord) async throws { + package func recordAccepted(_ record: AcceptedReviewRecord) async throws { let database = try preparedDatabase() let timestamp = now() let encoded = try ReviewHistoryRecordCodec.encodeStarted( @@ -68,6 +68,23 @@ package actor ReviewHistoryDatabase: ReviewHistoryPersistence { } } + package func recordExecutionStarted(id: String, at date: Date) async throws { + let database = try preparedDatabase() + let timestamp = ReviewHistoryTimestamp.encode(date) + try write(database) { db in + try ReviewRecordRow.where { $0.id.eq(id) && $0.phase.eq("queued") } + .update { + $0.phase = "active" + $0.startedAt = #bind(timestamp) + $0.updatedAt = #bind(timestamp) + } + .execute(db) + guard db.changesCount == 1 else { + throw ReviewHistoryDatabaseError.invalidRecord(id: id, reason: "review is not queued") + } + } + } + package func recordTerminal( _ record: TerminalReviewRecord, retentionPolicy: ReviewHistoryRetentionPolicy @@ -79,7 +96,7 @@ package actor ReviewHistoryDatabase: ReviewHistoryPersistence { else { throw ReviewHistoryDatabaseError.reviewNotFound(record.id) } - guard existing.phase == "active" else { + guard existing.phase == "active" || existing.phase == "queued" else { throw ReviewHistoryDatabaseError.invalidRecord( id: record.id, reason: "terminal state can only be recorded once" @@ -264,7 +281,7 @@ package actor ReviewHistoryDatabase: ReviewHistoryPersistence { in db: Database, committedAt: Date ) throws { - for row in try ReviewRecordRow.fetchAll(db) where row.phase == "active" { + for row in try ReviewRecordRow.fetchAll(db) where row.phase == "active" || row.phase == "queued" { let terminal = try TerminalReviewRecord( id: row.id, model: nil, @@ -441,7 +458,7 @@ package actor ReviewHistoryDatabase: ReviewHistoryPersistence { _ row: ReviewRecordRow, in db: Database ) throws { - try ReviewRecordRow.find(row.id).where { $0.phase.eq("active") }.update { + try ReviewRecordRow.find(row.id).where { $0.phase.eq("active") || $0.phase.eq("queued") }.update { $0.phase = #bind(row.phase) $0.terminalModel = #bind(row.terminalModel) $0.reviewThreadID = #bind(row.reviewThreadID) @@ -571,6 +588,7 @@ private extension ReviewRecordRow { if targetCommitTitle != other.targetCommitTitle { columns.append("targetCommitTitle") } if targetInstructions != other.targetInstructions { columns.append("targetInstructions") } if startedModel != other.startedModel { columns.append("startedModel") } + if acceptedAt != other.acceptedAt { columns.append("acceptedAt") } if startedAt != other.startedAt { columns.append("startedAt") } if phase != other.phase { columns.append("phase") } if terminalModel != other.terminalModel { columns.append("terminalModel") } diff --git a/Sources/CodexReviewPersistence/ReviewHistoryRecordCodec.swift b/Sources/CodexReviewPersistence/ReviewHistoryRecordCodec.swift index bcc26fdc..5724bbbe 100644 --- a/Sources/CodexReviewPersistence/ReviewHistoryRecordCodec.swift +++ b/Sources/CodexReviewPersistence/ReviewHistoryRecordCodec.swift @@ -12,13 +12,13 @@ struct EncodedTerminalReview { } enum DecodedReviewHistoryRow { - case active(StartedReviewRecord) + case active(AcceptedReviewRecord) case terminal(RestoredReviewRecord) } enum ReviewHistoryRecordCodec { static func encodeStarted( - _ record: StartedReviewRecord, + _ record: AcceptedReviewRecord, createdAt: Date, updatedAt: Date ) throws -> EncodedStartedReview { @@ -42,8 +42,9 @@ enum ReviewHistoryRecordCodec { targetCommitTitle: targetColumns.commitTitle, targetInstructions: targetColumns.instructions, startedModel: record.model, - startedAt: ReviewHistoryTimestamp.encode(record.startedAt), - phase: "active", + acceptedAt: ReviewHistoryTimestamp.encode(record.acceptedAt), + startedAt: record.startedAt.map(ReviewHistoryTimestamp.encode), + phase: record.startedAt == nil ? "queued" : "active", terminalModel: nil, terminalKind: nil, interruptionKind: nil, @@ -114,7 +115,7 @@ enum ReviewHistoryRecordCodec { findings: [ReviewFindingRow] ) throws -> DecodedReviewHistoryRow { switch row.phase { - case "active": + case "queued", "active": guard findings.isEmpty else { throw invalid(row.id, "active row contains terminal findings") } @@ -133,15 +134,15 @@ enum ReviewHistoryRecordCodec { static func decodeStarted( _ row: ReviewRecordRow, workspace: ReviewWorkspaceRow - ) throws -> StartedReviewRecord { + ) throws -> AcceptedReviewRecord { guard row.cwd == workspace.cwd else { throw invalid(row.id, "workspace foreign key does not match loaded workspace") } let target = try decodeTarget(row) let workspaceMetadata = try decodeWorkspaceMetadata(workspace, reviewID: row.id) - let started: StartedReviewRecord + let started: AcceptedReviewRecord do { - started = try StartedReviewRecord( + started = try AcceptedReviewRecord( id: row.id, cwd: row.cwd, workspaceMetadata: workspaceMetadata, @@ -149,7 +150,8 @@ enum ReviewHistoryRecordCodec { sortOrder: row.sortOrder, target: target, model: row.startedModel, - startedAt: ReviewHistoryTimestamp.decode(row.startedAt) + acceptedAt: ReviewHistoryTimestamp.decode(row.acceptedAt), + startedAt: row.startedAt.map(ReviewHistoryTimestamp.decode) ) } catch { throw invalid(row.id, error.localizedDescription) @@ -158,7 +160,7 @@ enum ReviewHistoryRecordCodec { throw invalid(row.id, "started model is not canonical") } - if row.phase == "active" { + if row.phase == "active" || row.phase == "queued" { let reencoded = try encodeStarted( started, createdAt: ReviewHistoryTimestamp.decode(row.createdAt), diff --git a/Sources/CodexReviewPersistence/ReviewHistorySchema.swift b/Sources/CodexReviewPersistence/ReviewHistorySchema.swift index 6de3c867..6f134c29 100644 --- a/Sources/CodexReviewPersistence/ReviewHistorySchema.swift +++ b/Sources/CodexReviewPersistence/ReviewHistorySchema.swift @@ -25,7 +25,8 @@ struct ReviewRecordRow: Equatable, Sendable { var targetInstructions: String? var startedModel: String? - var startedAt: Double + var acceptedAt: Double + var startedAt: Double? var phase: String var terminalModel: String? @@ -575,6 +576,226 @@ enum ReviewHistorySchema { .execute(db) } + migrator.registerMigration("v6_queue_review_starts") { db in + try #sql( + """ + CREATE TABLE "review_records_v6" ( + "id" TEXT NOT NULL PRIMARY KEY, + "cwd" TEXT NOT NULL REFERENCES "review_workspaces"("cwd") ON DELETE CASCADE, + "sortOrder" REAL NOT NULL, + + "targetKind" TEXT NOT NULL, + "targetBranch" TEXT, + "targetCommitSHA" TEXT, + "targetCommitTitle" TEXT, + "targetInstructions" TEXT, + + "startedModel" TEXT, + "acceptedAt" REAL NOT NULL, + "startedAt" REAL, + + "phase" TEXT NOT NULL, + "terminalModel" TEXT, + "terminalKind" TEXT, + "interruptionKind" TEXT, + "cancellationSource" TEXT, + "cancellationMessage" TEXT, + "terminalMessage" TEXT, + "endedAt" REAL, + "summary" TEXT, + "canonicalReview" TEXT, + "parsedState" TEXT, + "parsedFindingCount" INTEGER, + "parsedSource" TEXT, + "parserVersion" INTEGER, + "terminalCommittedAt" REAL, + + "createdAt" REAL NOT NULL, + "updatedAt" REAL NOT NULL, + "reviewThreadID" TEXT, + "threadID" TEXT, + + CHECK (("phase" != 'queued' OR "startedAt" IS NULL) + AND ("phase" != 'active' OR "startedAt" IS NOT NULL)), + CHECK (COALESCE( + ("targetKind" = 'uncommittedChanges' + AND "targetBranch" IS NULL + AND "targetCommitSHA" IS NULL + AND "targetCommitTitle" IS NULL + AND "targetInstructions" IS NULL) + OR ("targetKind" = 'baseBranch' + AND "targetBranch" IS NOT NULL + AND length(trim("targetBranch")) > 0 + AND "targetCommitSHA" IS NULL + AND "targetCommitTitle" IS NULL + AND "targetInstructions" IS NULL) + OR ("targetKind" = 'commit' + AND "targetBranch" IS NULL + AND "targetCommitSHA" IS NOT NULL + AND length(trim("targetCommitSHA")) > 0 + AND "targetInstructions" IS NULL) + OR ("targetKind" = 'custom' + AND "targetBranch" IS NULL + AND "targetCommitSHA" IS NULL + AND "targetCommitTitle" IS NULL + AND "targetInstructions" IS NOT NULL + AND length(trim("targetInstructions")) > 0) + , 0)), + CHECK (COALESCE( + ("phase" IN ('queued', 'active') + AND "terminalModel" IS NULL + AND "terminalKind" IS NULL + AND "interruptionKind" IS NULL + AND "cancellationSource" IS NULL + AND "cancellationMessage" IS NULL + AND "terminalMessage" IS NULL + AND "endedAt" IS NULL + AND "summary" IS NULL + AND "canonicalReview" IS NULL + AND "parsedState" IS NULL + AND "parsedFindingCount" IS NULL + AND "parsedSource" IS NULL + AND "parserVersion" IS NULL + AND "terminalCommittedAt" IS NULL) + OR ("phase" = 'terminal' + AND "summary" IS NOT NULL + AND "terminalCommittedAt" IS NOT NULL + AND "terminalKind" IS NOT NULL + AND ( + ("terminalKind" = 'completed' + AND "endedAt" IS NOT NULL + AND "interruptionKind" IS NULL + AND "cancellationSource" IS NULL + AND "cancellationMessage" IS NULL + AND "terminalMessage" IS NULL + AND "canonicalReview" IS NOT NULL + AND length(trim("canonicalReview")) > 0 + AND "parsedState" IS NOT NULL + AND "parsedSource" IS NOT NULL + AND "parserVersion" IS NOT NULL + AND "parserVersion" > 0 + AND ( + ("parsedState" = 'hasFindings' + AND "parsedFindingCount" > 0 + AND "parsedSource" = 'parsedFinalReviewText') + OR ("parsedState" = 'noFindings' + AND "parsedFindingCount" = 0 + AND "parsedSource" = 'parsedFinalReviewText') + OR ("parsedState" = 'unknown' + AND "parsedFindingCount" IS NULL + AND "parsedSource" IN ('unrecognizedFindingBlock', 'notAvailable')) + )) + OR ("terminalKind" = 'interrupted' + AND "interruptionKind" = 'requested' + AND "endedAt" IS NOT NULL + AND "cancellationSource" IN + ('userInterface', 'mcpClient', 'sessionClosed', 'system') + AND "cancellationMessage" IS NOT NULL + AND "terminalMessage" IS NULL + AND "canonicalReview" IS NULL + AND "parsedState" IS NULL + AND "parsedFindingCount" IS NULL + AND "parsedSource" IS NULL + AND "parserVersion" IS NULL) + OR ("terminalKind" = 'interrupted' + AND "interruptionKind" = 'server' + AND "endedAt" IS NOT NULL + AND "cancellationSource" IS NULL + AND "cancellationMessage" IS NULL + AND "canonicalReview" IS NULL + AND "parsedState" IS NULL + AND "parsedFindingCount" IS NULL + AND "parsedSource" IS NULL + AND "parserVersion" IS NULL) + OR ("terminalKind" = 'interrupted' + AND "interruptionKind" = 'transport' + AND "endedAt" IS NOT NULL + AND "cancellationSource" IS NULL + AND "cancellationMessage" IS NULL + AND "terminalMessage" IS NOT NULL + AND "canonicalReview" IS NULL + AND "parsedState" IS NULL + AND "parsedFindingCount" IS NULL + AND "parsedSource" IS NULL + AND "parserVersion" IS NULL) + OR ("terminalKind" = 'interrupted' + AND "interruptionKind" = 'previousProcessExit' + AND "cancellationSource" IS NULL + AND "cancellationMessage" IS NULL + AND "terminalMessage" IS NULL + AND "canonicalReview" IS NULL + AND "parsedState" IS NULL + AND "parsedFindingCount" IS NULL + AND "parsedSource" IS NULL + AND "parserVersion" IS NULL) + OR ("terminalKind" = 'failed' + AND "endedAt" IS NOT NULL + AND "interruptionKind" IS NULL + AND "cancellationSource" IS NULL + AND "cancellationMessage" IS NULL + AND "canonicalReview" IS NULL + AND "parsedState" IS NULL + AND "parsedFindingCount" IS NULL + AND "parsedSource" IS NULL + AND "parserVersion" IS NULL) + )) + , 0)) + ) STRICT + """ + ) + .execute(db) + + try #sql( + """ + INSERT INTO "review_records_v6" ( + "id", "cwd", "sortOrder", + "targetKind", "targetBranch", "targetCommitSHA", "targetCommitTitle", + "targetInstructions", "startedModel", "acceptedAt", "startedAt", "phase", + "terminalModel", "terminalKind", "interruptionKind", + "cancellationSource", "cancellationMessage", "terminalMessage", "endedAt", + "summary", "canonicalReview", "parsedState", "parsedFindingCount", + "parsedSource", "parserVersion", "terminalCommittedAt", "createdAt", "updatedAt", "reviewThreadID", "threadID" + ) + SELECT + "id", "cwd", "sortOrder", + "targetKind", "targetBranch", "targetCommitSHA", "targetCommitTitle", + "targetInstructions", "startedModel", "startedAt", "startedAt", "phase", + "terminalModel", "terminalKind", "interruptionKind", + "cancellationSource", "cancellationMessage", "terminalMessage", "endedAt", + "summary", "canonicalReview", "parsedState", "parsedFindingCount", + "parsedSource", "parserVersion", "terminalCommittedAt", "createdAt", "updatedAt", "reviewThreadID", "threadID" + FROM "review_records" + """ + ) + .execute(db) + + // GRDB owns deferred foreign-key validation for this migration. Leaving the child + // table in place preserves finding rows while the parent table is replaced. + try #sql("DROP TABLE \"review_records\"").execute(db) + try #sql( + """ + ALTER TABLE "review_records_v6" + RENAME TO "review_records" + """ + ) + .execute(db) + + try #sql( + """ + CREATE INDEX "review_records_order" + ON "review_records" ("sortOrder", "id") + """ + ) + .execute(db) + try #sql( + """ + CREATE INDEX "review_records_retention" + ON "review_records" ("phase", "terminalCommittedAt", "id") + """ + ) + .execute(db) + } + return migrator } diff --git a/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift b/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift index 50d4f4e6..473e2858 100644 --- a/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift +++ b/Tests/CodexReviewMCPServerTests/CodexReviewMCPHTTPServerTests.swift @@ -10,6 +10,38 @@ import CodexReviewTesting @Suite("MCP Streamable HTTP server") @MainActor struct CodexReviewMCPHTTPServerTests { + @Test func reviewStartRemainsPendingAcrossKitQueueWithoutClientResubmission() async throws { + let backend = FakeCodexReviewBackend() + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend), + idGenerator: .init(next: { "queued-job" }) + ) + store.suspendReviewStarts() + try await withHTTPServer(store: store) { server in + let endpoint = await server.url + let sessionID = try await initializeSession(endpoint: endpoint) + let body = try makeReviewStartBody(id: 2) + var returned = false + let request = Task { + defer { returned = true } + return try await postJSONRPCData(endpoint: endpoint, sessionID: sessionID, bodyData: body) + } + defer { request.cancel() } + try #require(await waitUntil(timeout: .seconds(2)) { + store.job(id: "queued-job")?.core.lifecycle.status == .queued + }) + #expect(returned == false) + #expect(store.activeJobIDs(for: sessionID) == ["queued-job"]) + #expect(await backend.recordedCommands().contains { if case .startReview = $0 { true } else { false } } == false) + store.resumeReviewStarts() + try await backend.waitForStartReview(timeout: .seconds(2)) + await backend.yield(.completed(summary: "Done", result: "No findings.")) + let result = try decodeSSEJSON(from: try await request.value) + #expect(result.value(for: ["result", "structuredContent", "jobId"]) as? String == "queued-job") + #expect(result.value(for: ["result", "structuredContent", "lifecycle", "status"]) as? String == "succeeded") + } + } + @Test func streamableHTTPInitializesAndListsTools() async throws { let backend = FakeCodexReviewBackend() let store = CodexReviewStore.makeTestingStore( diff --git a/Tests/CodexReviewPersistenceTests/ReviewHistoryDatabaseTests.swift b/Tests/CodexReviewPersistenceTests/ReviewHistoryDatabaseTests.swift index 1f373e6c..58e43f4f 100644 --- a/Tests/CodexReviewPersistenceTests/ReviewHistoryDatabaseTests.swift +++ b/Tests/CodexReviewPersistenceTests/ReviewHistoryDatabaseTests.swift @@ -6,6 +6,47 @@ import Testing @Suite("ReviewHistoryDatabase") struct ReviewHistoryDatabaseTests { + @Test("persists queue acceptance independently of execution start") + @MainActor + func queuedReviewPersistsWithoutAnExecutionTimestamp() async throws { + let (database, writer) = try ReviewHistoryTestSupport.database() + let acceptedAt = ReviewHistoryTestSupport.startedAt + let queued = try AcceptedReviewRecord( + id: "queued", cwd: "/tmp/queued", workspaceSortOrder: 0, sortOrder: 0, + target: .uncommittedChanges, model: "gpt-5", acceptedAt: acceptedAt + ) + try await database.recordAccepted(queued) + let row = try await writer.read { db in try ReviewRecordRow.find("queued").fetchOne(db) } + #expect(row?.phase == "queued") + #expect(row?.startedAt == nil) + #expect(row?.acceptedAt == ReviewHistoryTimestamp.encode(acceptedAt)) + + let executionDate = acceptedAt.addingTimeInterval(120) + try await database.recordExecutionStarted(id: queued.id, at: executionDate) + _ = try await database.recordTerminal( + ReviewHistoryTestSupport.completed(id: queued.id, endedAt: executionDate.addingTimeInterval(10)), + retentionPolicy: .default + ) + let restored = try #require(try await database.load(retentionPolicy: .default).first) + #expect(restored.started.acceptedAt == acceptedAt) + #expect(restored.started.startedAt == executionDate) + #expect(restored.makeRestoredJob().core.lifecycle.startedAt == executionDate) + } + + @Test("restores an undispatched queue entry as interrupted without inventing execution") + @MainActor + func queuedReviewIsNotAutomaticallyReplayedAfterProcessExit() async throws { + let (database, _) = try ReviewHistoryTestSupport.database() + try await database.recordAccepted(AcceptedReviewRecord( + id: "queued", cwd: "/tmp/queued", workspaceSortOrder: 0, sortOrder: 0, + target: .uncommittedChanges, model: nil, acceptedAt: ReviewHistoryTestSupport.startedAt + )) + let restored = try #require(try await database.load(retentionPolicy: .default).first) + #expect(restored.terminal.terminal == .interrupted(.previousProcessExit)) + #expect(restored.started.startedAt == nil) + #expect(restored.makeRestoredJob().core.lifecycle.startedAt == nil) + } + @Test("preserves deep-link identities across a file-backed relaunch", arguments: [ (nil, nil), (nil, "thread"), @@ -338,7 +379,7 @@ struct ReviewHistoryDatabaseTests { @Test("converts abandoned active rows without inventing an end time") func orphanConversion() async throws { let (database, _) = try ReviewHistoryTestSupport.database() - try await database.recordStarted(ReviewHistoryTestSupport.started(id: "orphan")) + try await database.recordAccepted(ReviewHistoryTestSupport.started(id: "orphan")) let restored = try await database.load(retentionPolicy: .default) let orphan = try #require(restored.first) @@ -351,14 +392,14 @@ struct ReviewHistoryDatabaseTests { @Test("rejects duplicate application-wide order on start insertion") func duplicateStartOrder() async throws { let (database, writer) = try ReviewHistoryTestSupport.database() - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "first", cwd: "/tmp/first", sortOrder: 0 )) await #expect(throws: ReviewHistoryDatabaseError.self) { - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "duplicate", cwd: "/tmp/duplicate", sortOrder: 0 @@ -499,7 +540,7 @@ struct ReviewHistoryDatabaseTests { @Test("terminal commit preserves admission identity and latest manual order") func terminalPreservesStartedFields() async throws { let (database, _) = try ReviewHistoryTestSupport.database() - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "ordered", cwd: "/tmp/original", workspaceSortOrder: 1, @@ -527,7 +568,7 @@ struct ReviewHistoryDatabaseTests { @Test("recording a start preserves current workspace order") func startPreservesWorkspaceOrder() async throws { let (database, writer) = try ReviewHistoryTestSupport.database() - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "existing", cwd: "/tmp/shared", workspaceSortOrder: 1, @@ -538,13 +579,13 @@ struct ReviewHistoryDatabaseTests { reviews: [.init(id: "existing", sortOrder: 30)] )) - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "new-shared", cwd: "/tmp/shared", workspaceSortOrder: 1, sortOrder: 31 )) - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "new-workspace", cwd: "/tmp/new", workspaceSortOrder: 40, @@ -659,13 +700,13 @@ struct ReviewHistoryDatabaseTests { displayTitle: "New", kind: .linkedWorktree ) - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "old-generation", cwd: "/tmp/reused", workspaceMetadata: oldMetadata, workspaceSortOrder: 20 )) - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "new-generation", cwd: "/tmp/reused", workspaceMetadata: newMetadata, @@ -715,7 +756,7 @@ struct ReviewHistoryDatabaseTests { @Test("terminal-only batch deletes preserve active rows and return exact membership") func terminalDeletionSemantics() async throws { let (database, writer) = try ReviewHistoryTestSupport.database() - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "active", cwd: "/tmp/active", sortOrder: 0 @@ -788,7 +829,7 @@ struct ReviewHistoryDatabaseTests { @Test("delete-all removes only terminal rows and reports every removed ID") func deleteAllTerminalReviews() async throws { let (database, writer) = try ReviewHistoryTestSupport.database() - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: "active", sortOrder: 0 )) @@ -822,7 +863,7 @@ struct ReviewHistoryDatabaseTests { _ = try await database.load(retentionPolicy: .default) } await #expect(throws: ReviewHistoryDatabaseError.closed) { - try await database.recordStarted(ReviewHistoryTestSupport.started(id: "closed")) + try await database.recordAccepted(ReviewHistoryTestSupport.started(id: "closed")) } await #expect(throws: ReviewHistoryDatabaseError.closed) { _ = try await database.recordTerminal( diff --git a/Tests/CodexReviewPersistenceTests/ReviewHistorySchemaTests.swift b/Tests/CodexReviewPersistenceTests/ReviewHistorySchemaTests.swift index 9203ad52..7e4a1674 100644 --- a/Tests/CodexReviewPersistenceTests/ReviewHistorySchemaTests.swift +++ b/Tests/CodexReviewPersistenceTests/ReviewHistorySchemaTests.swift @@ -84,6 +84,7 @@ struct ReviewHistorySchemaTests { "targetCommitTitle", "targetInstructions", "startedModel", + "acceptedAt", "startedAt", "phase", "terminalModel", @@ -314,7 +315,17 @@ struct ReviewHistorySchemaTests { let recordsBefore = try await writer.read { db in try #sql( - "SELECT *, NULL AS \"reviewThreadID\", NULL AS \"threadID\" FROM \"review_records\"", + """ + SELECT "id", "cwd", "sortOrder", "targetKind", + "targetBranch", "targetCommitSHA", "targetCommitTitle", "targetInstructions", + "startedModel", "startedAt" AS "acceptedAt", "startedAt", "phase", + "terminalModel", "terminalKind", "interruptionKind", "cancellationSource", + "cancellationMessage", "terminalMessage", "endedAt", "summary", + "canonicalReview", "parsedState", "parsedFindingCount", "parsedSource", + "parserVersion", "terminalCommittedAt", "createdAt", "updatedAt", + NULL AS "reviewThreadID", NULL AS "threadID" + FROM "review_records" + """, as: ReviewRecordRow.self ) .fetchAll(db).sorted { $0.id < $1.id } @@ -410,7 +421,7 @@ struct ReviewHistorySchemaTests { "missing-cancellation-source", ] for (index, id) in validActiveIDs.enumerated() { - try await database.recordStarted(ReviewHistoryTestSupport.started( + try await database.recordAccepted(ReviewHistoryTestSupport.started( id: id, sortOrder: Double(index) )) @@ -426,10 +437,10 @@ struct ReviewHistorySchemaTests { try #sql( """ INSERT INTO review_records ( - id, cwd, sortOrder, targetKind, startedAt, phase, createdAt, updatedAt + id, cwd, sortOrder, targetKind, acceptedAt, startedAt, phase, createdAt, updatedAt ) VALUES ( \(bind: id), \(bind: "/tmp/workspace"), 0, \(bind: targetKind), - 0, 'active', 0, 0 + 0, 0, 'active', 0, 0 ) """ ) @@ -763,7 +774,7 @@ struct ReviewHistorySchemaTests { let started = try ReviewHistoryTestSupport.started(id: "live-review") let first = ReviewHistoryDatabase(databaseURL: url) - try await first.recordStarted(started) + try await first.recordAccepted(started) let second = ReviewHistoryDatabase(databaseURL: url) await #expect(throws: ReviewHistoryDatabaseError.databaseInUse) { _ = try await second.load(retentionPolicy: .default) diff --git a/Tests/CodexReviewPersistenceTests/ReviewHistoryTestSupport.swift b/Tests/CodexReviewPersistenceTests/ReviewHistoryTestSupport.swift index 0d3c1149..d6f7f85a 100644 --- a/Tests/CodexReviewPersistenceTests/ReviewHistoryTestSupport.swift +++ b/Tests/CodexReviewPersistenceTests/ReviewHistoryTestSupport.swift @@ -29,8 +29,8 @@ enum ReviewHistoryTestSupport { target: CodexReviewAPI.Target = .uncommittedChanges, model: String? = "gpt-5.6-sol", startedAt: Date = startedAt - ) throws -> StartedReviewRecord { - try StartedReviewRecord( + ) throws -> AcceptedReviewRecord { + try AcceptedReviewRecord( id: id, cwd: cwd, workspaceMetadata: workspaceMetadata, @@ -38,6 +38,7 @@ enum ReviewHistoryTestSupport { sortOrder: sortOrder, target: target, model: model, + acceptedAt: startedAt, startedAt: startedAt ) } @@ -86,12 +87,12 @@ enum ReviewHistoryTestSupport { } static func record( - started: StartedReviewRecord, + started: AcceptedReviewRecord, terminal: TerminalReviewRecord, in database: ReviewHistoryDatabase, retentionPolicy: ReviewHistoryRetentionPolicy = .default ) async throws -> ReviewHistoryMutationResult { - try await database.recordStarted(started) + try await database.recordAccepted(started) return try await database.recordTerminal( terminal, retentionPolicy: retentionPolicy diff --git a/Tests/CodexReviewTests/CodexReviewStoreCommandTests.swift b/Tests/CodexReviewTests/CodexReviewStoreCommandTests.swift index bf8d9753..1790ed97 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreCommandTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreCommandTests.swift @@ -1,11 +1,179 @@ import Foundation import Testing @_spi(Testing) @testable import CodexReview +@_spi(ApplicationHostSupport) import CodexReview import CodexReviewTesting @Suite("Codex review store", .serialized) @MainActor struct CodexReviewStoreCommandTests { + @Test func shutdownCancelsQueuedReviewsAndFinishesTheirOriginalCalls() async throws { + let backend = FakeCodexReviewBackend() + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend) + ) + await store.start() + store.suspendReviewStarts() + let request = Task { + try await store.startReview( + sessionID: "session", request: .init(cwd: "/tmp/queued", target: .uncommittedChanges) + ) + } + try #require(await waitUntil { store.jobs.first?.core.lifecycle.status == .queued }) + await store.shutdown() + #expect(try await request.value.core.lifecycle.status == .cancelled) + #expect(store.queuedReviewStarts.isEmpty) + #expect(await backend.recordedCommands().contains { if case .startReview = $0 { true } else { false } } == false) + } + + @Test func queuedReviewKeepsTheOriginalStartCallWaitingUntilCompletion() async throws { + let backend = FakeCodexReviewBackend() + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend) + ) + store.suspendReviewStarts() + try await withStoreCommandTestCleanup(backend: backend, store: store) { + var returned = false + let request = Task { + defer { returned = true } + return try await store.startReview( + sessionID: "queued-client", + request: .init(cwd: "/tmp/queued", target: .uncommittedChanges) + ) + } + defer { request.cancel() } + try #require(await waitUntil { store.jobs.first?.core.lifecycle.status == .queued }) + let queued = try #require(store.jobs.first) + #expect(queued.core.lifecycle.startedAt == nil) + #expect(returned == false) + #expect(await backend.recordedCommands().contains { if case .startReview = $0 { true } else { false } } == false) + + store.resumeReviewStarts() + try await backend.waitForStartReview(timeout: .seconds(2)) + await backend.yield(.completed(summary: "Succeeded.", result: "No findings.")) + let result = try await request.value + #expect(result.jobID == queued.id) + #expect(result.core.lifecycle.status == .succeeded) + #expect(result.core.lifecycle.startedAt != nil) + #expect(returned) + } + } + + @Test func diagnosticsPublishRunningReviewWhileBackendStartIsHeld() async throws { + let directory = FileManager.default.temporaryDirectory.appendingPathComponent(UUID().uuidString) + defer { try? FileManager.default.removeItem(at: directory) } + let diagnosticsURL = directory.appendingPathComponent("diagnostics.json") + let backend = FakeCodexReviewBackend() + await backend.holdStartReview(with: AsyncGate()) + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend), + diagnosticsURL: diagnosticsURL + ) + store.suspendReviewStarts() + try await withStoreCommandTestCleanup(backend: backend, store: store) { + let queued = try await store.startReview( + sessionID: "session", request: .init(cwd: "/tmp/queued", target: .uncommittedChanges), + waitTimeout: .zero + ) + func diagnosticJob() throws -> [String: Any] { + let snapshot = try #require( + JSONSerialization.jsonObject(with: Data(contentsOf: diagnosticsURL)) as? [String: Any] + ) + return try #require((snapshot["jobs"] as? [[String: Any]])?.first) + } + #expect(try diagnosticJob()["status"] as? String == "queued") + #expect(try diagnosticJob()["startedAt"] == nil) + + store.resumeReviewStarts() + try await backend.waitForStartReview(timeout: .seconds(2)) + let runningSnapshot = Result { try diagnosticJob() } + try await store.cancelAllRunningJobs() + let running = try runningSnapshot.get() + #expect(running["id"] as? String == queued.jobID) + #expect(running["status"] as? String == "running") + #expect(running["startedAt"] != nil) + #expect(running["summary"] as? String == "Review started.") + } + } + + @Test func queuedReviewCanBeReadAwaitedAndCancelledBeforeDispatch() async throws { + let backend = FakeCodexReviewBackend() + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend) + ) + store.suspendReviewStarts() + let queued = try await store.startReview( + sessionID: "owner", request: .init(cwd: "/tmp/queued", target: .uncommittedChanges), + waitTimeout: .zero + ) + #expect(queued.core.lifecycle.status == .queued) + #expect(queued.elapsedSeconds == nil) + #expect(queued.cancellable) + #expect(throws: CodexReviewAPI.Error.self) { + try store.readReview(sessionID: "other", jobID: queued.jobID) + } + let waiting = try await store.awaitReview(sessionID: "owner", jobID: queued.jobID, timeout: .zero) + #expect(waiting.core.lifecycle.status == .queued) + let cancelled = try await store.cancelReview(jobID: queued.jobID, sessionID: "owner") + #expect(cancelled.cancelled) + #expect(cancelled.core.lifecycle.startedAt == nil) + #expect(store.queuedReviewStarts.isEmpty) + store.resumeReviewStarts() + #expect(await backend.recordedCommands().isEmpty) + } + + @Test func closingOneSessionOnlyCancelsItsQueuedReviews() async throws { + let backend = FakeCodexReviewBackend() + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend) + ) + store.suspendReviewStarts() + let first = try await store.startReview( + sessionID: "first", request: .init(cwd: "/tmp/first", target: .uncommittedChanges), + waitTimeout: .zero + ) + let second = try await store.startReview( + sessionID: "second", request: .init(cwd: "/tmp/second", target: .uncommittedChanges), + waitTimeout: .zero + ) + await store.closeSession("first") + #expect(try store.readReview(jobID: first.jobID).core.lifecycle.status == .cancelled) + #expect(try store.readReview(jobID: second.jobID).core.lifecycle.status == .queued) + await store.closeSession("second") + #expect(store.queuedReviewStarts.isEmpty) + #expect(await backend.recordedCommands().isEmpty) + } + + @Test func queuedReviewsStartInAcceptanceOrderWithoutSerializingCompletion() async throws { + let backend = FakeCodexReviewBackend() + await backend.holdStartReview(with: AsyncGate()) + let store = CodexReviewStore.makeTestingStore( + backend: TestingCodexReviewStoreBackend(reviewBackend: backend) + ) + store.suspendReviewStarts() + try await withStoreCommandTestCleanup(backend: backend, store: store) { + let first = try await store.startReview( + sessionID: "first", request: .init(cwd: "/tmp/first", target: .uncommittedChanges), + waitTimeout: .zero + ) + let second = try await store.startReview( + sessionID: "second", request: .init(cwd: "/tmp/second", target: .uncommittedChanges), + waitTimeout: .zero + ) + store.resumeReviewStarts() + try #require(await waitUntil { + await backend.recordedCommands().filter { if case .startReview = $0 { true } else { false } }.count == 2 + }) + let starts = await backend.recordedCommands().compactMap { command -> String? in + guard case .startReview(let request) = command else { return nil } + return request.jobID + } + #expect(starts == [first.jobID, second.jobID]) + #expect(store.jobs.allSatisfy { $0.core.lifecycle.status == .running }) + try await store.cancelAllRunningJobs() + } + } + @Test func storeBackendForwardsExplicitAdmission() async throws { let reviewBackend = FakeCodexReviewBackend() let storeBackend = TestingCodexReviewStoreBackend(reviewBackend: reviewBackend) diff --git a/Tests/CodexReviewTests/CodexReviewStoreDiagnosticsTests.swift b/Tests/CodexReviewTests/CodexReviewStoreDiagnosticsTests.swift index 70a36659..d3ce2c7a 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreDiagnosticsTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreDiagnosticsTests.swift @@ -16,13 +16,14 @@ struct CodexReviewStoreDiagnosticsTests { ) defer { try? FileManager.default.removeItem(at: directory) } - let started = try StartedReviewRecord( + let started = try AcceptedReviewRecord( id: "review-1", cwd: "/tmp/workspace", workspaceSortOrder: 0, sortOrder: 0, target: .uncommittedChanges, model: "gpt-5.6-sol", + acceptedAt: Date(timeIntervalSince1970: 100), startedAt: Date(timeIntervalSince1970: 100) ) let parsedResult = try PersistedParsedReviewResult( diff --git a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift index 554ef154..0a59dd55 100644 --- a/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift +++ b/Tests/CodexReviewTests/CodexReviewStoreHistoryTests.swift @@ -7,6 +7,45 @@ import CodexReviewTesting @Suite("review history store", .serialized) @MainActor struct CodexReviewStoreHistoryTests { + @Test func executionPersistenceFailureReturnsTheAcceptedJobWithoutDispatch() async throws { + let history = ReviewHistoryPersistenceProbe(executionWriteFailure: "Execution start could not be saved.") + let backend = FakeCodexReviewBackend() + let store = makeStore(history: history, backend: backend) + let result = try await store.startReview( + sessionID: "session", request: .init(cwd: "/tmp/queued", target: .uncommittedChanges) + ) + #expect(result.core.lifecycle.status == .failed) + #expect(result.core.lifecycle.errorMessage == "Execution start could not be saved.") + #expect(result.core.lifecycle.startedAt == nil) + #expect(result.elapsedSeconds == nil) + #expect(await history.startedRecords().first?.startedAt == nil) + #expect(store.queuedReviewStarts.isEmpty) + #expect(await backend.recordedCommands().isEmpty) + #expect(await history.terminalRecords().map(\.id) == [result.jobID]) + } + + @Test func suspensionDuringAcceptancePersistenceKeepsTheRequestQueued() async throws { + let entered = AsyncGate() + let release = AsyncGate() + let history = ReviewHistoryPersistenceProbe(startedWriteEntered: entered, startedWriteRelease: release) + let backend = FakeCodexReviewBackend() + let store = makeStore(history: history, backend: backend) + let request = Task { + try await store.startReview( + sessionID: "session", request: .init(cwd: "/tmp/queued", target: .uncommittedChanges), + waitTimeout: .zero + ) + } + await entered.wait() + store.suspendReviewStarts() + await release.open() + let queued = try await request.value + #expect(queued.core.lifecycle.status == .queued) + #expect(queued.core.lifecycle.startedAt == nil) + #expect(await backend.recordedCommands().isEmpty) + _ = try await store.cancelReview(jobID: queued.jobID, sessionID: "session") + } + @Test func loadOnceRestoresCompactApplicationHistoryWithoutSessionAuthority() async throws { let startedAt = Date(timeIntervalSince1970: 100) let history = ReviewHistoryPersistenceProbe(records: [ @@ -143,7 +182,7 @@ struct CodexReviewStoreHistoryTests { await backend.waitForStartReview() let startedRecord = try #require(await history.startedRecords().first) #expect(startedRecord.target == .commit(sha: "abc123", title: "Persist me")) - #expect(startedRecord.startedAt.timeIntervalSince1970 > 0) + #expect(try #require(startedRecord.startedAt).timeIntervalSince1970 > 0) await backend.yield(.completed(summary: "Done", result: "No findings.")) #expect(try await review.value.core.lifecycle.status == .succeeded) @@ -1076,6 +1115,7 @@ struct CodexReviewStoreHistoryTests { waitTimeout: .zero ) await backend.waitForStartReview() + try #require(await waitForHistoryTestCondition { store.reviewAttemptOwnerships[first.jobID]?.run != nil }) let firstRun = try #require(store.reviewAttemptOwnerships[first.jobID]?.run) await backend.yield( .completed(summary: "First", result: "First result"), @@ -1159,6 +1199,7 @@ struct CodexReviewStoreHistoryTests { waitTimeout: .zero ) await backend.waitForStartReview() + try #require(await waitForHistoryTestCondition { store.reviewAttemptOwnerships[first.jobID]?.run != nil }) let firstRun = try #require(store.reviewAttemptOwnerships[first.jobID]?.run) await backend.yield( .completed(summary: "First", result: "First result"), @@ -1858,7 +1899,7 @@ struct CodexReviewStoreHistoryTests { else { throw ReviewHistoryRecordError("History test fixture requires a terminal and start time.") } - let started = try StartedReviewRecord( + let started = try AcceptedReviewRecord( id: id, cwd: cwd, workspaceMetadata: workspaceMetadata, @@ -1866,6 +1907,7 @@ struct CodexReviewStoreHistoryTests { sortOrder: sortOrder, target: target, model: "gpt-5", + acceptedAt: startedAt, startedAt: startedAt ) let completed = terminal == .completed @@ -1950,11 +1992,12 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { private let orderingWriteRelease: AsyncGate? private let loadFailure: String? private let startedWriteFailure: String? + private let executionWriteFailure: String? private let terminalWriteFailure: String? private let orderingWriteFailure: String? private var terminalMutations: [ReviewHistoryMutationResult] private var loadCalls = 0 - private var started: [StartedReviewRecord] = [] + private var started: [AcceptedReviewRecord] = [] private var terminals: [TerminalReviewRecord] = [] private var savedOrderings: [ReviewHistoryOrdering] = [] private var mutationOperationLog: [String] = [] @@ -1973,6 +2016,7 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { orderingWriteRelease: AsyncGate? = nil, loadFailure: String? = nil, startedWriteFailure: String? = nil, + executionWriteFailure: String? = nil, terminalWriteFailure: String? = nil, orderingWriteFailure: String? = nil, terminalMutation: ReviewHistoryMutationResult = .init(), @@ -1989,6 +2033,7 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { self.orderingWriteRelease = orderingWriteRelease self.loadFailure = loadFailure self.startedWriteFailure = startedWriteFailure + self.executionWriteFailure = executionWriteFailure self.terminalWriteFailure = terminalWriteFailure self.orderingWriteFailure = orderingWriteFailure self.terminalMutations = terminalMutations ?? [terminalMutation] @@ -2006,7 +2051,7 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { return records } - func recordStarted(_ record: StartedReviewRecord) async throws { + func recordAccepted(_ record: AcceptedReviewRecord) async throws { await startedWriteEntered?.open() await startedWriteRelease?.waitIgnoringCancellation() if let startedWriteFailure { @@ -2016,6 +2061,15 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { started.append(record) } + func recordExecutionStarted(id: String, at date: Date) async throws { + if let executionWriteFailure { + throw ReviewHistoryPersistenceProbeError(message: executionWriteFailure) + } + if let index = started.firstIndex(where: { $0.id == id }) { + started[index].startedAt = date + } + } + func recordTerminal( _ record: TerminalReviewRecord, retentionPolicy _: ReviewHistoryRetentionPolicy @@ -2073,7 +2127,7 @@ private actor ReviewHistoryPersistenceProbe: ReviewHistoryPersistence { } func loadCallCount() -> Int { loadCalls } - func startedRecords() -> [StartedReviewRecord] { started } + func startedRecords() -> [AcceptedReviewRecord] { started } func terminalRecords() -> [TerminalReviewRecord] { terminals } func orderings() -> [ReviewHistoryOrdering] { savedOrderings } func deleteAllCallCount() -> Int { deleteAllCalls } diff --git a/Tests/CodexReviewTests/ReviewHistoryRecordTests.swift b/Tests/CodexReviewTests/ReviewHistoryRecordTests.swift index 89d18951..dd1560b1 100644 --- a/Tests/CodexReviewTests/ReviewHistoryRecordTests.swift +++ b/Tests/CodexReviewTests/ReviewHistoryRecordTests.swift @@ -7,13 +7,14 @@ import Testing struct ReviewHistoryRecordTests { @Test func phaseSpecificRecordsRejectIncompatiblePayloads() throws { #expect(throws: ReviewHistoryRecordError.self) { - try StartedReviewRecord( + try AcceptedReviewRecord( id: "", cwd: "/tmp/project", workspaceSortOrder: 0, sortOrder: 0, target: .uncommittedChanges, model: nil, + acceptedAt: Date(timeIntervalSince1970: 1), startedAt: Date(timeIntervalSince1970: 1) ) } @@ -89,13 +90,14 @@ struct ReviewHistoryRecordTests { let endedAt = Date(timeIntervalSince1970: 2) let parsed = ParsedReviewResult.parse(finalReviewText: "No findings.") let restored = try RestoredReviewRecord( - started: StartedReviewRecord( + started: AcceptedReviewRecord( id: "review-1", cwd: "/tmp/project", workspaceSortOrder: 3, sortOrder: 4, target: .baseBranch("main"), model: "gpt-5", + acceptedAt: startedAt, startedAt: startedAt ), terminal: TerminalReviewRecord( @@ -125,13 +127,14 @@ struct ReviewHistoryRecordTests { @Test func previousProcessExitRestoresUnknownEndWithoutLiveTimerState() throws { let startedAt = Date(timeIntervalSince1970: 1) let restored = try RestoredReviewRecord( - started: StartedReviewRecord( + started: AcceptedReviewRecord( id: "review-1", cwd: "/tmp/project", workspaceSortOrder: 0, sortOrder: 0, target: .uncommittedChanges, model: nil, + acceptedAt: startedAt, startedAt: startedAt ), terminal: TerminalReviewRecord( @@ -157,13 +160,14 @@ struct ReviewHistoryRecordTests { @Test func requestedCancellationRestoresTypedTerminalWithoutSyntheticLog() throws { let cancellation = ReviewCancellation.mcpClient(message: "Stop review.") let restored = try RestoredReviewRecord( - started: StartedReviewRecord( + started: AcceptedReviewRecord( id: "review-1", cwd: "/tmp/project", workspaceSortOrder: 0, sortOrder: 0, target: .uncommittedChanges, model: nil, + acceptedAt: Date(timeIntervalSince1970: 1), startedAt: Date(timeIntervalSince1970: 1) ), terminal: TerminalReviewRecord( diff --git a/Tests/ReviewUITests/ReviewMonitorCopyTests.swift b/Tests/ReviewUITests/ReviewMonitorCopyTests.swift index 725aade3..550f4f0e 100644 --- a/Tests/ReviewUITests/ReviewMonitorCopyTests.swift +++ b/Tests/ReviewUITests/ReviewMonitorCopyTests.swift @@ -85,7 +85,7 @@ struct ReviewMonitorCopyTests { let job = try RestoredReviewRecord( started: .init( id: "job", cwd: "/tmp/repo", workspaceSortOrder: 0, sortOrder: 0, - target: .uncommittedChanges, model: nil, startedAt: .distantPast + target: .uncommittedChanges, model: nil, acceptedAt: .distantPast, startedAt: .distantPast ), terminal: .init( id: "job", model: nil, reviewThreadID: reviewThreadID, threadID: threadID, diff --git a/Tests/ReviewUITests/ReviewMonitorLogProjectionTests.swift b/Tests/ReviewUITests/ReviewMonitorLogProjectionTests.swift index 8da22eb1..c1df3a79 100644 --- a/Tests/ReviewUITests/ReviewMonitorLogProjectionTests.swift +++ b/Tests/ReviewUITests/ReviewMonitorLogProjectionTests.swift @@ -137,13 +137,14 @@ struct ReviewMonitorLogProjectionTests { let cancellation = ReviewCancellation.sessionClosed(message: " \t") let summary = "Session closed." let restored = try RestoredReviewRecord( - started: StartedReviewRecord( + started: AcceptedReviewRecord( id: "review-1", cwd: "/tmp/project", workspaceSortOrder: 0, sortOrder: 0, target: .uncommittedChanges, model: "gpt-5", + acceptedAt: Date(timeIntervalSince1970: 1), startedAt: Date(timeIntervalSince1970: 1) ), terminal: TerminalReviewRecord( diff --git a/Tests/ReviewUITests/ReviewUIHistoryTests.swift b/Tests/ReviewUITests/ReviewUIHistoryTests.swift index 2a8cecd7..27304f38 100644 --- a/Tests/ReviewUITests/ReviewUIHistoryTests.swift +++ b/Tests/ReviewUITests/ReviewUIHistoryTests.swift @@ -200,13 +200,14 @@ struct ReviewUIHistoryTests { id: String, sortOrder: Double = 0 ) throws -> RestoredReviewRecord { - let started = try StartedReviewRecord( + let started = try AcceptedReviewRecord( id: id, cwd: "/tmp/workspace", workspaceSortOrder: 0, sortOrder: sortOrder, target: .uncommittedChanges, model: "gpt-5.6-sol", + acceptedAt: Date(timeIntervalSince1970: 100), startedAt: Date(timeIntervalSince1970: 100) ) let terminal = try TerminalReviewRecord( @@ -241,7 +242,9 @@ private actor ReviewUIHistoryPersistence: ReviewHistoryPersistence { records } - func recordStarted(_: StartedReviewRecord) async throws {} + func recordAccepted(_: AcceptedReviewRecord) async throws {} + + func recordExecutionStarted(id: String, at date: Date) async throws {} func recordTerminal( _: TerminalReviewRecord, diff --git a/Tests/ReviewUITests/ReviewUITests.swift b/Tests/ReviewUITests/ReviewUITests.swift index 28f8a1ea..bef28cc7 100644 --- a/Tests/ReviewUITests/ReviewUITests.swift +++ b/Tests/ReviewUITests/ReviewUITests.swift @@ -3764,13 +3764,14 @@ struct ReviewUITests { let cancellation = ReviewCancellation.sessionClosed(message: " \t") let summary = "Session closed." let restored = try RestoredReviewRecord( - started: StartedReviewRecord( + started: AcceptedReviewRecord( id: "job-restored-cancellation", cwd: "/tmp/workspace-alpha", workspaceSortOrder: 0, sortOrder: 0, target: .uncommittedChanges, model: "gpt-5", + acceptedAt: Date(timeIntervalSince1970: 200), startedAt: Date(timeIntervalSince1970: 200) ), terminal: TerminalReviewRecord( diff --git a/scripts/review-history-e2e/run.sh b/scripts/review-history-e2e/run.sh index aafdb66d..70f375c9 100755 --- a/scripts/review-history-e2e/run.sh +++ b/scripts/review-history-e2e/run.sh @@ -418,7 +418,7 @@ assert_database_contract() { /usr/bin/jq -e '[.[].name] == [ "id", "cwd", "sortOrder", "targetKind", "targetBranch", "targetCommitSHA", "targetCommitTitle", "targetInstructions", - "startedModel", "startedAt", "phase", "terminalModel", "terminalKind", + "startedModel", "acceptedAt", "startedAt", "phase", "terminalModel", "terminalKind", "interruptionKind", "cancellationSource", "cancellationMessage", "terminalMessage", "endedAt", "summary", "canonicalReview", "parsedState", "parsedFindingCount", "parsedSource", "parserVersion", @@ -430,7 +430,7 @@ assert_database_contract() { ]' "$finding_columns_path" >/dev/null || die "finding schema inventory changed" /usr/bin/sqlite3 -json "$history_path" \ - "SELECT id, cwd, targetKind, startedModel, startedAt, phase, terminalModel, terminalKind, endedAt, summary, canonicalReview, parsedState, parsedFindingCount, parsedSource, parserVersion FROM review_records" \ + "SELECT id, cwd, targetKind, startedModel, acceptedAt, startedAt, phase, terminalModel, terminalKind, endedAt, summary, canonicalReview, parsedState, parsedFindingCount, parsedSource, parserVersion FROM review_records" \ >"$semantic_record_path" /usr/bin/sqlite3 -json "$history_path" \ "SELECT reviewID, ordinal, priority, title, body, path, startLine, endLine FROM review_findings ORDER BY ordinal" \ @@ -441,6 +441,7 @@ assert_database_contract() { and .[0].id == $job_id and .[0].cwd == $cwd and .[0].targetKind == "uncommittedChanges" + and (.[0].acceptedAt | type == "number") and (.[0].startedAt | type == "number") and .[0].phase == "terminal" and .[0].terminalKind == "completed"