From e36f9898060d540167d7cd1b3ca807720efacfa8 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 11 Aug 2026 18:07:11 +0300 Subject: [PATCH 1/3] Require reviewed workflow source hash on start Todos: 40083146-e1a1-4c59-a7b6-a164e27ff37b Agent: agent-chief-shipping --- .../src/protocol/common.rs | 2 + .../src/protocol/v2/workflow.rs | 1 + .../thread_workflow_processor.rs | 4 + .../src/workflow_provider_credit.rs | 3 +- .../app-server/tests/suite/v2/workflow.rs | 165 ++++ codex-rs/ext/workflows/src/activation.rs | 209 ++++- .../src/activation/finance_acceptance.rs | 3 +- codex-rs/ext/workflows/src/manager_output.rs | 1 + codex-rs/ext/workflows/src/manager_tool.rs | 279 +++++- ...0071_workflow_run_source_yaml_snapshot.sql | 5 + codex-rs/state/src/model/workflow.rs | 4 + codex-rs/state/src/runtime.rs | 253 +++++- .../state/src/runtime/workflow_automation.rs | 9 +- .../runtime/workflow_goal_plan_projections.rs | 3 +- .../src/runtime/workflow_orchestrator.rs | 6 +- .../src/runtime/workflow_provider_credit.rs | 3 +- .../state/src/runtime/workflow_verifiers.rs | 6 +- codex-rs/state/src/runtime/workflows.rs | 860 +++++++++++++++++- .../tui/src/app/thread_workflow_actions.rs | 11 +- codex-rs/tui/src/app_event.rs | 29 +- codex-rs/tui/src/app_server_session.rs | 2 + codex-rs/tui/src/chatwidget/slash_dispatch.rs | 6 +- .../src/chatwidget/tests/slash_commands.rs | 5 +- .../tui/src/chatwidget/workflow_display.rs | 5 +- .../tui/src/chatwidget/workflow_manager.rs | 2 + codex-rs/tui/src/chatwidget/workflow_slash.rs | 76 +- 26 files changed, 1853 insertions(+), 99 deletions(-) create mode 100644 codex-rs/state/migrations/0071_workflow_run_source_yaml_snapshot.sql diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index 13352a771f..b726d94156 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -4597,6 +4597,8 @@ mod tests { params: v2::ThreadWorkflowRunStartParams { thread_id: "thr_123".to_string(), workflow_record_id: "workflow_123".to_string(), + expected_source_yaml_sha256: + "0123456789abcdef0123456789abcdef0123456789abcdef0123456789abcdef".to_string(), idempotency_key: None, }, }; diff --git a/codex-rs/app-server-protocol/src/protocol/v2/workflow.rs b/codex-rs/app-server-protocol/src/protocol/v2/workflow.rs index f52e5c016e..eaf2e78312 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/workflow.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/workflow.rs @@ -311,6 +311,7 @@ pub struct ThreadWorkflowRunGetResponse { pub struct ThreadWorkflowRunStartParams { pub thread_id: String, pub workflow_record_id: String, + pub expected_source_yaml_sha256: String, #[ts(optional = nullable)] pub idempotency_key: Option, } diff --git a/codex-rs/app-server/src/request_processors/thread_workflow_processor.rs b/codex-rs/app-server/src/request_processors/thread_workflow_processor.rs index 7dcbc53c7a..049249e835 100644 --- a/codex-rs/app-server/src/request_processors/thread_workflow_processor.rs +++ b/codex-rs/app-server/src/request_processors/thread_workflow_processor.rs @@ -270,6 +270,9 @@ impl ThreadWorkflowRequestProcessor { let state_db = self.state_db_for_materialized_thread(thread_id).await?; let idempotency_key = params.idempotency_key.and_then(normalize_optional_string); let workflow_record_id = params.workflow_record_id; + let expected_source_yaml_sha256 = + normalize_optional_string(params.expected_source_yaml_sha256) + .ok_or_else(|| invalid_request("expectedSourceYamlSha256 is required"))?; let workflow = retry_transient_sqlite_busy("read thread workflow before run start", || { state_db .workflows() @@ -287,6 +290,7 @@ impl ThreadWorkflowRequestProcessor { let start_request = codex_workflows_extension::WorkflowStartRequest { workflow_record_id, source_thread_id: thread_id, + expected_source_yaml_sha256, idempotency_key: idempotency_key.clone(), activation_config: crate::extensions::workflow_activation_config(&self.config), }; diff --git a/codex-rs/app-server/src/workflow_provider_credit.rs b/codex-rs/app-server/src/workflow_provider_credit.rs index e197f0524e..45a92b7754 100644 --- a/codex-rs/app-server/src/workflow_provider_credit.rs +++ b/codex-rs/app-server/src/workflow_provider_credit.rs @@ -1345,7 +1345,8 @@ mod tests { state .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: None, idempotency_key: None, }) diff --git a/codex-rs/app-server/tests/suite/v2/workflow.rs b/codex-rs/app-server/tests/suite/v2/workflow.rs index 0bf0b9de3f..99c22d49fd 100644 --- a/codex-rs/app-server/tests/suite/v2/workflow.rs +++ b/codex-rs/app-server/tests/suite/v2/workflow.rs @@ -359,6 +359,7 @@ async fn workflow_run_lifecycle_projects_tasks_and_returns_sanitized_state() -> Some(json!({ "threadId": thread_id.as_str(), "workflowRecordId": workflow.workflow_record_id.as_str(), + "expectedSourceYamlSha256": workflow.source_yaml_sha256.as_str(), "idempotencyKey": "run-lifecycle", })), ) @@ -521,6 +522,158 @@ async fn workflow_run_lifecycle_projects_tasks_and_returns_sanitized_state() -> Ok(()) } +#[tokio::test] +async fn workflow_create_canonicalizes_terminal_line_endings_in_shared_store() -> Result<()> { + let server = create_mock_responses_server_sequence_unchecked(Vec::new()).await; + let codex_home = TempDir::new()?; + create_config_toml(codex_home.path(), &server.uri(), WorkflowsFeature::Enabled)?; + let thread_id = create_materialized_thread(codex_home.path(), "workflow source canonical")?; + + let mut mcp = TestAppServer::new_without_managed_config(codex_home.path()).await?; + initialize(&mut mcp, ExperimentalApiCapability::Enabled).await?; + + for (index, case) in [ + "no-final-line-ending", + "final-lf", + "final-crlf", + "two-final-lfs", + ] + .into_iter() + .enumerate() + { + let marker = codex_home + .path() + .join(format!("workflow-source-canonical-{index}")); + let base = valid_workflow_yaml(&marker, &format!("wf_app_server_canonical_{index}")); + let base = base.trim_end_matches(['\r', '\n']).to_string(); + let (source_yaml, expected_source_yaml) = match case { + "no-final-line-ending" => (base.clone(), base.clone()), + "final-lf" => (format!("{base}\n"), base.clone()), + "final-crlf" => (format!("{base}\r\n"), base.clone()), + "two-final-lfs" => (format!("{base}\n\n"), format!("{base}\n")), + _ => unreachable!("case list is exhaustive"), + }; + + let request_id = + send_workflow_create(&mut mcp, thread_id.as_str(), source_yaml.as_str()).await?; + let response = read_response(&mut mcp, request_id).await?; + let ThreadWorkflowCreateResponse { workflow } = to_response(response)?; + let runtime = open_state_runtime(codex_home.path()).await?; + let stored = runtime + .workflows() + .get_thread_workflow_spec( + parse_thread_id(thread_id.as_str())?, + workflow.workflow_record_id.as_str(), + ) + .await? + .ok_or_else(|| anyhow::anyhow!("{case} workflow spec was not persisted"))?; + + assert_eq!( + expected_source_yaml.as_bytes(), + stored.source_yaml.as_bytes() + ); + assert_eq!(stored.source_yaml_sha256, workflow.source_yaml_sha256); + } + + Ok(()) +} + +#[tokio::test] +async fn workflow_run_start_rejects_stale_expected_source_sha_before_state_changes() -> Result<()> { + let server = create_mock_responses_server_sequence_unchecked(Vec::new()).await; + let codex_home = TempDir::new()?; + create_config_toml(codex_home.path(), &server.uri(), WorkflowsFeature::Enabled)?; + let thread_id = create_materialized_thread(codex_home.path(), "workflow stale source")?; + let marker = codex_home.path().join("workflow-stale-source-command-ran"); + let yaml = valid_workflow_yaml(&marker, "wf_app_server_stale_source"); + + let mut mcp = TestAppServer::new_without_managed_config(codex_home.path()).await?; + initialize(&mut mcp, ExperimentalApiCapability::Enabled).await?; + + let create_id = send_workflow_create(&mut mcp, thread_id.as_str(), yaml.as_str()).await?; + let create_resp = read_response(&mut mcp, create_id).await?; + let ThreadWorkflowCreateResponse { workflow } = + to_response::(create_resp)?; + let updated_yaml = yaml.replace( + "Build a serious workflow without leaking", + "Build an updated serious workflow without leaking", + ); + let update_id = + send_workflow_create(&mut mcp, thread_id.as_str(), updated_yaml.as_str()).await?; + let update_resp = read_response(&mut mcp, update_id).await?; + let ThreadWorkflowCreateResponse { workflow: updated } = + to_response::(update_resp)?; + assert_eq!(workflow.workflow_record_id, updated.workflow_record_id); + assert_ne!(workflow.source_yaml_sha256, updated.source_yaml_sha256); + + let stale_id = mcp + .send_raw_request( + "thread/workflow/run/start", + Some(json!({ + "threadId": thread_id.as_str(), + "workflowRecordId": workflow.workflow_record_id.as_str(), + "expectedSourceYamlSha256": workflow.source_yaml_sha256.as_str(), + "idempotencyKey": "stale-source", + })), + ) + .await?; + let stale = read_error(&mut mcp, stale_id).await?; + assert_eq!(stale.error.code, -32600); + assert!( + stale.error.message.contains("source YAML SHA mismatch"), + "unexpected stale-start error: {}", + stale.error.message + ); + assert_no_workflow_runs(codex_home.path(), parse_thread_id(thread_id.as_str())?).await?; + assert_no_execution_side_effects(codex_home.path(), parse_thread_id(thread_id.as_str())?) + .await?; + + let missing_id = mcp + .send_raw_request( + "thread/workflow/run/start", + Some(json!({ + "threadId": thread_id.as_str(), + "workflowRecordId": workflow.workflow_record_id.as_str(), + "idempotencyKey": "missing-source", + })), + ) + .await?; + let missing = read_error(&mut mcp, missing_id).await?; + assert_eq!(missing.error.code, -32600); + assert!( + missing + .error + .message + .contains("Invalid request: missing field") + && missing.error.message.contains("expectedSourceYamlSha256"), + "unexpected missing-source error: {}", + missing.error.message + ); + assert_no_workflow_runs(codex_home.path(), parse_thread_id(thread_id.as_str())?).await?; + + let matching_id = mcp + .send_raw_request( + "thread/workflow/run/start", + Some(json!({ + "threadId": thread_id.as_str(), + "workflowRecordId": updated.workflow_record_id.as_str(), + "expectedSourceYamlSha256": updated.source_yaml_sha256.as_str(), + "idempotencyKey": "matching-source", + })), + ) + .await?; + let matching_resp = read_response(&mut mcp, matching_id).await?; + let matching = to_response::(matching_resp)?; + assert_eq!(ThreadWorkflowRunStatus::Running, matching.run.run.status); + assert_eq!( + updated.source_yaml_sha256, + matching.run.run.source_yaml_sha256 + ); + assert!(!marker.exists()); + + Ok(()) +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn workflow_start_activates_paired_real_workers_and_verifiers() -> Result<()> { let codex_home = TempDir::new()?; @@ -589,6 +742,7 @@ async fn workflow_start_activates_paired_real_workers_and_verifiers() -> Result< Some(json!({ "threadId": thread_id.as_str(), "workflowRecordId": workflow.workflow_record_id.as_str(), + "expectedSourceYamlSha256": workflow.source_yaml_sha256.as_str(), "idempotencyKey": "actual-worker-activation", })), ) @@ -906,6 +1060,17 @@ async fn assert_no_execution_side_effects(codex_home: &Path, thread_id: ThreadId Ok(()) } +async fn assert_no_workflow_runs(codex_home: &Path, thread_id: ThreadId) -> Result<()> { + let runtime = open_state_runtime(codex_home).await?; + let page = runtime + .workflows() + .list_thread_workflow_runs_page(thread_id, /*cursor*/ None, /*limit*/ 10) + .await?; + assert_eq!(Vec::::new(), page.data); + assert_eq!(None, page.next_cursor); + Ok(()) +} + async fn assert_initial_execution_side_effects( codex_home: &Path, thread_id: ThreadId, diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index 0d5db25db6..c17c477ec8 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -162,6 +162,7 @@ enum WorkflowStateOperation { ListThreadRuns, ReconcileFailedProjection, LoadSnapshot, + FailInvalidSourceSnapshot, ClaimRun, HeartbeatRun, AdvanceRun, @@ -174,12 +175,13 @@ enum WorkflowStateOperation { impl WorkflowStateOperation { #[cfg(test)] - const ALL: [Self; 13] = [ + const ALL: [Self; 14] = [ Self::CreateRun, Self::ProjectGoalPlan, Self::ListThreadRuns, Self::ReconcileFailedProjection, Self::LoadSnapshot, + Self::FailInvalidSourceSnapshot, Self::ClaimRun, Self::HeartbeatRun, Self::AdvanceRun, @@ -197,6 +199,7 @@ impl WorkflowStateOperation { Self::ListThreadRuns => "list active thread workflow runs", Self::ReconcileFailedProjection => "reconcile failed workflow goal plan projection", Self::LoadSnapshot => "load workflow run snapshot", + Self::FailInvalidSourceSnapshot => "fail invalid workflow run source snapshot", Self::ClaimRun => "claim workflow run", Self::HeartbeatRun => "heartbeat workflow run", Self::AdvanceRun => "advance workflow run", @@ -262,6 +265,7 @@ impl Default for WorkflowActivationConfig { pub struct WorkflowStartRequest { pub workflow_record_id: String, pub source_thread_id: ThreadId, + pub expected_source_yaml_sha256: String, pub idempotency_key: Option, pub activation_config: WorkflowActivationConfig, } @@ -335,16 +339,10 @@ impl WorkflowActivationService { &self, request: WorkflowStartRequest, ) -> anyhow::Result { - let spec_record = self - .state_db - .workflows() - .get_workflow_spec(request.workflow_record_id.as_str()) - .await? - .ok_or_else(|| anyhow::anyhow!("workflow spec record not found"))?; - let spec = parse_workflow_yaml(spec_record.source_yaml.as_str())?; let create_params = WorkflowRunCreateParams { workflow_record_id: request.workflow_record_id, source_thread_id: Some(request.source_thread_id), + expected_source_yaml_sha256: request.expected_source_yaml_sha256, idempotency_key: request.idempotency_key.clone(), }; let snapshot = if create_params.idempotency_key.is_some() { @@ -360,6 +358,7 @@ impl WorkflowActivationService { .create_workflow_run(create_params) .await? }; + let spec = parse_run_source_yaml_snapshot(&snapshot)?; let projection_params = WorkflowGoalPlanProjectionParams { workflow_run_id: snapshot.run.run_id.clone(), thread_id: request.source_thread_id, @@ -891,6 +890,48 @@ impl WorkflowActivationService { { return Ok(()); } + let spec = if snapshot.run.status.is_terminal() { + None + } else { + match parse_run_source_yaml_snapshot(&snapshot) { + Ok(spec) => Some(spec), + Err(validation_error) => { + let Some(outcome) = retry_workflow_state( + WorkflowStateOperation::FailInvalidSourceSnapshot, + || { + self.state_db + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(run_id) + }, + ) + .await? + else { + return Ok(()); + }; + if outcome.changed { + tracing::warn!( + workflow_run_id = %run_id, + "workflow activation failed closed on an invalid source snapshot: {validation_error}" + ); + } + if !outcome.snapshot.run.status.is_terminal() { + continue; + } + self.reconcile_failed_goal_plan_projection(&outcome.snapshot) + .await?; + if outcome.snapshot.run.status.is_terminal() + && outcome + .snapshot + .provider_credit + .as_ref() + .is_none_or(|credit| credit.status.is_terminal()) + { + return Ok(()); + } + continue; + } + } + }; let claim_params = WorkflowRunClaimParams { run_id: run_id.to_string(), owner_id: self.owner_instance_id.to_string(), @@ -917,13 +958,9 @@ impl WorkflowActivationService { .await?; return Ok(()); } - let spec_record = self - .state_db - .workflows() - .get_workflow_spec(snapshot.run.workflow_record_id.as_str()) - .await? - .ok_or_else(|| anyhow::anyhow!("workflow spec record not found"))?; - let spec = parse_workflow_yaml(spec_record.source_yaml.as_str())?; + let spec = spec.ok_or_else(|| { + anyhow::anyhow!("non-terminal workflow run source snapshot was not loaded") + })?; config.route_runtime.credit_control = self .provider_credit_authority .reserve_or_restore(&fence, &spec, &config.route_runtime) @@ -1555,6 +1592,24 @@ impl WorkflowActivationService { } } +fn parse_run_source_yaml_snapshot(snapshot: &WorkflowRunSnapshot) -> anyhow::Result { + let observed_source_yaml_sha256 = format!( + "{:x}", + Sha256::digest(snapshot.run.source_yaml_snapshot.as_bytes()) + ); + if observed_source_yaml_sha256 != snapshot.run.source_yaml_sha256 { + anyhow::bail!( + "workflow run {} source YAML snapshot mismatch: expected {}, observed {}", + snapshot.run.run_id, + snapshot.run.source_yaml_sha256, + observed_source_yaml_sha256 + ); + } + Ok(parse_workflow_yaml( + snapshot.run.source_yaml_snapshot.as_str(), + )?) +} + fn initial_admission_is_observable(snapshot: &WorkflowRunSnapshot) -> bool { if snapshot.run.status.is_terminal() || snapshot.run.status == WorkflowRunStatus::Paused { return true; @@ -2276,6 +2331,7 @@ mod tests { use pretty_assertions::assert_eq; use serde_json::json; use std::sync::Arc; + use std::sync::Mutex as StdMutex; use std::sync::atomic::AtomicUsize; use tokio::sync::Notify; @@ -2295,6 +2351,34 @@ mod tests { recovered: Notify, } + struct RecordingFailingProviderAuthority { + observed_models: Arc>>, + } + + #[async_trait] + impl WorkflowProviderCreditAuthority for RecordingFailingProviderAuthority { + async fn reserve_or_restore( + &self, + _fence: &WorkflowRunFenceParams, + spec: &WorkflowSpec, + _runtime: &WorkflowRouteRuntime, + ) -> anyhow::Result { + self.observed_models + .lock() + .expect("observed models should not be poisoned") + .push(spec.execution_defaults.model.clone()); + anyhow::bail!("stop after route snapshot") + } + + async fn reconcile_terminal( + &self, + _fence: &WorkflowRunFenceParams, + _status: WorkflowProviderCreditReservationStatus, + ) -> anyhow::Result<()> { + Ok(()) + } + } + #[async_trait] impl WorkflowProviderCreditAuthority for BlockingFirstReservationAuthority { async fn reserve_or_restore( @@ -2563,6 +2647,7 @@ cleanup: ..Default::default() }; let workflow_record_id = spec.workflow_record_id; + let source_yaml_sha256 = spec.source_yaml_sha256; let idempotency_key = format!( "non-repo-start-{}", if has_active_goal_plan { @@ -2575,6 +2660,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: workflow_record_id.clone(), source_thread_id, + expected_source_yaml_sha256: source_yaml_sha256.clone(), idempotency_key: Some(idempotency_key.clone()), activation_config: activation_config.clone(), }) @@ -2635,6 +2721,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id, source_thread_id, + expected_source_yaml_sha256: source_yaml_sha256, idempotency_key: Some(idempotency_key), activation_config: activation_config.clone(), }) @@ -2803,6 +2890,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: spec.workflow_record_id, source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256, idempotency_key: Some("provider-fence-order".to_string()), activation_config: WorkflowActivationConfig { route_runtime: supported_route_runtime(), @@ -2860,6 +2948,7 @@ cleanup: let request = WorkflowStartRequest { workflow_record_id: spec.workflow_record_id, source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256, idempotency_key: Some("concurrent-idempotent-start".to_string()), activation_config: WorkflowActivationConfig { route_runtime: supported_route_runtime(), @@ -2973,6 +3062,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: spec.workflow_record_id, source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256, idempotency_key: Some("initial-admission-failure-cleanup".to_string()), activation_config: WorkflowActivationConfig { route_runtime: supported_route_runtime(), @@ -3046,6 +3136,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: spec.workflow_record_id, source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256, idempotency_key: Some("post-admission-supervisor-recovery".to_string()), activation_config: activation_config.clone(), }) @@ -3079,6 +3170,90 @@ cleanup: assert_eq!(WorkflowRunStatus::Cancelled, cancelled.run.status); } + #[tokio::test] + async fn supervisor_activation_uses_run_source_snapshot_after_spec_upsert() { + let temp_dir = tempfile::tempdir().expect("create state home"); + let codex_home = temp_dir.path().join("codex-home"); + let source_cwd = temp_dir.path().join("project-workspace"); + std::fs::create_dir_all(&source_cwd).expect("create non-repository source directory"); + let state_db = StateRuntime::init(codex_home.clone(), "activation-test".to_string()) + .await + .expect("state db should initialize"); + let source_thread_id = ThreadId::new(); + let mut thread = ThreadMetadataBuilder::new( + source_thread_id, + codex_home.join(format!("rollout-{source_thread_id}.jsonl")), + Utc::now(), + SessionSource::Cli, + ); + thread.cwd = source_cwd; + state_db + .upsert_thread(&thread.build("activation-test")) + .await + .expect("source thread should be persisted"); + let spec_a = state_db + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(source_thread_id), + source_yaml: ROUTE_ACTIVATION_WORKFLOW_YAML.to_string(), + }) + .await + .expect("workflow spec A should save"); + let run = state_db + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: spec_a.workflow_record_id.clone(), + source_thread_id: Some(source_thread_id), + expected_source_yaml_sha256: spec_a.source_yaml_sha256.clone(), + idempotency_key: Some("snapshot-bound-supervisor".to_string()), + }) + .await + .expect("workflow run should snapshot spec A"); + assert_eq!(spec_a.source_yaml, run.run.source_yaml_snapshot); + + let updated_yaml = spec_a.source_yaml.replace("\"gpt-5.4\"", "\"gpt-5.3\""); + let spec_b = state_db + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(source_thread_id), + source_yaml: updated_yaml, + }) + .await + .expect("workflow spec B should save"); + assert_eq!(spec_a.workflow_record_id, spec_b.workflow_record_id); + assert_ne!(spec_a.source_yaml_sha256, spec_b.source_yaml_sha256); + + let observed_models = Arc::new(StdMutex::new(Vec::new())); + let authority = Arc::new(RecordingFailingProviderAuthority { + observed_models: Arc::clone(&observed_models), + }); + let service = WorkflowActivationService::new_with_provider_credit_authority( + Arc::clone(&state_db), + authority, + ); + let activation_config = WorkflowActivationConfig { + route_runtime: supported_route_runtime(), + max_active_background_agent_runs: Some(10), + ..Default::default() + }; + let err = service + .run_supervisor(run.run.run_id.as_str(), activation_config) + .await + .expect_err("recording authority should stop after reading the route snapshot"); + assert!( + err.to_string().contains("stop after route snapshot"), + "unexpected supervisor error: {err}" + ); + assert_eq!( + vec!["gpt-5.4".to_string()], + observed_models + .lock() + .expect("observed models should not be poisoned") + .clone(), + "supervisor must use spec A from the run snapshot, not mutable spec B" + ); + } + #[tokio::test] async fn workflow_start_stays_queued_at_capacity_and_admits_after_release() { let temp_dir = tempfile::tempdir().expect("create state home"); @@ -3118,6 +3293,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: spec.workflow_record_id.clone(), source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), idempotency_key: Some("workflow-capacity-holder".to_string()), activation_config: activation_config.clone(), }) @@ -3144,6 +3320,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: spec.workflow_record_id.clone(), source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), idempotency_key: Some("workflow-start-at-capacity".to_string()), activation_config: activation_config.clone(), }) @@ -3177,6 +3354,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: spec.workflow_record_id, source_thread_id, + expected_source_yaml_sha256: spec.source_yaml_sha256, idempotency_key: Some("workflow-start-at-capacity".to_string()), activation_config: activation_config.clone(), }) @@ -3363,6 +3541,7 @@ cleanup: "list active thread workflow runs", "reconcile failed workflow goal plan projection", "load workflow run snapshot", + "fail invalid workflow run source snapshot", "claim workflow run", "heartbeat workflow run", "advance workflow run", diff --git a/codex-rs/ext/workflows/src/activation/finance_acceptance.rs b/codex-rs/ext/workflows/src/activation/finance_acceptance.rs index 38f0bfb12c..c5cac1119c 100644 --- a/codex-rs/ext/workflows/src/activation/finance_acceptance.rs +++ b/codex-rs/ext/workflows/src/activation/finance_acceptance.rs @@ -257,7 +257,8 @@ impl SyntheticFinanceHarness { let run = runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: saved.workflow_record_id, + workflow_record_id: saved.workflow_record_id.clone(), + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some(format!("{workflow_id}-run")), }) diff --git a/codex-rs/ext/workflows/src/manager_output.rs b/codex-rs/ext/workflows/src/manager_output.rs index 59a267eaa3..4e0f83e8d5 100644 --- a/codex-rs/ext/workflows/src/manager_output.rs +++ b/codex-rs/ext/workflows/src/manager_output.rs @@ -214,6 +214,7 @@ mod tests { spec_workflow_id: "wf_context_bound".to_string(), schema_version: "1".to_string(), source_yaml_sha256: "sha".to_string(), + source_yaml_snapshot: "schema_version: 1".to_string(), status: codex_state::WorkflowRunStatus::Running, status_reason: None, reason_code: None, diff --git a/codex-rs/ext/workflows/src/manager_tool.rs b/codex-rs/ext/workflows/src/manager_tool.rs index 705d816b62..bd20b4a721 100644 --- a/codex-rs/ext/workflows/src/manager_tool.rs +++ b/codex-rs/ext/workflows/src/manager_tool.rs @@ -110,6 +110,7 @@ impl ManageWorkflowTool { struct ManageWorkflowArgs { action: ManageWorkflowAction, workflow_record_id: Option, + expected_source_yaml_sha256: Option, run_id: Option, step_id: Option, yaml: Option, @@ -186,6 +187,12 @@ impl ToolExecutor for ManageWorkflowTool { "workflow_record_id".to_string(), nullable_string("Workflow record id for read/start actions."), ), + ( + "expected_source_yaml_sha256".to_string(), + nullable_string( + "Required for start: expected source YAML SHA-256 from the reviewed workflow.", + ), + ), ( "run_id".to_string(), nullable_string( @@ -336,7 +343,7 @@ impl ManageWorkflowTool { async fn create_workflow(&self, args: ManageWorkflowArgs) -> Result { let (state_db, thread_id) = self.runtime()?; - let yaml = required_field(args.yaml, "yaml", "create")?; + let yaml = required_workflow_yaml_field(args.yaml)?; let workflow = state_db .workflows() .save_workflow_spec_yaml(codex_state::WorkflowSpecCreateParams { @@ -372,6 +379,11 @@ impl ManageWorkflowTool { self.activation_runtime()?; let workflow_record_id = required_field(args.workflow_record_id, "workflow_record_id", "start")?; + let expected_source_yaml_sha256 = required_field( + args.expected_source_yaml_sha256, + "expected_source_yaml_sha256", + "start", + )?; let idempotency_key = normalize_optional_string(args.idempotency_key); if state_db .workflows() @@ -391,6 +403,7 @@ impl ManageWorkflowTool { .start_workflow_run(WorkflowStartRequest { workflow_record_id, source_thread_id: thread_id, + expected_source_yaml_sha256, idempotency_key, activation_config: activation_config.clone(), }) @@ -627,6 +640,12 @@ fn required_field( }) } +fn required_workflow_yaml_field(value: Option) -> Result { + value + .filter(|value| !value.trim().is_empty()) + .ok_or_else(|| respond("yaml is required for manage_workflow action create")) +} + fn normalize_optional_string(value: Option) -> Option { let value = value?.trim().to_string(); (!value.is_empty()).then_some(value) @@ -681,6 +700,8 @@ mod tests { use pretty_assertions::assert_eq; use serde_json::Value; use serde_json::json; + use sha2::Digest; + use sha2::Sha256; use super::MANAGE_WORKFLOW_TOOL_NAME; use super::ManageWorkflowTool; @@ -836,6 +857,43 @@ cleanup: ) } + fn workflow_record_id(value: &Value) -> String { + value["workflow"]["workflowRecordId"] + .as_str() + .expect("workflow id") + .to_string() + } + + fn source_yaml_sha256(value: &Value) -> String { + value["workflow"]["sourceYamlSha256"] + .as_str() + .expect("source YAML SHA-256") + .to_string() + } + + fn sha256_hex(value: &str) -> String { + format!("{:x}", Sha256::digest(value.as_bytes())) + } + + fn large_manage_workflow_yaml_without_trailing_lf() -> String { + let prompt_line = + r#"source_prompt: "Exercise workflow lifecycle operations without external effects.""#; + let empty_prompt_line = r#"source_prompt: """#; + let template = MANAGE_WORKFLOW_TEST_YAML + .trim() + .replacen(prompt_line, empty_prompt_line, 1); + let target_len = 44_066; + let padding_len = target_len - template.as_bytes().len(); + let padding = "x".repeat(padding_len); + let yaml = MANAGE_WORKFLOW_TEST_YAML.trim().replacen( + prompt_line, + &format!(r#"source_prompt: "{padding}""#), + 1, + ); + assert_eq!(target_len, yaml.as_bytes().len()); + yaml + } + fn workflow_activation_config() -> WorkflowActivationConfig { WorkflowActivationConfig { route_runtime: codex_workflows::WorkflowRouteRuntime { @@ -936,10 +994,8 @@ cleanup: ) .await; assert_eq!(create["action"], "create"); - let workflow_record_id = create["workflow"]["workflowRecordId"] - .as_str() - .expect("workflow id") - .to_string(); + let workflow_record_id = workflow_record_id(&create); + let source_yaml_sha256 = source_yaml_sha256(&create); let serialized = create.to_string(); assert!(!serialized.contains("source_prompt")); @@ -948,6 +1004,7 @@ cleanup: json!({ "action": "start", "workflow_record_id": workflow_record_id, + "expected_source_yaml_sha256": source_yaml_sha256, "idempotency_key": " run-1 ", }), ) @@ -1133,16 +1190,15 @@ cleanup: }), ) .await; - let workflow_record_id = create["workflow"]["workflowRecordId"] - .as_str() - .expect("workflow id") - .to_string(); + let workflow_record_id = workflow_record_id(&create); + let source_yaml_sha256 = source_yaml_sha256(&create); let first = call_tool( &tool, json!({ "action": "start", "workflow_record_id": workflow_record_id, + "expected_source_yaml_sha256": source_yaml_sha256, "idempotency_key": "heterogeneous-run", }), ) @@ -1152,6 +1208,7 @@ cleanup: json!({ "action": "start", "workflow_record_id": create["workflow"]["workflowRecordId"], + "expected_source_yaml_sha256": create["workflow"]["sourceYamlSha256"], "idempotency_key": "heterogeneous-run", }), ) @@ -1221,6 +1278,7 @@ cleanup: json!({ "action": "start", "workflow_record_id": create["workflow"]["workflowRecordId"], + "expected_source_yaml_sha256": create["workflow"]["sourceYamlSha256"], "idempotency_key": "recoverable-start-failure", }), ) @@ -1304,6 +1362,7 @@ cleanup: json!({ "action": "start", "workflow_record_id": create["workflow"]["workflowRecordId"], + "expected_source_yaml_sha256": create["workflow"]["sourceYamlSha256"], "idempotency_key": "concrete-admission-failure", }), ) @@ -1366,15 +1425,13 @@ cleanup: }), ) .await; - let workflow_record_id = create["workflow"]["workflowRecordId"] - .as_str() - .expect("workflow id") - .to_string(); + let workflow_record_id = workflow_record_id(&create); let start = call_tool( &tool, json!({ "action": "start", "workflow_record_id": workflow_record_id, + "expected_source_yaml_sha256": create["workflow"]["sourceYamlSha256"], }), ) .await; @@ -1431,6 +1488,180 @@ cleanup: ); } + #[tokio::test] + async fn manage_workflow_create_canonicalizes_final_lf_through_shared_store() { + let tempdir = tempfile::tempdir().expect("tempdir"); + let state_db = codex_state::StateRuntime::init( + tempdir.path().to_path_buf(), + "test-provider".to_string(), + ) + .await + .expect("state runtime should initialize"); + let thread_id = codex_protocol::ThreadId::new(); + let mut thread = codex_state::ThreadMetadataBuilder::new( + thread_id, + state_db.codex_home().join("rollout.jsonl"), + chrono::Utc::now(), + codex_protocol::protocol::SessionSource::Cli, + ); + thread.cwd = tempdir.path().to_path_buf(); + state_db + .upsert_thread(&thread.build("test-provider")) + .await + .expect("thread metadata should insert"); + let tool = + ManageWorkflowTool::new(Arc::new(AtomicBool::new(true)), state_db.clone(), thread_id); + let expected_stored_yaml = large_manage_workflow_yaml_without_trailing_lf(); + let input_yaml = format!("{expected_stored_yaml}\n"); + assert_eq!(44_067, input_yaml.as_bytes().len()); + assert_eq!(Some(&0x0a), input_yaml.as_bytes().last()); + + let create = call_tool( + &tool, + json!({ + "action": "create", + "yaml": input_yaml, + }), + ) + .await; + let workflow_record_id = workflow_record_id(&create); + let stored = state_db + .workflows() + .get_thread_workflow_spec(thread_id, workflow_record_id.as_str()) + .await + .expect("workflow spec should read") + .expect("workflow spec should exist"); + + assert_eq!(44_066, stored.source_yaml.as_bytes().len()); + assert_eq!( + expected_stored_yaml.as_bytes(), + stored.source_yaml.as_bytes() + ); + assert_eq!( + sha256_hex(expected_stored_yaml.as_str()), + stored.source_yaml_sha256 + ); + assert_eq!(stored.source_yaml_sha256, source_yaml_sha256(&create)); + } + + #[tokio::test] + async fn manage_workflow_start_rejects_stale_expected_source_sha_before_side_effects() { + let tempdir = tempfile::tempdir().expect("tempdir"); + let state_db = codex_state::StateRuntime::init( + tempdir.path().to_path_buf(), + "test-provider".to_string(), + ) + .await + .expect("state runtime should initialize"); + let thread_id = codex_protocol::ThreadId::new(); + let mut thread = codex_state::ThreadMetadataBuilder::new( + thread_id, + state_db.codex_home().join("rollout.jsonl"), + chrono::Utc::now(), + codex_protocol::protocol::SessionSource::Cli, + ); + thread.cwd = tempdir.path().to_path_buf(); + state_db + .upsert_thread(&thread.build("test-provider")) + .await + .expect("thread metadata should insert"); + let tool = + ManageWorkflowTool::new(Arc::new(AtomicBool::new(true)), state_db.clone(), thread_id); + + let create = call_tool( + &tool, + json!({ + "action": "create", + "yaml": MANAGE_WORKFLOW_TEST_YAML, + }), + ) + .await; + let created_workflow_record_id = workflow_record_id(&create); + let reviewed_source_yaml_sha256 = source_yaml_sha256(&create); + let updated = call_tool( + &tool, + json!({ + "action": "create", + "yaml": MANAGE_WORKFLOW_TEST_YAML.replace( + "Exercise workflow lifecycle operations without external effects.", + "Updated workflow source after review." + ), + }), + ) + .await; + assert_eq!(created_workflow_record_id, workflow_record_id(&updated)); + assert_ne!(reviewed_source_yaml_sha256, source_yaml_sha256(&updated)); + + let stale = call_tool( + &tool, + json!({ + "action": "start", + "workflow_record_id": created_workflow_record_id.clone(), + "expected_source_yaml_sha256": reviewed_source_yaml_sha256, + "idempotency_key": "stale-source", + }), + ) + .await; + + assert_eq!(stale["action"], "start"); + assert_eq!(stale["run"], Value::Null); + assert_eq!(stale["goalPlan"], Value::Null); + assert_eq!(stale["error"]["code"], "workflow_start_failed"); + assert_eq!(stale["error"]["stage"], "activation"); + assert_eq!(stale["error"]["recovered"], false); + assert!( + stale["error"]["message"] + .as_str() + .expect("stale source error") + .contains("workflow activation failed before durable recovery was available") + ); + let runs = state_db + .workflows() + .list_thread_workflow_runs_page(thread_id, /*cursor*/ None, /*limit*/ 10) + .await + .expect("workflow runs should list"); + assert_eq!(Vec::::new(), runs.data); + let plans = state_db + .thread_goals() + .list_thread_goal_plans(thread_id) + .await + .expect("workflow goal plans should list"); + assert!(plans.is_empty()); + + let missing = call_tool_error( + &tool, + json!({ + "action": "start", + "workflow_record_id": created_workflow_record_id.clone(), + "idempotency_key": "missing-source", + }), + ) + .await; + assert_eq!( + missing, + FunctionCallError::RespondToModel( + "expected_source_yaml_sha256 is required for manage_workflow action start" + .to_string() + ) + ); + + let matching = call_tool( + &tool, + json!({ + "action": "start", + "workflow_record_id": create["workflow"]["workflowRecordId"], + "expected_source_yaml_sha256": updated["workflow"]["sourceYamlSha256"], + "idempotency_key": "matching-source", + }), + ) + .await; + assert_eq!(matching["run"]["run"]["status"], "running"); + assert_eq!( + updated["workflow"]["sourceYamlSha256"], + matching["run"]["run"]["sourceYamlSha256"] + ); + } + #[tokio::test] async fn manage_workflow_is_model_only_and_uses_single_action_schema() { let tempdir = tempfile::tempdir().expect("tempdir"); @@ -1517,4 +1748,26 @@ cleanup: }; serde_json::from_str(&text).expect("output should be json") } + + async fn call_tool_error(tool: &ManageWorkflowTool, args: Value) -> FunctionCallError { + let payload = ToolPayload::Function { + arguments: args.to_string(), + }; + match tool + .handle(codex_tools::ToolCall { + turn_id: "turn".to_string(), + call_id: "call-workflow".to_string(), + tool_name: codex_tools::ToolName::plain(MANAGE_WORKFLOW_TOOL_NAME), + model: "test-model".to_string(), + truncation_policy: TruncationPolicy::Bytes(1024 * 64), + conversation_history: ConversationHistory::default(), + turn_item_emitter: Arc::new(NoopTurnItemEmitter), + payload, + }) + .await + { + Ok(_) => panic!("workflow manager should return a model-facing error"), + Err(error) => error, + } + } } diff --git a/codex-rs/state/migrations/0071_workflow_run_source_yaml_snapshot.sql b/codex-rs/state/migrations/0071_workflow_run_source_yaml_snapshot.sql new file mode 100644 index 0000000000..460791ae93 --- /dev/null +++ b/codex-rs/state/migrations/0071_workflow_run_source_yaml_snapshot.sql @@ -0,0 +1,5 @@ +ALTER TABLE workflow_runs + ADD COLUMN source_yaml_snapshot TEXT NOT NULL DEFAULT ''; + +-- Legacy rows stay empty until runtime repair proves the current canonical +-- bytes still match the hash recorded on that exact run. diff --git a/codex-rs/state/src/model/workflow.rs b/codex-rs/state/src/model/workflow.rs index 68e8d0579c..f514f03552 100644 --- a/codex-rs/state/src/model/workflow.rs +++ b/codex-rs/state/src/model/workflow.rs @@ -303,6 +303,7 @@ pub struct WorkflowRun { pub spec_workflow_id: String, pub schema_version: String, pub source_yaml_sha256: String, + pub source_yaml_snapshot: String, pub status: WorkflowRunStatus, pub status_reason: Option, pub reason_code: Option, @@ -485,6 +486,7 @@ pub(crate) struct WorkflowRunRow { pub spec_workflow_id: String, pub schema_version: String, pub source_yaml_sha256: String, + pub source_yaml_snapshot: String, pub status: String, pub status_reason: Option, pub reason_code: Option, @@ -520,6 +522,7 @@ impl WorkflowRunRow { spec_workflow_id: row.try_get("spec_workflow_id")?, schema_version: row.try_get("schema_version")?, source_yaml_sha256: row.try_get("source_yaml_sha256")?, + source_yaml_snapshot: row.try_get("source_yaml_snapshot")?, status: row.try_get("status")?, status_reason: row.try_get("status_reason")?, reason_code: row.try_get("reason_code")?, @@ -562,6 +565,7 @@ impl TryFrom for WorkflowRun { spec_workflow_id: row.spec_workflow_id, schema_version: row.schema_version, source_yaml_sha256: row.source_yaml_sha256, + source_yaml_snapshot: row.source_yaml_snapshot, status: WorkflowRunStatus::try_from(row.status.as_str())?, status_reason: row.status_reason, reason_code: row.reason_code, diff --git a/codex-rs/state/src/runtime.rs b/codex-rs/state/src/runtime.rs index 05470cc5dc..34ca718161 100644 --- a/codex-rs/state/src/runtime.rs +++ b/codex-rs/state/src/runtime.rs @@ -2457,6 +2457,253 @@ INSERT INTO background_agent_worktree_leases ( let _ = tokio::fs::remove_dir_all(codex_home).await; } + #[tokio::test] + async fn workflow_run_source_snapshot_migration_leaves_legacy_rows_unhydrated() { + let codex_home = unique_temp_dir(); + tokio::fs::create_dir_all(&codex_home) + .await + .expect("create codex home"); + let state_path = state_db_path(codex_home.as_path()); + let pool = SqlitePool::connect_with( + SqliteConnectOptions::new() + .filename(&state_path) + .create_if_missing(true), + ) + .await + .expect("open pre-snapshot state db"); + migrator_through(&STATE_MIGRATOR, /*version*/ 70) + .run(&pool) + .await + .expect("apply state schema before workflow source snapshots"); + + sqlx::query( + r#" +INSERT INTO workflow_specs ( + workflow_record_id, + spec_workflow_id, + source_thread_id, + schema_version, + display_name, + status, + source_yaml, + source_yaml_sha256, + agent_count, + step_count, + parallel_group_count, + verifier_count, + run_command_verifier_count, + model_routed_step_count, + created_at_ms, + updated_at_ms +) VALUES + ( + 'workflow-matching', + 'matching', + NULL, + '1', + 'Matching', + 'draft', + 'spec-a', + '2e47289bc38fc584af29a86da0eaa1795f53a7fe925b4b029a6abde544cf06f5', + 0, + 0, + 0, + 0, + 0, + 0, + 0, + 0 + ), + ( + 'workflow-changed', + 'changed', + NULL, + '1', + 'Changed', + 'draft', + 'spec-b', + '9097f20087e39d3b23d445383761940e267b16f518c26c358cc5f1e9858b7107', + 0, + 0, + 0, + 0, + 0, + 0, + 0, + 0 + ) + "#, + ) + .execute(&pool) + .await + .expect("insert legacy workflow specs"); + sqlx::query( + r#" +INSERT INTO workflow_runs ( + run_id, + workflow_record_id, + source_thread_id, + idempotency_key, + spec_workflow_id, + schema_version, + source_yaml_sha256, + status, + generation, + last_event_seq, + agents_json, + execution_defaults_json, + limits_json, + approvals_json, + artifacts_json, + cleanup_json, + created_at_ms, + updated_at_ms +) VALUES + ( + 'run-matching', + 'workflow-matching', + NULL, + 'matching', + 'matching', + '1', + '2e47289bc38fc584af29a86da0eaa1795f53a7fe925b4b029a6abde544cf06f5', + 'pending', + 0, + 0, + '[]', + '{}', + '{}', + '[]', + '[]', + '[]', + 0, + 0 + ), + ( + 'run-changed', + 'workflow-changed', + NULL, + 'changed', + 'changed', + '1', + '2e47289bc38fc584af29a86da0eaa1795f53a7fe925b4b029a6abde544cf06f5', + 'pending', + 0, + 0, + '[]', + '{}', + '{}', + '[]', + '[]', + '[]', + 0, + 0 + ) + "#, + ) + .execute(&pool) + .await + .expect("insert legacy workflow runs"); + + STATE_MIGRATOR + .run(&pool) + .await + .expect("apply workflow source snapshot migration"); + let migrated: Vec<(String, String, String)> = sqlx::query_as( + r#" +SELECT run_id, source_yaml_sha256, source_yaml_snapshot +FROM workflow_runs +ORDER BY run_id + "#, + ) + .fetch_all(&pool) + .await + .expect("read migrated workflow source snapshots"); + assert_eq!( + vec![ + ( + "run-changed".to_string(), + "2e47289bc38fc584af29a86da0eaa1795f53a7fe925b4b029a6abde544cf06f5".to_string(), + String::new(), + ), + ( + "run-matching".to_string(), + "2e47289bc38fc584af29a86da0eaa1795f53a7fe925b4b029a6abde544cf06f5".to_string(), + String::new(), + ), + ], + migrated + ); + + sqlx::query( + r#" +INSERT INTO workflow_runs ( + run_id, + workflow_record_id, + source_thread_id, + idempotency_key, + spec_workflow_id, + schema_version, + source_yaml_sha256, + status, + generation, + last_event_seq, + agents_json, + execution_defaults_json, + limits_json, + approvals_json, + artifacts_json, + cleanup_json, + created_at_ms, + updated_at_ms +) VALUES ( + 'run-rollback-writer', + 'workflow-changed', + NULL, + 'rollback-writer', + 'changed', + '1', + '9097f20087e39d3b23d445383761940e267b16f518c26c358cc5f1e9858b7107', + 'pending', + 0, + 0, + '[]', + '{}', + '{}', + '[]', + '[]', + '[]', + 0, + 0 +) + "#, + ) + .execute(&pool) + .await + .expect("simulate rollback writer omitting the new snapshot column"); + let rollback_writer: (String, i64, String) = sqlx::query_as( + r#" +SELECT source_yaml_sha256, LENGTH(source_yaml_snapshot), source_yaml_snapshot +FROM workflow_runs +WHERE run_id = 'run-rollback-writer' + "#, + ) + .fetch_one(&pool) + .await + .expect("read rollback-writer workflow source snapshot"); + assert_eq!( + ( + "9097f20087e39d3b23d445383761940e267b16f518c26c358cc5f1e9858b7107".to_string(), + 0, + String::new(), + ), + rollback_writer + ); + + pool.close().await; + let _ = tokio::fs::remove_dir_all(codex_home).await; + } + #[tokio::test] async fn workflow_automation_migration_upgrades_pre_0045_state_db() { let codex_home = unique_temp_dir(); @@ -2491,7 +2738,8 @@ INSERT INTO background_agent_worktree_leases ( let run = runtime .workflows() .create_workflow_run(crate::WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: None, idempotency_key: Some("automation-migration-run".to_string()), }) @@ -2583,7 +2831,8 @@ INSERT INTO background_agent_worktree_leases ( let run = runtime .workflows() .create_workflow_run(crate::WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some("projection-migration-run".to_string()), }) diff --git a/codex-rs/state/src/runtime/workflow_automation.rs b/codex-rs/state/src/runtime/workflow_automation.rs index c4fd630da8..74542f592c 100644 --- a/codex-rs/state/src/runtime/workflow_automation.rs +++ b/codex-rs/state/src/runtime/workflow_automation.rs @@ -863,7 +863,8 @@ cleanup: runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some("automation-run".to_string()), }) @@ -1348,7 +1349,7 @@ WHERE timer_id = ? thread_id: test_thread_id(), monitor_id: monitor.monitor_id.clone(), stream: crate::ThreadMonitorEventStream::Stdout, - text: "SECRET_TOKEN=do-not-copy".to_string(), + text: "WORKFLOW_FIXTURE_MARKER=[fixture-marker]".to_string(), }) .await .expect("monitor event should create"); @@ -1435,7 +1436,7 @@ WHERE timer_id = ? thread_id: test_thread_id(), monitor_id: monitor.monitor_id.clone(), stream: crate::ThreadMonitorEventStream::Stdout, - text: format!("event-{index} SECRET_TOKEN=do-not-copy"), + text: format!("event-{index} WORKFLOW_FIXTURE_MARKER=[fixture-marker]"), }) .await .expect("monitor event should create"); @@ -1536,7 +1537,7 @@ WHERE timer_id = ? thread_id: test_thread_id(), monitor_id: monitor.monitor_id, stream: crate::ThreadMonitorEventStream::Stdout, - text: "event-after-stop SECRET_TOKEN=do-not-copy".to_string(), + text: "event-after-stop WORKFLOW_FIXTURE_MARKER=[fixture-marker]".to_string(), }) .await .expect("post-stop monitor event should create"); diff --git a/codex-rs/state/src/runtime/workflow_goal_plan_projections.rs b/codex-rs/state/src/runtime/workflow_goal_plan_projections.rs index 609735f6fa..88345b92a6 100644 --- a/codex-rs/state/src/runtime/workflow_goal_plan_projections.rs +++ b/codex-rs/state/src/runtime/workflow_goal_plan_projections.rs @@ -984,7 +984,8 @@ cleanup: runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some(format!("{workflow_id}-run")), }) diff --git a/codex-rs/state/src/runtime/workflow_orchestrator.rs b/codex-rs/state/src/runtime/workflow_orchestrator.rs index b979a979f1..6739206f46 100644 --- a/codex-rs/state/src/runtime/workflow_orchestrator.rs +++ b/codex-rs/state/src/runtime/workflow_orchestrator.rs @@ -4757,7 +4757,8 @@ agents: runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some(format!("{workflow_id}-run")), }) @@ -4786,7 +4787,8 @@ agents: let run = runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some(format!("{workflow_id}-run")), }) diff --git a/codex-rs/state/src/runtime/workflow_provider_credit.rs b/codex-rs/state/src/runtime/workflow_provider_credit.rs index 293cfcf4c9..3a1802a969 100644 --- a/codex-rs/state/src/runtime/workflow_provider_credit.rs +++ b/codex-rs/state/src/runtime/workflow_provider_credit.rs @@ -929,7 +929,8 @@ mod tests { runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: None, idempotency_key: Some("provider-credit-state-test".to_string()), }) diff --git a/codex-rs/state/src/runtime/workflow_verifiers.rs b/codex-rs/state/src/runtime/workflow_verifiers.rs index 24c14b1cfd..27c2eff51b 100644 --- a/codex-rs/state/src/runtime/workflow_verifiers.rs +++ b/codex-rs/state/src/runtime/workflow_verifiers.rs @@ -1170,7 +1170,8 @@ cleanup: let snapshot = runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some(format!("{workflow_id}-run")), }) @@ -1379,7 +1380,8 @@ artifacts:"#, let snapshot = runtime .workflows() .create_workflow_run(WorkflowRunCreateParams { - workflow_record_id: spec.workflow_record_id, + workflow_record_id: spec.workflow_record_id.clone(), + expected_source_yaml_sha256: spec.source_yaml_sha256.clone(), source_thread_id: Some(thread_id), idempotency_key: Some(format!("{workflow_id}-run")), }) diff --git a/codex-rs/state/src/runtime/workflows.rs b/codex-rs/state/src/runtime/workflows.rs index 02f4bb7b2b..94c4e2917c 100644 --- a/codex-rs/state/src/runtime/workflows.rs +++ b/codex-rs/state/src/runtime/workflows.rs @@ -25,6 +25,41 @@ pub const WORKFLOW_STEP_APPROVAL_PENDING: &str = "pending"; pub const WORKFLOW_STEP_APPROVAL_APPROVED: &str = "approved"; /// A gated step has been rejected by the user and will be skipped. pub const WORKFLOW_STEP_APPROVAL_REJECTED: &str = "rejected"; +const WORKFLOW_SOURCE_SNAPSHOT_MISSING_REASON: &str = "workflow run source snapshot is missing"; +const WORKFLOW_SOURCE_SNAPSHOT_MISSING_REASON_CODE: &str = "workflow_source_snapshot_missing"; +const WORKFLOW_SOURCE_SNAPSHOT_HASH_MISMATCH_REASON: &str = + "workflow run source snapshot hash does not match the recorded source hash"; +const WORKFLOW_SOURCE_SNAPSHOT_HASH_MISMATCH_REASON_CODE: &str = + "workflow_source_snapshot_hash_mismatch"; +const WORKFLOW_SOURCE_SNAPSHOT_INVALID_YAML_REASON: &str = + "workflow run source snapshot is not valid workflow YAML"; +const WORKFLOW_SOURCE_SNAPSHOT_INVALID_YAML_REASON_CODE: &str = + "workflow_source_snapshot_invalid_yaml"; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +enum WorkflowRunSourceSnapshotFailure { + Missing, + HashMismatch, + InvalidYaml, +} + +impl WorkflowRunSourceSnapshotFailure { + fn reason(self) -> &'static str { + match self { + Self::Missing => WORKFLOW_SOURCE_SNAPSHOT_MISSING_REASON, + Self::HashMismatch => WORKFLOW_SOURCE_SNAPSHOT_HASH_MISMATCH_REASON, + Self::InvalidYaml => WORKFLOW_SOURCE_SNAPSHOT_INVALID_YAML_REASON, + } + } + + fn reason_code(self) -> &'static str { + match self { + Self::Missing => WORKFLOW_SOURCE_SNAPSHOT_MISSING_REASON_CODE, + Self::HashMismatch => WORKFLOW_SOURCE_SNAPSHOT_HASH_MISMATCH_REASON_CODE, + Self::InvalidYaml => WORKFLOW_SOURCE_SNAPSHOT_INVALID_YAML_REASON_CODE, + } + } +} #[derive(Clone)] pub struct WorkflowStore { @@ -76,6 +111,7 @@ pub struct WorkflowRunStatusMutationOutcome { pub struct WorkflowRunCreateParams { pub workflow_record_id: String, pub source_thread_id: Option, + pub expected_source_yaml_sha256: String, pub idempotency_key: Option, } @@ -158,6 +194,7 @@ impl WorkflowStore { source_thread_id, source_yaml, } = params; + let source_yaml = canonicalize_workflow_source_yaml(source_yaml); let spec = codex_workflows::parse_workflow_yaml(&source_yaml)?; let source_yaml_sha256 = workflow_source_sha256(&source_yaml); let metadata = metadata_from_spec(&spec, source_yaml_sha256)?; @@ -383,23 +420,25 @@ LIMIT ? OFFSET ? &self, params: WorkflowRunCreateParams, ) -> anyhow::Result { - let spec_record = self - .get_workflow_spec(params.workflow_record_id.as_str()) + let mut tx = self.pool.begin().await?; + let spec_record = workflow_spec_by_id_in_tx(&mut tx, params.workflow_record_id.as_str()) .await? .ok_or_else(|| { anyhow::anyhow!("workflow spec {} does not exist", params.workflow_record_id) })?; + ensure_expected_workflow_source_sha256(¶ms, &spec_record)?; let source_thread_id = resolve_run_source_thread_id(¶ms, &spec_record)?; - let source_workspace = self.resolve_run_source_workspace(source_thread_id).await?; + let source_workspace = + resolve_run_source_workspace_in_tx(&mut tx, source_thread_id).await?; let spec = codex_workflows::parse_workflow_yaml(&spec_record.source_yaml)?; let now_ms = datetime_to_epoch_millis(Utc::now()); - let mut tx = self.pool.begin().await?; if let Some(idempotency_key) = params.idempotency_key.as_deref() && let Some(snapshot) = workflow_run_snapshot_by_idempotency_in_tx( &mut tx, spec_record.workflow_record_id.as_str(), idempotency_key, + params.expected_source_yaml_sha256.as_str(), ) .await? { @@ -436,6 +475,7 @@ LIMIT ? OFFSET ? &mut tx, spec_record.workflow_record_id.as_str(), idempotency_key, + params.expected_source_yaml_sha256.as_str(), ) .await? .ok_or_else(|| anyhow::anyhow!("idempotent workflow run was not found"))?; @@ -590,6 +630,78 @@ WHERE run_id = ? Ok(Some(WorkflowRunStatusMutationOutcome { snapshot, changed })) } + pub async fn fail_workflow_run_if_source_snapshot_invalid( + &self, + run_id: &str, + ) -> anyhow::Result> { + let now_ms = datetime_to_epoch_millis(Utc::now()); + let mut tx = self.pool.begin().await?; + hydrate_workflow_run_source_snapshot_in_tx(&mut tx, run_id).await?; + let Some(initial_snapshot) = maybe_snapshot_workflow_run_in_tx(&mut tx, run_id).await? + else { + tx.commit().await?; + return Ok(None); + }; + let failure = workflow_run_source_snapshot_failure(&initial_snapshot.run); + let mut changed = false; + + if let Some(failure) = failure + && !initial_snapshot.run.status.is_terminal() + { + let result = sqlx::query( + r#" +UPDATE workflow_runs +SET + status = ?, + status_reason = ?, + reason_code = ?, + owner_id = NULL, + owner_instance_id = NULL, + lease_expires_at_ms = NULL, + heartbeat_at_ms = NULL, + generation = generation + 1, + completed_at_ms = COALESCE(completed_at_ms, ?), + updated_at_ms = ? +WHERE run_id = ? + AND status NOT IN ('completed', 'complete', 'failed', 'cancelled') + "#, + ) + .bind(crate::WorkflowRunStatus::Failed.as_str()) + .bind(redact_state_string(failure.reason())) + .bind(failure.reason_code()) + .bind(now_ms) + .bind(now_ms) + .bind(run_id) + .execute(&mut *tx) + .await?; + changed = result.rows_affected() > 0; + if changed { + append_workflow_run_event_in_tx( + &mut tx, + run_id, + WorkflowRunEventAppend { + event_type: "run_status_changed", + actor_kind: "system", + actor_id: None, + step_run_id: None, + verifier_run_id: None, + visibility: "internal", + payload: json!({ + "status": crate::WorkflowRunStatus::Failed.as_str(), + "reasonCode": failure.reason_code(), + }), + now_ms, + }, + ) + .await?; + } + } + + let snapshot = snapshot_workflow_run_in_tx(&mut tx, run_id).await?; + tx.commit().await?; + Ok(Some(WorkflowRunStatusMutationOutcome { snapshot, changed })) + } + pub async fn pause_workflow_run( &self, params: WorkflowRunPauseParams, @@ -933,25 +1045,23 @@ fn resolve_run_source_thread_id( Ok(params.source_thread_id.or(spec_record.source_thread_id)) } -impl WorkflowStore { - async fn resolve_run_source_workspace( - &self, - source_thread_id: Option, - ) -> anyhow::Result> { - let Some(source_thread_id) = source_thread_id else { - return Ok(None); - }; - let cwd: Option = sqlx::query_scalar("SELECT cwd FROM threads WHERE id = ?") - .bind(source_thread_id.to_string()) - .fetch_optional(self.pool.as_ref()) - .await?; - let cwd = cwd.ok_or_else(|| { - anyhow::anyhow!("workflow source thread {source_thread_id} does not exist") - })?; - let cwd = PathBuf::from(cwd); - let repo_path = codex_git_utils::get_git_repo_root(cwd.as_path()); - Ok(Some(WorkflowRunSourceWorkspace { cwd, repo_path })) - } +async fn resolve_run_source_workspace_in_tx( + tx: &mut sqlx::Transaction<'_, Sqlite>, + source_thread_id: Option, +) -> anyhow::Result> { + let Some(source_thread_id) = source_thread_id else { + return Ok(None); + }; + let cwd: Option = sqlx::query_scalar("SELECT cwd FROM threads WHERE id = ?") + .bind(source_thread_id.to_string()) + .fetch_optional(&mut **tx) + .await?; + let cwd = cwd.ok_or_else(|| { + anyhow::anyhow!("workflow source thread {source_thread_id} does not exist") + })?; + let cwd = PathBuf::from(cwd); + let repo_path = codex_git_utils::get_git_repo_root(cwd.as_path()); + Ok(Some(WorkflowRunSourceWorkspace { cwd, repo_path })) } fn sanitized_workflow_pause_reason() -> &'static str { @@ -978,6 +1088,7 @@ INSERT INTO workflow_runs ( spec_workflow_id, schema_version, source_yaml_sha256, + source_yaml_snapshot, status, generation, last_event_seq, @@ -991,7 +1102,7 @@ INSERT INTO workflow_runs ( cleanup_json, created_at_ms, updated_at_ms -) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) +) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(workflow_record_id, idempotency_key) DO NOTHING RETURNING run_id "#, @@ -1017,6 +1128,7 @@ RETURNING run_id .bind(params.spec_record.spec_workflow_id.as_str()) .bind(params.spec_record.schema_version.as_str()) .bind(params.spec_record.source_yaml_sha256.as_str()) + .bind(params.spec_record.source_yaml.as_str()) .bind(crate::WorkflowRunStatus::Pending.as_str()) .bind(0_i64) .bind(0_i64) @@ -1419,10 +1531,54 @@ INSERT INTO workflow_run_monitor_links ( Ok(()) } +async fn workflow_spec_by_id_in_tx( + tx: &mut sqlx::Transaction<'_, Sqlite>, + workflow_record_id: &str, +) -> anyhow::Result> { + let sql = workflow_spec_select_by( + r#" +SELECT +"#, + "workflow_record_id = ?", + ); + let row = sqlx::query(sqlx::AssertSqlSafe(sql)) + .bind(workflow_record_id) + .fetch_optional(&mut **tx) + .await?; + + row.map(|row| workflow_spec_from_row(&row)).transpose() +} + +fn ensure_expected_workflow_source_sha256( + params: &WorkflowRunCreateParams, + spec_record: &crate::WorkflowSpecRecord, +) -> anyhow::Result<()> { + if params.expected_source_yaml_sha256 != spec_record.source_yaml_sha256 { + anyhow::bail!( + "workflow spec {} source YAML SHA mismatch: expected {}, current {}", + spec_record.workflow_record_id, + params.expected_source_yaml_sha256, + spec_record.source_yaml_sha256 + ); + } + Ok(()) +} + +fn canonicalize_workflow_source_yaml(mut source_yaml: String) -> String { + if source_yaml.ends_with('\n') { + source_yaml.pop(); + if source_yaml.ends_with('\r') { + source_yaml.pop(); + } + } + source_yaml +} + async fn workflow_run_snapshot_by_idempotency_in_tx( tx: &mut sqlx::Transaction<'_, Sqlite>, workflow_record_id: &str, idempotency_key: &str, + expected_source_yaml_sha256: &str, ) -> anyhow::Result> { let run_id = sqlx::query_scalar::<_, String>( r#" @@ -1437,7 +1593,18 @@ WHERE workflow_record_id = ? AND idempotency_key = ? .await?; if let Some(run_id) = run_id { - maybe_snapshot_workflow_run_in_tx(tx, run_id.as_str()).await + let snapshot = maybe_snapshot_workflow_run_in_tx(tx, run_id.as_str()).await?; + if let Some(snapshot) = &snapshot + && snapshot.run.source_yaml_sha256 != expected_source_yaml_sha256 + { + anyhow::bail!( + "workflow run {} idempotency replay source YAML SHA mismatch: expected {}, existing {}", + snapshot.run.run_id, + expected_source_yaml_sha256, + snapshot.run.source_yaml_sha256 + ); + } + Ok(snapshot) } else { Ok(None) } @@ -1478,6 +1645,68 @@ pub(super) async fn snapshot_workflow_run_in_tx( }) } +async fn hydrate_workflow_run_source_snapshot_in_tx( + tx: &mut sqlx::Transaction<'_, Sqlite>, + run_id: &str, +) -> anyhow::Result<()> { + let row: Option<(String, String, String, Option, Option)> = sqlx::query_as( + r#" +SELECT + workflow_runs.workflow_record_id, + workflow_runs.source_yaml_sha256, + workflow_runs.source_yaml_snapshot, + workflow_specs.source_yaml, + workflow_specs.source_yaml_sha256 +FROM workflow_runs +LEFT JOIN workflow_specs + ON workflow_specs.workflow_record_id = workflow_runs.workflow_record_id +WHERE workflow_runs.run_id = ? + "#, + ) + .bind(run_id) + .fetch_optional(&mut **tx) + .await?; + let Some((workflow_record_id, run_sha256, snapshot, source_yaml, spec_sha256)) = row else { + return Ok(()); + }; + if !snapshot.is_empty() { + return Ok(()); + } + let Some(source_yaml) = source_yaml else { + return Ok(()); + }; + let computed_sha256 = workflow_source_sha256(source_yaml.as_str()); + if computed_sha256 != run_sha256 || spec_sha256.as_deref() != Some(run_sha256.as_str()) { + return Ok(()); + } + + sqlx::query( + r#" +UPDATE workflow_runs +SET source_yaml_snapshot = ? +WHERE run_id = ? + AND source_yaml_snapshot = '' + AND source_yaml_sha256 = ? + AND EXISTS ( + SELECT 1 + FROM workflow_specs + WHERE workflow_record_id = ? + AND source_yaml = ? + AND source_yaml_sha256 = ? + ) + "#, + ) + .bind(source_yaml.as_str()) + .bind(run_id) + .bind(run_sha256.as_str()) + .bind(workflow_record_id.as_str()) + .bind(source_yaml.as_str()) + .bind(computed_sha256.as_str()) + .execute(&mut **tx) + .await?; + Ok(()) +} + async fn get_workflow_run_in_tx( tx: &mut sqlx::Transaction<'_, Sqlite>, run_id: &str, @@ -1697,6 +1926,7 @@ fn workflow_run_select_columns() -> &'static str { spec_workflow_id, schema_version, source_yaml_sha256, + source_yaml_snapshot, status, status_reason, reason_code, @@ -1797,6 +2027,21 @@ pub(super) fn workflow_source_sha256(source_yaml: &str) -> String { format!("{:x}", Sha256::digest(source_yaml.as_bytes())) } +fn workflow_run_source_snapshot_failure( + run: &crate::WorkflowRun, +) -> Option { + if run.source_yaml_snapshot.is_empty() { + return Some(WorkflowRunSourceSnapshotFailure::Missing); + } + if workflow_source_sha256(run.source_yaml_snapshot.as_str()) != run.source_yaml_sha256 { + return Some(WorkflowRunSourceSnapshotFailure::HashMismatch); + } + if codex_workflows::parse_workflow_yaml(run.source_yaml_snapshot.as_str()).is_err() { + return Some(WorkflowRunSourceSnapshotFailure::InvalidYaml); + } + None +} + fn sanitized_workflow_cancel_reason() -> &'static str { "user requested workflow cancellation" } @@ -1876,6 +2121,8 @@ mod tests { let runtime = test_runtime().await; let thread_id = test_thread_id(/*id*/ 1); upsert_test_thread(&runtime, thread_id).await; + let expected_source_yaml = + canonicalize_workflow_source_yaml(DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string()); let saved = runtime .workflows() @@ -1894,12 +2141,9 @@ mod tests { assert!(saved.verifier_count >= 6); assert!(saved.run_command_verifier_count >= 4); assert!(saved.model_routed_step_count >= 12); + assert_eq!(expected_source_yaml, saved.source_yaml); assert_eq!( - DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML, - saved.source_yaml.as_str() - ); - assert_eq!( - workflow_source_sha256(DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML), + workflow_source_sha256(expected_source_yaml.as_str()), saved.source_yaml_sha256 ); @@ -1962,6 +2206,41 @@ mod tests { ); } + #[tokio::test] + async fn save_workflow_spec_yaml_canonicalizes_one_terminal_line_ending() { + let runtime = test_runtime().await; + let thread_id = test_thread_id(/*id*/ 1); + upsert_test_thread(&runtime, thread_id).await; + let base = + canonicalize_workflow_source_yaml(DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string()); + let cases = [ + ("no-final-line-ending", base.clone(), base.clone()), + ("final-lf", format!("{base}\n"), base.clone()), + ("final-crlf", format!("{base}\r\n"), base.clone()), + ("two-final-lfs", format!("{base}\n\n"), format!("{base}\n")), + ]; + + for (case, source_yaml, expected_source_yaml) in cases { + let saved = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(thread_id), + source_yaml, + }) + .await + .unwrap_or_else(|err| panic!("{case} should save: {err}")); + + assert_eq!( + expected_source_yaml.as_bytes(), + saved.source_yaml.as_bytes() + ); + assert_eq!( + workflow_source_sha256(expected_source_yaml.as_str()), + saved.source_yaml_sha256 + ); + } + } + #[tokio::test] async fn delete_thread_workflow_spec_removes_spec_and_reports_missing() { let runtime = test_runtime().await; @@ -2061,6 +2340,7 @@ mod tests { .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), source_thread_id: Some(thread_id), + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: Some("delete-guard".to_string()), }) .await @@ -2153,6 +2433,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id, source_thread_id: Some(thread_id), + expected_source_yaml_sha256: workflow_source_sha256(stored_invalid_yaml.as_str()), idempotency_key: Some("stored-invalid-cwd".to_string()), }) .await @@ -2308,6 +2589,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), source_thread_id: Some(thread_id), + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: Some("run-key-1".to_string()), }) .await @@ -2317,6 +2599,7 @@ WHERE workflow_record_id = ? assert_eq!(saved.spec_workflow_id, snapshot.run.spec_workflow_id); assert_eq!(saved.schema_version, snapshot.run.schema_version); assert_eq!(saved.source_yaml_sha256, snapshot.run.source_yaml_sha256); + assert_eq!(saved.source_yaml, snapshot.run.source_yaml_snapshot); assert_eq!(Some(thread_id), snapshot.run.source_thread_id); assert_eq!(Some("run-key-1".to_string()), snapshot.run.idempotency_key); assert_eq!(crate::WorkflowRunStatus::Pending, snapshot.run.status); @@ -2401,6 +2684,432 @@ WHERE workflow_record_id = ? ); } + #[tokio::test] + async fn invalid_workflow_run_source_snapshots_fail_once_and_stay_terminal() { + let runtime = test_runtime().await; + let saved = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: None, + source_yaml: DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string(), + }) + .await + .expect("workflow spec should save"); + let cases = [("mismatched", "spec-b")]; + + for (case, source_yaml_snapshot) in cases { + let run = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: saved.workflow_record_id.clone(), + source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), + idempotency_key: Some(case.to_string()), + }) + .await + .unwrap_or_else(|err| panic!("{case} workflow run should create: {err}")); + let valid = runtime + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(run.run.run_id.as_str()) + .await + .unwrap_or_else(|err| panic!("{case} valid snapshot should recheck: {err}")) + .unwrap_or_else(|| panic!("{case} workflow run should exist")); + assert!( + !valid.changed, + "{case} valid snapshot must remain unchanged" + ); + assert_eq!(crate::WorkflowRunStatus::Pending, valid.snapshot.run.status); + sqlx::query( + r#" +UPDATE workflow_runs +SET source_yaml_snapshot = ? +WHERE run_id = ? + "#, + ) + .bind(source_yaml_snapshot) + .bind(run.run.run_id.as_str()) + .execute(runtime.workflows().pool.as_ref()) + .await + .unwrap_or_else(|err| panic!("{case} source snapshot fixture should update: {err}")); + + let failed = runtime + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(run.run.run_id.as_str()) + .await + .unwrap_or_else(|err| panic!("{case} invalid snapshot should fail closed: {err}")) + .unwrap_or_else(|| panic!("{case} workflow run should still exist")); + assert!(failed.changed, "{case} should transition exactly once"); + assert_eq!(crate::WorkflowRunStatus::Failed, failed.snapshot.run.status); + assert_eq!( + Some("workflow_source_snapshot_hash_mismatch"), + failed.snapshot.run.reason_code.as_deref() + ); + assert!( + failed.snapshot.run.completed_at.is_some(), + "{case} should become terminal" + ); + + let repeated = runtime + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(run.run.run_id.as_str()) + .await + .unwrap_or_else(|err| panic!("{case} terminal recheck should succeed: {err}")) + .unwrap_or_else(|| panic!("{case} workflow run should remain")); + assert!(!repeated.changed, "{case} must not emit a retry transition"); + assert_eq!( + crate::WorkflowRunStatus::Failed, + repeated.snapshot.run.status + ); + let failure_event_count: i64 = sqlx::query_scalar( + r#" +SELECT COUNT(*) +FROM workflow_run_events +WHERE run_id = ? + AND event_type = 'run_status_changed' + AND json_extract(event_payload_json, '$.data.reasonCode') = 'workflow_source_snapshot_hash_mismatch' + "#, + ) + .bind(run.run.run_id.as_str()) + .fetch_one(runtime.workflows().pool.as_ref()) + .await + .unwrap_or_else(|err| panic!("{case} failure events should count: {err}")); + assert_eq!( + 1, failure_event_count, + "{case} should emit one failure event" + ); + } + } + + #[tokio::test] + async fn matching_legacy_empty_source_snapshot_hydrates_from_verified_current_bytes() { + let runtime = test_runtime().await; + let saved = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: None, + source_yaml: DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string(), + }) + .await + .expect("workflow spec should save"); + let run = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: saved.workflow_record_id.clone(), + source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), + idempotency_key: Some("matching-legacy-empty".to_string()), + }) + .await + .expect("workflow run should create"); + sqlx::query("UPDATE workflow_runs SET source_yaml_snapshot = '' WHERE run_id = ?") + .bind(run.run.run_id.as_str()) + .execute(runtime.workflows().pool.as_ref()) + .await + .expect("legacy empty snapshot fixture should persist"); + + let hydrated = runtime + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(run.run.run_id.as_str()) + .await + .expect("legacy workflow run should repair") + .expect("legacy workflow run should exist"); + + assert!(!hydrated.changed); + assert_eq!( + saved.source_yaml, + hydrated.snapshot.run.source_yaml_snapshot + ); + assert_eq!( + saved.source_yaml_sha256, + hydrated.snapshot.run.source_yaml_sha256 + ); + assert_eq!( + crate::WorkflowRunStatus::Pending, + hydrated.snapshot.run.status + ); + assert_eq!(0, hydrated.snapshot.run.generation); + } + + #[tokio::test] + async fn changed_spec_does_not_hydrate_empty_snapshot_and_fails_terminally() { + let runtime = test_runtime().await; + let saved = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: None, + source_yaml: DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string(), + }) + .await + .expect("workflow spec A should save"); + let run = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: saved.workflow_record_id.clone(), + source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), + idempotency_key: Some("legacy-changed-spec".to_string()), + }) + .await + .expect("legacy workflow run should create"); + sqlx::query("UPDATE workflow_runs SET source_yaml_snapshot = '' WHERE run_id = ?") + .bind(run.run.run_id.as_str()) + .execute(runtime.workflows().pool.as_ref()) + .await + .expect("legacy empty snapshot fixture should persist"); + let changed_yaml = saved + .source_yaml + .replace("Dental Lead SaaS", "Changed Spec"); + let changed = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: None, + source_yaml: changed_yaml, + }) + .await + .expect("workflow spec B should save"); + assert_eq!(saved.workflow_record_id, changed.workflow_record_id); + assert_ne!(saved.source_yaml_sha256, changed.source_yaml_sha256); + + let unhydrated = runtime + .workflows() + .get_workflow_run_snapshot(run.run.run_id.as_str()) + .await + .expect("changed-spec workflow run should load") + .expect("changed-spec workflow run should exist"); + assert!(unhydrated.run.source_yaml_snapshot.is_empty()); + assert_eq!(saved.source_yaml_sha256, unhydrated.run.source_yaml_sha256); + assert_ne!(changed.source_yaml, unhydrated.run.source_yaml_snapshot); + + let failed = runtime + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(run.run.run_id.as_str()) + .await + .expect("changed-spec workflow run should fail closed") + .expect("changed-spec workflow run should exist"); + assert!(failed.changed); + assert_eq!(crate::WorkflowRunStatus::Failed, failed.snapshot.run.status); + assert_eq!( + Some("workflow_source_snapshot_missing"), + failed.snapshot.run.reason_code.as_deref() + ); + assert!(failed.snapshot.run.completed_at.is_some()); + assert_eq!(1, failed.snapshot.run.generation); + assert_eq!(None, failed.snapshot.run.owner_id); + } + + #[tokio::test] + async fn old_writer_after_migration_hydrates_omitted_snapshot_from_verified_bytes() { + let runtime = test_runtime().await; + let saved = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: None, + source_yaml: DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string(), + }) + .await + .expect("workflow spec should save"); + let template = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: saved.workflow_record_id.clone(), + source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), + idempotency_key: Some("old-writer-template".to_string()), + }) + .await + .expect("template workflow run should create"); + let old_writer_run_id = "old-writer-after-migration"; + sqlx::query( + r#" +INSERT INTO workflow_runs ( + run_id, + workflow_record_id, + source_thread_id, + source_cwd, + source_repo_path, + idempotency_key, + spec_workflow_id, + schema_version, + source_yaml_sha256, + status, + status_reason, + reason_code, + generation, + owner_id, + owner_instance_id, + lease_expires_at_ms, + heartbeat_at_ms, + last_event_seq, + agents_json, + execution_defaults_json, + limits_json, + approvals_json, + loops_json, + monitor_links_json, + artifacts_json, + cleanup_json, + created_at_ms, + updated_at_ms, + started_at_ms, + completed_at_ms +) +SELECT + ?, + workflow_record_id, + source_thread_id, + source_cwd, + source_repo_path, + ?, + spec_workflow_id, + schema_version, + source_yaml_sha256, + status, + status_reason, + reason_code, + generation, + owner_id, + owner_instance_id, + lease_expires_at_ms, + heartbeat_at_ms, + 0, + agents_json, + execution_defaults_json, + limits_json, + approvals_json, + loops_json, + monitor_links_json, + artifacts_json, + cleanup_json, + created_at_ms, + updated_at_ms, + started_at_ms, + completed_at_ms +FROM workflow_runs +WHERE run_id = ? + "#, + ) + .bind(old_writer_run_id) + .bind("old-writer-after-migration") + .bind(template.run.run_id.as_str()) + .execute(runtime.workflows().pool.as_ref()) + .await + .expect("old writer should insert while omitting the snapshot column"); + let raw_snapshot: (String, i64) = sqlx::query_as( + "SELECT source_yaml_sha256, LENGTH(source_yaml_snapshot) FROM workflow_runs WHERE run_id = ?", + ) + .bind(old_writer_run_id) + .fetch_one(runtime.workflows().pool.as_ref()) + .await + .expect("old-writer row should load before hydration"); + assert_eq!((saved.source_yaml_sha256.clone(), 0), raw_snapshot); + + let hydrated = runtime + .workflows() + .fail_workflow_run_if_source_snapshot_invalid(old_writer_run_id) + .await + .expect("old-writer workflow run should repair") + .expect("old-writer workflow run should exist"); + + assert!(!hydrated.changed); + assert_eq!( + saved.source_yaml, + hydrated.snapshot.run.source_yaml_snapshot + ); + assert_eq!( + saved.source_yaml_sha256, + hydrated.snapshot.run.source_yaml_sha256 + ); + assert_eq!( + crate::WorkflowRunStatus::Pending, + hydrated.snapshot.run.status + ); + assert_eq!(0, hydrated.snapshot.run.generation); + } + + #[tokio::test] + async fn create_workflow_run_rejects_stale_expected_source_sha_without_effects() { + let runtime = test_runtime().await; + let thread_id = test_thread_id(/*id*/ 1); + upsert_test_thread(&runtime, thread_id).await; + let marker = unique_temp_dir().join("workflow-stale-sha-command-ran"); + let marker_command = format!("touch {}", marker.display()); + let yaml = DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.replace( + "\"npm test -- --runInBand\"", + &yaml_single_quoted(&marker_command), + ); + let reviewed = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(thread_id), + source_yaml: yaml.clone(), + }) + .await + .expect("reviewed workflow spec should save"); + let update_anchor = + "build me a saas that collects leads to dentists and sells them to these"; + assert!( + yaml.contains(update_anchor), + "stale-source positive control anchor must exist in the fixture" + ); + let updated_yaml = yaml.replace( + update_anchor, + "build me an updated saas that collects leads to dentists and sells them to these", + ); + let updated = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(thread_id), + source_yaml: updated_yaml, + }) + .await + .expect("updated workflow spec should save"); + assert_eq!(reviewed.workflow_record_id, updated.workflow_record_id); + assert_ne!(reviewed.source_yaml_sha256, updated.source_yaml_sha256); + + let err = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: reviewed.workflow_record_id.clone(), + source_thread_id: Some(thread_id), + expected_source_yaml_sha256: reviewed.source_yaml_sha256, + idempotency_key: Some("stale-source-sha".to_string()), + }) + .await + .expect_err("stale reviewed source SHA must refuse before run creation"); + assert!( + err.to_string().contains("source YAML SHA mismatch"), + "unexpected stale-source error: {err}" + ); + let run_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM workflow_runs") + .fetch_one(runtime.workflows().pool.as_ref()) + .await + .expect("workflow run count should read"); + assert_eq!(0, run_count, "stale source SHA must create no run"); + let event_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM workflow_run_events") + .fetch_one(runtime.workflows().pool.as_ref()) + .await + .expect("workflow run event count should read"); + assert_eq!(0, event_count, "stale source SHA must create no events"); + + let snapshot = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: updated.workflow_record_id.clone(), + source_thread_id: Some(thread_id), + expected_source_yaml_sha256: updated.source_yaml_sha256.clone(), + idempotency_key: Some("matching-source-sha".to_string()), + }) + .await + .expect("matching source SHA should create workflow run"); + assert_eq!(updated.workflow_record_id, snapshot.run.workflow_record_id); + assert_eq!(updated.source_yaml_sha256, snapshot.run.source_yaml_sha256); + assert_eq!(updated.source_yaml, snapshot.run.source_yaml_snapshot); + assert!( + !marker.exists(), + "verifier command must not execute at start" + ); + } + #[tokio::test] async fn keyed_workflow_run_create_is_idempotent_but_unkeyed_runs_are_distinct() { let runtime = test_runtime().await; @@ -2418,6 +3127,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: Some("same-key".to_string()), }) .await @@ -2427,11 +3137,13 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: Some("same-key".to_string()), }) .await .expect("second keyed workflow run should reuse first run"); assert_eq!(first.run.run_id, second.run.run_id); + assert_eq!(saved.source_yaml, second.run.source_yaml_snapshot); assert_eq!(1, second.events.len()); let unkeyed_first = runtime @@ -2439,6 +3151,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: None, }) .await @@ -2448,6 +3161,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id, source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256, idempotency_key: None, }) .await @@ -2455,6 +3169,89 @@ WHERE workflow_record_id = ? assert_ne!(unkeyed_first.run.run_id, unkeyed_second.run.run_id); } + #[tokio::test] + async fn idempotent_replay_rejects_same_key_for_different_reviewed_source() { + let runtime = test_runtime().await; + let thread_id = test_thread_id(/*id*/ 1); + upsert_test_thread(&runtime, thread_id).await; + let saved_a = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(thread_id), + source_yaml: DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string(), + }) + .await + .expect("workflow spec A should save"); + let first = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: saved_a.workflow_record_id.clone(), + source_thread_id: Some(thread_id), + expected_source_yaml_sha256: saved_a.source_yaml_sha256.clone(), + idempotency_key: Some("same-key".to_string()), + }) + .await + .expect("first keyed workflow run should save"); + + let update_anchor = + "build me a saas that collects leads to dentists and sells them to these"; + assert!( + saved_a.source_yaml.contains(update_anchor), + "idempotency positive control anchor must exist in the fixture" + ); + let yaml_b = saved_a.source_yaml.replace( + update_anchor, + "build me a second saas that collects leads to dentists and sells them to these", + ); + let saved_b = runtime + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(thread_id), + source_yaml: yaml_b, + }) + .await + .expect("workflow spec B should save"); + assert_eq!(saved_a.workflow_record_id, saved_b.workflow_record_id); + assert_ne!(saved_a.source_yaml_sha256, saved_b.source_yaml_sha256); + + let err = runtime + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: saved_b.workflow_record_id.clone(), + source_thread_id: Some(thread_id), + expected_source_yaml_sha256: saved_b.source_yaml_sha256.clone(), + idempotency_key: Some("same-key".to_string()), + }) + .await + .expect_err("same idempotency key must not replay a different reviewed source"); + assert!( + err.to_string() + .contains("idempotency replay source YAML SHA mismatch"), + "unexpected cross-source idempotency error: {err}" + ); + + let runs = runtime + .workflows() + .list_thread_workflow_runs_page( + thread_id, + /*cursor*/ None, + DEFAULT_THREAD_WORKFLOW_RUN_LIST_LIMIT, + ) + .await + .expect("thread workflow runs should list"); + assert_eq!(1, runs.data.len(), "mismatch must not create a second run"); + assert_eq!(first.run.run_id, runs.data[0].run.run_id); + assert_eq!( + saved_a.source_yaml_sha256, + runs.data[0].run.source_yaml_sha256 + ); + let event_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM workflow_run_events") + .fetch_one(runtime.workflows().pool.as_ref()) + .await + .expect("workflow run event count should read"); + assert_eq!(1, event_count, "mismatch must not append run events"); + } + #[tokio::test] async fn workflow_run_cancel_request_appends_one_transactional_event() { let runtime = test_runtime().await; @@ -2471,6 +3268,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id, source_thread_id: None, + expected_source_yaml_sha256: saved.source_yaml_sha256, idempotency_key: Some("cancel-key".to_string()), }) .await @@ -2550,6 +3348,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), source_thread_id: Some(thread_id), + expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: Some("snapshot-key".to_string()), }) .await @@ -2601,6 +3400,7 @@ WHERE workflow_record_id = ? .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id, source_thread_id: Some(second_thread_id), + expected_source_yaml_sha256: saved.source_yaml_sha256, idempotency_key: Some("wrong-thread".to_string()), }) .await diff --git a/codex-rs/tui/src/app/thread_workflow_actions.rs b/codex-rs/tui/src/app/thread_workflow_actions.rs index 2305b2213c..9387885bad 100644 --- a/codex-rs/tui/src/app/thread_workflow_actions.rs +++ b/codex-rs/tui/src/app/thread_workflow_actions.rs @@ -68,8 +68,15 @@ impl App { .thread_workflow_run_get(thread_id, run_id) .await .map(|response| ThreadWorkflowDisplayResponse::RunShow(Box::new(response))), - ThreadWorkflowAction::RunStart { workflow_record_id } => app_server - .thread_workflow_run_start(thread_id, workflow_record_id) + ThreadWorkflowAction::RunStart { + workflow_record_id, + expected_source_yaml_sha256, + } => app_server + .thread_workflow_run_start( + thread_id, + workflow_record_id, + expected_source_yaml_sha256, + ) .await .map(|response| ThreadWorkflowDisplayResponse::RunStart(Box::new(response))), ThreadWorkflowAction::RunPause { run_id } => app_server diff --git a/codex-rs/tui/src/app_event.rs b/codex-rs/tui/src/app_event.rs index 4be6ca4897..2d3b03548d 100644 --- a/codex-rs/tui/src/app_event.rs +++ b/codex-rs/tui/src/app_event.rs @@ -1918,14 +1918,29 @@ impl AppEvent { #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum ThreadWorkflowAction { List, - Show { workflow_record_id: String }, - Delete { workflow_record_id: String }, + Show { + workflow_record_id: String, + }, + Delete { + workflow_record_id: String, + }, RunList, - RunShow { run_id: String }, - RunStart { workflow_record_id: String }, - RunPause { run_id: String }, - RunResume { run_id: String }, - RunCancel { run_id: String }, + RunShow { + run_id: String, + }, + RunStart { + workflow_record_id: String, + expected_source_yaml_sha256: String, + }, + RunPause { + run_id: String, + }, + RunResume { + run_id: String, + }, + RunCancel { + run_id: String, + }, } /// Named profile selection to apply after any required UI guardrails complete. diff --git a/codex-rs/tui/src/app_server_session.rs b/codex-rs/tui/src/app_server_session.rs index 72eb6fc9aa..a4914751c1 100644 --- a/codex-rs/tui/src/app_server_session.rs +++ b/codex-rs/tui/src/app_server_session.rs @@ -1424,6 +1424,7 @@ impl AppServerSession { &mut self, thread_id: ThreadId, workflow_record_id: String, + expected_source_yaml_sha256: String, ) -> Result { let request_id = self.next_request_id(); self.client @@ -1432,6 +1433,7 @@ impl AppServerSession { params: ThreadWorkflowRunStartParams { thread_id: thread_id.to_string(), workflow_record_id, + expected_source_yaml_sha256, idempotency_key: None, }, }) diff --git a/codex-rs/tui/src/chatwidget/slash_dispatch.rs b/codex-rs/tui/src/chatwidget/slash_dispatch.rs index bd80df3c85..b1dac23320 100644 --- a/codex-rs/tui/src/chatwidget/slash_dispatch.rs +++ b/codex-rs/tui/src/chatwidget/slash_dispatch.rs @@ -2935,8 +2935,12 @@ fn workflow_slash_command_to_action(command: WorkflowSlashCommand<'_>) -> Thread WorkflowSlashCommand::RunShow { run_id } => ThreadWorkflowAction::RunShow { run_id: run_id.to_string(), }, - WorkflowSlashCommand::RunStart { workflow_record_id } => ThreadWorkflowAction::RunStart { + WorkflowSlashCommand::RunStart { + workflow_record_id, + expected_source_yaml_sha256, + } => ThreadWorkflowAction::RunStart { workflow_record_id: workflow_record_id.to_string(), + expected_source_yaml_sha256: expected_source_yaml_sha256.to_string(), }, WorkflowSlashCommand::RunPause { run_id } => ThreadWorkflowAction::RunPause { run_id: run_id.to_string(), diff --git a/codex-rs/tui/src/chatwidget/tests/slash_commands.rs b/codex-rs/tui/src/chatwidget/tests/slash_commands.rs index a80e8d0906..a53c6dc16a 100644 --- a/codex-rs/tui/src/chatwidget/tests/slash_commands.rs +++ b/codex-rs/tui/src/chatwidget/tests/slash_commands.rs @@ -2624,9 +2624,10 @@ async fn workflow_run_slash_commands_emit_management_events() { }, ), ( - "/workflow run start workflow-1", + "/workflow run start workflow-1 sha-1", crate::app_event::ThreadWorkflowAction::RunStart { workflow_record_id: "workflow-1".to_string(), + expected_source_yaml_sha256: "sha-1".to_string(), }, ), ( @@ -2787,7 +2788,7 @@ async fn workflow_commands_are_inert_when_feature_is_disabled() { "/workflow delete workflow-1", "/workflow draft build a SaaS that collects leads for dentists", "/workflow run", - "/workflow run start workflow-1", + "/workflow run start workflow-1 sha-1", ] { let (mut chat, mut rx, mut op_rx) = make_chatwidget_manual(/*model_override*/ None).await; chat.thread_id = Some(ThreadId::new()); diff --git a/codex-rs/tui/src/chatwidget/workflow_display.rs b/codex-rs/tui/src/chatwidget/workflow_display.rs index 02cfb286ce..fc057bb049 100644 --- a/codex-rs/tui/src/chatwidget/workflow_display.rs +++ b/codex-rs/tui/src/chatwidget/workflow_display.rs @@ -69,7 +69,10 @@ impl ChatWidget { if response.data.is_empty() { self.add_info_message( "No workflow runs for this thread.".to_string(), - Some("/workflow run start ".to_string()), + Some( + "/workflow run start " + .to_string(), + ), ); return; } diff --git a/codex-rs/tui/src/chatwidget/workflow_manager.rs b/codex-rs/tui/src/chatwidget/workflow_manager.rs index 154da5d717..35467979ec 100644 --- a/codex-rs/tui/src/chatwidget/workflow_manager.rs +++ b/codex-rs/tui/src/chatwidget/workflow_manager.rs @@ -342,6 +342,7 @@ fn thread_workflow_actions_params( ) -> SelectionViewParams { let inspect_workflow_id = workflow.workflow_record_id.clone(); let run_workflow_id = workflow.workflow_record_id.clone(); + let run_source_yaml_sha256 = workflow.source_yaml_sha256.clone(); let mut items = vec![ workflow_action_item( "Inspect spec", @@ -364,6 +365,7 @@ fn thread_workflow_actions_params( thread_id, action: ThreadWorkflowAction::RunStart { workflow_record_id: run_workflow_id.clone(), + expected_source_yaml_sha256: run_source_yaml_sha256.clone(), }, }, ), diff --git a/codex-rs/tui/src/chatwidget/workflow_slash.rs b/codex-rs/tui/src/chatwidget/workflow_slash.rs index 3545d45bd9..57870fac7f 100644 --- a/codex-rs/tui/src/chatwidget/workflow_slash.rs +++ b/codex-rs/tui/src/chatwidget/workflow_slash.rs @@ -1,25 +1,44 @@ pub(super) const WORKFLOW_USAGE: &str = concat!( "Usage: /workflow [list|show |delete |draft |", - "run [list|show |start |pause |resume |cancel ]]" + "run [list|show |start |", + "pause |resume |cancel ]]" ); pub(super) const WORKFLOW_USAGE_HINT: &str = concat!( "Examples: /workflow list, /workflow show , ", "/workflow delete , ", - "/workflow run, /workflow run start , /workflow run pause " + "/workflow run, /workflow run start , ", + "/workflow run pause " ); #[derive(Debug, PartialEq, Eq)] pub(super) enum WorkflowSlashCommand<'a> { List, - Show { workflow_record_id: &'a str }, - Delete { workflow_record_id: &'a str }, - Draft { request: &'a str }, + Show { + workflow_record_id: &'a str, + }, + Delete { + workflow_record_id: &'a str, + }, + Draft { + request: &'a str, + }, RunList, - RunShow { run_id: &'a str }, - RunStart { workflow_record_id: &'a str }, - RunPause { run_id: &'a str }, - RunResume { run_id: &'a str }, - RunCancel { run_id: &'a str }, + RunShow { + run_id: &'a str, + }, + RunStart { + workflow_record_id: &'a str, + expected_source_yaml_sha256: &'a str, + }, + RunPause { + run_id: &'a str, + }, + RunResume { + run_id: &'a str, + }, + RunCancel { + run_id: &'a str, + }, } pub(super) fn parse_workflow_slash_args(trimmed: &str) -> Result, String> { @@ -60,8 +79,14 @@ fn parse_workflow_run_slash_args(trimmed: &str) -> Result { required_id(value, "run_id").map(|run_id| WorkflowSlashCommand::RunShow { run_id }) } - "start" => required_id(value, "workflow_record_id") - .map(|workflow_record_id| WorkflowSlashCommand::RunStart { workflow_record_id }), + "start" => { + required_start_ids(value).map(|(workflow_record_id, expected_source_yaml_sha256)| { + WorkflowSlashCommand::RunStart { + workflow_record_id, + expected_source_yaml_sha256, + } + }) + } "pause" => { required_id(value, "run_id").map(|run_id| WorkflowSlashCommand::RunPause { run_id }) } @@ -86,6 +111,20 @@ fn required_id<'a>(value: &'a str, name: &str) -> Result<&'a str, String> { Ok(id) } +fn required_start_ids(value: &str) -> Result<(&str, &str), String> { + let mut parts = value.split_whitespace(); + let Some(workflow_record_id) = parts.next() else { + return Err(WORKFLOW_USAGE.to_string()); + }; + let Some(expected_source_yaml_sha256) = parts.next() else { + return Err(WORKFLOW_USAGE.to_string()); + }; + if parts.next().is_some() { + return Err("Expected workflow_record_id and expected_source_yaml_sha256.".to_string()); + } + Ok((workflow_record_id, expected_source_yaml_sha256)) +} + pub(super) fn workflow_generation_prompt(request: &str) -> String { format!( r#"Create a Codewith workflow YAML document for this request: @@ -146,9 +185,10 @@ mod tests { Ok(WorkflowSlashCommand::RunShow { run_id: "run-1" }) ); assert_eq!( - parse_workflow_slash_args("run start workflow-1"), + parse_workflow_slash_args("run start workflow-1 sha-1"), Ok(WorkflowSlashCommand::RunStart { - workflow_record_id: "workflow-1" + workflow_record_id: "workflow-1", + expected_source_yaml_sha256: "sha-1" }) ); assert_eq!( @@ -183,6 +223,14 @@ mod tests { ); } + #[test] + fn rejects_run_start_without_expected_source_sha() { + assert_eq!( + parse_workflow_slash_args("run start workflow-1"), + Err(WORKFLOW_USAGE.to_string()) + ); + } + #[test] fn parses_explicit_draft_command() { assert_eq!( From 476c5d02b852c520c09842b7cfb7dacf9a61bc43 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Wed, 12 Aug 2026 11:05:31 +0300 Subject: [PATCH 2/3] Fix workflow snapshot remediation checks Todos: 40083146-e1a1-4c59-a7b6-a164e27ff37b Agent: agent-chief-shipping --- codex-rs/state/src/runtime/workflows.rs | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/codex-rs/state/src/runtime/workflows.rs b/codex-rs/state/src/runtime/workflows.rs index 94c4e2917c..d16e707c98 100644 --- a/codex-rs/state/src/runtime/workflows.rs +++ b/codex-rs/state/src/runtime/workflows.rs @@ -43,6 +43,9 @@ enum WorkflowRunSourceSnapshotFailure { InvalidYaml, } +type WorkflowRunSourceSnapshotHydrationRow = + (String, String, String, Option, Option); + impl WorkflowRunSourceSnapshotFailure { fn reason(self) -> &'static str { match self { @@ -1649,7 +1652,7 @@ async fn hydrate_workflow_run_source_snapshot_in_tx( tx: &mut sqlx::Transaction<'_, Sqlite>, run_id: &str, ) -> anyhow::Result<()> { - let row: Option<(String, String, String, Option, Option)> = sqlx::query_as( + let row: Option = sqlx::query_as( r#" SELECT workflow_runs.workflow_record_id, @@ -2833,10 +2836,12 @@ WHERE run_id = ? #[tokio::test] async fn changed_spec_does_not_hydrate_empty_snapshot_and_fails_terminally() { let runtime = test_runtime().await; + let thread_id = test_thread_id(/*id*/ 13); + upsert_test_thread(&runtime, thread_id).await; let saved = runtime .workflows() .save_workflow_spec_yaml(WorkflowSpecCreateParams { - source_thread_id: None, + source_thread_id: Some(thread_id), source_yaml: DENTAL_LEAD_SAAS_WORKFLOW_EXAMPLE_YAML.to_string(), }) .await @@ -2845,7 +2850,7 @@ WHERE run_id = ? .workflows() .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: saved.workflow_record_id.clone(), - source_thread_id: None, + source_thread_id: Some(thread_id), expected_source_yaml_sha256: saved.source_yaml_sha256.clone(), idempotency_key: Some("legacy-changed-spec".to_string()), }) @@ -2862,7 +2867,7 @@ WHERE run_id = ? let changed = runtime .workflows() .save_workflow_spec_yaml(WorkflowSpecCreateParams { - source_thread_id: None, + source_thread_id: Some(thread_id), source_yaml: changed_yaml, }) .await From 6962716c5a415ecbf2f038195c5bd98affaf2bcd Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Wed, 12 Aug 2026 11:53:07 +0300 Subject: [PATCH 3/3] Fix workflow YAML length clippy findings Todos: 40083146-e1a1-4c59-a7b6-a164e27ff37b Agent: agent-chief-shipping --- codex-rs/ext/workflows/src/manager_tool.rs | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/codex-rs/ext/workflows/src/manager_tool.rs b/codex-rs/ext/workflows/src/manager_tool.rs index bd20b4a721..61dc95f7f6 100644 --- a/codex-rs/ext/workflows/src/manager_tool.rs +++ b/codex-rs/ext/workflows/src/manager_tool.rs @@ -883,14 +883,14 @@ cleanup: .trim() .replacen(prompt_line, empty_prompt_line, 1); let target_len = 44_066; - let padding_len = target_len - template.as_bytes().len(); + let padding_len = target_len - template.len(); let padding = "x".repeat(padding_len); let yaml = MANAGE_WORKFLOW_TEST_YAML.trim().replacen( prompt_line, &format!(r#"source_prompt: "{padding}""#), 1, ); - assert_eq!(target_len, yaml.as_bytes().len()); + assert_eq!(target_len, yaml.len()); yaml } @@ -1513,7 +1513,7 @@ cleanup: ManageWorkflowTool::new(Arc::new(AtomicBool::new(true)), state_db.clone(), thread_id); let expected_stored_yaml = large_manage_workflow_yaml_without_trailing_lf(); let input_yaml = format!("{expected_stored_yaml}\n"); - assert_eq!(44_067, input_yaml.as_bytes().len()); + assert_eq!(44_067, input_yaml.len()); assert_eq!(Some(&0x0a), input_yaml.as_bytes().last()); let create = call_tool( @@ -1532,7 +1532,7 @@ cleanup: .expect("workflow spec should read") .expect("workflow spec should exist"); - assert_eq!(44_066, stored.source_yaml.as_bytes().len()); + assert_eq!(44_066, stored.source_yaml.len()); assert_eq!( expected_stored_yaml.as_bytes(), stored.source_yaml.as_bytes()