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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 13 additions & 5 deletions codex-rs/core/src/guardian/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,11 +33,13 @@ pub(crate) use review::is_guardian_reviewer_source;
pub(crate) use review::new_guardian_review_id;
#[cfg(test)]
pub(crate) use review::record_guardian_denial_for_test;
pub(crate) use review::reset_auto_review_rejection_circuit_breaker;
pub(crate) use review::review_approval_request;
#[cfg(test)]
pub(crate) use review::review_approval_request_with_cancel;
pub(crate) use review::routes_approval_to_guardian;
pub(crate) use review::spawn_approval_request_review;
pub(crate) use review::take_pending_auto_review_escalation;
pub(crate) use review_session::GuardianReviewSessionManager;

const GUARDIAN_PREFERRED_MODEL: &str = "codex-auto-review";
Expand Down Expand Up @@ -79,13 +81,13 @@ pub(crate) struct GuardianRejectionCircuitBreaker {
struct GuardianRejectionCircuitBreakerTurn {
consecutive_denials: u32,
total_denials: u32,
interrupt_triggered: bool,
pending_auto_review_escalation: bool,
}

#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum GuardianRejectionCircuitBreakerAction {
Continue,
InterruptTurn {
EscalateToUser {
consecutive_denials: u32,
total_denials: u32,
},
Expand All @@ -100,12 +102,12 @@ impl GuardianRejectionCircuitBreaker {
let turn = self.turns.entry(turn_id.to_string()).or_default();
turn.consecutive_denials = turn.consecutive_denials.saturating_add(1);
turn.total_denials = turn.total_denials.saturating_add(1);
if !turn.interrupt_triggered
if !turn.pending_auto_review_escalation
&& (turn.consecutive_denials >= MAX_CONSECUTIVE_GUARDIAN_DENIALS_PER_TURN
|| turn.total_denials >= MAX_TOTAL_GUARDIAN_DENIALS_PER_TURN)
{
turn.interrupt_triggered = true;
GuardianRejectionCircuitBreakerAction::InterruptTurn {
turn.pending_auto_review_escalation = true;
GuardianRejectionCircuitBreakerAction::EscalateToUser {
Comment on lines +109 to +110

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

P2 Badge Bind escalation to the denied review

The circuit breaker records only a per-turn boolean, so any later denied review on the same turn can consume the pending escalation. With concurrent approvals (or delegated/helper reviews that do not consume it), the request that tripped the breaker can still be returned as an Auto-review denial while an unrelated request is escalated to the user.

Useful? React with 👍 / 👎.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

@won-openai can you respond to this one?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is concerned that say for example there are 2 parallel auto-review going on, whichever one gets denied first as the 3rd/10th one will be sent for a user approval and the other will be denied like normal (which is okay)

consecutive_denials: turn.consecutive_denials,
total_denials: turn.total_denials,
}
Expand All @@ -118,6 +120,12 @@ impl GuardianRejectionCircuitBreaker {
let turn = self.turns.entry(turn_id.to_string()).or_default();
turn.consecutive_denials = 0;
}

pub(crate) fn take_pending_auto_review_escalation(&mut self, turn_id: &str) -> bool {
self.turns
.get_mut(turn_id)
.is_some_and(|turn| std::mem::take(&mut turn.pending_auto_review_escalation))
}
}

#[cfg(test)]
Expand Down
36 changes: 25 additions & 11 deletions codex-rs/core/src/guardian/review.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@ use codex_protocol::protocol::GuardianRiskLevel;
use codex_protocol::protocol::GuardianUserAuthorization;
use codex_protocol::protocol::ReviewDecision;
use codex_protocol::protocol::SubAgentSource;
use codex_protocol::protocol::TurnAbortReason;
use codex_protocol::protocol::WarningEvent;
use tokio::sync::oneshot;
use tokio_util::sync::CancellationToken;
Expand Down Expand Up @@ -185,7 +184,7 @@ async fn record_guardian_denial(session: &Arc<Session>, turn: &Arc<TurnContext>,
.lock()
.await
.record_denial(turn_id);
let GuardianRejectionCircuitBreakerAction::InterruptTurn {
let GuardianRejectionCircuitBreakerAction::EscalateToUser {
consecutive_denials,
total_denials,
} = action
Expand All @@ -202,20 +201,35 @@ async fn record_guardian_denial(session: &Arc<Session>, turn: &Arc<TurnContext>,
turn.as_ref(),
EventMsg::GuardianWarning(WarningEvent {
message: format!(
"Automatic approval review rejected too many approval requests for this turn ({consecutive_denials} consecutive, {total_denials} total); interrupting the turn."
"Automatic approval review rejected too many approval requests for this turn ({consecutive_denials} consecutive, {total_denials} total); requesting manual approval for the current request."
),
}),
)
.await;
}

let runtime_handle = session.services.runtime_handle.clone();
let session = Arc::clone(session);
let turn_id = turn_id.to_string();
let _abort_task = runtime_handle.spawn(async move {
session
.abort_turn_if_active(&turn_id, TurnAbortReason::Interrupted)
.await;
});
pub(crate) async fn reset_auto_review_rejection_circuit_breaker(
session: &Arc<Session>,
turn_id: &str,
) {
session
.services
.guardian_rejection_circuit_breaker
.lock()
.await
.clear_turn(turn_id);
}

pub(crate) async fn take_pending_auto_review_escalation(
session: &Arc<Session>,
turn_id: &str,
) -> bool {
session
.services
.guardian_rejection_circuit_breaker
.lock()
.await
.take_pending_auto_review_escalation(turn_id)
}

#[cfg(test)]
Expand Down
30 changes: 25 additions & 5 deletions codex-rs/core/src/guardian/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -72,7 +72,7 @@ fn fixed_guardian_parent_session_id() -> ThreadId {
}

#[test]
fn guardian_rejection_circuit_breaker_interrupts_after_three_consecutive_denials() {
fn guardian_rejection_circuit_breaker_escalates_after_three_consecutive_denials() {
let mut circuit_breaker = GuardianRejectionCircuitBreaker::default();
assert_eq!(
circuit_breaker.record_denial("turn-1"),
Expand All @@ -84,7 +84,7 @@ fn guardian_rejection_circuit_breaker_interrupts_after_three_consecutive_denials
);
assert_eq!(
circuit_breaker.record_denial("turn-1"),
GuardianRejectionCircuitBreakerAction::InterruptTurn {
GuardianRejectionCircuitBreakerAction::EscalateToUser {
consecutive_denials: 3,
total_denials: 3,
}
Expand Down Expand Up @@ -113,15 +113,15 @@ fn guardian_rejection_circuit_breaker_resets_consecutive_denials_on_non_denial()
);
assert_eq!(
circuit_breaker.record_denial("turn-1"),
GuardianRejectionCircuitBreakerAction::InterruptTurn {
GuardianRejectionCircuitBreakerAction::EscalateToUser {
consecutive_denials: 3,
total_denials: 4,
}
);
}

#[test]
fn guardian_rejection_circuit_breaker_interrupts_after_ten_total_denials() {
fn guardian_rejection_circuit_breaker_escalates_after_ten_total_denials() {
let mut circuit_breaker = GuardianRejectionCircuitBreaker::default();
for _ in 0..9 {
assert_eq!(
Expand All @@ -132,13 +132,33 @@ fn guardian_rejection_circuit_breaker_interrupts_after_ten_total_denials() {
}
assert_eq!(
circuit_breaker.record_denial("turn-1"),
GuardianRejectionCircuitBreakerAction::InterruptTurn {
GuardianRejectionCircuitBreakerAction::EscalateToUser {
consecutive_denials: 1,
total_denials: 10,
}
);
}

#[test]
fn guardian_rejection_circuit_breaker_tracks_pending_escalation_once() {
let mut circuit_breaker = GuardianRejectionCircuitBreaker::default();
for _ in 0..2 {
assert_eq!(
circuit_breaker.record_denial("turn-1"),
GuardianRejectionCircuitBreakerAction::Continue
);
}
assert_eq!(
circuit_breaker.record_denial("turn-1"),
GuardianRejectionCircuitBreakerAction::EscalateToUser {
consecutive_denials: 3,
total_denials: 3,
}
);
assert!(circuit_breaker.take_pending_auto_review_escalation("turn-1"));
assert!(!circuit_breaker.take_pending_auto_review_escalation("turn-1"));
}

async fn guardian_test_session_and_turn(
server: &wiremock::MockServer,
) -> (Arc<Session>, Arc<TurnContext>) {
Expand Down
64 changes: 33 additions & 31 deletions codex-rs/core/src/mcp_tool_call.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,8 +23,10 @@ use crate::guardian::guardian_approval_request_to_json;
use crate::guardian::guardian_rejection_message;
use crate::guardian::guardian_timeout_message;
use crate::guardian::new_guardian_review_id;
use crate::guardian::reset_auto_review_rejection_circuit_breaker;
use crate::guardian::review_approval_request;
use crate::guardian::routes_approval_to_guardian;
use crate::guardian::take_pending_auto_review_escalation;
use crate::hook_runtime::run_permission_request_hooks;
use crate::mcp_openai_file::rewrite_mcp_tool_arguments_for_openai_files;
use crate::mcp_tool_approval_templates::RenderedMcpToolApprovalParam;
Expand Down Expand Up @@ -1229,6 +1231,7 @@ async fn maybe_request_mcp_tool_approval(
.features
.enabled(Feature::ToolCallMcpElicitation);

let mut escalated_from_auto_review = false;
if routes_approval_to_guardian(turn_context) {
let review_id = new_guardian_review_id();
let decision = review_approval_request(
Expand All @@ -1239,16 +1242,21 @@ async fn maybe_request_mcp_tool_approval(
monitor_reason.clone(),
)
.await;
let decision = mcp_tool_approval_decision_from_guardian(sess, &review_id, decision).await;
apply_mcp_tool_approval_decision(
sess,
turn_context,
&decision,
session_approval_key,
persistent_approval_key,
)
.await;
return Some(decision);
escalated_from_auto_review = matches!(decision, ReviewDecision::Denied)
&& take_pending_auto_review_escalation(sess, &turn_context.sub_id).await;
if !escalated_from_auto_review {
let decision =
mcp_tool_approval_decision_from_guardian(sess, &review_id, decision).await;
apply_mcp_tool_approval_decision(
sess,
turn_context,
&decision,
session_approval_key,
persistent_approval_key,
)
.await;
return Some(decision);
}
}

let prompt_options = mcp_tool_approval_prompt_options(
Expand Down Expand Up @@ -1280,7 +1288,7 @@ async fn maybe_request_mcp_tool_approval(
);
question.question =
mcp_tool_approval_question_text(question.question, monitor_reason.as_deref());
if tool_call_mcp_elicitation_enabled {
let decision = if tool_call_mcp_elicitation_enabled {
let request_id = rmcp::model::RequestId::String(
format!("{MCP_TOOL_APPROVAL_QUESTION_ID_PREFIX}_{call_id}").into(),
);
Expand Down Expand Up @@ -1309,28 +1317,22 @@ async fn maybe_request_mcp_tool_approval(
.await,
&question_id,
);
let decision = normalize_approval_decision_for_mode(decision, approval_mode);
apply_mcp_tool_approval_decision(
sess,
turn_context,
&decision,
session_approval_key,
persistent_approval_key,
normalize_approval_decision_for_mode(decision, approval_mode)
} else {
let args = RequestUserInputArgs {
questions: vec![question],
};
let response = sess
.request_user_input(turn_context.as_ref(), call_id.to_string(), args)
.await;
normalize_approval_decision_for_mode(
parse_mcp_tool_approval_response(response, &question_id),
approval_mode,
)
.await;
return Some(decision);
}

let args = RequestUserInputArgs {
questions: vec![question],
};
let response = sess
.request_user_input(turn_context.as_ref(), call_id.to_string(), args)
.await;
let decision = normalize_approval_decision_for_mode(
parse_mcp_tool_approval_response(response, &question_id),
approval_mode,
);
if escalated_from_auto_review {
reset_auto_review_rejection_circuit_breaker(sess, &turn_context.sub_id).await;
}
apply_mcp_tool_approval_decision(
sess,
turn_context,
Comment thread
won-openai marked this conversation as resolved.
Expand Down
Loading
Loading