Skip to content
Merged
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
50 changes: 6 additions & 44 deletions codex-rs/core/src/agent/control.rs
Original file line number Diff line number Diff line change
Expand Up @@ -50,8 +50,6 @@ pub(crate) use self::execution::AgentExecutionGuard;
use self::execution::AgentExecutionLimiter;
use self::residency::V2Residency;

const ROOT_LAST_TASK_MESSAGE: &str = "Main thread";

mod execution;
mod legacy;
mod residency;
Expand Down Expand Up @@ -82,7 +80,6 @@ pub(crate) struct LiveAgent {
pub(crate) struct ListedAgent {
pub(crate) agent_name: String,
pub(crate) agent_status: AgentStatus,
pub(crate) last_task_message: Option<String>,
}

/// Control-plane handle for multi-agent operations.
Expand Down Expand Up @@ -156,23 +153,12 @@ impl AgentControl {
state: &Arc<ThreadManagerState>,
input: Vec<UserInput>,
) -> CodexResult<String> {
let last_task_message = non_empty_task_message(render_input_preview(&input));
let result = self
.handle_thread_request_result(
agent_id,
state,
state.send_op(agent_id, input.into()).await,
)
.await;
if result.is_ok() {
match last_task_message {
Some(last_task_message) => self
.state
.update_last_task_message(agent_id, last_task_message),
None => self.state.clear_last_task_message(agent_id),
}
}
result
self.handle_thread_request_result(
agent_id,
state,
state.send_op(agent_id, input.into()).await,
)
.await
}

pub(crate) async fn send_inter_agent_communication(
Expand Down Expand Up @@ -211,7 +197,6 @@ impl AgentControl {
communication: InterAgentCommunication,
context: AgentCommunicationContext,
) -> CodexResult<String> {
let last_task_message = last_task_message_from_communication(&communication);
let communication_for_log =
crate::agent_communication::logging_enabled().then(|| communication.clone());
let result = self
Expand All @@ -233,14 +218,6 @@ impl AgentControl {
agent_id,
);
}
if result.is_ok() {
match last_task_message {
Some(last_task_message) => self
.state
.update_last_task_message(agent_id, last_task_message),
None => self.state.clear_last_task_message(agent_id),
}
}
result
}

Expand Down Expand Up @@ -417,7 +394,6 @@ impl AgentControl {
agents.push(ListedAgent {
agent_name: root_path.to_string(),
agent_status: root_thread.agent_status().await,
last_task_message: Some(ROOT_LAST_TASK_MESSAGE.to_string()),
});
}

Expand All @@ -440,11 +416,9 @@ impl AgentControl {
.as_ref()
.map(ToString::to_string)
.unwrap_or_else(|| thread_id.to_string());
let last_task_message = metadata.last_task_message.clone();
agents.push(ListedAgent {
agent_name,
agent_status: thread.agent_status().await,
last_task_message,
});
}

Expand Down Expand Up @@ -562,7 +536,6 @@ impl AgentControl {
agent_path,
agent_nickname,
agent_role,
last_task_message: None,
})
}

Expand Down Expand Up @@ -787,17 +760,6 @@ pub(crate) fn render_input_preview(input: &[UserInput]) -> String {
.join("\n")
}

fn last_task_message_from_communication(communication: &InterAgentCommunication) -> Option<String> {
if communication.encrypted_content.is_some() {
return None;
}
non_empty_task_message(communication.content.clone())
}

fn non_empty_task_message(message: String) -> Option<String> {
(!message.is_empty()).then_some(message)
}

fn thread_spawn_depth(session_source: &SessionSource) -> Option<i32> {
match session_source {
SessionSource::SubAgent(SubAgentSource::ThreadSpawn { depth, .. }) => Some(*depth),
Expand Down
60 changes: 0 additions & 60 deletions codex-rs/core/src/agent/control_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -775,66 +775,6 @@ async fn resume_agent_from_rollout_does_not_reopen_v2_descendants() {
assert_thread_not_loaded(&resumed_manager, reviewer_thread_id).await;
}

#[tokio::test]
async fn encrypted_inter_agent_communication_clears_existing_last_task_message() {
let harness = AgentControlHarness::new().await;
let (parent_thread_id, _) = harness.start_thread().await;
let agent_path = AgentPath::try_from("/root/worker").expect("agent path");
let spawned_agent = harness
.control
.spawn_agent_with_metadata(
harness.config.clone(),
text_input("old plaintext task"),
Some(SessionSource::SubAgent(SubAgentSource::ThreadSpawn {
parent_thread_id,
depth: 1,
agent_path: Some(agent_path.clone()),
agent_nickname: None,
agent_role: None,
})),
SpawnAgentOptions {
parent_thread_id: Some(parent_thread_id),
..Default::default()
},
)
.await
.expect("spawn_agent should succeed");
assert_eq!(
harness
.control
.state
.agent_metadata_for_thread(spawned_agent.thread_id)
.and_then(|metadata| metadata.last_task_message),
Some("old plaintext task".to_string())
);

let communication = InterAgentCommunication::new_encrypted(
AgentPath::root(),
agent_path,
Vec::new(),
"encrypted-task".to_string(),
/*trigger_turn*/ true,
);
harness
.control
.send_inter_agent_communication(
spawned_agent.thread_id,
communication,
AgentCommunicationContext::new(AgentCommunicationKind::Followup, ThreadId::new()),
)
.await
.expect("send_inter_agent_communication should succeed");

assert_eq!(
harness
.control
.state
.agent_metadata_for_thread(spawned_agent.thread_id)
.and_then(|metadata| metadata.last_task_message),
None
);
}

#[tokio::test]
async fn spawn_agent_creates_thread_and_sends_prompt() {
let harness = AgentControlHarness::new().await;
Expand Down
29 changes: 0 additions & 29 deletions codex-rs/core/src/agent/registry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,6 @@ pub(crate) struct AgentMetadata {
pub(crate) agent_path: Option<AgentPath>,
pub(crate) agent_nickname: Option<String>,
pub(crate) agent_role: Option<String>,
pub(crate) last_task_message: Option<String>,
}

fn format_agent_nickname(name: &str, nickname_reset_count: usize) -> String {
Expand Down Expand Up @@ -166,34 +165,6 @@ impl AgentRegistry {
.collect()
}

pub(crate) fn update_last_task_message(&self, thread_id: ThreadId, last_task_message: String) {
let mut active_agents = self
.active_agents
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(metadata) = active_agents
.agent_tree
.values_mut()
.find(|metadata| metadata.agent_id == Some(thread_id))
{
metadata.last_task_message = Some(last_task_message);
}
}

pub(crate) fn clear_last_task_message(&self, thread_id: ThreadId) {
let mut active_agents = self
.active_agents
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(metadata) = active_agents
.agent_tree
.values_mut()
.find(|metadata| metadata.agent_id == Some(thread_id))
{
metadata.last_task_message = None;
}
}

fn register_spawned_thread(&self, agent_metadata: AgentMetadata) {
let Some(thread_id) = agent_metadata.agent_id else {
return;
Expand Down
6 changes: 1 addition & 5 deletions codex-rs/core/src/tools/handlers/multi_agents_spec.rs
Original file line number Diff line number Diff line change
Expand Up @@ -460,13 +460,9 @@ fn list_agents_output_schema() -> Value {
"agent_status": {
"description": "Last known status of the agent.",
"allOf": [agent_status_output_schema()]
},
"last_task_message": {
"type": ["string", "null"],
"description": "Most recent user or inter-agent instruction received by the agent, when available."
}
},
"required": ["agent_name", "agent_status", "last_task_message"],
"required": ["agent_name", "agent_status"],
"additionalProperties": false
},
"description": "Live agents visible in the current root thread tree."
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -443,7 +443,7 @@ fn list_agents_tool_includes_path_prefix_and_agent_fields() {
);
assert_eq!(
output_schema.expect("list_agents output schema")["properties"]["agents"]["items"]["required"],
json!(["agent_name", "agent_status", "last_task_message"])
json!(["agent_name", "agent_status"])
);
}

Expand Down
19 changes: 1 addition & 18 deletions codex-rs/core/src/tools/handlers/multi_agents_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,6 @@ struct ListAgentsResult {
struct ListedAgentResult {
agent_name: String,
agent_status: serde_json::Value,
last_task_message: Option<String>,
}

#[derive(Debug, Deserialize)]
Expand Down Expand Up @@ -1546,7 +1545,7 @@ async fn multi_agent_v2_followup_task_rejects_root_target_from_child() {
}

#[tokio::test]
async fn multi_agent_v2_list_agents_returns_completed_status_without_encrypted_spawn_preview() {
async fn multi_agent_v2_list_agents_returns_completed_status() {
let (mut session, mut turn) = make_session_and_context().await;
let manager = thread_manager();
let root = manager
Expand Down Expand Up @@ -1622,19 +1621,12 @@ async fn multi_agent_v2_list_agents_returns_completed_status_without_encrypted_s
.map(|agent| agent.agent_name.as_str())
.collect::<Vec<_>>();
assert_eq!(agent_names, vec!["/root", "/root/worker"]);
let root_agent = result
.agents
.iter()
.find(|agent| agent.agent_name == "/root")
.expect("root agent should be listed");
assert_eq!(root_agent.last_task_message.as_deref(), Some("Main thread"));
let worker = result
.agents
.iter()
.find(|agent| agent.agent_name == "/root/worker")
.expect("worker agent should be listed");
assert_eq!(worker.agent_status, json!({"completed": "done"}));
assert_eq!(worker.last_task_message, None);
assert_eq!(success, Some(true));
}

Expand Down Expand Up @@ -1720,7 +1712,6 @@ async fn multi_agent_v2_list_agents_filters_by_relative_path_prefix() {

assert_eq!(result.agents.len(), 1);
assert_eq!(result.agents[0].agent_name, worker_path.as_str());
assert_eq!(result.agents[0].last_task_message.as_deref(), Some("build"));
}

#[tokio::test]
Expand Down Expand Up @@ -1781,10 +1772,6 @@ async fn multi_agent_v2_list_agents_omits_closed_agents() {

assert_eq!(result.agents.len(), 1);
assert_eq!(result.agents[0].agent_name, "/root");
assert_eq!(
result.agents[0].last_task_message.as_deref(),
Some("Main thread")
);
}

#[tokio::test]
Expand Down Expand Up @@ -1856,10 +1843,6 @@ async fn multi_agent_v2_list_agents_keeps_interrupted_resident_agents() {

assert_eq!(result.agents.len(), 2);
assert_eq!(result.agents[0].agent_name, "/root");
assert_eq!(
result.agents[0].last_task_message.as_deref(),
Some("Main thread")
);
assert_eq!(result.agents[1].agent_name, agent_path.as_str());
}

Expand Down
Loading