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
33 changes: 32 additions & 1 deletion codex-rs/app-server-protocol/src/protocol/thread_history.rs
Original file line number Diff line number Diff line change
Expand Up @@ -600,6 +600,7 @@ impl ThreadHistoryBuilder {
let should_upsert = match item {
codex_protocol::items::TurnItem::Plan(plan) => !plan.text.is_empty(),
codex_protocol::items::TurnItem::Sleep(_)
| codex_protocol::items::TurnItem::HookPrompt(_)
| codex_protocol::items::TurnItem::CommandExecution(_)
| codex_protocol::items::TurnItem::DynamicToolCall(_)
| codex_protocol::items::TurnItem::CollabAgentToolCall(_)
Expand All @@ -608,7 +609,6 @@ impl ThreadHistoryBuilder {
| codex_protocol::items::TurnItem::EnteredReviewMode(_)
| codex_protocol::items::TurnItem::ExitedReviewMode(_) => true,
codex_protocol::items::TurnItem::UserMessage(_)
| codex_protocol::items::TurnItem::HookPrompt(_)
| codex_protocol::items::TurnItem::AgentMessage(_)
| codex_protocol::items::TurnItem::Reasoning(_)
| codex_protocol::items::TurnItem::WebSearch(_)
Expand Down Expand Up @@ -4027,6 +4027,37 @@ mod tests {
);
}

#[test]
fn canonical_hook_prompt_completion_updates_turn_history() {
let hook_prompt = CoreTurnItem::HookPrompt(codex_protocol::items::HookPromptItem {
id: "hook-prompt-1".into(),
fragments: vec![CoreHookPromptFragment::from_single_hook(
"Retry with tests.",
"hook-run-1",
)],
});
let expected_item = ThreadItem::from(hook_prompt.clone());
let mut builder = ThreadHistoryBuilder::new();
builder.handle_event(&EventMsg::TurnStarted(TurnStartedEvent {
turn_id: "turn-a".into(),
trace_id: None,
started_at: None,
model_context_window: None,
collaboration_mode_kind: Default::default(),
}));
builder.handle_event(&EventMsg::ItemCompleted(ItemCompletedEvent {
thread_id: ThreadId::new(),
turn_id: "turn-a".into(),
item: hook_prompt,
completed_at_ms: 0,
}));

assert_eq!(
builder.active_turn_snapshot().expect("active turn").items,
vec![expected_item]
);
}

#[test]
fn ignores_plain_user_response_items_in_rollout_replay() {
let items = vec![
Expand Down
100 changes: 0 additions & 100 deletions codex-rs/app-server/src/bespoke_event_handling.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,6 @@ use codex_core::ThreadManager;
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;
use codex_protocol::plan_tool::UpdatePlanArgs;
use codex_protocol::protocol::CodexErrorInfo as CoreCodexErrorInfo;
Expand Down Expand Up @@ -1022,13 +1021,6 @@ pub(crate) async fn apply_bespoke_event_handling(
.await;
}
EventMsg::RawResponseItem(raw_response_item_event) => {
maybe_emit_hook_prompt_item_completed(
conversation_id,
&event_turn_id,
&raw_response_item_event.item,
&outgoing,
)
.await;
maybe_emit_raw_response_item_completed(
conversation_id,
&event_turn_id,
Expand Down Expand Up @@ -1406,45 +1398,6 @@ async fn maybe_emit_raw_response_item_completed(
.await;
}

pub(crate) async fn maybe_emit_hook_prompt_item_completed(
conversation_id: ThreadId,
turn_id: &str,
item: &codex_protocol::models::ResponseItem,
outgoing: &ThreadScopedOutgoingMessageSender,
) {
let codex_protocol::models::ResponseItem::Message {
role, content, id, ..
} = item
else {
return;
};

if role != "user" {
return;
}

let Some(hook_prompt) = parse_hook_prompt_message(id.as_ref(), content) else {
return;
};

let notification = ItemCompletedNotification {
thread_id: conversation_id.to_string(),
turn_id: turn_id.to_string(),
completed_at_ms: now_unix_timestamp_ms(),
item: ThreadItem::HookPrompt {
id: hook_prompt.id,
fragments: hook_prompt
.fragments
.into_iter()
.map(codex_app_server_protocol::HookPromptFragment::from)
.collect(),
},
};
outgoing
.send_server_notification(ServerNotification::ItemCompleted(notification))
.await;
}

async fn find_and_remove_turn_summary(
_conversation_id: ThreadId,
thread_state: &Arc<Mutex<ThreadState>>,
Expand Down Expand Up @@ -2104,10 +2057,8 @@ mod tests {
use codex_protocol::AgentPath;
use codex_protocol::items::DynamicToolCallItem;
use codex_protocol::items::DynamicToolCallStatus as CoreDynamicToolCallStatus;
use codex_protocol::items::HookPromptFragment;
use codex_protocol::items::SubAgentActivityItem;
use codex_protocol::items::TurnItem as CoreTurnItem;
use codex_protocol::items::build_hook_prompt_message;
use codex_protocol::models::FileSystemPermissions as CoreFileSystemPermissions;
use codex_protocol::models::NetworkPermissions as CoreNetworkPermissions;
use codex_protocol::models::PermissionProfile;
Expand Down Expand Up @@ -3993,55 +3944,4 @@ mod tests {
assert!(rx.try_recv().is_err(), "no extra messages expected");
Ok(())
}

#[tokio::test]
async fn test_hook_prompt_raw_response_emits_item_completed() -> Result<()> {
let (tx, mut rx) = mpsc::channel(CHANNEL_CAPACITY);
let outgoing = Arc::new(OutgoingMessageSender::new(
tx,
codex_analytics::AnalyticsEventsClient::disabled(),
));
let conversation_id = ThreadId::new();
let outgoing = ThreadScopedOutgoingMessageSender::new(
outgoing,
vec![ConnectionId(1)],
conversation_id,
);
let item = build_hook_prompt_message(&[
HookPromptFragment::from_single_hook("Retry with tests.", "hook-run-1"),
HookPromptFragment::from_single_hook("Then summarize cleanly.", "hook-run-2"),
])
.expect("hook prompt message");

maybe_emit_hook_prompt_item_completed(conversation_id, "turn-1", &item, &outgoing).await;

let msg = recv_broadcast_message(&mut rx).await?;
match msg {
OutgoingMessage::AppServerNotification(ServerNotification::ItemCompleted(
notification,
)) => {
assert_eq!(notification.thread_id, conversation_id.to_string());
assert_eq!(notification.turn_id, "turn-1");
assert_eq!(
notification.item,
ThreadItem::HookPrompt {
id: notification.item.id().to_string(),
fragments: vec![
codex_app_server_protocol::HookPromptFragment {
text: "Retry with tests.".into(),
hook_run_id: "hook-run-1".into(),
},
codex_app_server_protocol::HookPromptFragment {
text: "Then summarize cleanly.".into(),
hook_run_id: "hook-run-2".into(),
},
],
}
);
}
other => bail!("unexpected message: {other:?}"),
}
assert!(rx.try_recv().is_err(), "no extra messages expected");
Ok(())
}
}
1 change: 0 additions & 1 deletion codex-rs/app-server/src/request_processors.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,4 @@
use crate::bespoke_event_handling::apply_bespoke_event_handling;
use crate::bespoke_event_handling::maybe_emit_hook_prompt_item_completed;
use crate::command_exec::CommandExecManager;
use crate::command_exec::StartCommandExecParams;
use crate::config_manager::ConfigManager;
Expand Down
16 changes: 3 additions & 13 deletions codex-rs/app-server/src/request_processors/thread_lifecycle.rs
Original file line number Diff line number Diff line change
Expand Up @@ -316,6 +316,9 @@ pub(super) async fn ensure_listener_task_running(
thread_state.track_current_turn_event(&event.id, &event.msg);
thread_state.experimental_raw_events
};
if matches!(&event.msg, EventMsg::RawResponseItem(_)) && !raw_events_enabled {
continue;
}
let subscribed_connection_ids = thread_state_manager
.subscribed_connection_ids(conversation_id)
.await;
Expand All @@ -325,19 +328,6 @@ pub(super) async fn ensure_listener_task_running(
conversation_id,
);

if let EventMsg::RawResponseItem(raw_response_item_event) = &event.msg
&& !raw_events_enabled
{
maybe_emit_hook_prompt_item_completed(
conversation_id,
&event.id,
&raw_response_item_event.item,
&thread_outgoing,
)
.await;
continue;
}

apply_bespoke_event_handling(
event.clone(),
conversation_id,
Expand Down
40 changes: 40 additions & 0 deletions codex-rs/core/src/session/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,8 @@ use codex_otel::TelemetryAuthMode;
use codex_protocol::config_types::CollaborationMode;
use codex_protocol::config_types::ModeKind;
use codex_protocol::config_types::Settings;
use codex_protocol::items::HookPromptFragment;
use codex_protocol::items::build_hook_prompt_message;
use codex_protocol::models::BaseInstructions;
use codex_protocol::models::ContentItem;
use codex_protocol::models::InternalChatMessageMetadataPassthrough;
Expand Down Expand Up @@ -1763,6 +1765,44 @@ async fn record_conversation_items_stamps_missing_turn_id_and_preserves_existing
);
}

#[tokio::test]
async fn record_response_item_and_emit_turn_item_emits_hook_prompt_lifecycle() {
let (session, turn_context, rx) = make_session_and_context_with_rx().await;
let response_item = build_hook_prompt_message(&[HookPromptFragment::from_single_hook(
"Retry with tests.",
"hook-run-1",
)])
.expect("hook prompt message");
let response_item_id = response_item.id().expect("hook prompt id").to_string();

session
.record_response_item_and_emit_turn_item(&turn_context, response_item)
.await;

let raw_response = rx.recv().await.expect("raw response item event");
assert!(matches!(raw_response.msg, EventMsg::RawResponseItem(_)));

let started = rx.recv().await.expect("started hook prompt event");
assert!(matches!(
started.msg,
EventMsg::ItemStarted(ItemStartedEvent {
item: TurnItem::HookPrompt(item),
..
}) if item.id == response_item_id
));

let completed = rx.recv().await.expect("completed hook prompt event");
assert!(matches!(
completed.msg,
EventMsg::ItemCompleted(ItemCompletedEvent {
item: TurnItem::HookPrompt(item),
..
}) if item.id == response_item_id
));

assert!(rx.try_recv().is_err(), "no extra events expected");
}

#[tokio::test]
async fn record_inter_agent_communication_sets_turn_id_in_rollout_and_resume() {
let (mut session, turn_context) = make_session_and_context().await;
Expand Down
4 changes: 2 additions & 2 deletions codex-rs/core/src/session/turn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -382,9 +382,9 @@ pub(crate) async fn run_turn(
if let Some(hook_prompt_message) =
build_hook_prompt_message(&stop_outcome.continuation_fragments)
{
sess.record_conversation_items(
sess.record_response_item_and_emit_turn_item(
&turn_context,
std::slice::from_ref(&hook_prompt_message),
hook_prompt_message,
)
.await;
stop_hook_active = true;
Expand Down
Loading