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
36 changes: 15 additions & 21 deletions codex-rs/app-server/src/bespoke_event_handling.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,7 @@ use codex_core::ThreadManager;
use codex_core::review_format::format_review_findings_block;
use codex_core::review_prompts;
use codex_protocol::ThreadId;
use codex_protocol::items::CollabAgentTool as CoreCollabAgentTool;
use codex_protocol::items::TurnItem as CoreTurnItem;
use codex_protocol::items::parse_hook_prompt_message;
use codex_protocol::models::AdditionalPermissionProfile as CoreAdditionalPermissionProfile;
Expand Down Expand Up @@ -831,15 +832,20 @@ pub(crate) async fn apply_bespoke_event_handling(
// Deprecated MCP tool-call events are still fanned out for legacy clients.
// App-server v2 receives the canonical TurnItem::McpToolCall lifecycle instead.
}
msg @ (EventMsg::CollabAgentSpawnBegin(_)
EventMsg::CollabAgentSpawnBegin(_)
| EventMsg::CollabAgentSpawnEnd(_)
| EventMsg::CollabAgentInteractionBegin(_)
| EventMsg::CollabAgentInteractionEnd(_)
| EventMsg::CollabWaitingBegin(_)
| EventMsg::CollabWaitingEnd(_)
| EventMsg::CollabCloseBegin(_)
| EventMsg::CollabCloseEnd(_)
| EventMsg::CollabResumeBegin(_)
| EventMsg::CollabResumeEnd(_)
| EventMsg::CollabResumeEnd(_) => {
// Deprecated non-wait collaboration events are still fanned out for raw-event and
// rollout compatibility consumers. App-server v2 receives the canonical
// CollabAgentToolCall item lifecycle instead.
}
msg @ (EventMsg::CollabWaitingBegin(_)
| EventMsg::CollabWaitingEnd(_)
| EventMsg::AgentMessageContentDelta(_)
| EventMsg::PlanDelta(_)
| EventMsg::ReasoningContentDelta(_)
Expand All @@ -857,23 +863,6 @@ pub(crate) async fn apply_bespoke_event_handling(
// rollout compatibility consumers. App-server v2 receives the canonical
// SubAgentActivity item lifecycle instead.
}
EventMsg::CollabCloseEnd(end_event) => {
if thread_manager
.get_thread(end_event.receiver_thread_id)
.await
.is_err()
{
thread_watch_manager
.remove_thread(&end_event.receiver_thread_id.to_string())
.await;
}
let notification = item_event_to_server_notification(
EventMsg::CollabCloseEnd(end_event),
&conversation_id.to_string(),
&event_turn_id,
);
outgoing.send_server_notification(notification).await;
}
EventMsg::ContextCompacted(..) => {
// Core still fans out this deprecated event for legacy clients;
// v2 clients receive the canonical ContextCompaction item instead.
Expand Down Expand Up @@ -1351,6 +1340,11 @@ async fn apply_canonical_item_completed_side_effects(
)
.await;
}
CoreTurnItem::CollabAgentToolCall(item) if item.tool == CoreCollabAgentTool::CloseAgent => {
for thread_id in &item.receiver_thread_ids {
remove_missing_thread_watch(thread_manager, thread_watch_manager, *thread_id).await;
}
}
_ => {}
}
}
Expand Down
25 changes: 15 additions & 10 deletions codex-rs/core/src/tools/handlers/multi_agents.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,8 +8,6 @@
use crate::agent::AgentStatus;
use crate::agent::exceeds_thread_spawn_depth_limit;
use crate::function_tool::FunctionCallError;
use crate::session::session::Session;
use crate::session::turn_context::TurnContext;
use crate::tools::context::ToolInvocation;
use crate::tools::context::ToolOutput;
use crate::tools::context::ToolPayload;
Expand All @@ -20,17 +18,13 @@ use crate::tools::handlers::parse_arguments;
use crate::tools::registry::CoreToolRuntime;
use crate::tools::registry::ToolExecutor;
use codex_protocol::ThreadId;
use codex_protocol::items::CollabAgentTool;
use codex_protocol::items::CollabAgentToolCallItem;
use codex_protocol::items::CollabAgentToolCallStatus;
use codex_protocol::items::TurnItem;
use codex_protocol::models::ResponseInputItem;
use codex_protocol::openai_models::ReasoningEffort;
use codex_protocol::protocol::CollabAgentInteractionBeginEvent;
use codex_protocol::protocol::CollabAgentInteractionEndEvent;
use codex_protocol::protocol::CollabAgentRef;
use codex_protocol::protocol::CollabAgentSpawnBeginEvent;
use codex_protocol::protocol::CollabAgentSpawnEndEvent;
use codex_protocol::protocol::CollabCloseBeginEvent;
use codex_protocol::protocol::CollabCloseEndEvent;
use codex_protocol::protocol::CollabResumeBeginEvent;
use codex_protocol::protocol::CollabResumeEndEvent;
use codex_protocol::protocol::CollabWaitingBeginEvent;
use codex_protocol::protocol::CollabWaitingEndEvent;
use codex_protocol::user_input::UserInput;
Expand Down Expand Up @@ -91,6 +85,17 @@ mod send_input;
mod spawn;
pub(crate) mod wait;

pub(crate) fn collab_tool_call_status(
status: &AgentStatus,
receiver_thread_id: Option<ThreadId>,
) -> CollabAgentToolCallStatus {
match status {
AgentStatus::Errored(_) | AgentStatus::NotFound => CollabAgentToolCallStatus::Failed,
_ if receiver_thread_id.is_some() => CollabAgentToolCallStatus::Completed,
_ => CollabAgentToolCallStatus::Failed,
}
}

#[cfg(test)]
#[path = "multi_agents_tests.rs"]
mod tests;
72 changes: 44 additions & 28 deletions codex-rs/core/src/tools/handlers/multi_agents/close_agent.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,5 @@
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;

Expand Down Expand Up @@ -44,15 +43,20 @@ async fn handle_close_agent(
let known_agent = receiver_agent.is_some();
let receiver_agent = receiver_agent.unwrap_or_default();
session
.send_event(
.emit_turn_item_started(
&turn,
CollabCloseBeginEvent {
call_id: call_id.clone(),
started_at_ms: now_unix_timestamp_ms(),
&TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id.clone(),
tool: CollabAgentTool::CloseAgent,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: session.thread_id,
receiver_thread_id: agent_id,
}
.into(),
receiver_thread_ids: vec![agent_id],
receiver_agents: Vec::new(),
prompt: None,
model: None,
reasoning_effort: None,
agents_states: Default::default(),
}),
)
.await;
let status = match session
Expand All @@ -68,18 +72,24 @@ async fn handle_close_agent(
Err(err) => {
let status = session.services.agent_control.get_status(agent_id).await;
session
.send_event(
.emit_turn_item_completed(
&turn,
CollabCloseEndEvent {
call_id: call_id.clone(),
completed_at_ms: now_unix_timestamp_ms(),
TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id.clone(),
tool: CollabAgentTool::CloseAgent,
status: collab_tool_call_status(&status, Some(agent_id)),
sender_thread_id: session.thread_id(),
receiver_thread_id: agent_id,
receiver_agent_nickname: receiver_agent.agent_nickname.clone(),
receiver_agent_role: receiver_agent.agent_role.clone(),
status,
}
.into(),
receiver_thread_ids: vec![agent_id],
receiver_agents: vec![CollabAgentRef {
thread_id: agent_id,
agent_nickname: receiver_agent.agent_nickname.clone(),
agent_role: receiver_agent.agent_role.clone(),
}],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: [(agent_id, status)].into_iter().collect(),
}),
)
.await;
return Err(collab_agent_error(agent_id, err));
Expand All @@ -90,18 +100,24 @@ async fn handle_close_agent(
.map_err(|err| collab_agent_error(agent_id, err))
.map(|_| ());
session
.send_event(
.emit_turn_item_completed(
&turn,
CollabCloseEndEvent {
call_id,
completed_at_ms: now_unix_timestamp_ms(),
TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id,
tool: CollabAgentTool::CloseAgent,
status: collab_tool_call_status(&status, Some(agent_id)),
sender_thread_id: session.thread_id,
receiver_thread_id: agent_id,
receiver_agent_nickname: receiver_agent.agent_nickname,
receiver_agent_role: receiver_agent.agent_role,
status: status.clone(),
}
.into(),
receiver_thread_ids: vec![agent_id],
receiver_agents: vec![CollabAgentRef {
thread_id: agent_id,
agent_nickname: receiver_agent.agent_nickname,
agent_role: receiver_agent.agent_role,
}],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: [(agent_id, status.clone())].into_iter().collect(),
}),
)
.await;
result?;
Expand Down
54 changes: 34 additions & 20 deletions codex-rs/core/src/tools/handlers/multi_agents/resume_agent.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,8 @@
use super::*;
use crate::agent::next_thread_spawn_depth;
use crate::session::session::Session;
use crate::session::turn_context::TurnContext;
use crate::tools::handlers::multi_agents_spec::create_resume_agent_tool;
use crate::turn_timing::now_unix_timestamp_ms;
use codex_tools::ToolSpec;
use std::sync::Arc;

Expand Down Expand Up @@ -57,17 +58,24 @@ async fn handle_resume_agent(
}

session
.send_event(
.emit_turn_item_started(
&turn,
CollabResumeBeginEvent {
call_id: call_id.clone(),
started_at_ms: now_unix_timestamp_ms(),
&TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id.clone(),
tool: CollabAgentTool::ResumeAgent,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: session.thread_id,
receiver_thread_id,
receiver_agent_nickname: receiver_agent.agent_nickname.clone(),
receiver_agent_role: receiver_agent.agent_role.clone(),
}
.into(),
receiver_thread_ids: vec![receiver_thread_id],
receiver_agents: vec![CollabAgentRef {
thread_id: receiver_thread_id,
agent_nickname: receiver_agent.agent_nickname.clone(),
agent_role: receiver_agent.agent_role.clone(),
}],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: Default::default(),
}),
)
.await;

Expand Down Expand Up @@ -113,18 +121,24 @@ async fn handle_resume_agent(
(receiver_agent, None)
};
session
.send_event(
.emit_turn_item_completed(
&turn,
CollabResumeEndEvent {
call_id,
completed_at_ms: now_unix_timestamp_ms(),
TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id,
tool: CollabAgentTool::ResumeAgent,
status: collab_tool_call_status(&status, Some(receiver_thread_id)),
sender_thread_id: session.thread_id(),
receiver_thread_id,
receiver_agent_nickname: receiver_agent.agent_nickname,
receiver_agent_role: receiver_agent.agent_role,
status: status.clone(),
}
.into(),
receiver_thread_ids: vec![receiver_thread_id],
receiver_agents: vec![CollabAgentRef {
thread_id: receiver_thread_id,
agent_nickname: receiver_agent.agent_nickname,
agent_role: receiver_agent.agent_role,
}],
prompt: None,
model: None,
reasoning_effort: None,
agents_states: [(receiver_thread_id, status.clone())].into_iter().collect(),
}),
)
.await;

Expand Down
48 changes: 28 additions & 20 deletions codex-rs/core/src/tools/handlers/multi_agents/send_input.rs
Original file line number Diff line number Diff line change
@@ -1,7 +1,6 @@
use super::*;
use crate::agent::control::render_input_preview;
use crate::tools::handlers::multi_agents_spec::create_send_input_tool_v1;
use crate::turn_timing::now_unix_timestamp_ms;
use codex_tools::ToolSpec;

pub(crate) struct Handler;
Expand Down Expand Up @@ -67,16 +66,20 @@ impl Handler {
.map_err(|err| collab_agent_error(receiver_thread_id, err))?;
}
session
.send_event(
.emit_turn_item_started(
&turn,
CollabAgentInteractionBeginEvent {
call_id: call_id.clone(),
started_at_ms: now_unix_timestamp_ms(),
&TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id.clone(),
tool: CollabAgentTool::SendInput,
status: CollabAgentToolCallStatus::InProgress,
sender_thread_id: session.thread_id,
receiver_thread_id,
prompt: prompt.clone(),
}
.into(),
receiver_thread_ids: vec![receiver_thread_id],
receiver_agents: Vec::new(),
prompt: Some(prompt.clone()),
model: None,
reasoning_effort: None,
agents_states: Default::default(),
}),
)
.await;
let agent_control = session.services.agent_control.clone();
Expand All @@ -90,19 +93,24 @@ impl Handler {
.get_status(receiver_thread_id)
.await;
session
.send_event(
.emit_turn_item_completed(
&turn,
CollabAgentInteractionEndEvent {
call_id,
completed_at_ms: now_unix_timestamp_ms(),
TurnItem::CollabAgentToolCall(CollabAgentToolCallItem {
id: call_id,
tool: CollabAgentTool::SendInput,
status: collab_tool_call_status(&status, Some(receiver_thread_id)),
sender_thread_id: session.thread_id,
receiver_thread_id,
receiver_agent_nickname: receiver_agent.agent_nickname,
receiver_agent_role: receiver_agent.agent_role,
prompt,
status,
}
.into(),
receiver_thread_ids: vec![receiver_thread_id],
receiver_agents: vec![CollabAgentRef {
thread_id: receiver_thread_id,
agent_nickname: receiver_agent.agent_nickname,
agent_role: receiver_agent.agent_role,
}],
prompt: Some(prompt),
model: None,
reasoning_effort: None,
agents_states: [(receiver_thread_id, status)].into_iter().collect(),
}),
)
.await;
let submission_id = result?;
Expand Down
Loading
Loading