mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
Reap stale multi-agent slots (#24903)
## Summary
- Let `close_agent` clean up an agent that is still registered in
`AgentRegistry` even when its underlying thread is already missing.
- Preserve the explicit-close boundary: for known stale thread-spawn
agents, mark the persisted spawn edge `Closed`, then treat
`ThreadNotFound` / `InternalAgentDied` as a successful close so the
registry slot can be released.
- Add a regression for MultiAgentV2 task-name targets where
`close_agent("worker")` succeeds after the worker thread has already
disappeared.
## Motivation
A worker can disappear from `ThreadManager` while its metadata still
exists in the root `AgentRegistry`. Before this change, the close tool
failed while trying to subscribe to the missing thread status, so it
never reached the cleanup path that releases the registered agent slot.
With `agents.max_threads = 1`, an explicit close of that stale task-name
agent could fail and leave the session unable to spawn a replacement.
## Scope
This PR intentionally does not add automatic stale-agent reaping to
`spawn_agent`, `resume_agent`, or `list_agents`. A thread being missing
from `ThreadManager` is not the same as an explicit close: persisted
open spawn edges are still the durable source of truth for resume and
task-name ownership until `close_agent` is called.
## Validation
- `just test -p codex-core -E
'test(multi_agent_v2_close_agent_reaps_stale_task_name_target) |
test(resume_agent_from_rollout_reopens_open_descendants_after_manager_shutdown)'`
- `just fix -p codex-core`
This commit is contained in:
committed by
GitHub
Unverified
parent
8a827d6426
commit
e2551a5e36
@@ -772,15 +772,45 @@ impl AgentControl {
|
||||
/// agent and any live descendants reached from the in-memory tree.
|
||||
pub(crate) async fn close_agent(&self, agent_id: ThreadId) -> CodexResult<String> {
|
||||
let state = self.upgrade()?;
|
||||
if let Ok(thread) = state.get_thread(agent_id).await
|
||||
&& let Some(state_db_ctx) = thread.state_db()
|
||||
&& let Err(err) = state_db_ctx
|
||||
.set_thread_spawn_edge_status(agent_id, DirectionalThreadSpawnEdgeStatus::Closed)
|
||||
.await
|
||||
{
|
||||
warn!("failed to persist thread-spawn edge status for {agent_id}: {err}");
|
||||
let known_agent = self.state.agent_metadata_for_thread(agent_id).is_some();
|
||||
match state.get_thread(agent_id).await {
|
||||
Ok(thread) => {
|
||||
if let Some(state_db_ctx) = thread.state_db()
|
||||
&& let Err(err) = state_db_ctx
|
||||
.set_thread_spawn_edge_status(
|
||||
agent_id,
|
||||
DirectionalThreadSpawnEdgeStatus::Closed,
|
||||
)
|
||||
.await
|
||||
{
|
||||
warn!("failed to persist thread-spawn edge status for {agent_id}: {err}");
|
||||
}
|
||||
}
|
||||
Err(CodexErr::ThreadNotFound(_)) if known_agent => {
|
||||
if let Some(state_db_ctx) = state.state_db()
|
||||
&& let Err(err) = state_db_ctx
|
||||
.set_thread_spawn_edge_status(
|
||||
agent_id,
|
||||
DirectionalThreadSpawnEdgeStatus::Closed,
|
||||
)
|
||||
.await
|
||||
{
|
||||
return Err(CodexErr::Fatal(format!(
|
||||
"failed to persist stale thread-spawn edge status for {agent_id}: {err}"
|
||||
)));
|
||||
}
|
||||
}
|
||||
Err(CodexErr::ThreadNotFound(_)) => {}
|
||||
Err(err) => {
|
||||
warn!("failed to inspect agent before close {agent_id}: {err}");
|
||||
}
|
||||
}
|
||||
match Box::pin(self.shutdown_agent_tree(agent_id)).await {
|
||||
Err(CodexErr::ThreadNotFound(_)) | Err(CodexErr::InternalAgentDied) if known_agent => {
|
||||
Ok(String::new())
|
||||
}
|
||||
result => result,
|
||||
}
|
||||
Box::pin(self.shutdown_agent_tree(agent_id)).await
|
||||
}
|
||||
|
||||
/// Shut down `agent_id` and any live descendants reachable from the in-memory spawn tree.
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
use super::*;
|
||||
use crate::tools::handlers::multi_agents_spec::create_close_agent_tool_v1;
|
||||
use crate::turn_timing::now_unix_timestamp_ms;
|
||||
use codex_protocol::error::CodexErr;
|
||||
use codex_tools::ToolSpec;
|
||||
|
||||
pub(crate) struct Handler;
|
||||
@@ -36,11 +37,9 @@ async fn handle_close_agent(
|
||||
let arguments = function_arguments(payload)?;
|
||||
let args: CloseAgentArgs = parse_arguments(&arguments)?;
|
||||
let agent_id = parse_agent_id_target(&args.target)?;
|
||||
let receiver_agent = session
|
||||
.services
|
||||
.agent_control
|
||||
.get_agent_metadata(agent_id)
|
||||
.unwrap_or_default();
|
||||
let receiver_agent = session.services.agent_control.get_agent_metadata(agent_id);
|
||||
let known_agent = receiver_agent.is_some();
|
||||
let receiver_agent = receiver_agent.unwrap_or_default();
|
||||
session
|
||||
.send_event(
|
||||
&turn,
|
||||
@@ -60,6 +59,9 @@ async fn handle_close_agent(
|
||||
.await
|
||||
{
|
||||
Ok(mut status_rx) => status_rx.borrow_and_update().clone(),
|
||||
Err(CodexErr::ThreadNotFound(_)) if known_agent => {
|
||||
session.services.agent_control.get_status(agent_id).await
|
||||
}
|
||||
Err(err) => {
|
||||
let status = session.services.agent_control.get_status(agent_id).await;
|
||||
session
|
||||
|
||||
@@ -52,6 +52,7 @@ use codex_protocol::protocol::TurnAbortReason;
|
||||
use codex_protocol::protocol::TurnAbortedEvent;
|
||||
use codex_protocol::protocol::TurnCompleteEvent;
|
||||
use codex_protocol::user_input::UserInput;
|
||||
use codex_state::DirectionalThreadSpawnEdgeStatus;
|
||||
use core_test_support::TempDirExt;
|
||||
use pretty_assertions::assert_eq;
|
||||
use serde::Deserialize;
|
||||
@@ -3751,6 +3752,126 @@ async fn multi_agent_v2_close_agent_accepts_task_name_target() {
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn multi_agent_v2_close_agent_reaps_stale_task_name_target() {
|
||||
let (mut session, mut turn) = make_session_and_context().await;
|
||||
let mut config = (*turn.config).clone();
|
||||
config.agent_max_threads = Some(1);
|
||||
config
|
||||
.features
|
||||
.enable(Feature::MultiAgentV2)
|
||||
.expect("test config should allow feature update");
|
||||
config
|
||||
.features
|
||||
.enable(Feature::Sqlite)
|
||||
.expect("test config should allow sqlite");
|
||||
let state_db = init_state_db(&config)
|
||||
.await
|
||||
.expect("sqlite state db should initialize");
|
||||
let manager = ThreadManager::with_models_provider_home_and_state_for_tests(
|
||||
CodexAuth::from_api_key("dummy"),
|
||||
config.model_provider.clone(),
|
||||
config.codex_home.to_path_buf(),
|
||||
Arc::new(codex_exec_server::EnvironmentManager::default_for_tests()),
|
||||
Some(state_db.clone()),
|
||||
);
|
||||
let root = manager
|
||||
.start_thread(config.clone())
|
||||
.await
|
||||
.expect("root thread should start");
|
||||
session.services.agent_control = manager.agent_control();
|
||||
session.conversation_id = root.thread_id;
|
||||
turn.config = Arc::new(config.clone());
|
||||
|
||||
let session = Arc::new(session);
|
||||
let turn = Arc::new(turn);
|
||||
SpawnAgentHandlerV2::default()
|
||||
.handle(invocation(
|
||||
session.clone(),
|
||||
turn.clone(),
|
||||
"spawn_agent",
|
||||
function_payload(json!({
|
||||
"message": "inspect this repo",
|
||||
"task_name": "worker"
|
||||
})),
|
||||
))
|
||||
.await
|
||||
.expect("spawn_agent should succeed");
|
||||
|
||||
let agent_id = session
|
||||
.services
|
||||
.agent_control
|
||||
.resolve_agent_reference(session.conversation_id, &turn.session_source, "worker")
|
||||
.await
|
||||
.expect("worker path should resolve");
|
||||
let stale_thread = manager
|
||||
.remove_thread(&agent_id)
|
||||
.await
|
||||
.expect("worker thread should be loaded before removal");
|
||||
stale_thread
|
||||
.submit(Op::Shutdown {})
|
||||
.await
|
||||
.expect("removed worker thread should still accept shutdown");
|
||||
stale_thread.wait_until_terminated().await;
|
||||
|
||||
let output = CloseAgentHandlerV2
|
||||
.handle(invocation(
|
||||
session.clone(),
|
||||
turn.clone(),
|
||||
"close_agent",
|
||||
function_payload(json!({"target": "worker"})),
|
||||
))
|
||||
.await
|
||||
.expect("close_agent should reap stale v2 task names");
|
||||
let (content, success) = expect_text_output(output);
|
||||
let result: close_agent::CloseAgentResult =
|
||||
serde_json::from_str(&content).expect("close_agent result should be json");
|
||||
assert_eq!(result.previous_status, AgentStatus::NotFound);
|
||||
assert_eq!(success, Some(true));
|
||||
|
||||
let open_children = state_db
|
||||
.list_thread_spawn_children_with_status(
|
||||
root.thread_id,
|
||||
DirectionalThreadSpawnEdgeStatus::Open,
|
||||
)
|
||||
.await
|
||||
.expect("open children should load");
|
||||
assert_eq!(open_children, Vec::<ThreadId>::new());
|
||||
let closed_children = state_db
|
||||
.list_thread_spawn_children_with_status(
|
||||
root.thread_id,
|
||||
DirectionalThreadSpawnEdgeStatus::Closed,
|
||||
)
|
||||
.await
|
||||
.expect("closed children should load");
|
||||
assert_eq!(closed_children, vec![agent_id]);
|
||||
|
||||
SpawnAgentHandlerV2::default()
|
||||
.handle(invocation(
|
||||
session.clone(),
|
||||
turn.clone(),
|
||||
"spawn_agent",
|
||||
function_payload(json!({
|
||||
"message": "inspect this repo again",
|
||||
"task_name": "replacement"
|
||||
})),
|
||||
))
|
||||
.await
|
||||
.expect("spawn_agent should succeed after stale close releases the slot");
|
||||
let replacement_id = session
|
||||
.services
|
||||
.agent_control
|
||||
.resolve_agent_reference(session.conversation_id, &turn.session_source, "replacement")
|
||||
.await
|
||||
.expect("replacement path should resolve");
|
||||
let _ = session
|
||||
.services
|
||||
.agent_control
|
||||
.shutdown_live_agent(replacement_id)
|
||||
.await
|
||||
.expect("replacement should shut down");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn multi_agent_v2_close_agent_rejects_root_target_and_id() {
|
||||
let (mut session, mut turn) = make_session_and_context().await;
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
use super::*;
|
||||
use crate::tools::handlers::multi_agents_spec::create_close_agent_tool_v2;
|
||||
use crate::turn_timing::now_unix_timestamp_ms;
|
||||
use codex_protocol::error::CodexErr;
|
||||
use codex_tools::ToolSpec;
|
||||
|
||||
pub(crate) struct Handler;
|
||||
@@ -36,11 +37,9 @@ async fn handle_close_agent(
|
||||
let arguments = function_arguments(payload)?;
|
||||
let args: CloseAgentArgs = parse_arguments(&arguments)?;
|
||||
let agent_id = resolve_agent_target(&session, &turn, &args.target).await?;
|
||||
let receiver_agent = session
|
||||
.services
|
||||
.agent_control
|
||||
.get_agent_metadata(agent_id)
|
||||
.unwrap_or_default();
|
||||
let receiver_agent = session.services.agent_control.get_agent_metadata(agent_id);
|
||||
let known_agent = receiver_agent.is_some();
|
||||
let receiver_agent = receiver_agent.unwrap_or_default();
|
||||
if receiver_agent
|
||||
.agent_path
|
||||
.as_ref()
|
||||
@@ -69,6 +68,9 @@ async fn handle_close_agent(
|
||||
.await
|
||||
{
|
||||
Ok(mut status_rx) => status_rx.borrow_and_update().clone(),
|
||||
Err(CodexErr::ThreadNotFound(_)) if known_agent => {
|
||||
session.services.agent_control.get_status(agent_id).await
|
||||
}
|
||||
Err(err) => {
|
||||
let status = session.services.agent_control.get_status(agent_id).await;
|
||||
session
|
||||
|
||||
Reference in New Issue
Block a user