From 6f98e769a123eded7c0e2562f78d23c1910290e3 Mon Sep 17 00:00:00 2001 From: Owen Lin Date: Wed, 8 Jul 2026 12:30:51 -0700 Subject: [PATCH] feat(core): emit canonical hook prompt items --- .../src/protocol/thread_history.rs | 33 +++++- .../app-server/src/bespoke_event_handling.rs | 100 ------------------ codex-rs/app-server/src/request_processors.rs | 1 - .../request_processors/thread_lifecycle.rs | 16 +-- codex-rs/core/src/session/tests.rs | 40 +++++++ codex-rs/core/src/session/turn.rs | 4 +- 6 files changed, 77 insertions(+), 117 deletions(-) diff --git a/codex-rs/app-server-protocol/src/protocol/thread_history.rs b/codex-rs/app-server-protocol/src/protocol/thread_history.rs index 9f7ffa4447e..ca8487b792c 100644 --- a/codex-rs/app-server-protocol/src/protocol/thread_history.rs +++ b/codex-rs/app-server-protocol/src/protocol/thread_history.rs @@ -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(_) @@ -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(_) @@ -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![ diff --git a/codex-rs/app-server/src/bespoke_event_handling.rs b/codex-rs/app-server/src/bespoke_event_handling.rs index db5e505eb49..829aa3f1265 100644 --- a/codex-rs/app-server/src/bespoke_event_handling.rs +++ b/codex-rs/app-server/src/bespoke_event_handling.rs @@ -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; @@ -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, @@ -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>, @@ -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; @@ -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(()) - } } diff --git a/codex-rs/app-server/src/request_processors.rs b/codex-rs/app-server/src/request_processors.rs index 42ec2063b49..a7a1f6ce819 100644 --- a/codex-rs/app-server/src/request_processors.rs +++ b/codex-rs/app-server/src/request_processors.rs @@ -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; diff --git a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs index 82a49bae0ba..784cd51add3 100644 --- a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs +++ b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs @@ -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; @@ -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, diff --git a/codex-rs/core/src/session/tests.rs b/codex-rs/core/src/session/tests.rs index ac8f9d049e2..90f43f2bc35 100644 --- a/codex-rs/core/src/session/tests.rs +++ b/codex-rs/core/src/session/tests.rs @@ -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; @@ -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; diff --git a/codex-rs/core/src/session/turn.rs b/codex-rs/core/src/session/turn.rs index e231f6af867..1fe4b15b85e 100644 --- a/codex-rs/core/src/session/turn.rs +++ b/codex-rs/core/src/session/turn.rs @@ -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;