Skip to content

Commit 5348fb6

Browse files
Fix completed subagent quota reclamation
Reclaim spawn quota slots when subagents reach final states, retry spawn after opportunistic cleanup, and keep interrupted subagents visible as active quota holders. This prevents completed subagents from leaking thread quota while preserving the ability to resume interrupted agents. Co-authored-by: Open Codex <hff582580@gmail.com>
1 parent ba54960 commit 5348fb6

5 files changed

Lines changed: 401 additions & 5 deletions

File tree

codex-rs/core/src/agent/control.rs

Lines changed: 72 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -200,7 +200,9 @@ impl AgentControl {
200200
options: SpawnAgentOptions,
201201
) -> CodexResult<LiveAgent> {
202202
let state = self.upgrade()?;
203-
let mut reservation = self.state.reserve_spawn_slot(config.agent_max_threads)?;
203+
let mut reservation = self
204+
.reserve_spawn_slot_with_final_reclaim(&state, &config)
205+
.await?;
204206
let inherited_shell_snapshot = self
205207
.inherited_shell_snapshot_for_source(&state, session_source.as_ref())
206208
.await;
@@ -538,7 +540,9 @@ impl AgentControl {
538540
}
539541
let state = self.upgrade()?;
540542
let state_db_ctx = state.state_db();
541-
let mut reservation = self.state.reserve_spawn_slot(config.agent_max_threads)?;
543+
let mut reservation = self
544+
.reserve_spawn_slot_with_final_reclaim(&state, &config)
545+
.await?;
542546
let (session_source, agent_metadata) = match session_source {
543547
SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
544548
parent_thread_id,
@@ -638,6 +642,7 @@ impl AgentControl {
638642
) -> CodexResult<String> {
639643
let last_task_message = render_input_preview(&initial_operation);
640644
let state = self.upgrade()?;
645+
self.ensure_agent_slot_for_input(agent_id, &state).await?;
641646
let result = self
642647
.handle_thread_request_result(
643648
agent_id,
@@ -675,6 +680,7 @@ impl AgentControl {
675680
) -> CodexResult<String> {
676681
let last_task_message = communication.content.clone();
677682
let state = self.upgrade()?;
683+
self.ensure_agent_slot_for_input(agent_id, &state).await?;
678684
let result = self
679685
.handle_thread_request_result(
680686
agent_id,
@@ -710,6 +716,66 @@ impl AgentControl {
710716
result
711717
}
712718

719+
async fn reserve_spawn_slot_with_final_reclaim(
720+
&self,
721+
state: &Arc<ThreadManagerState>,
722+
config: &crate::config::Config,
723+
) -> CodexResult<crate::agent::registry::SpawnReservation> {
724+
match self.state.reserve_spawn_slot(config.agent_max_threads) {
725+
Ok(reservation) => Ok(reservation),
726+
Err(CodexErr::AgentLimitReached { .. }) => {
727+
self.reclaim_final_agent_slots(state).await;
728+
self.state.reserve_spawn_slot(config.agent_max_threads)
729+
}
730+
Err(err) => Err(err),
731+
}
732+
}
733+
734+
async fn ensure_agent_slot_for_input(
735+
&self,
736+
agent_id: ThreadId,
737+
state: &Arc<ThreadManagerState>,
738+
) -> CodexResult<()> {
739+
if !self.state.agent_slot_is_reclaimed(agent_id) {
740+
return Ok(());
741+
}
742+
self.reclaim_final_agent_slots(state).await;
743+
let thread = state.get_thread(agent_id).await?;
744+
let config = thread.codex.session.get_config().await;
745+
self.state
746+
.occupy_reclaimed_thread_slot(agent_id, config.agent_max_threads)?;
747+
Ok(())
748+
}
749+
750+
async fn reclaim_final_agent_slots(&self, state: &Arc<ThreadManagerState>) {
751+
for agent_id in self.state.occupied_agent_ids() {
752+
let status = match state.get_thread(agent_id).await {
753+
Ok(thread) => thread.agent_status().await,
754+
Err(_) => AgentStatus::NotFound,
755+
};
756+
self.reclaim_final_agent_slot(agent_id, state, &status)
757+
.await;
758+
}
759+
}
760+
761+
async fn reclaim_final_agent_slot(
762+
&self,
763+
agent_id: ThreadId,
764+
state: &Arc<ThreadManagerState>,
765+
status: &AgentStatus,
766+
) {
767+
if !is_final(status) {
768+
return;
769+
}
770+
if let Ok(thread) = state.get_thread(agent_id).await {
771+
thread.codex.session.ensure_rollout_materialized().await;
772+
if let Err(err) = thread.codex.session.flush_rollout().await {
773+
warn!("failed to flush final subagent rollout before reclaiming slot: {err}");
774+
}
775+
}
776+
self.state.reclaim_spawned_thread_slot(agent_id);
777+
}
778+
713779
/// Submit a shutdown request for a live agent without marking it explicitly closed in
714780
/// persisted spawn-edge state.
715781
pub(crate) async fn shutdown_live_agent(&self, agent_id: ThreadId) -> CodexResult<String> {
@@ -978,6 +1044,9 @@ impl AgentControl {
9781044
};
9791045
let child_thread = state.get_thread(child_thread_id).await.ok();
9801046
let message = format_subagent_notification_message(child_reference.as_str(), &status);
1047+
control
1048+
.reclaim_final_agent_slot(child_thread_id, &state, &status)
1049+
.await;
9811050
if child_agent_path.is_some()
9821051
&& child_thread
9831052
.as_ref()
@@ -1051,6 +1120,7 @@ impl AgentControl {
10511120
agent_nickname,
10521121
agent_role,
10531122
last_task_message: None,
1123+
slot_occupied: false,
10541124
};
10551125
Ok((session_source, agent_metadata))
10561126
}

codex-rs/core/src/agent/control_tests.rs

Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1052,6 +1052,134 @@ async fn spawn_agent_releases_slot_after_shutdown() {
10521052
.expect("shutdown agent");
10531053
}
10541054

1055+
#[tokio::test]
1056+
async fn spawn_agent_reclaims_completed_slot_before_limit_error() {
1057+
let max_threads = 1usize;
1058+
let (_home, config) = test_config_with_cli_overrides(vec![(
1059+
"agents.max_threads".to_string(),
1060+
TomlValue::Integer(max_threads as i64),
1061+
)])
1062+
.await;
1063+
let manager = ThreadManager::with_models_provider_and_home_for_tests(
1064+
CodexAuth::from_api_key("dummy"),
1065+
config.model_provider.clone(),
1066+
config.codex_home.to_path_buf(),
1067+
std::sync::Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()),
1068+
);
1069+
let control = manager.agent_control();
1070+
1071+
let completed_agent_id = control
1072+
.spawn_agent(
1073+
config.clone(),
1074+
text_input("complete then release"),
1075+
/*session_source*/ None,
1076+
)
1077+
.await
1078+
.expect("spawn_agent should succeed");
1079+
let completed_agent = manager
1080+
.get_thread(completed_agent_id)
1081+
.await
1082+
.expect("completed agent should be registered");
1083+
let turn = completed_agent.codex.session.new_default_turn().await;
1084+
completed_agent
1085+
.codex
1086+
.session
1087+
.send_event(
1088+
turn.as_ref(),
1089+
EventMsg::TurnComplete(TurnCompleteEvent {
1090+
turn_id: turn.sub_id.clone(),
1091+
last_agent_message: Some("done".to_string()),
1092+
completed_at: None,
1093+
duration_ms: None,
1094+
time_to_first_token_ms: None,
1095+
}),
1096+
)
1097+
.await;
1098+
1099+
let next_agent_id = control
1100+
.spawn_agent(
1101+
config.clone(),
1102+
text_input("new work after completed agent"),
1103+
/*session_source*/ None,
1104+
)
1105+
.await
1106+
.expect("completed agent should not block the spawn quota");
1107+
1108+
let _ = control
1109+
.shutdown_live_agent(completed_agent_id)
1110+
.await
1111+
.expect("shutdown completed agent");
1112+
let _ = control
1113+
.shutdown_live_agent(next_agent_id)
1114+
.await
1115+
.expect("shutdown next agent");
1116+
}
1117+
1118+
#[tokio::test]
1119+
async fn spawn_agent_does_not_reclaim_interrupted_slot() {
1120+
let max_threads = 1usize;
1121+
let (_home, config) = test_config_with_cli_overrides(vec![(
1122+
"agents.max_threads".to_string(),
1123+
TomlValue::Integer(max_threads as i64),
1124+
)])
1125+
.await;
1126+
let manager = ThreadManager::with_models_provider_and_home_for_tests(
1127+
CodexAuth::from_api_key("dummy"),
1128+
config.model_provider.clone(),
1129+
config.codex_home.to_path_buf(),
1130+
std::sync::Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()),
1131+
);
1132+
let control = manager.agent_control();
1133+
1134+
let interrupted_agent_id = control
1135+
.spawn_agent(
1136+
config.clone(),
1137+
text_input("interrupt and keep slot"),
1138+
/*session_source*/ None,
1139+
)
1140+
.await
1141+
.expect("spawn_agent should succeed");
1142+
let interrupted_agent = manager
1143+
.get_thread(interrupted_agent_id)
1144+
.await
1145+
.expect("interrupted agent should be registered");
1146+
let turn = interrupted_agent.codex.session.new_default_turn().await;
1147+
interrupted_agent
1148+
.codex
1149+
.session
1150+
.send_event(
1151+
turn.as_ref(),
1152+
EventMsg::TurnAborted(TurnAbortedEvent {
1153+
turn_id: Some("turn-1".to_string()),
1154+
reason: TurnAbortReason::Interrupted,
1155+
completed_at: None,
1156+
duration_ms: None,
1157+
}),
1158+
)
1159+
.await;
1160+
1161+
let err = control
1162+
.spawn_agent(
1163+
config,
1164+
text_input("should still be blocked"),
1165+
/*session_source*/ None,
1166+
)
1167+
.await
1168+
.expect_err("interrupted agent should still occupy the spawn quota");
1169+
let CodexErr::AgentLimitReached {
1170+
max_threads: seen_max_threads,
1171+
} = err
1172+
else {
1173+
panic!("expected CodexErr::AgentLimitReached");
1174+
};
1175+
assert_eq!(seen_max_threads, max_threads);
1176+
1177+
let _ = control
1178+
.shutdown_live_agent(interrupted_agent_id)
1179+
.await
1180+
.expect("shutdown interrupted agent");
1181+
}
1182+
10551183
#[tokio::test]
10561184
async fn spawn_agent_limit_shared_across_clones() {
10571185
let max_threads = 1usize;

0 commit comments

Comments
 (0)