@@ -1162,6 +1162,145 @@ async fn user_input_does_not_preempt_after_reasoning_item() {
11621162 server. shutdown ( ) . await ;
11631163}
11641164
1165+ #[ derive( Clone , Copy ) ]
1166+ enum ConditionalInterruptCase {
1167+ CurrentTurn ,
1168+ StaleTurn ,
1169+ PendingUserInput ,
1170+ PendingMailbox ,
1171+ AbandonedRequest ,
1172+ }
1173+
1174+ #[ test_case( ConditionalInterruptCase :: CurrentTurn ; "current_turn_without_pending_input" ) ]
1175+ #[ test_case( ConditionalInterruptCase :: StaleTurn ; "stale_turn" ) ]
1176+ #[ test_case( ConditionalInterruptCase :: PendingUserInput ; "pending_user_input" ) ]
1177+ #[ test_case( ConditionalInterruptCase :: PendingMailbox ; "pending_mailbox" ) ]
1178+ #[ test_case( ConditionalInterruptCase :: AbandonedRequest ; "abandoned_request" ) ]
1179+ #[ tokio:: test( flavor = "multi_thread" , worker_threads = 2 ) ]
1180+ async fn interrupt_if_no_pending_input_checks_turn_and_queue (
1181+ case : ConditionalInterruptCase ,
1182+ ) -> anyhow:: Result < ( ) > {
1183+ const INITIAL_PROMPT : & str = "first prompt" ;
1184+ const PENDING_PROMPT : & str = "preserve this pending input" ;
1185+ let ( release_response, response_gate) = oneshot:: channel ( ) ;
1186+ let first_chunks = vec ! [
1187+ chunk( ev_response_created( "resp-1" ) ) ,
1188+ chunk( ev_reasoning_item_added( "reason-1" , & [ "thinking" ] ) ) ,
1189+ gated_chunk(
1190+ response_gate,
1191+ vec![
1192+ ev_reasoning_item( "reason-1" , & [ "thinking" ] , & [ ] ) ,
1193+ ev_completed( "resp-1" ) ,
1194+ ] ,
1195+ ) ,
1196+ ] ;
1197+ let ( server, _completions) =
1198+ start_streaming_sse_server ( vec ! [ first_chunks, response_completed_chunks( "resp-2" ) ] ) . await ;
1199+ let config_server = responses:: start_mock_server ( ) . await ;
1200+ let base_url = format ! ( "{}/v1" , server. uri( ) ) ;
1201+ let test = test_codex ( )
1202+ . with_model ( "gpt-5.4" )
1203+ . with_config ( move |config| {
1204+ config. model_provider . base_url = Some ( base_url) ;
1205+ let _ = config. features . disable ( Feature :: EnableRequestCompression ) ;
1206+ } )
1207+ . build_with_auto_env ( & config_server)
1208+ . await ?;
1209+ let codex = & test. codex ;
1210+ let TurnInputSubmission :: Started { turn_id } = codex
1211+ . start_or_steer_turn ( TurnInputRequest :: user_input ( vec ! [ UserInput :: Text {
1212+ text: INITIAL_PROMPT . to_string( ) ,
1213+ text_elements: Vec :: new( ) ,
1214+ } ] ) )
1215+ . await ?
1216+ else {
1217+ panic ! ( "initial input should start a turn" ) ;
1218+ } ;
1219+ wait_for_reasoning_item_started ( codex) . await ;
1220+ if matches ! ( case, ConditionalInterruptCase :: PendingUserInput ) {
1221+ steer_user_input ( codex, PENDING_PROMPT ) . await ;
1222+ }
1223+ if matches ! ( case, ConditionalInterruptCase :: PendingMailbox ) {
1224+ submit_queue_only_agent_mail ( codex, PENDING_PROMPT ) . await ;
1225+ }
1226+ let expected_turn_id = match case {
1227+ ConditionalInterruptCase :: StaleTurn | ConditionalInterruptCase :: AbandonedRequest => {
1228+ format ! ( "stale-{turn_id}" )
1229+ }
1230+ ConditionalInterruptCase :: CurrentTurn
1231+ | ConditionalInterruptCase :: PendingUserInput
1232+ | ConditionalInterruptCase :: PendingMailbox => turn_id. clone ( ) ,
1233+ } ;
1234+ if matches ! ( case, ConditionalInterruptCase :: AbandonedRequest ) {
1235+ let ( reply, result) = oneshot:: channel ( ) ;
1236+ drop ( result) ;
1237+ codex
1238+ . submit ( Op :: InterruptIfNoPendingInput { turn_id, reply } )
1239+ . await ?;
1240+ // Wait for the abandoned request to be handled before releasing the response.
1241+ let ( reply, result) = oneshot:: channel ( ) ;
1242+ codex
1243+ . submit ( Op :: InterruptIfNoPendingInput {
1244+ turn_id : expected_turn_id. clone ( ) ,
1245+ reply,
1246+ } )
1247+ . await ?;
1248+ assert ! ( !tokio:: time:: timeout( std:: time:: Duration :: from_secs( /*secs*/ 10 ) , result) . await ??) ;
1249+ }
1250+ let should_abort = matches ! ( case, ConditionalInterruptCase :: CurrentTurn ) ;
1251+ assert_eq ! (
1252+ tokio:: time:: timeout(
1253+ std:: time:: Duration :: from_secs( /*secs*/ 10 ) ,
1254+ codex. interrupt_if_no_pending_input( & expected_turn_id) ,
1255+ )
1256+ . await ??,
1257+ should_abort,
1258+ ) ;
1259+
1260+ if should_abort {
1261+ wait_for_event ( codex, |event| {
1262+ assert ! ( !matches!( event, EventMsg :: TurnComplete ( _) ) ) ;
1263+ matches ! ( event, EventMsg :: TurnAborted ( _) )
1264+ } )
1265+ . await ;
1266+ let _ = release_response. send ( ( ) ) ;
1267+ } else {
1268+ release_response. send ( ( ) ) . expect ( "release model response" ) ;
1269+ wait_for_event ( codex, |event| {
1270+ assert ! ( !matches!( event, EventMsg :: TurnAborted ( _) ) ) ;
1271+ matches ! ( event, EventMsg :: TurnComplete ( _) )
1272+ } )
1273+ . await ;
1274+ }
1275+ let requests = server. requests ( ) . await ;
1276+ if matches ! ( case, ConditionalInterruptCase :: PendingUserInput ) {
1277+ assert_eq ! ( requests. len( ) , 2 ) ;
1278+ let second: Value = from_slice ( & requests[ 1 ] ) ?;
1279+ let prompts = message_input_texts ( & second, "user" )
1280+ . into_iter ( )
1281+ . filter ( |text| text == INITIAL_PROMPT || text == PENDING_PROMPT )
1282+ . collect :: < Vec < _ > > ( ) ;
1283+ assert_eq ! ( prompts, vec![ INITIAL_PROMPT , PENDING_PROMPT ] ) ;
1284+ } else if matches ! ( case, ConditionalInterruptCase :: PendingMailbox ) {
1285+ assert_eq ! ( requests. len( ) , 2 ) ;
1286+ let second: Value = from_slice ( & requests[ 1 ] ) ?;
1287+ let mail = second[ "input" ]
1288+ . as_array ( )
1289+ . expect ( "model input" )
1290+ . iter ( )
1291+ . find ( |item| item[ "type" ] == "agent_message" )
1292+ . expect ( "pending mailbox input" ) ;
1293+ assert_eq ! (
1294+ mail[ "content" ] ,
1295+ json!( [ { "type" : "input_text" , "text" : PENDING_PROMPT } ] )
1296+ ) ;
1297+ } else {
1298+ assert_eq ! ( requests. len( ) , 1 ) ;
1299+ }
1300+ server. shutdown ( ) . await ;
1301+ Ok ( ( ) )
1302+ }
1303+
11651304#[ derive( Clone , Copy , PartialEq , Eq ) ]
11661305enum CompactionFailurePoint {
11671306 PreTurn ,
0 commit comments