From ed12cc75d34f7cb5e3b08c8ac0c14e6bc7f67c4f Mon Sep 17 00:00:00 2001 From: jif Date: Sat, 19 Sep 2026 02:18:28 +0000 Subject: [PATCH] Keep Guardian reviews on the applied instruction snapshot (#46580) ## Why Shared thread instructions can change after an action is generated but before Guardian reviews it. Polling the live provider during review can give the reviewer different instructions from those applied to the action. ## What changed Initialize Guardian reviewers with the user and thread instructions captured in the review session reuse key, keeping the review tied to the applied snapshot. ## Testing Add a regression test that updates and removes shared instructions between action generation and review in a running descendant thread. Verify that the worker picks up each change, Guardian uses the applied instructions, and instruction changes create new reviewer sessions even while ancestor threads remain idle. GitOrigin-RevId: 4980762bc7b6646a1691a43be93c71a7dc611ec6 --- .../core/src/guardian/review_session_setup.rs | 8 +- .../suite/scenarios_shared_instructions.rs | 139 ++++++++++++++++++ 2 files changed, 146 insertions(+), 1 deletion(-) diff --git a/codex-rs/core/src/guardian/review_session_setup.rs b/codex-rs/core/src/guardian/review_session_setup.rs index c87b8cb622..5108bcd487 100644 --- a/codex-rs/core/src/guardian/review_session_setup.rs +++ b/codex-rs/core/src/guardian/review_session_setup.rs @@ -105,7 +105,13 @@ impl PreparedGuardianContext { auth_manager: Arc::clone(&self.parent.services.auth_manager), agent_control: self.parent.services.agent_control.clone(), originator: self.context.turn().originator.clone(), - inherited_instructions: Some(self.parent.inherited_instructions().await), + // Review the same applied instructions captured by the reuse key. + // A live provider could advance independently while reviewing this action. + inherited_instructions: Some(SessionInstructions { + user: self.key.user_instructions.clone(), + thread: self.key.thread_instructions.clone(), + ..Default::default() + }), }), initial_history: initial_history.unwrap_or(InitialHistory::New), environments: Some(self.context.environments().to_selections()), diff --git a/codex-rs/core/tests/suite/scenarios_shared_instructions.rs b/codex-rs/core/tests/suite/scenarios_shared_instructions.rs index 1e119aa551..d92e18b9b1 100644 --- a/codex-rs/core/tests/suite/scenarios_shared_instructions.rs +++ b/codex-rs/core/tests/suite/scenarios_shared_instructions.rs @@ -6,7 +6,14 @@ use super::super::agents_md::instruction_fragments; use super::super::agents_md::persisted_resume_history; use super::super::agents_md::submit_thread_turn; use super::*; +use codex_core::StartThreadOptions; +use codex_core::config::Constrained; use codex_extension_api::Instructions; +use codex_extension_api::ToolLifecycleContributor; +use codex_extension_api::ToolLifecycleFuture; +use codex_extension_api::ToolStartInput; +use codex_protocol::config_types::ApprovalsReviewer; +use codex_protocol::protocol::AskForApproval; use codex_protocol::protocol::SessionSource; use codex_protocol::protocol::SubAgentSource; use core_test_support::responses; @@ -154,3 +161,135 @@ async fn running_descendants_refresh_only_shared_thread_instructions(shared: boo } Ok(()) } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn guardian_tracks_shared_instruction_updates_in_running_descendants() -> Result<()> { + skip_if_no_network!(Ok(())); + + // Change the host's rules after each action is sampled, before its review. + // The worker must refresh on its next request; Guardian must review the + // applied snapshot, without independently polling the live provider. + struct UpdateRulesOnAction(Arc); + + impl ToolLifecycleContributor for UpdateRulesOnAction { + fn on_tool_start<'a>(&'a self, input: ToolStartInput<'a>) -> ToolLifecycleFuture<'a> { + Box::pin(async move { + let instructions = match input.call_id { + "initial-action" => Some(Instructions { + text: UPDATED.to_owned(), + source: None, + }), + "updated-action" => None, + _ => return, + }; + self.0.set_instructions(instructions); + }) + } + } + + let server = start_mock_server().await; + let provider = Arc::new(RecordingThreadInstructionsProvider::with_text(INITIAL).shared()); + let mut extensions = ExtensionRegistryBuilder::::new(); + extensions.tool_lifecycle_contributor(Arc::new(UpdateRulesOnAction(provider.clone()))); + let test = test_codex() + .with_extensions(Arc::new(extensions.build())) + .with_config(|config| { + config.permissions.approval_policy = Constrained::allow_any(AskForApproval::OnRequest); + config.approvals_reviewer = ApprovalsReviewer::AutoReview; + }) + .build_with_auto_env(&server) + .await?; + let root = test + .thread_manager + .start_thread(StartThreadOptions { + environments: Some(vec![test.executor_environment().selection().clone()]), + thread_instructions_provider: Some(provider), + ..StartThreadOptions::new(test.config.clone()) + }) + .await?; + let mut descendant = root; + for depth in 1..=2 { + descendant = test + .thread_manager + .start_thread(StartThreadOptions { + environments: Some(vec![test.executor_environment().selection().clone()]), + session_source: Some(SessionSource::SubAgent(SubAgentSource::ThreadSpawn { + parent_thread_id: descendant.thread_id, + depth, + agent_path: None, + agent_nickname: None, + agent_role: None, + })), + ..StartThreadOptions::new(test.config.clone()) + }) + .await?; + } + + let mut responses = Vec::new(); + for call_id in ["initial-action", "updated-action", "cleared-action"] { + responses.push(sse(vec![ + ev_response_created(call_id), + responses::ev_function_call( + call_id, + "exec_command", + &json!({ + "cmd": format!("echo {call_id}"), + "sandbox_permissions": "require_escalated", + "justification": "Review this action against the current rules.", + }) + .to_string(), + ), + ev_completed(call_id), + ])); + responses.push(sse(vec![ + ev_assistant_message("guardian", r#"{"outcome":"deny"}"#), + ev_completed(&format!("review-{call_id}")), + ])); + responses.push(responses::sse_completed(&format!("finished-{call_id}"))); + } + let requests = mount_sse_sequence(&server, responses).await; + // Neither ancestor takes another turn during the update and removal. + for prompt in [ + "Check the initial rules.", + "Check the updated rules.", + "Check after removal.", + ] { + submit_thread_turn(&descendant.thread, prompt).await?; + } + let requests = requests.requests(); + assert_eq!(requests.len(), 9); + let initial = format!("# AGENTS.md instructions\n\n\n{INITIAL}\n"); + let updated = format!("# AGENTS.md instructions\n\n\n{UPDATED}\n"); + let replacement = format!( + "# AGENTS.md instructions\n\n\nThese AGENTS.md instructions replace all previously provided AGENTS.md instructions.\n\n{UPDATED}\n" + ); + let cleared = "# AGENTS.md instructions\n\n\nThe previously provided AGENTS.md instructions no longer apply.\n".to_owned(); + assert_eq!( + requests + .iter() + .map(instruction_fragments) + .collect::>(), + vec![ + vec![initial.clone()], + vec![initial.clone()], + vec![initial.clone(), replacement.clone()], + vec![initial.clone(), replacement.clone()], + vec![updated], + vec![initial.clone(), replacement.clone(), cleared.clone()], + vec![initial.clone(), replacement.clone(), cleared.clone()], + vec![], + vec![initial, replacement, cleared], + ], + ); + let reviewers = [1, 4, 7].map(|index| { + let metadata = requests[index].body_json()["client_metadata"].clone(); + assert_eq!(metadata["x-openai-subagent"], "guardian"); + metadata["thread_id"] + .as_str() + .expect("reviewer thread id") + .to_owned() + }); + assert_ne!(reviewers[0], reviewers[1]); + assert_ne!(reviewers[1], reviewers[2]); + Ok(()) +}