From 06009efff33e49aef51975524f12733de0830f8f Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Sun, 2 Aug 2026 01:01:25 +0300 Subject: [PATCH 1/5] feat(review): add durable publisher bridge Agent: Herminia --- .../src/protocol/common.rs | 10 + .../src/protocol/v2/review.rs | 115 ++ codex-rs/app-server/README.md | 4 + .../app-server/src/bespoke_event_handling.rs | 21 + codex-rs/app-server/src/message_processor.rs | 20 + codex-rs/app-server/src/request_processors.rs | 4 + .../request_processors/review_publisher.rs | 768 ++++++++++++ .../src/request_processors/turn_processor.rs | 115 +- .../tests/suite/v2/client_metadata.rs | 1 + codex-rs/app-server/tests/suite/v2/review.rs | 10 + codex-rs/core/src/session/review.rs | 5 +- codex-rs/core/src/session/tests.rs | 2 +- codex-rs/core/src/tasks/review.rs | 33 +- codex-rs/core/tests/suite/codex_delegate.rs | 3 + codex-rs/core/tests/suite/review.rs | 21 +- codex-rs/exec/src/lib.rs | 2 + codex-rs/exec/src/lib_tests.rs | 3 + codex-rs/prompts/src/review_request.rs | 23 +- codex-rs/prompts/src/review_request_tests.rs | 45 + codex-rs/protocol/src/protocol.rs | 86 ++ .../0063_review_publisher_outbox.sql | 48 + codex-rs/state/src/lib.rs | 20 + codex-rs/state/src/model/mod.rs | 7 + codex-rs/state/src/model/review_publisher.rs | 112 ++ codex-rs/state/src/runtime.rs | 21 + .../state/src/runtime/review_publisher.rs | 1078 +++++++++++++++++ codex-rs/tui/src/app_server_session.rs | 1 + 27 files changed, 2552 insertions(+), 26 deletions(-) create mode 100644 codex-rs/app-server/src/request_processors/review_publisher.rs create mode 100644 codex-rs/state/migrations/0063_review_publisher_outbox.sql create mode 100644 codex-rs/state/src/model/review_publisher.rs create mode 100644 codex-rs/state/src/runtime/review_publisher.rs diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index 13352a771f..43dbf589b9 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -1424,6 +1424,16 @@ client_request_definitions! { serialization: thread_id(params.thread_id), response: v2::ReviewStartResponse, }, + ReviewPublisherStatusRead => "review/publisher/status/read" { + params: v2::ReviewPublisherStatusReadParams, + serialization: global_shared_read("review-publisher"), + response: v2::ReviewPublisherStatusReadResponse, + }, + ReviewPublisherReplay => "review/publisher/replay" { + params: v2::ReviewPublisherReplayParams, + serialization: global("review-publisher"), + response: v2::ReviewPublisherReplayResponse, + }, ModelList => "model/list" { params: v2::ModelListParams, diff --git a/codex-rs/app-server-protocol/src/protocol/v2/review.rs b/codex-rs/app-server-protocol/src/protocol/v2/review.rs index 82ec5b6f59..574fa07a3c 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/review.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/review.rs @@ -23,6 +23,13 @@ pub struct ReviewStartParams { #[serde(default)] #[ts(optional = nullable)] pub delivery: Option, + + /// Exact pull-request candidate metadata for authenticated check publishing. + /// Repository origin and implementer identity are deliberately absent: the + /// app-server derives them from the clean local Git checkout. + #[serde(default)] + #[ts(optional = nullable)] + pub publisher_context: Option, } #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] @@ -35,6 +42,114 @@ pub struct ReviewStartResponse { /// For inline reviews, this is the original thread id. /// For detached reviews, this is the id of the new review thread. pub review_thread_id: String, + /// Stable durable review run id when `publisherContext` was supplied. + pub review_run_id: Option, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherContext { + pub pull_request_number: u64, + pub base_ref: String, + pub reviewed_base_sha: String, + pub head_sha: String, + pub acceptance_scope_id: String, + pub acceptance_scope_sha256: String, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherStatusReadParams { + pub review_run_id: String, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherStatusReadResponse { + pub run: Option, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherReplayParams { + pub event_id: String, + pub payload_sha256: String, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherReplayResponse { + pub event: Option, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherRun { + pub review_run_id: String, + pub envelope_sha256: String, + pub status: ReviewPublisherRunStatus, + pub verdict: Option, + pub created_at: i64, + pub completed_at: Option, + pub events: Vec, +} + +#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase", export_to = "v2/")] +pub enum ReviewPublisherRunStatus { + Started, + Completed, +} + +#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +#[ts(rename_all = "SCREAMING_SNAKE_CASE", export_to = "v2/")] +pub enum ReviewPublisherVerdict { + Go, + NoGo, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ReviewPublisherOutboxEvent { + pub event_id: String, + pub event_kind: ReviewPublisherEventKind, + pub sequence: u8, + pub status: ReviewPublisherEventStatus, + pub payload_sha256: String, + pub attempt_count: u32, + pub next_attempt_at: i64, + pub lease_expires_at: Option, + pub receipt_id: Option, + pub last_error_code: Option, + pub created_at: i64, + pub delivered_at: Option, +} + +#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase", export_to = "v2/")] +pub enum ReviewPublisherEventKind { + Started, + Completed, +} + +#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase", export_to = "v2/")] +pub enum ReviewPublisherEventStatus { + Pending, + InFlight, + Delivered, + DeadLetter, } #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 5b7efc468a..beb025c93a 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -1631,6 +1631,10 @@ Example request/response: For a detached review, use `"delivery": "detached"`. The response is the same shape, but `reviewThreadId` will be the id of the new review thread (different from the original `threadId`). The server also emits a `thread/started` notification for that new thread before streaming the review turn. +An authenticated publisher can add `publisherContext` to a `baseBranch` review. This path is fail-closed: `baseRef` must be the same full Git ref used by the target, the worktree must be clean, `reviewedBaseSha` and `headSha` must resolve exactly, `git merge-tree --write-tree` must produce a clean result, and the head commit must carry exactly one `Agent:` trailer. The server derives the canonical origin and implementer from Git, persists an immutable `codewith-review-envelope-v1` start event before starting the turn, and returns its stable `reviewRunId`. It publishes the terminal `GO` or `NO_GO` event only from structured review output; missing output, unknown correctness, malformed priorities, or P0/P1 findings all map to `NO_GO`. + +The owner-only outbox dispatcher is enabled only when `CODEWITH_REVIEW_PUBLISHER_URL` names the HTTP endpoint and `CODEWITH_REVIEW_PUBLISHER_CREDENTIAL_ENV` names the environment variable holding its bearer credential. The credential itself is never accepted in RPC payloads or persisted. Delivery is ordered start-before-terminal, leases in-flight work, treats HTTP 409 as an idempotent receipt, retries timeouts/429/5xx, and dead-letters permanent failures. Inspect a run with `review/publisher/status/read`; replay only an exact immutable payload by passing both `eventId` and `payloadSha256` to `review/publisher/replay`. + Codewith streams the usual `turn/started` notification followed by an `item/started` with an `enteredReviewMode` item so clients can show progress: diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index 6b20d2545a..a3a05a9f80 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -1377,6 +1377,27 @@ pub(crate) async fn apply_bespoke_event_handling( .await; } EventMsg::ExitedReviewMode(review_event) => { + if let Some(envelope) = review_event.review_envelope.as_ref() { + if let Some(state_db) = conversation.state_db() { + let review_run_id = codex_state::review_run_id_from_envelope_sha256( + envelope.envelope_sha256.as_str(), + ); + if let Err(err) = state_db + .review_publisher() + .complete_review_run(codex_state::ReviewPublisherCompleteParams { + review_run_id, + envelope_sha256: envelope.envelope_sha256.clone(), + review_output: review_event.review_output.clone(), + terminal_reason_override: None, + }) + .await + { + warn!("failed to persist review publisher terminal event: {err}"); + } + } else { + warn!("review publisher terminal event has no durable state runtime"); + } + } let review = match review_event.review_output { Some(output) => render_review_output_text(&output), None => REVIEW_FALLBACK_MESSAGE.to_string(), diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index cb35840347..11b3ea9ec2 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -35,6 +35,8 @@ use crate::request_processors::PluginRequestProcessor; use crate::request_processors::ProcessExecRequestProcessor; use crate::request_processors::RemoteControlRequestProcessor; use crate::request_processors::RemoteDispatchRequestProcessor; +use crate::request_processors::ReviewPublisherDispatcherRuntime; +use crate::request_processors::ReviewPublisherRequestProcessor; use crate::request_processors::SearchRequestProcessor; use crate::request_processors::ThreadGoalRequestProcessor; use crate::request_processors::ThreadMailboxDispatcherRuntime; @@ -201,6 +203,8 @@ pub(crate) struct MessageProcessor { plugin_processor: PluginRequestProcessor, remote_control_processor: RemoteControlRequestProcessor, remote_dispatch_processor: RemoteDispatchRequestProcessor, + review_publisher_dispatcher_runtime: ReviewPublisherDispatcherRuntime, + review_publisher_processor: ReviewPublisherRequestProcessor, search_processor: SearchRequestProcessor, thread_goal_processor: ThreadGoalRequestProcessor, thread_mailbox_dispatcher_runtime: Option, @@ -504,6 +508,10 @@ impl MessageProcessor { let machine_registry_processor = MachineRegistryRequestProcessor::new(state_db.clone()); let remote_control_processor = RemoteControlRequestProcessor::new(remote_control_handle); let remote_dispatch_processor = RemoteDispatchRequestProcessor::new(state_db.clone()); + let review_publisher_dispatcher_runtime = + ReviewPublisherDispatcherRuntime::new(state_db.clone()); + review_publisher_dispatcher_runtime.start(); + let review_publisher_processor = ReviewPublisherRequestProcessor::new(state_db.clone()); let search_processor = SearchRequestProcessor::new(outgoing.clone()); let thread_goal_processor = ThreadGoalRequestProcessor::new( Arc::clone(&thread_manager), @@ -677,6 +685,8 @@ impl MessageProcessor { plugin_processor, remote_control_processor, remote_dispatch_processor, + review_publisher_dispatcher_runtime, + review_publisher_processor, search_processor, thread_goal_processor, thread_mailbox_dispatcher_runtime, @@ -709,6 +719,7 @@ impl MessageProcessor { if let Some(runtime) = self.thread_mailbox_dispatcher_runtime.as_ref() { runtime.shutdown(); } + self.review_publisher_dispatcher_runtime.shutdown(); self.thread_monitor_runtime.shutdown(); self.thread_schedule_runtime.shutdown(); } @@ -894,6 +905,9 @@ impl MessageProcessor { if let Some(runtime) = self.thread_mailbox_dispatcher_runtime.as_ref() { runtime.drain_background_tasks().await; } + self.review_publisher_dispatcher_runtime + .drain_background_tasks() + .await; self.thread_monitor_runtime.drain_background_tasks().await; self.thread_schedule_runtime.drain_background_tasks().await; self.thread_processor.drain_background_tasks().await; @@ -1874,6 +1888,12 @@ impl MessageProcessor { ClientRequest::ReviewStart { params, .. } => { self.turn_processor.review_start(&request_id, params).await } + ClientRequest::ReviewPublisherStatusRead { params, .. } => { + self.review_publisher_processor.status_read(params).await + } + ClientRequest::ReviewPublisherReplay { params, .. } => { + self.review_publisher_processor.replay(params).await + } ClientRequest::McpServerOauthLogin { params, .. } => { self.mcp_processor.mcp_server_oauth_login(params).await } diff --git a/codex-rs/app-server/src/request_processors.rs b/codex-rs/app-server/src/request_processors.rs index 873877c214..c5107b8c5f 100644 --- a/codex-rs/app-server/src/request_processors.rs +++ b/codex-rs/app-server/src/request_processors.rs @@ -673,6 +673,7 @@ mod plugins; mod process_exec_processor; mod remote_control_processor; mod remote_dispatch_processor; +mod review_publisher; mod search; mod sqlite_retry; mod thread_external_agent_processor; @@ -715,6 +716,9 @@ pub(crate) use plugins::PluginRequestProcessor; pub(crate) use process_exec_processor::ProcessExecRequestProcessor; pub(crate) use remote_control_processor::RemoteControlRequestProcessor; pub(crate) use remote_dispatch_processor::RemoteDispatchRequestProcessor; +pub(crate) use review_publisher::ReviewPublisherDispatcherRuntime; +pub(crate) use review_publisher::ReviewPublisherRequestProcessor; +pub(crate) use review_publisher::build_review_envelope; pub(crate) use search::SearchRequestProcessor; pub(crate) use thread_goal_processor::ThreadGoalRequestProcessor; pub(crate) use thread_mailbox_dispatcher_runtime::ThreadMailboxDispatcherRuntime; diff --git a/codex-rs/app-server/src/request_processors/review_publisher.rs b/codex-rs/app-server/src/request_processors/review_publisher.rs new file mode 100644 index 0000000000..43122c3c77 --- /dev/null +++ b/codex-rs/app-server/src/request_processors/review_publisher.rs @@ -0,0 +1,768 @@ +use super::*; +use crate::error_code::internal_error; +use crate::error_code::invalid_request; +use chrono::Utc; +use codex_app_server_protocol::ClientResponsePayload; +use codex_app_server_protocol::JSONRPCErrorError; +use codex_app_server_protocol::ReviewPublisherContext; +use codex_app_server_protocol::ReviewPublisherEventKind as ApiEventKind; +use codex_app_server_protocol::ReviewPublisherEventStatus as ApiEventStatus; +use codex_app_server_protocol::ReviewPublisherOutboxEvent as ApiOutboxEvent; +use codex_app_server_protocol::ReviewPublisherReplayParams; +use codex_app_server_protocol::ReviewPublisherReplayResponse; +use codex_app_server_protocol::ReviewPublisherRun as ApiRun; +use codex_app_server_protocol::ReviewPublisherRunStatus as ApiRunStatus; +use codex_app_server_protocol::ReviewPublisherStatusReadParams; +use codex_app_server_protocol::ReviewPublisherStatusReadResponse; +use codex_app_server_protocol::ReviewPublisherVerdict as ApiVerdict; +use codex_git_utils::canonicalize_git_remote_url; +use codex_protocol::protocol::ReviewEnvelope; +use codex_protocol::protocol::ReviewImplementerProvenance; +use codex_protocol::protocol::ReviewImplementerProvenanceSource; +use codex_rollout::StateDbHandle; +use codex_state::REVIEW_ENVELOPE_SCHEMA_VERSION; +use codex_state::ReviewPublisherClaimParams; +use codex_state::ReviewPublisherDeliveryAckParams; +use codex_state::ReviewPublisherDeliveryFailParams; +use codex_state::ReviewPublisherFailureDisposition; +use reqwest::StatusCode; +use std::path::Path; +use tokio::process::Command; +use tokio_util::sync::CancellationToken; +use tokio_util::task::TaskTracker; + +const GIT_PROBE_TIMEOUT: Duration = Duration::from_secs(10); +const DISPATCH_POLL_INTERVAL: Duration = Duration::from_secs(1); +const DISPATCH_LEASE_DURATION: Duration = Duration::from_secs(30); +const DISPATCH_HTTP_TIMEOUT: Duration = Duration::from_secs(10); +const DISPATCH_DRAIN_TIMEOUT: Duration = Duration::from_secs(15); +const DISPATCH_MAX_ATTEMPTS: u32 = 6; +const DISPATCH_LEASE_OWNER: &str = "app-server-review-publisher"; +const REVIEW_PUBLISHER_URL_ENV: &str = "CODEWITH_REVIEW_PUBLISHER_URL"; +const REVIEW_PUBLISHER_CREDENTIAL_ENV_ENV: &str = "CODEWITH_REVIEW_PUBLISHER_CREDENTIAL_ENV"; + +#[derive(Clone)] +pub(crate) struct ReviewPublisherRequestProcessor { + state_db: Option, +} + +impl ReviewPublisherRequestProcessor { + pub(crate) fn new(state_db: Option) -> Self { + Self { state_db } + } + + pub(crate) async fn status_read( + &self, + params: ReviewPublisherStatusReadParams, + ) -> Result, JSONRPCErrorError> { + let state_db = self.state_db()?; + let review_run_id = normalize_digest_bound_id("reviewRunId", params.review_run_id)?; + let run = state_db + .review_publisher() + .get_review_run(review_run_id.as_str()) + .await + .map_err(|err| { + internal_error(format!("failed to read review publisher status: {err}")) + })? + .map(api_run); + Ok(Some(ReviewPublisherStatusReadResponse { run }.into())) + } + + pub(crate) async fn replay( + &self, + params: ReviewPublisherReplayParams, + ) -> Result, JSONRPCErrorError> { + let state_db = self.state_db()?; + let event_id = normalize_digest_bound_id("eventId", params.event_id)?; + let payload_sha256 = normalize_sha256("payloadSha256", params.payload_sha256)?; + let event = state_db + .review_publisher() + .exact_replay(event_id.as_str(), payload_sha256.as_str(), Utc::now()) + .await + .map_err(|err| { + internal_error(format!("failed to replay review publisher event: {err}")) + })? + .map(api_event); + Ok(Some(ReviewPublisherReplayResponse { event }.into())) + } + + fn state_db(&self) -> Result<&StateDbHandle, JSONRPCErrorError> { + self.state_db + .as_ref() + .ok_or_else(|| internal_error("review publisher state is unavailable".to_string())) + } +} + +#[derive(Clone)] +pub(crate) struct ReviewPublisherDispatcherRuntime { + state_db: Option, + config: Option, + client: reqwest::Client, + cancel_token: CancellationToken, + tasks: TaskTracker, +} + +#[derive(Clone)] +struct ReviewPublisherDispatcherConfig { + endpoint: reqwest::Url, + credential_env: String, +} + +impl ReviewPublisherDispatcherRuntime { + pub(crate) fn new(state_db: Option) -> Self { + let config = dispatcher_config_from_env(); + let client = reqwest::Client::builder() + .timeout(DISPATCH_HTTP_TIMEOUT) + .build() + .unwrap_or_else(|_| reqwest::Client::new()); + Self { + state_db, + config, + client, + cancel_token: CancellationToken::new(), + tasks: TaskTracker::new(), + } + } + + pub(crate) fn start(&self) { + if self.state_db.is_none() || self.config.is_none() { + return; + } + let runtime = self.clone(); + self.tasks.spawn(async move { runtime.run().await }); + } + + pub(crate) fn shutdown(&self) { + self.cancel_token.cancel(); + } + + pub(crate) async fn drain_background_tasks(&self) { + self.shutdown(); + self.tasks.close(); + if tokio::time::timeout(DISPATCH_DRAIN_TIMEOUT, self.tasks.wait()) + .await + .is_err() + { + warn!("timed out waiting for review publisher dispatcher to drain"); + } + } + + async fn run(self) { + let mut interval = tokio::time::interval(DISPATCH_POLL_INTERVAL); + interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); + loop { + tokio::select! { + _ = self.cancel_token.cancelled() => break, + _ = interval.tick() => self.dispatch_one().await, + } + } + } + + async fn dispatch_one(&self) { + let (Some(state_db), Some(config)) = (self.state_db.as_ref(), self.config.as_ref()) else { + return; + }; + let claim = match state_db + .review_publisher() + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: DISPATCH_LEASE_OWNER.to_string(), + lease_duration: DISPATCH_LEASE_DURATION, + now: Utc::now(), + }) + .await + { + Ok(Some(claim)) => claim, + Ok(None) => return, + Err(err) => { + warn!("failed to claim review publisher event: {err}"); + return; + } + }; + + let result = self.send_event(config, &claim.event.payload).await; + let requested_dead_letter = matches!(&result, DispatchResult::DeadLetter { .. }); + let now = Utc::now(); + match result { + DispatchResult::Delivered { receipt_id } => { + if let Err(err) = state_db + .review_publisher() + .acknowledge_delivery(ReviewPublisherDeliveryAckParams { + event_id: claim.event.event_id, + lease_owner: DISPATCH_LEASE_OWNER.to_string(), + receipt_id, + now, + }) + .await + { + warn!("failed to acknowledge review publisher delivery: {err}"); + } + } + DispatchResult::Retry { error_code } | DispatchResult::DeadLetter { error_code } => { + let exhausted = claim.event.attempt_count >= DISPATCH_MAX_ATTEMPTS; + let disposition = if requested_dead_letter || exhausted { + ReviewPublisherFailureDisposition::DeadLetter + } else { + ReviewPublisherFailureDisposition::Retry + }; + let retry_seconds = + i64::from(2_u32.saturating_pow(claim.event.attempt_count.min(8))).clamp(2, 300); + if let Err(err) = state_db + .review_publisher() + .fail_delivery(ReviewPublisherDeliveryFailParams { + event_id: claim.event.event_id, + lease_owner: DISPATCH_LEASE_OWNER.to_string(), + error_code, + disposition, + retry_at: now + chrono::Duration::seconds(retry_seconds), + now, + }) + .await + { + warn!("failed to record review publisher delivery failure: {err}"); + } + } + } + } + + async fn send_event( + &self, + config: &ReviewPublisherDispatcherConfig, + event: &codex_protocol::protocol::ReviewPublisherEvent, + ) -> DispatchResult { + let token = match std::env::var(config.credential_env.as_str()) { + Ok(token) if !token.is_empty() => token, + _ => { + return DispatchResult::DeadLetter { + error_code: "credential_unavailable".to_string(), + }; + } + }; + let response = self + .client + .post(config.endpoint.clone()) + .bearer_auth(token) + .json(event) + .send() + .await; + match response { + Ok(response) => classify_http_response(&response), + Err(err) => classify_transport_error(&err), + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +enum DispatchResult { + Delivered { receipt_id: Option }, + Retry { error_code: String }, + DeadLetter { error_code: String }, +} + +fn classify_http_response(response: &reqwest::Response) -> DispatchResult { + let status = response.status(); + match classify_status(status) { + StatusDisposition::Delivered => { + let receipt_id = response + .headers() + .get("x-review-receipt-id") + .or_else(|| response.headers().get("x-github-request-id")) + .and_then(|value| value.to_str().ok()) + .and_then(normalize_receipt_id); + DispatchResult::Delivered { receipt_id } + } + StatusDisposition::Retry => DispatchResult::Retry { + error_code: format!("http_{}", status.as_u16()), + }, + StatusDisposition::DeadLetter => DispatchResult::DeadLetter { + error_code: format!("http_{}", status.as_u16()), + }, + } +} + +fn classify_transport_error(error: &reqwest::Error) -> DispatchResult { + DispatchResult::Retry { + error_code: if error.is_timeout() { + "http_timeout".to_string() + } else { + "http_transport".to_string() + }, + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum StatusDisposition { + Delivered, + Retry, + DeadLetter, +} + +fn classify_status(status: StatusCode) -> StatusDisposition { + if status.is_success() || status == StatusCode::CONFLICT { + StatusDisposition::Delivered + } else if status == StatusCode::TOO_MANY_REQUESTS || status.is_server_error() { + StatusDisposition::Retry + } else { + StatusDisposition::DeadLetter + } +} + +fn normalize_receipt_id(value: &str) -> Option { + let value = value.trim(); + if value.is_empty() + || value.len() > 128 + || !value + .chars() + .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.' | ':')) + { + return None; + } + Some(value.to_string()) +} + +fn dispatcher_config_from_env() -> Option { + let endpoint = std::env::var(REVIEW_PUBLISHER_URL_ENV).ok()?; + let credential_env = std::env::var(REVIEW_PUBLISHER_CREDENTIAL_ENV_ENV).ok()?; + let endpoint = reqwest::Url::parse(endpoint.trim()).ok()?; + if !matches!(endpoint.scheme(), "http" | "https") + || !endpoint.username().is_empty() + || endpoint.password().is_some() + || endpoint.query().is_some() + || endpoint.fragment().is_some() + || !valid_env_name(credential_env.trim()) + { + warn!("review publisher configuration is invalid; dispatcher is disabled"); + return None; + } + Some(ReviewPublisherDispatcherConfig { + endpoint, + credential_env: credential_env.trim().to_string(), + }) +} + +fn valid_env_name(value: &str) -> bool { + let mut chars = value.chars(); + matches!(chars.next(), Some(ch) if ch == '_' || ch.is_ascii_alphabetic()) + && chars.all(|ch| ch == '_' || ch.is_ascii_alphanumeric()) +} + +pub(crate) async fn build_review_envelope( + cwd: &Path, + context: ReviewPublisherContext, +) -> Result { + if context.pull_request_number == 0 { + return Err(invalid_request( + "publisherContext.pullRequestNumber must be positive".to_string(), + )); + } + let base_ref = normalize_ref(context.base_ref)?; + let reviewed_base_sha = normalize_git_sha("reviewedBaseSha", context.reviewed_base_sha)?; + let head_sha = normalize_git_sha("headSha", context.head_sha)?; + let acceptance_scope_id = normalize_scope_id(context.acceptance_scope_id)?; + let acceptance_scope_sha256 = normalize_sha256( + "publisherContext.acceptanceScopeSha256", + context.acceptance_scope_sha256, + )?; + + let status = git_stdout(cwd, &["status", "--porcelain=v1", "--untracked-files=all"]).await?; + if !status.is_empty() { + return Err(invalid_request( + "review publisher requires a clean working tree".to_string(), + )); + } + let actual_head = git_stdout(cwd, &["rev-parse", "--verify", "HEAD^{commit}"]).await?; + if actual_head != head_sha { + return Err(invalid_request( + "publisherContext.headSha does not match local HEAD".to_string(), + )); + } + let base_commit_expr = format!("{base_ref}^{{commit}}"); + let actual_base = git_stdout( + cwd, + &[ + "rev-parse", + "--verify", + "--end-of-options", + base_commit_expr.as_str(), + ], + ) + .await?; + if actual_base != reviewed_base_sha { + return Err(invalid_request( + "publisherContext.reviewedBaseSha does not match baseRef".to_string(), + )); + } + let remote = git_stdout(cwd, &["remote", "get-url", "origin"]).await?; + let repository_origin = canonicalize_git_remote_url(remote.as_str()).ok_or_else(|| { + invalid_request("origin remote is not a canonical repository URL".to_string()) + })?; + let merge_tree = git_stdout( + cwd, + &[ + "merge-tree", + "--write-tree", + reviewed_base_sha.as_str(), + head_sha.as_str(), + ], + ) + .await?; + let merge_result_tree_sha = merge_tree + .lines() + .next() + .map(str::trim) + .filter(|value| is_lower_hex(value, 40)) + .ok_or_else(|| invalid_request("git merge-tree returned no clean result tree".to_string()))? + .to_string(); + let message = git_stdout(cwd, &["show", "-s", "--format=%B", head_sha.as_str()]).await?; + let implementer_agent = parse_single_agent_trailer(message.as_str())?; + let candidate_sha256 = codex_state::review_candidate_sha256( + repository_origin.as_str(), + context.pull_request_number, + reviewed_base_sha.as_str(), + head_sha.as_str(), + merge_result_tree_sha.as_str(), + ) + .map_err(|err| internal_error(format!("failed to digest review candidate: {err}")))?; + let mut envelope = ReviewEnvelope { + schema_version: REVIEW_ENVELOPE_SCHEMA_VERSION.to_string(), + repository_origin, + pull_request_number: context.pull_request_number, + base_ref, + reviewed_base_sha, + head_sha: head_sha.clone(), + merge_result_tree_sha, + candidate_sha256, + acceptance_scope_id, + acceptance_scope_sha256, + implementer: ReviewImplementerProvenance { + source: ReviewImplementerProvenanceSource::GitAgentTrailer, + agent: implementer_agent, + commit_sha: head_sha, + }, + envelope_sha256: String::new(), + }; + envelope.envelope_sha256 = codex_state::review_envelope_sha256(&envelope) + .map_err(|err| internal_error(format!("failed to digest review envelope: {err}")))?; + Ok(envelope) +} + +async fn git_stdout(cwd: &Path, args: &[&str]) -> Result { + let output = tokio::time::timeout( + GIT_PROBE_TIMEOUT, + Command::new("git").args(args).current_dir(cwd).output(), + ) + .await + .map_err(|_| invalid_request("timed out while verifying review candidate".to_string()))? + .map_err(|_| { + invalid_request("failed to execute git while verifying review candidate".to_string()) + })?; + if !output.status.success() { + return Err(invalid_request(format!( + "git {} failed while verifying review candidate", + args.first().copied().unwrap_or("command") + ))); + } + String::from_utf8(output.stdout) + .map(|value| value.trim().to_string()) + .map_err(|_| invalid_request("git returned non-UTF-8 candidate metadata".to_string())) +} + +fn parse_single_agent_trailer(message: &str) -> Result { + let trailer_block = message + .rsplit_once("\n\n") + .map_or(message, |(_, tail)| tail); + let agents = trailer_block + .lines() + .filter_map(|line| line.strip_prefix("Agent:")) + .map(str::trim) + .filter(|agent| !agent.is_empty()) + .collect::>(); + if agents.len() != 1 + || agents[0].len() > 64 + || !agents[0] + .chars() + .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_')) + { + return Err(invalid_request( + "reviewed head commit must contain exactly one valid Agent trailer".to_string(), + )); + } + Ok(agents[0].to_string()) +} + +fn normalize_ref(value: String) -> Result { + let value = value.trim(); + if value.len() > 256 + || !value.starts_with("refs/") + || value.contains("..") + || value.contains("@{") + || value.chars().any(|ch| { + ch.is_ascii_control() + || ch.is_whitespace() + || matches!(ch, '~' | '^' | ':' | '?' | '*' | '[' | '\\') + }) + { + return Err(invalid_request( + "publisherContext.baseRef must be a full, safe Git ref".to_string(), + )); + } + Ok(value.to_string()) +} + +fn normalize_scope_id(value: String) -> Result { + let value = value.trim(); + if value.is_empty() + || value.len() > 128 + || !value + .chars() + .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.' | ':' | '/')) + { + return Err(invalid_request( + "publisherContext.acceptanceScopeId is invalid".to_string(), + )); + } + Ok(value.to_string()) +} + +fn normalize_git_sha(field: &str, value: String) -> Result { + let value = value.trim(); + if !is_lower_hex(value, 40) { + return Err(invalid_request(format!( + "publisherContext.{field} must be a full lowercase Git SHA" + ))); + } + Ok(value.to_string()) +} + +fn normalize_sha256(field: &str, value: String) -> Result { + let value = value.trim(); + if !is_lower_hex(value, 64) { + return Err(invalid_request(format!( + "{field} must be a lowercase SHA-256" + ))); + } + Ok(value.to_string()) +} + +fn normalize_digest_bound_id(field: &str, value: String) -> Result { + let value = value.trim(); + if value.is_empty() + || value.len() > 160 + || !value + .chars() + .all(|ch| ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_')) + { + return Err(invalid_request(format!("{field} is invalid"))); + } + Ok(value.to_string()) +} + +fn is_lower_hex(value: &str, len: usize) -> bool { + value.len() == len + && value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) +} + +fn api_run(snapshot: codex_state::ReviewPublisherRunSnapshot) -> ApiRun { + ApiRun { + review_run_id: snapshot.run.review_run_id, + envelope_sha256: snapshot.run.envelope_sha256, + status: match snapshot.run.status { + codex_state::ReviewPublisherRunStatus::Started => ApiRunStatus::Started, + codex_state::ReviewPublisherRunStatus::Completed => ApiRunStatus::Completed, + }, + verdict: snapshot.run.verdict.map(|verdict| match verdict { + codex_protocol::protocol::ReviewPublisherVerdict::Go => ApiVerdict::Go, + codex_protocol::protocol::ReviewPublisherVerdict::NoGo => ApiVerdict::NoGo, + }), + created_at: snapshot.run.created_at.timestamp(), + completed_at: snapshot.run.completed_at.map(|value| value.timestamp()), + events: snapshot.events.into_iter().map(api_event).collect(), + } +} + +fn api_event(event: codex_state::ReviewPublisherOutboxEvent) -> ApiOutboxEvent { + ApiOutboxEvent { + event_id: event.event_id, + event_kind: match event.event_kind { + codex_protocol::protocol::ReviewPublisherEventKind::Started => ApiEventKind::Started, + codex_protocol::protocol::ReviewPublisherEventKind::Completed => { + ApiEventKind::Completed + } + }, + sequence: event.sequence, + status: match event.status { + codex_state::ReviewPublisherOutboxStatus::Pending => ApiEventStatus::Pending, + codex_state::ReviewPublisherOutboxStatus::InFlight => ApiEventStatus::InFlight, + codex_state::ReviewPublisherOutboxStatus::Delivered => ApiEventStatus::Delivered, + codex_state::ReviewPublisherOutboxStatus::DeadLetter => ApiEventStatus::DeadLetter, + }, + payload_sha256: event.payload_sha256, + attempt_count: event.attempt_count, + next_attempt_at: event.next_attempt_at.timestamp(), + lease_expires_at: event.lease_expires_at.map(|value| value.timestamp()), + receipt_id: event.receipt_id, + last_error_code: event.last_error_code, + created_at: event.created_at.timestamp(), + delivered_at: event.delivered_at.map(|value| value.timestamp()), + } +} + +#[cfg(test)] +mod tests { + use super::*; + use std::fs; + use std::process::Command as StdCommand; + use wiremock::Mock; + use wiremock::MockServer; + use wiremock::ResponseTemplate; + use wiremock::matchers::method; + + #[test] + fn http_statuses_are_classified_fail_closed() { + assert_eq!( + classify_status(StatusCode::CONFLICT), + StatusDisposition::Delivered + ); + assert_eq!( + classify_status(StatusCode::UNAUTHORIZED), + StatusDisposition::DeadLetter + ); + assert_eq!( + classify_status(StatusCode::FORBIDDEN), + StatusDisposition::DeadLetter + ); + assert_eq!( + classify_status(StatusCode::TOO_MANY_REQUESTS), + StatusDisposition::Retry + ); + assert_eq!( + classify_status(StatusCode::INTERNAL_SERVER_ERROR), + StatusDisposition::Retry + ); + } + + #[tokio::test] + async fn http_timeout_is_retryable_without_persisting_response_content() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .respond_with(ResponseTemplate::new(200).set_delay(Duration::from_millis(100))) + .mount(&server) + .await; + let client = reqwest::Client::builder() + .timeout(Duration::from_millis(5)) + .build() + .expect("client"); + let error = client + .post(server.uri()) + .send() + .await + .expect_err("request must time out"); + assert_eq!( + classify_transport_error(&error), + DispatchResult::Retry { + error_code: "http_timeout".to_string() + } + ); + } + + #[test] + fn agent_provenance_comes_only_from_one_trailer() { + assert_eq!( + parse_single_agent_trailer("subject\n\nbody\n\nAgent: Herminia").unwrap(), + "Herminia" + ); + assert!(parse_single_agent_trailer("Agent: One\nAgent: Two").is_err()); + assert!(parse_single_agent_trailer("body names Agent: Caller").is_err()); + } + + #[tokio::test] + async fn envelope_rejects_candidate_mismatch_and_base_ref_movement() { + let repo = tempfile::tempdir().expect("temp repo"); + git(repo.path(), &["init", "-b", "main"]); + git(repo.path(), &["config", "user.name", "Test"]); + git(repo.path(), &["config", "user.email", "test@example.com"]); + git( + repo.path(), + &[ + "remote", + "add", + "origin", + "https://github.com/hasna/codewith.git", + ], + ); + fs::write(repo.path().join("file.txt"), "base\n").expect("base file"); + git(repo.path(), &["add", "file.txt"]); + git(repo.path(), &["commit", "-m", "base", "-m", "Agent: Base"]); + let base = git(repo.path(), &["rev-parse", "HEAD"]); + git(repo.path(), &["switch", "-c", "feature"]); + fs::write(repo.path().join("file.txt"), "base\nfeature\n").expect("feature file"); + git(repo.path(), &["add", "file.txt"]); + git( + repo.path(), + &["commit", "-m", "feature", "-m", "Agent: Herminia"], + ); + let head = git(repo.path(), &["rev-parse", "HEAD"]); + let context = publisher_context(base.clone(), head.clone()); + let envelope = build_review_envelope(repo.path(), context.clone()) + .await + .expect("verified envelope"); + assert_eq!(envelope.reviewed_base_sha, base); + assert_eq!(envelope.head_sha, head); + assert_eq!(envelope.implementer.agent, "Herminia"); + + let mut mismatched = context.clone(); + mismatched.head_sha = "c".repeat(40); + assert!( + build_review_envelope(repo.path(), mismatched) + .await + .is_err() + ); + + let base_tree = git(repo.path(), &["rev-parse", "refs/heads/main^{tree}"]); + let moved_base = git( + repo.path(), + &[ + "commit-tree", + base_tree.as_str(), + "-p", + base.as_str(), + "-m", + "moved base\n\nAgent: Base", + ], + ); + git( + repo.path(), + &["update-ref", "refs/heads/main", moved_base.as_str()], + ); + assert!(build_review_envelope(repo.path(), context).await.is_err()); + } + + fn publisher_context(reviewed_base_sha: String, head_sha: String) -> ReviewPublisherContext { + ReviewPublisherContext { + pull_request_number: 17, + base_ref: "refs/heads/main".to_string(), + reviewed_base_sha, + head_sha, + acceptance_scope_id: "codewith-review-envelope-v1".to_string(), + acceptance_scope_sha256: "d".repeat(64), + } + } + + fn git(cwd: &Path, args: &[&str]) -> String { + let output = StdCommand::new("git") + .args(args) + .current_dir(cwd) + .output() + .expect("git command"); + assert!( + output.status.success(), + "git command failed: {}", + args.join(" ") + ); + String::from_utf8(output.stdout) + .expect("git stdout") + .trim() + .to_string() + } +} diff --git a/codex-rs/app-server/src/request_processors/turn_processor.rs b/codex-rs/app-server/src/request_processors/turn_processor.rs index 54c9e840c4..245a8e819f 100644 --- a/codex-rs/app-server/src/request_processors/turn_processor.rs +++ b/codex-rs/app-server/src/request_processors/turn_processor.rs @@ -499,6 +499,7 @@ impl TurnRequestProcessor { let review_request = ReviewRequest { target: core_target, user_facing_hint: Some(hint.clone()), + review_envelope: None, }; Ok((review_request, hint)) @@ -1311,10 +1312,12 @@ impl TurnRequestProcessor { request_id: &ConnectionRequestId, turn: Turn, review_thread_id: String, + review_run_id: Option, ) { let response = ReviewStartResponse { turn, review_thread_id, + review_run_id, }; self.outgoing .send_response(request_id.clone(), response) @@ -1328,6 +1331,7 @@ impl TurnRequestProcessor { review_request: ReviewRequest, display_text: &str, parent_thread_id: String, + review_run_id: Option, ) -> std::result::Result<(), JSONRPCErrorError> { let turn_id = self .submit_core_op( @@ -1337,8 +1341,12 @@ impl TurnRequestProcessor { ) .await .map_err(|err| internal_error(format!("failed to start review: {err}")))?; + if let Some(review_run_id) = review_run_id.as_deref() { + self.bind_review_turn(review_run_id, turn_id.as_str()) + .await?; + } let turn = Self::build_review_turn(turn_id, display_text); - self.emit_review_started(request_id, turn, parent_thread_id) + self.emit_review_started(request_id, turn, parent_thread_id, review_run_id) .await; Ok(()) } @@ -1350,6 +1358,7 @@ impl TurnRequestProcessor { parent_thread: Arc, review_request: ReviewRequest, display_text: &str, + review_run_id: Option, ) -> std::result::Result<(), JSONRPCErrorError> { parent_thread.ensure_rollout_materialized().await; parent_thread.flush_rollout().await.map_err(|err| { @@ -1366,7 +1375,11 @@ impl TurnRequestProcessor { )) })?; - let mut config = self.config.as_ref().clone(); + let mut config = if review_request.review_envelope.is_some() { + parent_thread.config().await.as_ref().clone() + } else { + self.config.as_ref().clone() + }; if let Some(review_model) = &config.review_model { config.model = Some(review_model.clone()); } @@ -1446,9 +1459,13 @@ impl TurnRequestProcessor { internal_error(format!("failed to start detached review turn: {err}")) })?; + if let Some(review_run_id) = review_run_id.as_deref() { + self.bind_review_turn(review_run_id, turn_id.as_str()) + .await?; + } let turn = Self::build_review_turn(turn_id, display_text); let review_thread_id = thread_id.to_string(); - self.emit_review_started(request_id, turn, review_thread_id) + self.emit_review_started(request_id, turn, review_thread_id, review_run_id) .await; Ok(()) @@ -1463,11 +1480,47 @@ impl TurnRequestProcessor { thread_id, target, delivery, + publisher_context, } = params; let (parent_thread_id, parent_thread) = self.load_thread(&thread_id).await?; - let (review_request, display_text) = Self::review_request_from_target(target)?; - match delivery.unwrap_or(ApiReviewDelivery::Inline).to_core() { + if let Some(context) = publisher_context.as_ref() { + match &target { + ApiReviewTarget::BaseBranch { branch } + if branch.trim() == context.base_ref.trim() => {} + _ => { + return Err(invalid_request( + "publisherContext requires a baseBranch target matching baseRef" + .to_string(), + )); + } + } + } + let (mut review_request, display_text) = Self::review_request_from_target(target)?; + let review_run = if let Some(context) = publisher_context { + let cwd = parent_thread.config_snapshot().await.cwd; + let envelope = super::build_review_envelope(cwd.as_path(), context).await?; + let state_db = self.state_db.as_ref().ok_or_else(|| { + internal_error("review publisher requires durable state".to_string()) + })?; + let snapshot = state_db + .review_publisher() + .start_review_run(codex_state::ReviewPublisherStartParams { + thread_id: parent_thread_id.to_string(), + envelope: envelope.clone(), + }) + .await + .map_err(|err| { + internal_error(format!("failed to persist review publisher start: {err}")) + })?; + let run = Some((snapshot.run.review_run_id, envelope.envelope_sha256.clone())); + review_request.review_envelope = Some(envelope); + run + } else { + None + }; + let review_run_id = review_run.as_ref().map(|(run_id, _)| run_id.clone()); + let result = match delivery.unwrap_or(ApiReviewDelivery::Inline).to_core() { CoreReviewDelivery::Inline => { self.start_inline_review( request_id, @@ -1475,8 +1528,9 @@ impl TurnRequestProcessor { review_request, &display_text, thread_id, + review_run_id, ) - .await?; + .await } CoreReviewDelivery::Detached => { self.start_detached_review( @@ -1485,11 +1539,56 @@ impl TurnRequestProcessor { parent_thread, review_request, &display_text, + review_run_id, ) - .await?; + .await + } + }; + if result.is_err() { + if let Some((review_run_id, envelope_sha256)) = review_run { + self.terminalize_review_start_failure(review_run_id, envelope_sha256) + .await; } } - Ok(()) + result + } + + async fn bind_review_turn( + &self, + review_run_id: &str, + turn_id: &str, + ) -> Result<(), JSONRPCErrorError> { + let state_db = self + .state_db + .as_ref() + .ok_or_else(|| internal_error("review publisher requires durable state".to_string()))?; + state_db + .review_publisher() + .bind_review_turn(review_run_id, turn_id) + .await + .map_err(|err| internal_error(format!("failed to bind review publisher turn: {err}"))) + } + + async fn terminalize_review_start_failure( + &self, + review_run_id: String, + envelope_sha256: String, + ) { + let Some(state_db) = self.state_db.as_ref() else { + return; + }; + if let Err(err) = state_db + .review_publisher() + .complete_review_run(codex_state::ReviewPublisherCompleteParams { + review_run_id, + envelope_sha256, + review_output: None, + terminal_reason_override: Some("review_start_failed".to_string()), + }) + .await + { + warn!("failed to terminalize review publisher start failure: {err}"); + } } async fn turn_interrupt_inner( diff --git a/codex-rs/app-server/tests/suite/v2/client_metadata.rs b/codex-rs/app-server/tests/suite/v2/client_metadata.rs index 412b0f3308..264560c877 100644 --- a/codex-rs/app-server/tests/suite/v2/client_metadata.rs +++ b/codex-rs/app-server/tests/suite/v2/client_metadata.rs @@ -249,6 +249,7 @@ async fn review_start_sends_parent_lineage_in_turn_metadata_for_thread_fork_v2() target: ReviewTarget::Custom { instructions: "Review the fork".to_string(), }, + publisher_context: None, }) .await?; let review_resp: JSONRPCResponse = timeout( diff --git a/codex-rs/app-server/tests/suite/v2/review.rs b/codex-rs/app-server/tests/suite/v2/review.rs index 427a9f288b..147880db96 100644 --- a/codex-rs/app-server/tests/suite/v2/review.rs +++ b/codex-rs/app-server/tests/suite/v2/review.rs @@ -72,6 +72,7 @@ async fn review_start_runs_review_turn_and_emits_code_review_item() -> Result<() sha: "1234567deadbeef".to_string(), title: Some("Tidy UI colors".to_string()), }, + publisher_context: None, }) .await?; let review_resp: JSONRPCResponse = timeout( @@ -82,7 +83,9 @@ async fn review_start_runs_review_turn_and_emits_code_review_item() -> Result<() let ReviewStartResponse { turn, review_thread_id, + review_run_id, } = to_response::(review_resp)?; + assert!(review_run_id.is_none()); assert_eq!(review_thread_id, thread_id.clone()); let turn_id = turn.id.clone(); assert_eq!(turn.status, TurnStatus::InProgress); @@ -186,6 +189,7 @@ async fn review_start_exec_approval_item_id_matches_command_execution_item() -> sha: "1234567deadbeef".to_string(), title: Some("Check review approvals".to_string()), }, + publisher_context: None, }) .await?; let review_resp: JSONRPCResponse = timeout( @@ -267,6 +271,7 @@ async fn review_start_rejects_empty_base_branch() -> Result<()> { target: ReviewTarget::BaseBranch { branch: " ".to_string(), }, + publisher_context: None, }) .await?; let error: JSONRPCError = timeout( @@ -312,6 +317,7 @@ async fn review_start_with_detached_delivery_returns_new_thread_id() -> Result<( target: ReviewTarget::Custom { instructions: "detached review".to_string(), }, + publisher_context: None, }) .await?; let review_resp: JSONRPCResponse = timeout( @@ -322,7 +328,9 @@ async fn review_start_with_detached_delivery_returns_new_thread_id() -> Result<( let ReviewStartResponse { turn, review_thread_id, + review_run_id, } = to_response::(review_resp)?; + assert!(review_run_id.is_none()); assert_eq!(turn.status, TurnStatus::InProgress); assert_eq!(turn.items_view, TurnItemsView::NotLoaded); @@ -389,6 +397,7 @@ async fn review_start_rejects_empty_commit_sha() -> Result<()> { sha: "\t".to_string(), title: None, }, + publisher_context: None, }) .await?; let error: JSONRPCError = timeout( @@ -423,6 +432,7 @@ async fn review_start_rejects_empty_custom_instructions() -> Result<()> { target: ReviewTarget::Custom { instructions: "\n\n".to_string(), }, + publisher_context: None, }) .await?; let error: JSONRPCError = timeout( diff --git a/codex-rs/core/src/session/review.rs b/codex-rs/core/src/session/review.rs index c564dcc2f7..4fcd21050f 100644 --- a/codex-rs/core/src/session/review.rs +++ b/codex-rs/core/src/session/review.rs @@ -39,6 +39,7 @@ pub(super) async fn spawn_review_thread( ); let review_prompt = resolved.prompt.clone(); + let review_envelope = resolved.review_envelope.clone(); let provider = parent_turn_context.provider.clone(); let auth_manager = parent_turn_context.auth_manager.clone(); let model_info = review_model_info.clone(); @@ -183,12 +184,14 @@ pub(super) async fn spawn_review_thread( // TODO(ccunningham): Review turns currently rely on `spawn_task` for TurnComplete but do not // emit a parent TurnStarted. Consider giving review a full parent turn lifecycle // (TurnStarted + TurnComplete) for consistency with other standalone tasks. - sess.spawn_task(tc.clone(), input, ReviewTask::new()).await; + sess.spawn_task(tc.clone(), input, ReviewTask::new(review_envelope.clone())) + .await; // Announce entering review mode so UIs can switch modes. let review_request = ReviewRequest { target: resolved.target, user_facing_hint: Some(resolved.user_facing_hint), + review_envelope, }; sess.send_event(&tc, EventMsg::EnteredReviewMode(review_request)) .await; diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index 712a398a32..1a9f1f47af 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -11839,7 +11839,7 @@ async fn abort_review_task_emits_exited_then_aborted_and_records_history() { }], client_id: None, }]; - sess.spawn_task(Arc::clone(&tc), input, ReviewTask::new()) + sess.spawn_task(Arc::clone(&tc), input, ReviewTask::new(None)) .await; sess.abort_all_tasks(TurnAbortReason::Interrupted).await; diff --git a/codex-rs/core/src/tasks/review.rs b/codex-rs/core/src/tasks/review.rs index 9f7eba9573..137109d841 100644 --- a/codex-rs/core/src/tasks/review.rs +++ b/codex-rs/core/src/tasks/review.rs @@ -12,6 +12,7 @@ use codex_protocol::protocol::Event; use codex_protocol::protocol::EventMsg; use codex_protocol::protocol::ExitedReviewModeEvent; use codex_protocol::protocol::ItemCompletedEvent; +use codex_protocol::protocol::ReviewEnvelope; use codex_protocol::protocol::ReviewOutputEvent; use codex_protocol::protocol::SubAgentSource; use tokio_util::sync::CancellationToken; @@ -30,12 +31,14 @@ use codex_protocol::user_input::UserInput; use super::SessionTask; use super::SessionTaskContext; -#[derive(Clone, Copy)] -pub(crate) struct ReviewTask; +#[derive(Clone)] +pub(crate) struct ReviewTask { + review_envelope: Option, +} impl ReviewTask { - pub(crate) fn new() -> Self { - Self + pub(crate) fn new(review_envelope: Option) -> Self { + Self { review_envelope } } } @@ -82,13 +85,25 @@ impl SessionTask for ReviewTask { None => None, }; if !cancellation_token.is_cancelled() { - exit_review_mode(session.clone_session(), output.clone(), ctx.clone()).await; + exit_review_mode( + session.clone_session(), + output.clone(), + self.review_envelope.clone(), + ctx.clone(), + ) + .await; } None } async fn abort(&self, session: Arc, ctx: Arc) { - exit_review_mode(session.clone_session(), /*review_output*/ None, ctx).await; + exit_review_mode( + session.clone_session(), + /*review_output*/ None, + self.review_envelope.clone(), + ctx, + ) + .await; } } @@ -214,6 +229,7 @@ fn parse_review_output_event(text: &str) -> ReviewOutputEvent { pub(crate) async fn exit_review_mode( session: Arc, review_output: Option, + review_envelope: Option, ctx: Arc, ) { const REVIEW_USER_MESSAGE_ID: &str = "review_rollout_user"; @@ -254,7 +270,10 @@ pub(crate) async fn exit_review_mode( session .send_event( ctx.as_ref(), - EventMsg::ExitedReviewMode(ExitedReviewModeEvent { review_output }), + EventMsg::ExitedReviewMode(ExitedReviewModeEvent { + review_output, + review_envelope, + }), ) .await; session diff --git a/codex-rs/core/tests/suite/codex_delegate.rs b/codex-rs/core/tests/suite/codex_delegate.rs index f5db9856df..6b10fb4739 100644 --- a/codex-rs/core/tests/suite/codex_delegate.rs +++ b/codex-rs/core/tests/suite/codex_delegate.rs @@ -79,6 +79,7 @@ async fn codex_delegate_forwards_exec_approval_and_proceeds_on_approval() { instructions: "Please review".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -163,6 +164,7 @@ async fn codex_delegate_forwards_patch_approval_and_proceeds_on_decision() { instructions: "Please review".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -222,6 +224,7 @@ async fn codex_delegate_ignores_legacy_deltas() { instructions: "Please review".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await diff --git a/codex-rs/core/tests/suite/review.rs b/codex-rs/core/tests/suite/review.rs index 4ba711c99e..4e2bf15011 100644 --- a/codex-rs/core/tests/suite/review.rs +++ b/codex-rs/core/tests/suite/review.rs @@ -77,6 +77,7 @@ async fn review_op_emits_lifecycle_and_review_output() { instructions: "Please review my changes".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -223,6 +224,7 @@ async fn review_op_with_plain_text_emits_review_fallback() { instructions: "Plain text review".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -279,6 +281,7 @@ async fn review_filters_agent_message_related_events() { instructions: "Filter streaming events".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -353,6 +356,7 @@ async fn review_does_not_emit_agent_message_on_structured_output() { instructions: "check structured".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -410,6 +414,7 @@ async fn review_uses_custom_review_model_from_config() { instructions: "use custom model".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -421,7 +426,8 @@ async fn review_uses_custom_review_model_from_config() { matches!( ev, EventMsg::ExitedReviewMode(ExitedReviewModeEvent { - review_output: None + review_output: None, + review_envelope: None, }) ) }) @@ -460,6 +466,7 @@ async fn review_uses_session_model_when_review_model_unset() { instructions: "use session model".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -470,7 +477,8 @@ async fn review_uses_session_model_when_review_model_unset() { matches!( ev, EventMsg::ExitedReviewMode(ExitedReviewModeEvent { - review_output: None + review_output: None, + review_envelope: None, }) ) }) @@ -573,6 +581,7 @@ async fn review_input_isolated_from_parent_history() { instructions: review_prompt.clone(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -583,7 +592,8 @@ async fn review_input_isolated_from_parent_history() { matches!( ev, EventMsg::ExitedReviewMode(ExitedReviewModeEvent { - review_output: None + review_output: None, + review_envelope: None, }) ) }) @@ -685,6 +695,7 @@ async fn review_history_surfaces_in_parent_session() { instructions: "Start a review".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await @@ -694,7 +705,8 @@ async fn review_history_surfaces_in_parent_session() { matches!( ev, EventMsg::ExitedReviewMode(ExitedReviewModeEvent { - review_output: Some(_) + review_output: Some(_), + review_envelope: None, }) ) }) @@ -835,6 +847,7 @@ async fn review_uses_overridden_cwd_for_base_branch_merge_base() { branch: "main".to_string(), }, user_facing_hint: None, + review_envelope: None, }, }) .await diff --git a/codex-rs/exec/src/lib.rs b/codex-rs/exec/src/lib.rs index 022967b15c..e2010ad293 100644 --- a/codex-rs/exec/src/lib.rs +++ b/codex-rs/exec/src/lib.rs @@ -965,6 +965,7 @@ async fn run_exec_session(args: ExecRunArgs) -> anyhow::Result<()> { thread_id: primary_thread_id_for_span.clone(), target: review_target_to_api(review_request.target), delivery: None, + publisher_context: None, }, }, "review/start", @@ -2149,6 +2150,7 @@ fn build_review_request(args: &ReviewArgs) -> anyhow::Result { Ok(ReviewRequest { target, user_facing_hint: None, + review_envelope: None, }) } diff --git a/codex-rs/exec/src/lib_tests.rs b/codex-rs/exec/src/lib_tests.rs index 6bc057b62a..f2e39b4903 100644 --- a/codex-rs/exec/src/lib_tests.rs +++ b/codex-rs/exec/src/lib_tests.rs @@ -116,6 +116,7 @@ fn builds_uncommitted_review_request() { let expected = ReviewRequest { target: ReviewTarget::UncommittedChanges, user_facing_hint: None, + review_envelope: None, }; assert_eq!(request, expected); @@ -138,6 +139,7 @@ fn builds_commit_review_request_with_title() { title: Some("Add review command".to_string()), }, user_facing_hint: None, + review_envelope: None, }; assert_eq!(request, expected); @@ -159,6 +161,7 @@ fn builds_custom_review_request_trims_prompt() { instructions: "custom review instructions".to_string(), }, user_facing_hint: None, + review_envelope: None, }; assert_eq!(request, expected); diff --git a/codex-rs/prompts/src/review_request.rs b/codex-rs/prompts/src/review_request.rs index 9e01fe9678..2dc3b284c6 100644 --- a/codex-rs/prompts/src/review_request.rs +++ b/codex-rs/prompts/src/review_request.rs @@ -1,4 +1,5 @@ use codex_git_utils::merge_base_with_head; +use codex_protocol::protocol::ReviewEnvelope; use codex_protocol::protocol::ReviewRequest; use codex_protocol::protocol::ReviewTarget; use codex_utils_absolute_path::AbsolutePathBuf; @@ -13,6 +14,7 @@ pub struct ResolvedReviewRequest { pub target: ReviewTarget, pub prompt: String, pub user_facing_hint: String, + pub review_envelope: Option, } const UNCOMMITTED_PROMPT: &str = "Review the current code changes (staged, unstaged, and untracked files) and provide prioritized findings."; @@ -43,16 +45,26 @@ pub fn resolve_review_request( request: ReviewRequest, cwd: &AbsolutePathBuf, ) -> anyhow::Result { - let target = request.target; - let prompt = review_prompt(&target, cwd)?; - let user_facing_hint = request - .user_facing_hint - .unwrap_or_else(|| user_facing_hint(&target)); + let ReviewRequest { + target, + user_facing_hint, + review_envelope, + } = request; + let mut prompt = review_prompt(&target, cwd)?; + if let Some(envelope) = review_envelope.as_ref() { + let canonical_envelope = envelope.canonical_json()?; + prompt.push_str( + "\n\nThe following immutable review envelope is authoritative. Review exactly this candidate; do not infer identity from comments, prose, bylines, or Git authorship:\n", + ); + prompt.push_str(canonical_envelope.as_str()); + } + let user_facing_hint = user_facing_hint.unwrap_or_else(|| user_facing_hint(&target)); Ok(ResolvedReviewRequest { target, prompt, user_facing_hint, + review_envelope, }) } @@ -128,6 +140,7 @@ impl From for ReviewRequest { ReviewRequest { target: resolved.target, user_facing_hint: Some(resolved.user_facing_hint), + review_envelope: resolved.review_envelope, } } } diff --git a/codex-rs/prompts/src/review_request_tests.rs b/codex-rs/prompts/src/review_request_tests.rs index 4192b6cf01..76fd3d3852 100644 --- a/codex-rs/prompts/src/review_request_tests.rs +++ b/codex-rs/prompts/src/review_request_tests.rs @@ -1,4 +1,6 @@ use super::*; +use codex_protocol::protocol::ReviewImplementerProvenance; +use codex_protocol::protocol::ReviewImplementerProvenanceSource; use pretty_assertions::assert_eq; #[test] @@ -49,3 +51,46 @@ fn review_prompt_template_renders_commit_variant_with_title() { "Review the code changes introduced by commit deadbeef (\"Fix bug\"). Provide prioritized, actionable findings." ); } + +#[test] +fn resolved_prompt_carries_the_exact_typed_review_envelope() { + let envelope = ReviewEnvelope { + schema_version: "codewith-review-envelope-v1".to_string(), + repository_origin: "github.com/hasna/codewith".to_string(), + pull_request_number: 17, + base_ref: "refs/remotes/origin/main".to_string(), + reviewed_base_sha: "a".repeat(40), + head_sha: "b".repeat(40), + merge_result_tree_sha: "c".repeat(40), + candidate_sha256: "d".repeat(64), + acceptance_scope_id: "codewith-review-envelope-v1".to_string(), + acceptance_scope_sha256: "e".repeat(64), + implementer: ReviewImplementerProvenance { + source: ReviewImplementerProvenanceSource::GitAgentTrailer, + agent: "Herminia".to_string(), + commit_sha: "b".repeat(40), + }, + envelope_sha256: "f".repeat(64), + }; + let canonical = envelope.canonical_json().expect("canonical envelope"); + let resolved = resolve_review_request( + ReviewRequest { + target: ReviewTarget::Commit { + sha: envelope.head_sha.clone(), + title: None, + }, + user_facing_hint: None, + review_envelope: Some(envelope.clone()), + }, + &AbsolutePathBuf::current_dir().expect("cwd"), + ) + .expect("resolved review request"); + + assert_eq!(resolved.review_envelope, Some(envelope)); + assert_eq!(resolved.prompt.matches(canonical.as_str()).count(), 1); + assert!( + resolved + .prompt + .contains("do not infer identity from comments") + ); +} diff --git a/codex-rs/protocol/src/protocol.rs b/codex-rs/protocol/src/protocol.rs index 0fd2cc20c2..72c91e1f31 100644 --- a/codex-rs/protocol/src/protocol.rs +++ b/codex-rs/protocol/src/protocol.rs @@ -1900,6 +1900,9 @@ impl HasLegacyEvent for EventMsg { #[derive(Debug, Clone, Deserialize, Serialize, JsonSchema, TS)] pub struct ExitedReviewModeEvent { pub review_output: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + #[ts(optional)] + pub review_envelope: Option, } // Individual event payload types matching each `EventMsg` variant. @@ -3255,6 +3258,89 @@ pub struct ReviewRequest { #[serde(skip_serializing_if = "Option::is_none")] #[ts(optional)] pub user_facing_hint: Option, + /// Immutable, runtime-verified identity of the pull-request candidate. + /// + /// This is populated only by the trusted app-server review entrypoint. It + /// is never synthesized from review prose, comments, or model output. + #[serde(default, skip_serializing_if = "Option::is_none")] + #[ts(optional)] + pub review_envelope: Option, +} + +/// Immutable identity of an exact pull-request review candidate. +#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase")] +pub struct ReviewEnvelope { + pub schema_version: String, + pub repository_origin: String, + pub pull_request_number: u64, + pub base_ref: String, + pub reviewed_base_sha: String, + pub head_sha: String, + pub merge_result_tree_sha: String, + pub candidate_sha256: String, + pub acceptance_scope_id: String, + pub acceptance_scope_sha256: String, + pub implementer: ReviewImplementerProvenance, + pub envelope_sha256: String, +} + +impl ReviewEnvelope { + /// Serialize the envelope deterministically using its typed field order. + pub fn canonical_json(&self) -> serde_json::Result { + serde_json::to_string(self) + } +} + +/// Trusted implementer provenance derived from the reviewed head commit. +#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase")] +pub struct ReviewImplementerProvenance { + pub source: ReviewImplementerProvenanceSource, + pub agent: String, + pub commit_sha: String, +} + +#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase")] +pub enum ReviewImplementerProvenanceSource { + GitAgentTrailer, +} + +#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase")] +pub enum ReviewPublisherEventKind { + Started, + Completed, +} + +#[derive(Debug, Clone, Copy, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "SCREAMING_SNAKE_CASE")] +#[ts(rename_all = "SCREAMING_SNAKE_CASE")] +pub enum ReviewPublisherVerdict { + Go, + NoGo, +} + +/// Authenticated payload sent to the package-owned review publisher ingress. +#[derive(Debug, Clone, Deserialize, Serialize, PartialEq, Eq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(rename_all = "camelCase")] +pub struct ReviewPublisherEvent { + pub schema_version: String, + pub event_id: String, + pub event_kind: ReviewPublisherEventKind, + pub review_run_id: String, + pub sequence: u8, + pub envelope: ReviewEnvelope, + pub verdict: Option, + pub overall_correctness: Option, + pub finding_priorities: Vec, + pub terminal_reason: Option, } /// Structured review result produced by a child review session. diff --git a/codex-rs/state/migrations/0063_review_publisher_outbox.sql b/codex-rs/state/migrations/0063_review_publisher_outbox.sql new file mode 100644 index 0000000000..6fa2559ba0 --- /dev/null +++ b/codex-rs/state/migrations/0063_review_publisher_outbox.sql @@ -0,0 +1,48 @@ +CREATE TABLE review_publisher_runs ( + review_run_id TEXT PRIMARY KEY NOT NULL, + thread_id TEXT NOT NULL, + turn_id TEXT, + envelope_json TEXT NOT NULL, + envelope_sha256 TEXT NOT NULL UNIQUE, + status TEXT NOT NULL CHECK(status IN ('started', 'completed')), + verdict TEXT CHECK(verdict IS NULL OR verdict IN ('GO', 'NO_GO')), + terminal_reason TEXT, + created_at_ms INTEGER NOT NULL, + completed_at_ms INTEGER, + updated_at_ms INTEGER NOT NULL +); + +CREATE TABLE review_publisher_outbox_events ( + event_id TEXT PRIMARY KEY NOT NULL, + review_run_id TEXT NOT NULL REFERENCES review_publisher_runs(review_run_id) ON DELETE CASCADE, + event_kind TEXT NOT NULL CHECK(event_kind IN ('started', 'completed')), + sequence INTEGER NOT NULL CHECK(sequence IN (0, 1)), + status TEXT NOT NULL CHECK(status IN ('pending', 'in_flight', 'delivered', 'dead_letter')), + payload_json TEXT NOT NULL, + payload_sha256 TEXT NOT NULL, + attempt_count INTEGER NOT NULL DEFAULT 0 CHECK(attempt_count >= 0), + next_attempt_at_ms INTEGER NOT NULL, + lease_owner TEXT, + lease_expires_at_ms INTEGER, + receipt_id TEXT, + last_error_code TEXT, + created_at_ms INTEGER NOT NULL, + delivered_at_ms INTEGER, + updated_at_ms INTEGER NOT NULL, + UNIQUE(review_run_id, event_kind), + UNIQUE(review_run_id, sequence), + CHECK( + (event_kind = 'started' AND sequence = 0) + OR (event_kind = 'completed' AND sequence = 1) + ), + CHECK( + (status = 'in_flight' AND lease_owner IS NOT NULL AND lease_expires_at_ms IS NOT NULL) + OR (status != 'in_flight' AND lease_owner IS NULL AND lease_expires_at_ms IS NULL) + ) +); + +CREATE INDEX idx_review_publisher_outbox_due + ON review_publisher_outbox_events(status, next_attempt_at_ms, lease_expires_at_ms, sequence); + +CREATE INDEX idx_review_publisher_outbox_run + ON review_publisher_outbox_events(review_run_id, sequence); diff --git a/codex-rs/state/src/lib.rs b/codex-rs/state/src/lib.rs index 049d600741..5657c97c12 100644 --- a/codex-rs/state/src/lib.rs +++ b/codex-rs/state/src/lib.rs @@ -113,6 +113,12 @@ pub use model::PendingInteractionRespondParams; pub use model::PendingInteractionSourceKind; pub use model::PendingInteractionStatus; pub use model::PostGoalContextAction; +pub use model::ReviewPublisherOutboxClaim; +pub use model::ReviewPublisherOutboxEvent; +pub use model::ReviewPublisherOutboxStatus; +pub use model::ReviewPublisherRun; +pub use model::ReviewPublisherRunSnapshot; +pub use model::ReviewPublisherRunStatus; pub use model::SortDirection; pub use model::SortKey; pub use model::Stage1JobClaim; @@ -217,7 +223,17 @@ pub use runtime::MonitorStore; pub use runtime::PendingInteractionListParams; pub use runtime::PendingInteractionPage; pub use runtime::PendingInteractionRespondForSourceParams; +pub use runtime::REVIEW_ENVELOPE_SCHEMA_VERSION; +pub use runtime::REVIEW_PUBLISHER_EVENT_SCHEMA_VERSION; +pub use runtime::REVIEW_PUBLISHER_IMMUTABLE_CONFLICT; pub use runtime::RemoteControlEnrollmentRecord; +pub use runtime::ReviewPublisherClaimParams; +pub use runtime::ReviewPublisherCompleteParams; +pub use runtime::ReviewPublisherDeliveryAckParams; +pub use runtime::ReviewPublisherDeliveryFailParams; +pub use runtime::ReviewPublisherFailureDisposition; +pub use runtime::ReviewPublisherStartParams; +pub use runtime::ReviewPublisherStore; pub use runtime::RuntimeDbPath; pub use runtime::ScheduleStore; pub use runtime::ThreadFilterOptions; @@ -299,8 +315,12 @@ pub use runtime::goals_db_filename; pub use runtime::goals_db_path; pub use runtime::logs_db_filename; pub use runtime::logs_db_path; +pub use runtime::map_review_output; pub use runtime::memories_db_filename; pub use runtime::memories_db_path; +pub use runtime::review_candidate_sha256; +pub use runtime::review_envelope_sha256; +pub use runtime::review_run_id_from_envelope_sha256; pub use runtime::runtime_db_paths; pub use runtime::sqlite_integrity_check; pub use runtime::state_db_filename; diff --git a/codex-rs/state/src/model/mod.rs b/codex-rs/state/src/model/mod.rs index c5744df7b9..65e12c029c 100644 --- a/codex-rs/state/src/model/mod.rs +++ b/codex-rs/state/src/model/mod.rs @@ -8,6 +8,7 @@ mod mailbox; mod managed_worktree; mod memories; mod pending_interaction; +mod review_publisher; mod thread_goal; mod thread_metadata; mod thread_monitor; @@ -85,6 +86,12 @@ pub use pending_interaction::PendingInteractionKind; pub use pending_interaction::PendingInteractionRespondParams; pub use pending_interaction::PendingInteractionSourceKind; pub use pending_interaction::PendingInteractionStatus; +pub use review_publisher::ReviewPublisherOutboxClaim; +pub use review_publisher::ReviewPublisherOutboxEvent; +pub use review_publisher::ReviewPublisherOutboxStatus; +pub use review_publisher::ReviewPublisherRun; +pub use review_publisher::ReviewPublisherRunSnapshot; +pub use review_publisher::ReviewPublisherRunStatus; pub use thread_goal::PostGoalContextAction; pub use thread_goal::ThreadGoal; pub use thread_goal::ThreadGoalLineChangeStats; diff --git a/codex-rs/state/src/model/review_publisher.rs b/codex-rs/state/src/model/review_publisher.rs new file mode 100644 index 0000000000..f522f68dce --- /dev/null +++ b/codex-rs/state/src/model/review_publisher.rs @@ -0,0 +1,112 @@ +use chrono::DateTime; +use chrono::Utc; +use codex_protocol::protocol::ReviewEnvelope; +use codex_protocol::protocol::ReviewPublisherEvent; +use codex_protocol::protocol::ReviewPublisherEventKind; +use codex_protocol::protocol::ReviewPublisherVerdict; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ReviewPublisherRunStatus { + Started, + Completed, +} + +impl ReviewPublisherRunStatus { + pub fn as_str(self) -> &'static str { + match self { + Self::Started => "started", + Self::Completed => "completed", + } + } +} + +impl TryFrom<&str> for ReviewPublisherRunStatus { + type Error = anyhow::Error; + + fn try_from(value: &str) -> anyhow::Result { + match value { + "started" => Ok(Self::Started), + "completed" => Ok(Self::Completed), + other => anyhow::bail!("unknown review publisher run status `{other}`"), + } + } +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ReviewPublisherOutboxStatus { + Pending, + InFlight, + Delivered, + DeadLetter, +} + +impl ReviewPublisherOutboxStatus { + pub fn as_str(self) -> &'static str { + match self { + Self::Pending => "pending", + Self::InFlight => "in_flight", + Self::Delivered => "delivered", + Self::DeadLetter => "dead_letter", + } + } +} + +impl TryFrom<&str> for ReviewPublisherOutboxStatus { + type Error = anyhow::Error; + + fn try_from(value: &str) -> anyhow::Result { + match value { + "pending" => Ok(Self::Pending), + "in_flight" => Ok(Self::InFlight), + "delivered" => Ok(Self::Delivered), + "dead_letter" => Ok(Self::DeadLetter), + other => anyhow::bail!("unknown review publisher outbox status `{other}`"), + } + } +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ReviewPublisherRun { + pub review_run_id: String, + pub thread_id: String, + pub turn_id: Option, + pub envelope: ReviewEnvelope, + pub envelope_sha256: String, + pub status: ReviewPublisherRunStatus, + pub verdict: Option, + pub terminal_reason: Option, + pub created_at: DateTime, + pub completed_at: Option>, + pub updated_at: DateTime, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ReviewPublisherOutboxEvent { + pub event_id: String, + pub review_run_id: String, + pub event_kind: ReviewPublisherEventKind, + pub sequence: u8, + pub status: ReviewPublisherOutboxStatus, + pub payload: ReviewPublisherEvent, + pub payload_sha256: String, + pub attempt_count: u32, + pub next_attempt_at: DateTime, + pub lease_owner: Option, + pub lease_expires_at: Option>, + pub receipt_id: Option, + pub last_error_code: Option, + pub created_at: DateTime, + pub delivered_at: Option>, + pub updated_at: DateTime, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ReviewPublisherRunSnapshot { + pub run: ReviewPublisherRun, + pub events: Vec, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct ReviewPublisherOutboxClaim { + pub event: ReviewPublisherOutboxEvent, +} diff --git a/codex-rs/state/src/runtime.rs b/codex-rs/state/src/runtime.rs index 49739f1c84..f6e06efa4d 100644 --- a/codex-rs/state/src/runtime.rs +++ b/codex-rs/state/src/runtime.rs @@ -114,6 +114,7 @@ mod memories; mod monitors; mod pending_interactions; mod remote_control; +mod review_publisher; mod schedules; #[cfg(test)] mod test_support; @@ -243,6 +244,20 @@ pub use pending_interactions::PendingInteractionListParams; pub use pending_interactions::PendingInteractionPage; pub use pending_interactions::PendingInteractionRespondForSourceParams; pub use remote_control::RemoteControlEnrollmentRecord; +pub use review_publisher::REVIEW_ENVELOPE_SCHEMA_VERSION; +pub use review_publisher::REVIEW_PUBLISHER_EVENT_SCHEMA_VERSION; +pub use review_publisher::REVIEW_PUBLISHER_IMMUTABLE_CONFLICT; +pub use review_publisher::ReviewPublisherClaimParams; +pub use review_publisher::ReviewPublisherCompleteParams; +pub use review_publisher::ReviewPublisherDeliveryAckParams; +pub use review_publisher::ReviewPublisherDeliveryFailParams; +pub use review_publisher::ReviewPublisherFailureDisposition; +pub use review_publisher::ReviewPublisherStartParams; +pub use review_publisher::ReviewPublisherStore; +pub use review_publisher::map_review_output; +pub use review_publisher::review_candidate_sha256; +pub use review_publisher::review_envelope_sha256; +pub use review_publisher::review_run_id_from_envelope_sha256; pub use schedules::MAX_THREAD_SCHEDULE_NESTING_DEPTH; pub use schedules::ScheduleStore; pub use schedules::ThreadScheduleClaim; @@ -400,6 +415,7 @@ pub struct StateRuntime { thread_monitors: MonitorStore, local_active_sessions: LocalActiveSessionStore, webhook_events: WebhookEventStore, + review_publisher: ReviewPublisherStore, machine_registry: MachineRegistryStore, mailbox_messages: MailboxMessageStore, managed_worktrees: ManagedWorktreeStore, @@ -616,6 +632,7 @@ impl StateRuntime { thread_monitors: MonitorStore::new(Arc::clone(&pool)), local_active_sessions: LocalActiveSessionStore::new(Arc::clone(&pool)), webhook_events: WebhookEventStore::new(Arc::clone(&pool)), + review_publisher: ReviewPublisherStore::new(Arc::clone(&pool)), machine_registry: MachineRegistryStore::new(Arc::clone(&pool)), mailbox_messages: MailboxMessageStore::new(Arc::clone(&pool)), managed_worktrees: ManagedWorktreeStore::new(Arc::clone(&pool)), @@ -681,6 +698,10 @@ impl StateRuntime { &self.webhook_events } + pub fn review_publisher(&self) -> &ReviewPublisherStore { + &self.review_publisher + } + pub fn machine_registry(&self) -> &MachineRegistryStore { &self.machine_registry } diff --git a/codex-rs/state/src/runtime/review_publisher.rs b/codex-rs/state/src/runtime/review_publisher.rs new file mode 100644 index 0000000000..f21b0fc459 --- /dev/null +++ b/codex-rs/state/src/runtime/review_publisher.rs @@ -0,0 +1,1078 @@ +use super::*; +use codex_protocol::protocol::ReviewEnvelope; +use codex_protocol::protocol::ReviewOutputEvent; +use codex_protocol::protocol::ReviewPublisherEvent; +use codex_protocol::protocol::ReviewPublisherEventKind; +use codex_protocol::protocol::ReviewPublisherVerdict; +use sha2::Digest; +use sha2::Sha256; + +pub const REVIEW_ENVELOPE_SCHEMA_VERSION: &str = "codewith-review-envelope-v1"; +pub const REVIEW_PUBLISHER_EVENT_SCHEMA_VERSION: &str = "codewith-review-publisher-event-v1"; +pub const REVIEW_PUBLISHER_IMMUTABLE_CONFLICT: &str = "review publisher immutable payload conflict"; + +#[derive(Clone)] +pub struct ReviewPublisherStore { + pool: Arc, +} + +impl ReviewPublisherStore { + pub(crate) fn new(pool: Arc) -> Self { + Self { pool } + } +} + +#[derive(Debug, Clone)] +pub struct ReviewPublisherStartParams { + pub thread_id: String, + pub envelope: ReviewEnvelope, +} + +#[derive(Debug, Clone)] +pub struct ReviewPublisherCompleteParams { + pub review_run_id: String, + pub envelope_sha256: String, + pub review_output: Option, + pub terminal_reason_override: Option, +} + +#[derive(Debug, Clone)] +pub struct ReviewPublisherClaimParams { + pub lease_owner: String, + pub lease_duration: Duration, + pub now: DateTime, +} + +#[derive(Debug, Clone)] +pub struct ReviewPublisherDeliveryAckParams { + pub event_id: String, + pub lease_owner: String, + pub receipt_id: Option, + pub now: DateTime, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum ReviewPublisherFailureDisposition { + Retry, + DeadLetter, +} + +#[derive(Debug, Clone)] +pub struct ReviewPublisherDeliveryFailParams { + pub event_id: String, + pub lease_owner: String, + pub error_code: String, + pub disposition: ReviewPublisherFailureDisposition, + pub retry_at: DateTime, + pub now: DateTime, +} + +impl ReviewPublisherStore { + pub async fn start_review_run( + &self, + params: ReviewPublisherStartParams, + ) -> anyhow::Result { + verify_review_envelope(¶ms.envelope)?; + let review_run_id = review_run_id(params.envelope.envelope_sha256.as_str()); + let now_ms = datetime_to_epoch_millis(Utc::now()); + let envelope_json = params.envelope.canonical_json()?; + let start_event = publisher_event( + review_run_id.as_str(), + ReviewPublisherEventKind::Started, + params.envelope.clone(), + None, + None, + Vec::new(), + None, + ); + let start_event_json = serde_json::to_string(&start_event)?; + let start_payload_sha256 = sha256_hex(start_event_json.as_bytes()); + + let mut tx = self.pool.begin().await?; + if let Some(existing_sha) = sqlx::query_scalar::<_, String>( + "SELECT envelope_sha256 FROM review_publisher_runs WHERE review_run_id = ?", + ) + .bind(review_run_id.as_str()) + .fetch_optional(&mut *tx) + .await? + { + if existing_sha != params.envelope.envelope_sha256 { + anyhow::bail!(REVIEW_PUBLISHER_IMMUTABLE_CONFLICT); + } + tx.commit().await?; + return self + .get_review_run(review_run_id.as_str()) + .await? + .ok_or_else(|| anyhow::anyhow!("existing review run disappeared")); + } + + sqlx::query( + r#" +INSERT INTO review_publisher_runs ( + review_run_id, thread_id, turn_id, envelope_json, envelope_sha256, + status, verdict, terminal_reason, created_at_ms, completed_at_ms, updated_at_ms +) VALUES (?, ?, NULL, ?, ?, 'started', NULL, NULL, ?, NULL, ?) +"#, + ) + .bind(review_run_id.as_str()) + .bind(params.thread_id.as_str()) + .bind(envelope_json) + .bind(params.envelope.envelope_sha256.as_str()) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *tx) + .await?; + sqlx::query( + r#" +INSERT INTO review_publisher_outbox_events ( + event_id, review_run_id, event_kind, sequence, status, payload_json, + payload_sha256, attempt_count, next_attempt_at_ms, lease_owner, + lease_expires_at_ms, receipt_id, last_error_code, created_at_ms, + delivered_at_ms, updated_at_ms +) VALUES (?, ?, 'started', 0, 'pending', ?, ?, 0, ?, NULL, NULL, NULL, NULL, ?, NULL, ?) +"#, + ) + .bind(start_event.event_id.as_str()) + .bind(review_run_id.as_str()) + .bind(start_event_json) + .bind(start_payload_sha256) + .bind(now_ms) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *tx) + .await?; + tx.commit().await?; + + self.get_review_run(review_run_id.as_str()) + .await? + .ok_or_else(|| anyhow::anyhow!("created review run disappeared")) + } + + pub async fn bind_review_turn(&self, review_run_id: &str, turn_id: &str) -> anyhow::Result<()> { + let updated = sqlx::query( + r#" +UPDATE review_publisher_runs +SET turn_id = COALESCE(turn_id, ?), updated_at_ms = ? +WHERE review_run_id = ? AND (turn_id IS NULL OR turn_id = ?) +"#, + ) + .bind(turn_id) + .bind(datetime_to_epoch_millis(Utc::now())) + .bind(review_run_id) + .bind(turn_id) + .execute(self.pool.as_ref()) + .await?; + if updated.rows_affected() != 1 { + anyhow::bail!(REVIEW_PUBLISHER_IMMUTABLE_CONFLICT); + } + Ok(()) + } + + pub async fn complete_review_run( + &self, + params: ReviewPublisherCompleteParams, + ) -> anyhow::Result { + let mut tx = self.pool.begin().await?; + let row = sqlx::query( + "SELECT envelope_json, envelope_sha256, status FROM review_publisher_runs WHERE review_run_id = ?", + ) + .bind(params.review_run_id.as_str()) + .fetch_optional(&mut *tx) + .await? + .ok_or_else(|| anyhow::anyhow!("review publisher run not found"))?; + let envelope_sha256: String = row.try_get("envelope_sha256")?; + if envelope_sha256 != params.envelope_sha256 { + anyhow::bail!(REVIEW_PUBLISHER_IMMUTABLE_CONFLICT); + } + let envelope: ReviewEnvelope = serde_json::from_str(row.try_get("envelope_json")?)?; + verify_review_envelope(&envelope)?; + let (verdict, overall_correctness, finding_priorities, mapped_reason) = + map_review_output(params.review_output.as_ref()); + let terminal_reason = params + .terminal_reason_override + .as_deref() + .map(normalize_terminal_reason) + .transpose()? + .unwrap_or(mapped_reason); + let terminal_event = publisher_event( + params.review_run_id.as_str(), + ReviewPublisherEventKind::Completed, + envelope, + Some(verdict), + overall_correctness, + finding_priorities, + Some(terminal_reason.clone()), + ); + let terminal_event_json = serde_json::to_string(&terminal_event)?; + let terminal_payload_sha256 = sha256_hex(terminal_event_json.as_bytes()); + let status: String = row.try_get("status")?; + + if status == "completed" { + let existing_sha = sqlx::query_scalar::<_, String>( + "SELECT payload_sha256 FROM review_publisher_outbox_events WHERE review_run_id = ? AND sequence = 1", + ) + .bind(params.review_run_id.as_str()) + .fetch_optional(&mut *tx) + .await? + .ok_or_else(|| anyhow::anyhow!("completed review run has no terminal event"))?; + if existing_sha != terminal_payload_sha256 { + anyhow::bail!(REVIEW_PUBLISHER_IMMUTABLE_CONFLICT); + } + tx.commit().await?; + return self + .get_review_run(params.review_run_id.as_str()) + .await? + .ok_or_else(|| anyhow::anyhow!("completed review run disappeared")); + } + + let start_exists = sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM review_publisher_outbox_events WHERE review_run_id = ? AND sequence = 0", + ) + .bind(params.review_run_id.as_str()) + .fetch_one(&mut *tx) + .await?; + if start_exists != 1 { + anyhow::bail!("review publisher terminal event requires one start event"); + } + let now_ms = datetime_to_epoch_millis(Utc::now()); + sqlx::query( + r#" +UPDATE review_publisher_runs +SET status = 'completed', verdict = ?, terminal_reason = ?, completed_at_ms = ?, updated_at_ms = ? +WHERE review_run_id = ? AND status = 'started' +"#, + ) + .bind(verdict_str(verdict)) + .bind(terminal_reason.as_str()) + .bind(now_ms) + .bind(now_ms) + .bind(params.review_run_id.as_str()) + .execute(&mut *tx) + .await?; + sqlx::query( + r#" +INSERT INTO review_publisher_outbox_events ( + event_id, review_run_id, event_kind, sequence, status, payload_json, + payload_sha256, attempt_count, next_attempt_at_ms, lease_owner, + lease_expires_at_ms, receipt_id, last_error_code, created_at_ms, + delivered_at_ms, updated_at_ms +) VALUES (?, ?, 'completed', 1, 'pending', ?, ?, 0, ?, NULL, NULL, NULL, NULL, ?, NULL, ?) +"#, + ) + .bind(terminal_event.event_id.as_str()) + .bind(params.review_run_id.as_str()) + .bind(terminal_event_json) + .bind(terminal_payload_sha256) + .bind(now_ms) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *tx) + .await?; + tx.commit().await?; + + self.get_review_run(params.review_run_id.as_str()) + .await? + .ok_or_else(|| anyhow::anyhow!("completed review run disappeared")) + } + + pub async fn get_review_run( + &self, + review_run_id: &str, + ) -> anyhow::Result> { + let Some(row) = sqlx::query( + r#" +SELECT review_run_id, thread_id, turn_id, envelope_json, envelope_sha256, + status, verdict, terminal_reason, created_at_ms, completed_at_ms, updated_at_ms +FROM review_publisher_runs WHERE review_run_id = ? +"#, + ) + .bind(review_run_id) + .fetch_optional(self.pool.as_ref()) + .await? + else { + return Ok(None); + }; + let run = review_run_from_row(&row)?; + let query = format!( + "{} WHERE review_run_id = ? ORDER BY sequence ASC", + outbox_select("SELECT") + ); + let event_rows = sqlx::query(sqlx::AssertSqlSafe(query)) + .bind(review_run_id) + .fetch_all(self.pool.as_ref()) + .await?; + let events = event_rows + .iter() + .map(review_outbox_from_row) + .collect::>>()?; + Ok(Some(crate::ReviewPublisherRunSnapshot { run, events })) + } + + pub async fn claim_next_due_event( + &self, + params: ReviewPublisherClaimParams, + ) -> anyhow::Result> { + let now_ms = datetime_to_epoch_millis(params.now); + let lease_expires_at_ms = now_ms + .saturating_add(i64::try_from(params.lease_duration.as_millis()).unwrap_or(i64::MAX)); + let mut tx = self.pool.begin().await?; + let query = format!( + r#"{} +WHERE ( + (status = 'pending' AND next_attempt_at_ms <= ?) + OR (status = 'in_flight' AND lease_expires_at_ms <= ?) +) +AND ( + sequence = 0 + OR EXISTS ( + SELECT 1 FROM review_publisher_outbox_events start_event + WHERE start_event.review_run_id = review_publisher_outbox_events.review_run_id + AND start_event.sequence = 0 + AND start_event.status = 'delivered' + ) +) +ORDER BY sequence ASC, next_attempt_at_ms ASC, created_at_ms ASC +LIMIT 1 +"#, + outbox_select("SELECT") + ); + let Some(row) = sqlx::query(sqlx::AssertSqlSafe(query)) + .bind(now_ms) + .bind(now_ms) + .fetch_optional(&mut *tx) + .await? + else { + tx.commit().await?; + return Ok(None); + }; + let event_id: String = row.try_get("event_id")?; + let updated = sqlx::query( + r#" +UPDATE review_publisher_outbox_events +SET status = 'in_flight', attempt_count = attempt_count + 1, + lease_owner = ?, lease_expires_at_ms = ?, updated_at_ms = ? +WHERE event_id = ? + AND ( + (status = 'pending' AND next_attempt_at_ms <= ?) + OR (status = 'in_flight' AND lease_expires_at_ms <= ?) + ) +"#, + ) + .bind(params.lease_owner.as_str()) + .bind(lease_expires_at_ms) + .bind(now_ms) + .bind(event_id.as_str()) + .bind(now_ms) + .bind(now_ms) + .execute(&mut *tx) + .await?; + if updated.rows_affected() != 1 { + tx.rollback().await?; + return Ok(None); + } + let query = format!("{} WHERE event_id = ?", outbox_select("SELECT")); + let claimed_row = sqlx::query(sqlx::AssertSqlSafe(query)) + .bind(event_id.as_str()) + .fetch_one(&mut *tx) + .await?; + let event = review_outbox_from_row(&claimed_row)?; + tx.commit().await?; + Ok(Some(crate::ReviewPublisherOutboxClaim { event })) + } + + pub async fn acknowledge_delivery( + &self, + params: ReviewPublisherDeliveryAckParams, + ) -> anyhow::Result { + let now_ms = datetime_to_epoch_millis(params.now); + let updated = sqlx::query( + r#" +UPDATE review_publisher_outbox_events +SET status = 'delivered', lease_owner = NULL, lease_expires_at_ms = NULL, + receipt_id = ?, last_error_code = NULL, delivered_at_ms = ?, updated_at_ms = ? +WHERE event_id = ? AND status = 'in_flight' AND lease_owner = ? +"#, + ) + .bind(params.receipt_id.as_deref()) + .bind(now_ms) + .bind(now_ms) + .bind(params.event_id.as_str()) + .bind(params.lease_owner.as_str()) + .execute(self.pool.as_ref()) + .await?; + Ok(updated.rows_affected() == 1) + } + + pub async fn fail_delivery( + &self, + params: ReviewPublisherDeliveryFailParams, + ) -> anyhow::Result { + let error_code = normalize_error_code(params.error_code.as_str())?; + let now_ms = datetime_to_epoch_millis(params.now); + let mut tx = self.pool.begin().await?; + let status = match params.disposition { + ReviewPublisherFailureDisposition::Retry => "pending", + ReviewPublisherFailureDisposition::DeadLetter => "dead_letter", + }; + let updated = sqlx::query( + r#" +UPDATE review_publisher_outbox_events +SET status = ?, next_attempt_at_ms = ?, lease_owner = NULL, + lease_expires_at_ms = NULL, last_error_code = ?, updated_at_ms = ? +WHERE event_id = ? AND status = 'in_flight' AND lease_owner = ? +"#, + ) + .bind(status) + .bind(datetime_to_epoch_millis(params.retry_at)) + .bind(error_code.as_str()) + .bind(now_ms) + .bind(params.event_id.as_str()) + .bind(params.lease_owner.as_str()) + .execute(&mut *tx) + .await?; + if updated.rows_affected() != 1 { + tx.rollback().await?; + return Ok(false); + } + if params.disposition == ReviewPublisherFailureDisposition::DeadLetter { + let sequence = sqlx::query_scalar::<_, i64>( + "SELECT sequence FROM review_publisher_outbox_events WHERE event_id = ?", + ) + .bind(params.event_id.as_str()) + .fetch_one(&mut *tx) + .await?; + if sequence == 0 { + sqlx::query( + r#" +UPDATE review_publisher_outbox_events +SET status = 'dead_letter', last_error_code = 'blocked_by_start', updated_at_ms = ? +WHERE review_run_id = ( + SELECT review_run_id FROM review_publisher_outbox_events WHERE event_id = ? +) +AND sequence > 0 AND status = 'pending' +"#, + ) + .bind(now_ms) + .bind(params.event_id.as_str()) + .execute(&mut *tx) + .await?; + } + } + tx.commit().await?; + Ok(true) + } + + pub async fn exact_replay( + &self, + event_id: &str, + payload_sha256: &str, + now: DateTime, + ) -> anyhow::Result> { + let now_ms = datetime_to_epoch_millis(now); + let updated = sqlx::query( + r#" +UPDATE review_publisher_outbox_events +SET status = 'pending', attempt_count = 0, next_attempt_at_ms = ?, + lease_owner = NULL, lease_expires_at_ms = NULL, receipt_id = NULL, + last_error_code = NULL, delivered_at_ms = NULL, updated_at_ms = ? +WHERE event_id = ? AND payload_sha256 = ? AND status != 'in_flight' +"#, + ) + .bind(now_ms) + .bind(now_ms) + .bind(event_id) + .bind(payload_sha256) + .execute(self.pool.as_ref()) + .await?; + if updated.rows_affected() == 0 { + return Ok(None); + } + self.get_outbox_event(event_id).await + } + + pub async fn get_outbox_event( + &self, + event_id: &str, + ) -> anyhow::Result> { + let query = format!("{} WHERE event_id = ?", outbox_select("SELECT")); + let row = sqlx::query(sqlx::AssertSqlSafe(query)) + .bind(event_id) + .fetch_optional(self.pool.as_ref()) + .await?; + row.as_ref().map(review_outbox_from_row).transpose() + } +} + +pub fn review_envelope_sha256(envelope: &ReviewEnvelope) -> anyhow::Result { + let mut unsigned = envelope.clone(); + unsigned.envelope_sha256.clear(); + Ok(sha256_hex(unsigned.canonical_json()?.as_bytes())) +} + +pub fn review_candidate_sha256( + repository_origin: &str, + pull_request_number: u64, + reviewed_base_sha: &str, + head_sha: &str, + merge_result_tree_sha: &str, +) -> anyhow::Result { + #[derive(serde::Serialize)] + #[serde(rename_all = "camelCase")] + struct ReviewCandidateIdentity<'a> { + repository_origin: &'a str, + pull_request_number: u64, + reviewed_base_sha: &'a str, + head_sha: &'a str, + merge_result_tree_sha: &'a str, + } + + let identity = ReviewCandidateIdentity { + repository_origin, + pull_request_number, + reviewed_base_sha, + head_sha, + merge_result_tree_sha, + }; + Ok(sha256_hex(serde_json::to_vec(&identity)?.as_slice())) +} + +pub fn map_review_output( + output: Option<&ReviewOutputEvent>, +) -> (ReviewPublisherVerdict, Option, Vec, String) { + let Some(output) = output else { + return ( + ReviewPublisherVerdict::NoGo, + None, + Vec::new(), + "missing_or_interrupted_review_output".to_string(), + ); + }; + let correctness = match output.overall_correctness.as_str() { + "patch is correct" => Some("patch is correct".to_string()), + "patch is incorrect" => Some("patch is incorrect".to_string()), + _ => None, + }; + let priorities = output + .findings + .iter() + .map(|finding| finding.priority) + .collect::>(); + if correctness.as_deref() != Some("patch is correct") { + return ( + ReviewPublisherVerdict::NoGo, + correctness, + priorities, + "overall_correctness_not_exact_go".to_string(), + ); + } + if output + .findings + .iter() + .any(|finding| !(0..=3).contains(&finding.priority)) + { + return ( + ReviewPublisherVerdict::NoGo, + correctness, + priorities, + "unknown_finding_priority".to_string(), + ); + } + if output.findings.iter().any(|finding| { + finding_priority_from_title(finding.title.as_str()) != Some(finding.priority) + }) { + return ( + ReviewPublisherVerdict::NoGo, + correctness, + priorities, + "finding_priority_label_mismatch".to_string(), + ); + } + if priorities.iter().any(|priority| matches!(priority, 0 | 1)) { + return ( + ReviewPublisherVerdict::NoGo, + correctness, + priorities, + "blocking_finding_present".to_string(), + ); + } + ( + ReviewPublisherVerdict::Go, + correctness, + priorities, + "exact_go".to_string(), + ) +} + +fn verify_review_envelope(envelope: &ReviewEnvelope) -> anyhow::Result<()> { + if envelope.schema_version != REVIEW_ENVELOPE_SCHEMA_VERSION + || envelope.envelope_sha256.len() != 64 + || review_envelope_sha256(envelope)? != envelope.envelope_sha256 + { + anyhow::bail!(REVIEW_PUBLISHER_IMMUTABLE_CONFLICT); + } + Ok(()) +} + +fn publisher_event( + review_run_id: &str, + event_kind: ReviewPublisherEventKind, + envelope: ReviewEnvelope, + verdict: Option, + overall_correctness: Option, + finding_priorities: Vec, + terminal_reason: Option, +) -> ReviewPublisherEvent { + let sequence = match event_kind { + ReviewPublisherEventKind::Started => 0, + ReviewPublisherEventKind::Completed => 1, + }; + ReviewPublisherEvent { + schema_version: REVIEW_PUBLISHER_EVENT_SCHEMA_VERSION.to_string(), + event_id: format!("review-event-{}-{sequence}", envelope.envelope_sha256), + event_kind, + review_run_id: review_run_id.to_string(), + sequence, + envelope, + verdict, + overall_correctness, + finding_priorities, + terminal_reason, + } +} + +fn review_run_id(envelope_sha256: &str) -> String { + format!("review-run-{envelope_sha256}") +} + +pub fn review_run_id_from_envelope_sha256(envelope_sha256: &str) -> String { + review_run_id(envelope_sha256) +} + +fn finding_priority_from_title(title: &str) -> Option { + let bytes = title.as_bytes(); + if bytes.len() < 4 || bytes[0] != b'[' || bytes[1] != b'P' || bytes[3] != b']' { + return None; + } + char::from(bytes[2]).to_digit(10).map(|value| value as i32) +} + +fn normalize_terminal_reason(reason: &str) -> anyhow::Result { + let reason = reason.trim(); + if reason.is_empty() + || reason.len() > 128 + || !reason + .chars() + .all(|ch| ch.is_ascii_lowercase() || ch.is_ascii_digit() || ch == '_') + { + anyhow::bail!("invalid review publisher terminal reason"); + } + Ok(reason.to_string()) +} + +fn normalize_error_code(code: &str) -> anyhow::Result { + normalize_terminal_reason(code) +} + +fn verdict_str(verdict: ReviewPublisherVerdict) -> &'static str { + match verdict { + ReviewPublisherVerdict::Go => "GO", + ReviewPublisherVerdict::NoGo => "NO_GO", + } +} + +fn verdict_from_str(value: Option<&str>) -> anyhow::Result> { + value + .map(|value| match value { + "GO" => Ok(ReviewPublisherVerdict::Go), + "NO_GO" => Ok(ReviewPublisherVerdict::NoGo), + other => anyhow::bail!("unknown review publisher verdict `{other}`"), + }) + .transpose() +} + +fn event_kind_str(kind: ReviewPublisherEventKind) -> &'static str { + match kind { + ReviewPublisherEventKind::Started => "started", + ReviewPublisherEventKind::Completed => "completed", + } +} + +fn event_kind_from_str(value: &str) -> anyhow::Result { + match value { + "started" => Ok(ReviewPublisherEventKind::Started), + "completed" => Ok(ReviewPublisherEventKind::Completed), + other => anyhow::bail!("unknown review publisher event kind `{other}`"), + } +} + +fn review_run_from_row(row: &sqlx::sqlite::SqliteRow) -> anyhow::Result { + let verdict: Option = row.try_get("verdict")?; + Ok(crate::ReviewPublisherRun { + review_run_id: row.try_get("review_run_id")?, + thread_id: row.try_get("thread_id")?, + turn_id: row.try_get("turn_id")?, + envelope: serde_json::from_str(row.try_get("envelope_json")?)?, + envelope_sha256: row.try_get("envelope_sha256")?, + status: crate::ReviewPublisherRunStatus::try_from( + row.try_get::("status")?.as_str(), + )?, + verdict: verdict_from_str(verdict.as_deref())?, + terminal_reason: row.try_get("terminal_reason")?, + created_at: epoch_millis_to_datetime(row.try_get("created_at_ms")?)?, + completed_at: row + .try_get::, _>("completed_at_ms")? + .map(epoch_millis_to_datetime) + .transpose()?, + updated_at: epoch_millis_to_datetime(row.try_get("updated_at_ms")?)?, + }) +} + +fn outbox_select(prefix: &str) -> String { + format!( + "{prefix} event_id, review_run_id, event_kind, sequence, status, payload_json, \ + payload_sha256, attempt_count, next_attempt_at_ms, lease_owner, \ + lease_expires_at_ms, receipt_id, last_error_code, created_at_ms, \ + delivered_at_ms, updated_at_ms FROM review_publisher_outbox_events" + ) +} + +fn review_outbox_from_row( + row: &sqlx::sqlite::SqliteRow, +) -> anyhow::Result { + let event_kind = event_kind_from_str(row.try_get::("event_kind")?.as_str())?; + let event_id: String = row.try_get("event_id")?; + let review_run_id: String = row.try_get("review_run_id")?; + let sequence: u8 = row.try_get::("sequence")?.try_into()?; + let payload_json: String = row.try_get("payload_json")?; + let payload_sha256: String = row.try_get("payload_sha256")?; + let payload: ReviewPublisherEvent = serde_json::from_str(payload_json.as_str())?; + if payload.event_kind != event_kind + || event_kind_str(payload.event_kind) != event_kind_str(event_kind) + || matches!(event_kind, ReviewPublisherEventKind::Started) != (sequence == 0) + || payload.event_id != event_id + || payload.review_run_id != review_run_id + || payload.sequence != sequence + || sha256_hex(payload_json.as_bytes()) != payload_sha256 + || review_envelope_sha256(&payload.envelope)? != payload.envelope.envelope_sha256 + { + anyhow::bail!(REVIEW_PUBLISHER_IMMUTABLE_CONFLICT); + } + Ok(crate::ReviewPublisherOutboxEvent { + event_id, + review_run_id, + event_kind, + sequence, + status: crate::ReviewPublisherOutboxStatus::try_from( + row.try_get::("status")?.as_str(), + )?, + payload, + payload_sha256, + attempt_count: row.try_get::("attempt_count")?.try_into()?, + next_attempt_at: epoch_millis_to_datetime(row.try_get("next_attempt_at_ms")?)?, + lease_owner: row.try_get("lease_owner")?, + lease_expires_at: row + .try_get::, _>("lease_expires_at_ms")? + .map(epoch_millis_to_datetime) + .transpose()?, + receipt_id: row.try_get("receipt_id")?, + last_error_code: row.try_get("last_error_code")?, + created_at: epoch_millis_to_datetime(row.try_get("created_at_ms")?)?, + delivered_at: row + .try_get::, _>("delivered_at_ms")? + .map(epoch_millis_to_datetime) + .transpose()?, + updated_at: epoch_millis_to_datetime(row.try_get("updated_at_ms")?)?, + }) +} + +fn sha256_hex(bytes: &[u8]) -> String { + format!("{:x}", Sha256::digest(bytes)) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::runtime::test_support::unique_temp_dir; + use codex_protocol::protocol::ReviewCodeLocation; + use codex_protocol::protocol::ReviewFinding; + use codex_protocol::protocol::ReviewImplementerProvenance; + use codex_protocol::protocol::ReviewImplementerProvenanceSource; + use codex_protocol::protocol::ReviewLineRange; + use pretty_assertions::assert_eq; + + #[test] + fn candidate_digest_has_a_feature_independent_typed_order() { + let digest = review_candidate_sha256( + "github.com/hasna/codewith", + 475, + "06ff792de7303ba5867c8a329e0aca80cc65cbf2", + "f69adb218518d88f2192dccb670d01b2a5402239", + "fd2de011b53698b4b16dc2f4fcb3a2661760228c", + ) + .expect("candidate digest"); + + assert_eq!( + digest, + "a77396eb10afeb21407736ad90e43ccc391751e27c0f4254a2826bd7996f6b6a" + ); + } + + fn envelope() -> ReviewEnvelope { + let mut envelope = ReviewEnvelope { + schema_version: REVIEW_ENVELOPE_SCHEMA_VERSION.to_string(), + repository_origin: "github.com/hasna/codewith".to_string(), + pull_request_number: 1, + base_ref: "refs/remotes/origin/main".to_string(), + reviewed_base_sha: "a".repeat(40), + head_sha: "b".repeat(40), + merge_result_tree_sha: "c".repeat(40), + candidate_sha256: "d".repeat(64), + acceptance_scope_id: "codewith-review-envelope-v1".to_string(), + acceptance_scope_sha256: "e".repeat(64), + implementer: ReviewImplementerProvenance { + source: ReviewImplementerProvenanceSource::GitAgentTrailer, + agent: "Herminia".to_string(), + commit_sha: "b".repeat(40), + }, + envelope_sha256: String::new(), + }; + envelope.envelope_sha256 = review_envelope_sha256(&envelope).expect("digest"); + envelope + } + + fn finding(priority: i32, title_priority: i32) -> ReviewFinding { + ReviewFinding { + title: format!("[P{title_priority}] finding"), + body: "body".to_string(), + confidence_score: 1.0, + priority, + code_location: ReviewCodeLocation { + absolute_file_path: "/tmp/file.rs".into(), + line_range: ReviewLineRange { start: 1, end: 1 }, + }, + } + } + + fn output(correctness: &str, findings: Vec) -> ReviewOutputEvent { + ReviewOutputEvent { + findings, + overall_correctness: correctness.to_string(), + overall_explanation: "ignored untrusted prose".to_string(), + overall_confidence_score: 1.0, + } + } + + async fn runtime() -> Arc { + StateRuntime::init(unique_temp_dir(), "test-provider".to_string()) + .await + .expect("state runtime") + } + + #[test] + fn verdict_mapping_is_fail_closed() { + assert_eq!(ReviewPublisherVerdict::NoGo, map_review_output(None).0); + assert_eq!( + ReviewPublisherVerdict::NoGo, + map_review_output(Some(&output("patch is correct ", Vec::new()))).0 + ); + assert_eq!( + ReviewPublisherVerdict::NoGo, + map_review_output(Some(&output("patch is correct", vec![finding(4, 4)]))).0 + ); + assert_eq!( + ReviewPublisherVerdict::NoGo, + map_review_output(Some(&output("patch is correct", vec![finding(2, 3)]))).0 + ); + assert_eq!( + ReviewPublisherVerdict::NoGo, + map_review_output(Some(&output("patch is correct", vec![finding(1, 1)]))).0 + ); + assert_eq!( + ReviewPublisherVerdict::Go, + map_review_output(Some(&output("patch is correct", vec![finding(2, 2)]))).0 + ); + } + + #[tokio::test] + async fn outbox_orders_deduplicates_reclaims_and_replays_exact_payload() { + let runtime = runtime().await; + let store = runtime.review_publisher(); + let envelope = envelope(); + let started = store + .start_review_run(ReviewPublisherStartParams { + thread_id: "thread".to_string(), + envelope: envelope.clone(), + }) + .await + .expect("start"); + let duplicate = store + .start_review_run(ReviewPublisherStartParams { + thread_id: "thread".to_string(), + envelope: envelope.clone(), + }) + .await + .expect("duplicate start"); + assert_eq!(started, duplicate); + assert_eq!(started.events[0].payload.envelope, envelope); + store + .complete_review_run(ReviewPublisherCompleteParams { + review_run_id: started.run.review_run_id.clone(), + envelope_sha256: envelope.envelope_sha256.clone(), + review_output: Some(output("patch is correct", Vec::new())), + terminal_reason_override: None, + }) + .await + .expect("complete"); + + let now = Utc::now(); + let first = store + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: "owner-a".to_string(), + lease_duration: Duration::from_millis(1), + now, + }) + .await + .expect("claim") + .expect("start claim"); + assert_eq!(0, first.event.sequence); + assert!( + store + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: "owner-b".to_string(), + lease_duration: Duration::from_secs(1), + now, + }) + .await + .expect("claim while leased") + .is_none() + ); + let reclaimed = store + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: "owner-b".to_string(), + lease_duration: Duration::from_secs(1), + now: now + chrono::Duration::milliseconds(2), + }) + .await + .expect("reclaim") + .expect("expired lease claim"); + assert_eq!(2, reclaimed.event.attempt_count); + store + .acknowledge_delivery(ReviewPublisherDeliveryAckParams { + event_id: reclaimed.event.event_id.clone(), + lease_owner: "owner-b".to_string(), + receipt_id: Some("receipt".to_string()), + now, + }) + .await + .expect("ack"); + let terminal = store + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: "owner-b".to_string(), + lease_duration: Duration::from_secs(1), + now, + }) + .await + .expect("terminal claim") + .expect("terminal"); + assert_eq!(1, terminal.event.sequence); + assert!( + !serde_json::to_string(&terminal.event.payload) + .expect("event json") + .contains("ignored untrusted prose") + ); + store + .acknowledge_delivery(ReviewPublisherDeliveryAckParams { + event_id: terminal.event.event_id.clone(), + lease_owner: "owner-b".to_string(), + receipt_id: None, + now, + }) + .await + .expect("terminal ack"); + let replayed = store + .exact_replay( + terminal.event.event_id.as_str(), + terminal.event.payload_sha256.as_str(), + now, + ) + .await + .expect("replay") + .expect("exact event"); + assert_eq!(terminal.event.payload, replayed.payload); + assert!( + store + .exact_replay(terminal.event.event_id.as_str(), "wrong-digest", now,) + .await + .expect("mismatch") + .is_none() + ); + } + + #[tokio::test] + async fn dead_lettered_start_blocks_and_dead_letters_terminal_event() { + let runtime = runtime().await; + let store = runtime.review_publisher(); + let envelope = envelope(); + let started = store + .start_review_run(ReviewPublisherStartParams { + thread_id: "thread".to_string(), + envelope: envelope.clone(), + }) + .await + .expect("start"); + store + .complete_review_run(ReviewPublisherCompleteParams { + review_run_id: started.run.review_run_id.clone(), + envelope_sha256: envelope.envelope_sha256, + review_output: None, + terminal_reason_override: None, + }) + .await + .expect("complete"); + let now = Utc::now(); + let start = store + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: "owner".to_string(), + lease_duration: Duration::from_secs(30), + now, + }) + .await + .expect("claim") + .expect("start event"); + store + .fail_delivery(ReviewPublisherDeliveryFailParams { + event_id: start.event.event_id, + lease_owner: "owner".to_string(), + error_code: "http_401".to_string(), + disposition: ReviewPublisherFailureDisposition::DeadLetter, + retry_at: now, + now, + }) + .await + .expect("dead letter start"); + + let snapshot = store + .get_review_run(started.run.review_run_id.as_str()) + .await + .expect("read run") + .expect("run"); + assert_eq!(snapshot.events.len(), 2); + assert!( + snapshot + .events + .iter() + .all(|event| event.status == crate::ReviewPublisherOutboxStatus::DeadLetter) + ); + assert!( + store + .claim_next_due_event(ReviewPublisherClaimParams { + lease_owner: "owner".to_string(), + lease_duration: Duration::from_secs(30), + now, + }) + .await + .expect("claim after dead letter") + .is_none() + ); + } +} diff --git a/codex-rs/tui/src/app_server_session.rs b/codex-rs/tui/src/app_server_session.rs index 72eb6fc9aa..c790447b0b 100644 --- a/codex-rs/tui/src/app_server_session.rs +++ b/codex-rs/tui/src/app_server_session.rs @@ -2403,6 +2403,7 @@ impl AppServerSession { thread_id: thread_id.to_string(), target, delivery: Some(ReviewDelivery::Inline), + publisher_context: None, }, }) .await From 2f3c9aa3ee4eb6f9e8a8466c1367a540fa8b4ef6 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Sun, 2 Aug 2026 01:19:21 +0300 Subject: [PATCH 2/5] fix(review): close publisher gate failures Agent: unresolved-account001 --- codex-rs/app-server/README.md | 2 +- .../request_processors/review_publisher.rs | 78 +++++++++++++++---- codex-rs/prompts/src/review_request.rs | 5 +- .../state/src/runtime/review_publisher.rs | 6 +- 4 files changed, 72 insertions(+), 19 deletions(-) diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index beb025c93a..a32f03635e 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -1633,7 +1633,7 @@ For a detached review, use `"delivery": "detached"`. The response is the same sh An authenticated publisher can add `publisherContext` to a `baseBranch` review. This path is fail-closed: `baseRef` must be the same full Git ref used by the target, the worktree must be clean, `reviewedBaseSha` and `headSha` must resolve exactly, `git merge-tree --write-tree` must produce a clean result, and the head commit must carry exactly one `Agent:` trailer. The server derives the canonical origin and implementer from Git, persists an immutable `codewith-review-envelope-v1` start event before starting the turn, and returns its stable `reviewRunId`. It publishes the terminal `GO` or `NO_GO` event only from structured review output; missing output, unknown correctness, malformed priorities, or P0/P1 findings all map to `NO_GO`. -The owner-only outbox dispatcher is enabled only when `CODEWITH_REVIEW_PUBLISHER_URL` names the HTTP endpoint and `CODEWITH_REVIEW_PUBLISHER_CREDENTIAL_ENV` names the environment variable holding its bearer credential. The credential itself is never accepted in RPC payloads or persisted. Delivery is ordered start-before-terminal, leases in-flight work, treats HTTP 409 as an idempotent receipt, retries timeouts/429/5xx, and dead-letters permanent failures. Inspect a run with `review/publisher/status/read`; replay only an exact immutable payload by passing both `eventId` and `payloadSha256` to `review/publisher/replay`. +The owner-only outbox dispatcher is enabled only when `CODEWITH_REVIEW_PUBLISHER_URL` names an HTTPS endpoint (loopback HTTP is allowed for local development) and `CODEWITH_REVIEW_PUBLISHER_CREDENTIAL_ENV` names the environment variable holding its bearer credential. The credential itself is never accepted in RPC payloads or persisted. Delivery is ordered start-before-terminal, leases in-flight work, treats HTTP 409 as an idempotent receipt, retries timeouts/429/5xx, and dead-letters permanent failures. Inspect a run with `review/publisher/status/read`; replay only an exact immutable payload by passing both `eventId` and `payloadSha256` to `review/publisher/replay`. Codewith streams the usual `turn/started` notification followed by an `item/started` with an `enteredReviewMode` item so clients can show progress: diff --git a/codex-rs/app-server/src/request_processors/review_publisher.rs b/codex-rs/app-server/src/request_processors/review_publisher.rs index 43122c3c77..e113c9d296 100644 --- a/codex-rs/app-server/src/request_processors/review_publisher.rs +++ b/codex-rs/app-server/src/request_processors/review_publisher.rs @@ -323,13 +323,7 @@ fn dispatcher_config_from_env() -> Option { let endpoint = std::env::var(REVIEW_PUBLISHER_URL_ENV).ok()?; let credential_env = std::env::var(REVIEW_PUBLISHER_CREDENTIAL_ENV_ENV).ok()?; let endpoint = reqwest::Url::parse(endpoint.trim()).ok()?; - if !matches!(endpoint.scheme(), "http" | "https") - || !endpoint.username().is_empty() - || endpoint.password().is_some() - || endpoint.query().is_some() - || endpoint.fragment().is_some() - || !valid_env_name(credential_env.trim()) - { + if !valid_publisher_endpoint(&endpoint) || !valid_env_name(credential_env.trim()) { warn!("review publisher configuration is invalid; dispatcher is disabled"); return None; } @@ -339,6 +333,28 @@ fn dispatcher_config_from_env() -> Option { }) } +fn valid_publisher_endpoint(endpoint: &reqwest::Url) -> bool { + let transport_is_safe = match endpoint.scheme() { + "https" => endpoint.host_str().is_some(), + "http" => endpoint.host_str().is_some_and(|host| { + let host = host + .strip_prefix('[') + .and_then(|host| host.strip_suffix(']')) + .unwrap_or(host); + host.eq_ignore_ascii_case("localhost") + || host + .parse::() + .is_ok_and(|address| address.is_loopback()) + }), + _ => false, + }; + transport_is_safe + && endpoint.username().is_empty() + && endpoint.password().is_none() + && endpoint.query().is_none() + && endpoint.fragment().is_none() +} + fn valid_env_name(value: &str) -> bool { let mut chars = value.chars(); matches!(chars.next(), Some(ch) if ch == '_' || ch.is_ascii_alphabetic()) @@ -575,8 +591,8 @@ fn api_run(snapshot: codex_state::ReviewPublisherRunSnapshot) -> ApiRun { codex_protocol::protocol::ReviewPublisherVerdict::Go => ApiVerdict::Go, codex_protocol::protocol::ReviewPublisherVerdict::NoGo => ApiVerdict::NoGo, }), - created_at: snapshot.run.created_at.timestamp(), - completed_at: snapshot.run.completed_at.map(|value| value.timestamp()), + created_at: api_timestamp(snapshot.run.created_at), + completed_at: snapshot.run.completed_at.map(api_timestamp), events: snapshot.events.into_iter().map(api_event).collect(), } } @@ -599,15 +615,19 @@ fn api_event(event: codex_state::ReviewPublisherOutboxEvent) -> ApiOutboxEvent { }, payload_sha256: event.payload_sha256, attempt_count: event.attempt_count, - next_attempt_at: event.next_attempt_at.timestamp(), - lease_expires_at: event.lease_expires_at.map(|value| value.timestamp()), + next_attempt_at: api_timestamp(event.next_attempt_at), + lease_expires_at: event.lease_expires_at.map(api_timestamp), receipt_id: event.receipt_id, last_error_code: event.last_error_code, - created_at: event.created_at.timestamp(), - delivered_at: event.delivered_at.map(|value| value.timestamp()), + created_at: api_timestamp(event.created_at), + delivered_at: event.delivered_at.map(api_timestamp), } } +fn api_timestamp(value: chrono::DateTime) -> i64 { + value.timestamp() +} + #[cfg(test)] mod tests { use super::*; @@ -618,6 +638,38 @@ mod tests { use wiremock::ResponseTemplate; use wiremock::matchers::method; + #[test] + fn publisher_endpoint_requires_https_except_for_loopback_http() { + for endpoint in [ + "https://publisher.example.com/reviews", + "http://localhost:8080/reviews", + "http://127.0.0.1:8080/reviews", + "http://[::1]:8080/reviews", + ] { + assert!(valid_publisher_endpoint( + &reqwest::Url::parse(endpoint).expect("valid URL") + )); + } + for endpoint in [ + "http://publisher.example.com/reviews", + "http://10.0.0.8/reviews", + "https://user@publisher.example.com/reviews", + "https://publisher.example.com/reviews?key=value", + "https://publisher.example.com/reviews#fragment", + ] { + assert!(!valid_publisher_endpoint( + &reqwest::Url::parse(endpoint).expect("valid URL") + )); + } + } + + #[test] + fn publisher_api_timestamps_are_unix_seconds() { + let timestamp = chrono::DateTime::::from_timestamp_millis(1_700_000_000_123) + .expect("valid timestamp"); + assert_eq!(api_timestamp(timestamp), 1_700_000_000); + } + #[test] fn http_statuses_are_classified_fail_closed() { assert_eq!( diff --git a/codex-rs/prompts/src/review_request.rs b/codex-rs/prompts/src/review_request.rs index 2dc3b284c6..e387158147 100644 --- a/codex-rs/prompts/src/review_request.rs +++ b/codex-rs/prompts/src/review_request.rs @@ -47,7 +47,7 @@ pub fn resolve_review_request( ) -> anyhow::Result { let ReviewRequest { target, - user_facing_hint, + user_facing_hint: requested_user_facing_hint, review_envelope, } = request; let mut prompt = review_prompt(&target, cwd)?; @@ -58,7 +58,8 @@ pub fn resolve_review_request( ); prompt.push_str(canonical_envelope.as_str()); } - let user_facing_hint = user_facing_hint.unwrap_or_else(|| user_facing_hint(&target)); + let user_facing_hint = + requested_user_facing_hint.unwrap_or_else(|| user_facing_hint(&target)); Ok(ResolvedReviewRequest { target, diff --git a/codex-rs/state/src/runtime/review_publisher.rs b/codex-rs/state/src/runtime/review_publisher.rs index f21b0fc459..263e16198a 100644 --- a/codex-rs/state/src/runtime/review_publisher.rs +++ b/codex-rs/state/src/runtime/review_publisher.rs @@ -80,10 +80,10 @@ impl ReviewPublisherStore { review_run_id.as_str(), ReviewPublisherEventKind::Started, params.envelope.clone(), - None, - None, + /*verdict*/ None, + /*overall_correctness*/ None, Vec::new(), - None, + /*terminal_reason*/ None, ); let start_event_json = serde_json::to_string(&start_event)?; let start_payload_sha256 = sha256_hex(start_event_json.as_bytes()); From 670349ef7fdd4a01271f199afe09f94516542b28 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Sun, 2 Aug 2026 01:22:26 +0300 Subject: [PATCH 3/5] fix(review): apply rustfmt output Agent: unresolved-account001 --- codex-rs/prompts/src/review_request.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/codex-rs/prompts/src/review_request.rs b/codex-rs/prompts/src/review_request.rs index e387158147..16ab11654c 100644 --- a/codex-rs/prompts/src/review_request.rs +++ b/codex-rs/prompts/src/review_request.rs @@ -58,8 +58,7 @@ pub fn resolve_review_request( ); prompt.push_str(canonical_envelope.as_str()); } - let user_facing_hint = - requested_user_facing_hint.unwrap_or_else(|| user_facing_hint(&target)); + let user_facing_hint = requested_user_facing_hint.unwrap_or_else(|| user_facing_hint(&target)); Ok(ResolvedReviewRequest { target, From c39eb9c6a1361cf40e5bf89b9b4fd4a6c16d8a91 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Sun, 2 Aug 2026 01:27:06 +0300 Subject: [PATCH 4/5] fix(review): prevent credential-bearing redirects Agent: unresolved-account001 --- .../request_processors/review_publisher.rs | 45 ++++++++++++++++--- 1 file changed, 40 insertions(+), 5 deletions(-) diff --git a/codex-rs/app-server/src/request_processors/review_publisher.rs b/codex-rs/app-server/src/request_processors/review_publisher.rs index e113c9d296..c3e15bd990 100644 --- a/codex-rs/app-server/src/request_processors/review_publisher.rs +++ b/codex-rs/app-server/src/request_processors/review_publisher.rs @@ -110,11 +110,17 @@ struct ReviewPublisherDispatcherConfig { impl ReviewPublisherDispatcherRuntime { pub(crate) fn new(state_db: Option) -> Self { - let config = dispatcher_config_from_env(); - let client = reqwest::Client::builder() - .timeout(DISPATCH_HTTP_TIMEOUT) - .build() - .unwrap_or_else(|_| reqwest::Client::new()); + let mut config = dispatcher_config_from_env(); + let client = match build_dispatch_client() { + Ok(client) => client, + Err(err) => { + warn!( + "failed to build review publisher HTTP client; dispatcher is disabled: {err}" + ); + config = None; + reqwest::Client::new() + } + }; Self { state_db, config, @@ -251,6 +257,13 @@ impl ReviewPublisherDispatcherRuntime { } } +fn build_dispatch_client() -> Result { + reqwest::Client::builder() + .timeout(DISPATCH_HTTP_TIMEOUT) + .redirect(reqwest::redirect::Policy::none()) + .build() +} + #[derive(Debug, Clone, PartialEq, Eq)] enum DispatchResult { Delivered { receipt_id: Option }, @@ -718,6 +731,28 @@ mod tests { ); } + #[tokio::test] + async fn publisher_client_does_not_follow_redirects() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .respond_with( + ResponseTemplate::new(StatusCode::FOUND.as_u16()) + .append_header("location", server.uri()), + ) + .expect(1) + .mount(&server) + .await; + + let response = build_dispatch_client() + .expect("client") + .post(server.uri()) + .send() + .await + .expect("redirect response"); + + assert_eq!(response.status(), StatusCode::FOUND); + } + #[test] fn agent_provenance_comes_only_from_one_trailer() { assert_eq!( From 03c2cd3b35c12367ee25948b083732b4dc26ba02 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Sun, 2 Aug 2026 01:41:06 +0300 Subject: [PATCH 5/5] fix(review): expose reqwest to Bazel Agent: unresolved-account001 --- codex-rs/app-server/BUILD.bazel | 1 + 1 file changed, 1 insertion(+) diff --git a/codex-rs/app-server/BUILD.bazel b/codex-rs/app-server/BUILD.bazel index 6765141bdc..c0c672acab 100644 --- a/codex-rs/app-server/BUILD.bazel +++ b/codex-rs/app-server/BUILD.bazel @@ -3,6 +3,7 @@ load("//:defs.bzl", "codex_rust_crate") codex_rust_crate( name = "app-server", crate_name = "codex_app_server", + deps_extra = ["@crates//:reqwest"], integration_test_timeout = "long", test_shard_counts = { # Note app-server-all-test has a large number of integration tests, so