diff --git a/codex-rs/app-server/src/message_processor_schedule_tests.rs b/codex-rs/app-server/src/message_processor_schedule_tests.rs index 8f1dca059..ec086aee2 100644 --- a/codex-rs/app-server/src/message_processor_schedule_tests.rs +++ b/codex-rs/app-server/src/message_processor_schedule_tests.rs @@ -92,6 +92,7 @@ use codex_rollout::state_db::StateDbHandle; use core_test_support::responses; use pretty_assertions::assert_eq; use std::collections::BTreeMap; +use std::collections::VecDeque; use std::future::Future; use std::path::Path; use std::sync::Arc; @@ -117,6 +118,7 @@ struct ScheduleHarness { state_db: StateDbHandle, processor: Arc, outgoing_rx: mpsc::Receiver, + pending_notifications: VecDeque, session: Arc, next_request_id: i64, } @@ -153,6 +155,7 @@ impl ScheduleHarness { state_db, processor, outgoing_rx, + pending_notifications: VecDeque::new(), session: Arc::new(ConnectionSessionState::new(ConnectionOrigin::WebSocket)), next_request_id: 1, }; @@ -250,6 +253,7 @@ impl ScheduleHarness { params: ThreadStartParams { cwd: Some(self.workspace_cwd()), ephemeral: Some(ephemeral), + auth_profile: Some(None), ..ThreadStartParams::default() }, }) @@ -302,17 +306,9 @@ impl ScheduleHarness { .await .expect("timed out waiting for response") .expect("outgoing channel closed"); - let OutgoingEnvelope::ToConnection { - connection_id, - message, - .. - } = envelope - else { + let Some(message) = message_for_test_connection(envelope) else { continue; }; - if connection_id != TEST_CONNECTION_ID { - continue; - } match message { OutgoingMessage::Response(response) if response.id == RequestId::Integer(request_id) => @@ -323,6 +319,9 @@ impl ScheduleHarness { OutgoingMessage::Error(error) if error.id == RequestId::Integer(request_id) => { panic!("request {request_id} failed: {:?}", error.error); } + OutgoingMessage::AppServerNotification(notification) => { + self.pending_notifications.push_back(notification); + } _ => { continue; } @@ -337,17 +336,9 @@ impl ScheduleHarness { .await .expect("timed out waiting for error") .expect("outgoing channel closed"); - let OutgoingEnvelope::ToConnection { - connection_id, - message, - .. - } = envelope - else { + let Some(message) = message_for_test_connection(envelope) else { continue; }; - if connection_id != TEST_CONNECTION_ID { - continue; - } match message { OutgoingMessage::Response(response) if response.id == RequestId::Integer(request_id) => @@ -360,6 +351,9 @@ impl ScheduleHarness { OutgoingMessage::Error(error) if error.id == RequestId::Integer(request_id) => { return error.error; } + OutgoingMessage::AppServerNotification(notification) => { + self.pending_notifications.push_back(notification); + } _ => { continue; } @@ -423,6 +417,43 @@ impl ScheduleHarness { response.schedule } + async fn seed_schedule_failure(&self, schedule_id: &str) -> Result<()> { + let now = Utc::now(); + let local_active_fresh_after = self + .processor + .thread_schedule_runtime + .local_active_fresh_after(now); + let claim = self + .state_db + .thread_schedules() + .claim_thread_schedule_now_with_params(codex_state::ThreadScheduleNowClaimParams { + schedule_id, + now, + lease_id: "lease-fail", + lease_duration: std::time::Duration::from_secs(300), + local_active_owner_id: Some( + self.processor + .thread_schedule_runtime + .local_active_owner_id(), + ), + local_active_fresh_after: Some(local_active_fresh_after), + }) + .await? + .expect("schedule should claim for seeded failure"); + self.state_db + .thread_schedules() + .fail_thread_schedule_run( + schedule_id, + claim.run.run_id.as_str(), + claim.run.lease_id.as_str(), + now, + /*next_run_at*/ None, + "model unavailable".to_string(), + ) + .await?; + Ok(()) + } + async fn read_schedule_deleted(&mut self, thread_id: &str, schedule_id: &str) { loop { let notification = self.read_server_notification().await; @@ -706,6 +737,9 @@ impl ScheduleHarness { } async fn read_server_notification(&mut self) -> ServerNotification { + if let Some(notification) = self.pending_notifications.pop_front() { + return notification; + } loop { let envelope = tokio::time::timeout( std::time::Duration::from_secs(/*secs*/ 20), @@ -714,18 +748,8 @@ impl ScheduleHarness { .await .expect("timed out waiting for server notification") .expect("outgoing channel closed"); - let message = match envelope { - OutgoingEnvelope::ToConnection { - connection_id, - message, - .. - } => { - if connection_id != TEST_CONNECTION_ID { - continue; - } - message - } - OutgoingEnvelope::Broadcast { message } => message, + let Some(message) = message_for_test_connection(envelope) else { + continue; }; if let OutgoingMessage::AppServerNotification(notification) = message { return notification; @@ -734,6 +758,17 @@ impl ScheduleHarness { } } +fn message_for_test_connection(envelope: OutgoingEnvelope) -> Option { + match envelope { + OutgoingEnvelope::ToConnection { + connection_id, + message, + .. + } => (connection_id == TEST_CONNECTION_ID).then_some(message), + OutgoingEnvelope::Broadcast { message } => Some(message), + } +} + async fn create_mock_responses_server_unauthorized() -> MockServer { let server = MockServer::start().await; Mock::given(method("POST")) @@ -759,7 +794,10 @@ where .name("schedule-harness".to_string()) .stack_size(16 * 1024 * 1024) .spawn(|| { - tokio::runtime::Builder::new_current_thread() + // Message processing spawns runtime work that must remain runnable + // while Windows filesystem and SQLite operations block the caller. + tokio::runtime::Builder::new_multi_thread() + .worker_threads(2) .enable_all() .build() .expect("schedule harness runtime should build") @@ -1191,28 +1229,8 @@ fn thread_schedule_resume_recomputes_recurring_without_next_run_at() -> Result<( .await; harness.read_schedule_updated(&thread_id).await; - let claim = harness - .state_db - .thread_schedules() - .claim_thread_schedule_now( - create_response.schedule.schedule_id.as_str(), - Utc::now(), - "lease-fail", - std::time::Duration::from_secs(300), - ) - .await? - .expect("schedule should claim for seeded failure"); harness - .state_db - .thread_schedules() - .fail_thread_schedule_run( - create_response.schedule.schedule_id.as_str(), - claim.run.run_id.as_str(), - "lease-fail", - Utc::now(), - /*next_run_at*/ None, - "model unavailable".to_string(), - ) + .seed_schedule_failure(create_response.schedule.schedule_id.as_str()) .await?; let failed_schedule = harness .state_db @@ -1281,28 +1299,8 @@ fn thread_schedule_update_to_active_resets_failure_count() -> Result<()> { .await; harness.read_schedule_updated(&thread_id).await; - let claim = harness - .state_db - .thread_schedules() - .claim_thread_schedule_now( - create_response.schedule.schedule_id.as_str(), - Utc::now(), - "lease-fail", - std::time::Duration::from_secs(300), - ) - .await? - .expect("schedule should claim for seeded failure"); harness - .state_db - .thread_schedules() - .fail_thread_schedule_run( - create_response.schedule.schedule_id.as_str(), - claim.run.run_id.as_str(), - "lease-fail", - Utc::now(), - /*next_run_at*/ None, - "model unavailable".to_string(), - ) + .seed_schedule_failure(create_response.schedule.schedule_id.as_str()) .await?; let failed_schedule = harness .state_db @@ -1788,13 +1786,18 @@ fn thread_schedule_create_nests_loops_to_depth_five() -> Result<()> { let thread_id = thread.thread.id.clone(); let root = harness - .create_interval_thread_schedule(&thread_id, "root loop", 1, None) + .create_interval_thread_schedule( + &thread_id, + "root loop", + /*amount_minutes*/ 1, + /*parent_schedule_id*/ None, + ) .await; let level_2 = harness .create_interval_thread_schedule( &thread_id, "level 2 loop", - 2, + /*amount_minutes*/ 2, Some(root.schedule_id.clone()), ) .await; @@ -1802,7 +1805,7 @@ fn thread_schedule_create_nests_loops_to_depth_five() -> Result<()> { .create_interval_thread_schedule( &thread_id, "branch level 2 loop", - 3, + /*amount_minutes*/ 3, Some(root.schedule_id.clone()), ) .await; @@ -1810,7 +1813,7 @@ fn thread_schedule_create_nests_loops_to_depth_five() -> Result<()> { .create_interval_thread_schedule( &thread_id, "level 3 loop", - 3, + /*amount_minutes*/ 3, Some(level_2.schedule_id.clone()), ) .await; @@ -1818,7 +1821,7 @@ fn thread_schedule_create_nests_loops_to_depth_five() -> Result<()> { .create_interval_thread_schedule( &thread_id, "level 4 loop", - 4, + /*amount_minutes*/ 4, Some(level_3.schedule_id.clone()), ) .await; @@ -1826,7 +1829,7 @@ fn thread_schedule_create_nests_loops_to_depth_five() -> Result<()> { .create_interval_thread_schedule( &thread_id, "level 5 loop", - 5, + /*amount_minutes*/ 5, Some(level_4.schedule_id.clone()), ) .await; @@ -1889,13 +1892,18 @@ fn thread_schedule_delete_parent_emits_descendant_delete_notifications() -> Resu let thread = harness.start_materialized_thread().await; let thread_id = thread.thread.id.clone(); let root = harness - .create_interval_thread_schedule(&thread_id, "root loop", 1, None) + .create_interval_thread_schedule( + &thread_id, + "root loop", + /*amount_minutes*/ 1, + /*parent_schedule_id*/ None, + ) .await; let child = harness .create_interval_thread_schedule( &thread_id, "child loop", - 2, + /*amount_minutes*/ 2, Some(root.schedule_id.clone()), ) .await; @@ -1903,7 +1911,7 @@ fn thread_schedule_delete_parent_emits_descendant_delete_notifications() -> Resu .create_interval_thread_schedule( &thread_id, "grandchild loop", - 3, + /*amount_minutes*/ 3, Some(child.schedule_id.clone()), ) .await; @@ -2300,7 +2308,7 @@ fn schedule_create_materializes_fresh_thread_rollout_before_first_user_turn() -> let rollout_path = codex_rollout::find_thread_path_by_id_str( harness._codex_home.path(), &thread_id, - Option::<&codex_state::StateRuntime>::None, + /*state_db_ctx*/ Option::<&codex_state::StateRuntime>::None, ) .await? .expect("fresh scheduled thread should have a materialized rollout"); diff --git a/codex-rs/app-server/src/request_processors/active_session_processor.rs b/codex-rs/app-server/src/request_processors/active_session_processor.rs index 45d3a2353..914955582 100644 --- a/codex-rs/app-server/src/request_processors/active_session_processor.rs +++ b/codex-rs/app-server/src/request_processors/active_session_processor.rs @@ -571,7 +571,10 @@ mod tests { assert_eq!(named_api_peer.auth_profile.as_deref(), Some("work")); assert_eq!(named_api_peer.auth_profile_kind, AuthProfileKind::Named); - let default_api_peer = api_active_session_peer(test_active_peer(ThreadId::new(), None)); + let default_api_peer = api_active_session_peer(test_active_peer( + ThreadId::new(), + /*auth_profile*/ None, + )); assert_eq!(default_api_peer.auth_profile, None); assert_eq!(default_api_peer.auth_profile_kind, AuthProfileKind::Default); @@ -682,7 +685,7 @@ mod tests { auth_profile, process: None, capabilities: ActivePeerCapabilities::codewith_session(), - last_seen_at: LastSeenAt::from_unix_seconds(100), + last_seen_at: LastSeenAt::from_unix_seconds(/*seconds*/ 100), } } } diff --git a/codex-rs/app-server/src/request_processors/background_agent_live.rs b/codex-rs/app-server/src/request_processors/background_agent_live.rs index ae061bd73..a9a796cb0 100644 --- a/codex-rs/app-server/src/request_processors/background_agent_live.rs +++ b/codex-rs/app-server/src/request_processors/background_agent_live.rs @@ -2783,14 +2783,18 @@ async fn load_all_managed_worktrees( let page = state_db .managed_worktrees() .list_managed_worktrees_page( - Some(base_repo_path), + /*base_repo_path*/ None, /*include_deleted*/ false, cursor.as_deref(), codex_state::MAX_MANAGED_WORKTREE_LIST_LIMIT, ) .await .map_err(|err| internal_error(format!("failed to list managed worktrees: {err}")))?; - worktrees.extend(page.data); + worktrees.extend( + page.data + .into_iter() + .filter(|worktree| worktree_matches_base_repo(worktree, base_repo_path)), + ); let Some(next_cursor) = page.next_cursor else { return Ok(worktrees); }; @@ -5472,7 +5476,8 @@ mod tests { // context_length_exceeded): it must be reported as a failure, and the // recorded error must be consumed. let mut turn_error = Some("context_length_exceeded".to_string()); - let reason = turn_completion_failure_reason(None, &mut turn_error); + let reason = + turn_completion_failure_reason(/*last_agent_message*/ None, &mut turn_error); assert_eq!(reason.as_deref(), Some("context_length_exceeded")); assert_eq!(turn_error, None); } @@ -5482,7 +5487,8 @@ mod tests { // No error observed and no final message (e.g. a tool-only turn) is a // legitimate completion, not a failure. let mut turn_error = None; - let reason = turn_completion_failure_reason(None, &mut turn_error); + let reason = + turn_completion_failure_reason(/*last_agent_message*/ None, &mut turn_error); assert_eq!(reason, None); assert_eq!(turn_error, None); } @@ -5775,6 +5781,12 @@ done if !worker_process_group_exists(pgid)? { return Ok(()); } + // A dead descendant can remain as an orphaned zombie when the host's + // init process does not reap promptly. It cannot execute work or hold + // resources owned by the fixture, so the process group is drained. + if !worker_process_group_has_live_members(pgid)? { + return Ok(()); + } // PID/PGID reuse: a live process now leads group `pgid`. Our leader // was already reaped, so this must be a recycled id, not our group. if worker_process_group_id(pgid)? == Some(pgid) { @@ -5795,6 +5807,28 @@ done .success()) } + fn worker_process_group_has_live_members(pgid: u32) -> std::io::Result { + let output = std::process::Command::new("ps") + .args(["-eo", "pgid=,stat="]) + .stderr(std::process::Stdio::null()) + .output()?; + if !output.status.success() { + return Err(std::io::Error::other(format!( + "ps failed while inspecting process group {pgid}: {}", + output.status + ))); + } + Ok(String::from_utf8_lossy(&output.stdout) + .lines() + .filter_map(|line| { + let mut fields = line.split_whitespace(); + let process_group_id = fields.next()?.parse::().ok()?; + let state = fields.next()?; + Some((process_group_id, state)) + }) + .any(|(process_group_id, state)| process_group_id == pgid && !state.starts_with('Z'))) + } + /// Returns the process-group id of the live process `pid`, or `None` when no /// such process exists. Used to detect PID/PGID reuse in the group-reap wait. fn worker_process_group_id(pid: u32) -> std::io::Result> { @@ -6291,7 +6325,7 @@ done .list_background_agent_events_after( "initial-prompt-receipt", /*after_seq*/ None, - None, + /*limit*/ None, ) .await?; assert_eq!( @@ -6340,7 +6374,9 @@ done assert_eq!(goal.objective, "Investigate flaky test"); assert_eq!(goal.status, codex_state::ThreadGoalStatus::Active); let events = state_db - .list_background_agent_events_after("goal-run", /*after_seq*/ None, None) + .list_background_agent_events_after( + "goal-run", /*after_seq*/ None, /*limit*/ None, + ) .await?; let thread_id_string = thread_id.to_string(); let initial_goal_events = events @@ -6374,7 +6410,9 @@ done .await?; let events = state_db - .list_background_agent_events_after("goal-run", /*after_seq*/ None, None) + .list_background_agent_events_after( + "goal-run", /*after_seq*/ None, /*limit*/ None, + ) .await?; let initial_goal_event_count = events .iter() diff --git a/codex-rs/app-server/src/request_processors/background_agent_processor.rs b/codex-rs/app-server/src/request_processors/background_agent_processor.rs index 54e24a08a..eec9982cb 100644 --- a/codex-rs/app-server/src/request_processors/background_agent_processor.rs +++ b/codex-rs/app-server/src/request_processors/background_agent_processor.rs @@ -763,7 +763,8 @@ impl BackgroundAgentRequestProcessor { let state_db = self.state_db()?; let limit = params .limit - .unwrap_or(codex_state::DEFAULT_MANAGED_WORKTREE_LIST_LIMIT); + .unwrap_or(codex_state::DEFAULT_MANAGED_WORKTREE_LIST_LIMIT) + .clamp(1, codex_state::MAX_MANAGED_WORKTREE_LIST_LIMIT); let include_deleted = params.include_deleted.unwrap_or(false); let base_repo_path = params .base_repo_path @@ -778,23 +779,72 @@ impl BackgroundAgentRequestProcessor { "worktree/list baseRepoPath must be absolute", )); } - let page = state_db - .managed_worktrees() - .list_managed_worktrees_page( - base_repo_path.as_deref(), - include_deleted, - params.cursor.as_deref(), - limit, + let (worktrees, next_cursor) = if let Some(base_repo_path) = base_repo_path.as_deref() { + let cursor = params.cursor.as_deref().map(str::trim).unwrap_or_default(); + let offset = if cursor.is_empty() { + 0 + } else { + cursor.parse::().map_err(|_| { + internal_error(format!( + "failed to list worktrees: invalid managed worktree list cursor `{cursor}`" + )) + })? + }; + let page_end = offset.saturating_add(limit); + let required_matches = page_end.saturating_add(1) as usize; + let mut matching_worktrees = Vec::new(); + let mut scan_cursor = None; + loop { + let page = state_db + .managed_worktrees() + .list_managed_worktrees_page( + /*base_repo_path*/ None, + include_deleted, + scan_cursor.as_deref(), + codex_state::MAX_MANAGED_WORKTREE_LIST_LIMIT, + ) + .await + .map_err(|err| internal_error(format!("failed to list worktrees: {err}")))?; + matching_worktrees.extend(page.data.into_iter().filter(|worktree| { + paths_equivalent(worktree.base_repo_path.as_path(), base_repo_path) + })); + if matching_worktrees.len() >= required_matches { + break; + } + let Some(next_cursor) = page.next_cursor else { + break; + }; + scan_cursor = Some(next_cursor); + } + let has_more = matching_worktrees.len() > page_end as usize; + ( + matching_worktrees + .into_iter() + .skip(offset as usize) + .take(limit as usize) + .collect(), + has_more.then(|| page_end.to_string()), ) - .await - .map_err(|err| internal_error(format!("failed to list worktrees: {err}")))?; - let mut data = Vec::with_capacity(page.data.len()); - for worktree in page.data { + } else { + let page = state_db + .managed_worktrees() + .list_managed_worktrees_page( + /*base_repo_path*/ None, + include_deleted, + params.cursor.as_deref(), + limit, + ) + .await + .map_err(|err| internal_error(format!("failed to list worktrees: {err}")))?; + (page.data, page.next_cursor) + }; + let mut data = Vec::with_capacity(worktrees.len()); + for worktree in worktrees { data.push(api_worktree_from_state(state_db.as_ref(), worktree).await?); } Ok(WorktreeListResponse { data, - next_cursor: page.next_cursor, + next_cursor, policy, }) } diff --git a/codex-rs/app-server/src/request_processors/local_session_directory.rs b/codex-rs/app-server/src/request_processors/local_session_directory.rs index 3a4c1038d..229de95d1 100644 --- a/codex-rs/app-server/src/request_processors/local_session_directory.rs +++ b/codex-rs/app-server/src/request_processors/local_session_directory.rs @@ -510,8 +510,8 @@ mod tests { let local_session = api_local_session( thread, - None, - None, + /*model*/ None, + /*thread_agent_path*/ None, &live_overlay, &auth_profile_account_labels, &HashSet::new(), 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 e10970dc9..3ac34b588 100644 --- a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs +++ b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs @@ -472,6 +472,14 @@ async fn heartbeat_local_active_session( let Some(state_db) = listener_task_context.state_db.as_ref() else { return; }; + if let Err(err) = + materialize_local_active_session_thread(state_db, conversation_id, conversation).await + { + tracing::warn!( + thread_id = %conversation_id, + "failed to materialize thread metadata before local active-session heartbeat: {err}" + ); + } let session_id = conversation.session_configured().session_id.to_string(); if let Err(err) = state_db .local_active_sessions() @@ -491,6 +499,41 @@ async fn heartbeat_local_active_session( } } +async fn materialize_local_active_session_thread( + state_db: &StateDbHandle, + conversation_id: ThreadId, + conversation: &CodexThread, +) -> anyhow::Result<()> { + if state_db.get_thread(conversation_id).await?.is_some() { + return Ok(()); + } + let config_snapshot = conversation.config_snapshot().await; + if config_snapshot.ephemeral { + return Ok(()); + } + let Some(rollout_path) = conversation.rollout_path() else { + return Ok(()); + }; + let mut builder = ThreadMetadataBuilder::new( + conversation_id, + rollout_path, + Utc::now(), + config_snapshot.session_source.clone(), + ); + builder.thread_source = config_snapshot.thread_source; + builder.agent_nickname = config_snapshot.session_source.get_nickname(); + builder.agent_role = config_snapshot.session_source.get_agent_role(); + builder.model_provider = Some(config_snapshot.model_provider_id.clone()); + builder.cwd = config_snapshot.cwd.to_path_buf(); + builder.cli_version = Some(env!("CARGO_PKG_VERSION").to_string()); + builder.approval_mode = config_snapshot.approval_policy; + let mut metadata = builder.build(config_snapshot.model_provider_id.as_str()); + metadata.model = Some(config_snapshot.model); + metadata.sandbox_policy = serde_json::to_string(&config_snapshot.permission_profile)?; + state_db.insert_thread_if_absent(&metadata).await?; + Ok(()) +} + pub(super) async fn wait_for_thread_shutdown(thread: &Arc) -> ThreadShutdownResult { match tokio::time::timeout(Duration::from_secs(10), thread.shutdown_and_wait()).await { Ok(Ok(())) => ThreadShutdownResult::Complete, diff --git a/codex-rs/app-server/src/request_processors/thread_monitor_runtime.rs b/codex-rs/app-server/src/request_processors/thread_monitor_runtime.rs index 753f18138..5df05078f 100644 --- a/codex-rs/app-server/src/request_processors/thread_monitor_runtime.rs +++ b/codex-rs/app-server/src/request_processors/thread_monitor_runtime.rs @@ -823,7 +823,7 @@ mod tests { state_db .upsert_thread(&builder.build("test-provider")) .await?; - let monitor = test_monitor_for_thread(thread_id, None); + let monitor = test_monitor_for_thread(thread_id, /*cwd*/ None); assert_eq!(monitor_thread_cwd(&state_db, &monitor).await?, thread_cwd); assert_ne!( @@ -854,7 +854,7 @@ mod tests { #[test] fn monitor_output_communication_uses_wake_if_idle_mailbox_shape() { - let monitor = test_monitor(None); + let monitor = test_monitor(/*cwd*/ None); let communication = monitor_output_communication(&monitor, "new conversations message".to_string()); diff --git a/codex-rs/app-server/src/request_processors/thread_processor.rs b/codex-rs/app-server/src/request_processors/thread_processor.rs index 86906ad54..2cc894ea1 100644 --- a/codex-rs/app-server/src/request_processors/thread_processor.rs +++ b/codex-rs/app-server/src/request_processors/thread_processor.rs @@ -1700,6 +1700,15 @@ impl ThreadRequestProcessor { { Ok(thread) => { if thread.archived_at.is_none() { + if let Some(rollout_path) = thread.rollout_path.as_deref() + && codex_rollout::existing_rollout_path(rollout_path) + .await + .is_none() + { + return Err(invalid_request(format!( + "no rollout found for thread id {thread_id}" + ))); + } archive_thread_ids.push(thread_id); } } @@ -1717,6 +1726,17 @@ impl ThreadRequestProcessor { { Ok(thread) => { if thread.archived_at.is_none() { + if let Some(rollout_path) = thread.rollout_path.as_deref() + && codex_rollout::existing_rollout_path(rollout_path) + .await + .is_none() + { + warn!( + "skipping unmaterialized spawned descendant thread \ + {descendant_thread_id} while archiving {thread_id}" + ); + continue; + } archive_thread_ids.push(descendant_thread_id); } } diff --git a/codex-rs/app-server/src/request_processors/thread_processor_tests.rs b/codex-rs/app-server/src/request_processors/thread_processor_tests.rs index 60c4c4669..e7e275b72 100644 --- a/codex-rs/app-server/src/request_processors/thread_processor_tests.rs +++ b/codex-rs/app-server/src/request_processors/thread_processor_tests.rs @@ -1032,7 +1032,7 @@ mod thread_processor_behavior_tests { permission_profile: PermissionProfile, ) -> RolloutItem { let RolloutItem::TurnContext(mut turn_context) = - turn_context_with_auth_profile(thread_id, None) + turn_context_with_auth_profile(thread_id, /*auth_profile*/ None) else { unreachable!("helper returns turn context") }; @@ -1046,7 +1046,7 @@ mod thread_processor_behavior_tests { workspace_roots: Vec, ) -> RolloutItem { let RolloutItem::TurnContext(mut turn_context) = - turn_context_with_auth_profile(thread_id, None) + turn_context_with_auth_profile(thread_id, /*auth_profile*/ None) else { unreachable!("helper returns turn context") }; @@ -1060,7 +1060,7 @@ mod thread_processor_behavior_tests { approval_policy: AskForApproval, ) -> RolloutItem { let RolloutItem::TurnContext(mut turn_context) = - turn_context_with_auth_profile(thread_id, None) + turn_context_with_auth_profile(thread_id, /*auth_profile*/ None) else { unreachable!("helper returns turn context") }; @@ -1320,7 +1320,11 @@ mod thread_processor_behavior_tests { }); let mut typesafe_overrides = ConfigOverrides::default(); - merge_persisted_permission_profile_from_history(&mut typesafe_overrides, None, &history); + merge_persisted_permission_profile_from_history( + &mut typesafe_overrides, + /*request_overrides*/ None, + &history, + ); assert_eq!( typesafe_overrides.permission_profile, @@ -1344,7 +1348,11 @@ mod thread_processor_behavior_tests { ..Default::default() }; - merge_persisted_permission_profile_from_history(&mut typesafe_overrides, None, &history); + merge_persisted_permission_profile_from_history( + &mut typesafe_overrides, + /*request_overrides*/ None, + &history, + ); assert_eq!(typesafe_overrides.permission_profile, None); assert_eq!( @@ -1370,7 +1378,11 @@ mod thread_processor_behavior_tests { }); let mut typesafe_overrides = ConfigOverrides::default(); - merge_persisted_approval_settings_from_history(&mut typesafe_overrides, None, &history); + merge_persisted_approval_settings_from_history( + &mut typesafe_overrides, + /*request_overrides*/ None, + &history, + ); assert_eq!( typesafe_overrides.approval_policy, @@ -1400,7 +1412,11 @@ mod thread_processor_behavior_tests { ..Default::default() }; - merge_persisted_approval_settings_from_history(&mut typesafe_overrides, None, &history); + merge_persisted_approval_settings_from_history( + &mut typesafe_overrides, + /*request_overrides*/ None, + &history, + ); assert_eq!( typesafe_overrides.approval_policy, @@ -1499,7 +1515,7 @@ mod thread_processor_behavior_tests { merge_persisted_cwd_and_workspace_roots_from_history( &mut typesafe_overrides, - None, + /*request_overrides*/ None, &history, ); diff --git a/codex-rs/app-server/src/request_processors/thread_schedule_runtime.rs b/codex-rs/app-server/src/request_processors/thread_schedule_runtime.rs index ddf0b54ae..278ca500c 100644 --- a/codex-rs/app-server/src/request_processors/thread_schedule_runtime.rs +++ b/codex-rs/app-server/src/request_processors/thread_schedule_runtime.rs @@ -2992,8 +2992,8 @@ mod tests { &schedule.schedule_id, &retry_claim.run.run_id, "lease-retry", - None, - None, + /*goal_id*/ None, + /*error*/ None, completed_at, ) .await @@ -3075,8 +3075,8 @@ mod tests { &schedule.schedule_id, &claim.run.run_id, "lease-run", - None, - None, + /*goal_id*/ None, + /*error*/ None, completed_at, ) .await @@ -3171,7 +3171,7 @@ mod tests { &claim.run.run_id, "lease-run", Some(goal.goal_id.as_str()), - None, + /*error*/ None, completed_at, ) .await @@ -3294,7 +3294,7 @@ mod tests { recovered.run_id.as_str(), recovered.lease_id.as_str(), recovered.goal_id.as_deref(), - None, + /*error*/ None, completed_at, ) .await @@ -3638,8 +3638,8 @@ mod tests { &schedule.schedule_id, &claim.run.run_id, "lease-run", - None, - None, + /*goal_id*/ None, + /*error*/ None, completed_at, ) .await @@ -3720,7 +3720,7 @@ mod tests { unit: codex_state::ThreadScheduleIntervalUnit::Minutes, }), "UTC", - None, + /*scheduled_for*/ None, at(/*seconds*/ 1_700_000_300), ) .expect("next interval should compute") diff --git a/codex-rs/app-server/src/request_processors/usage_profile_broker.rs b/codex-rs/app-server/src/request_processors/usage_profile_broker.rs index 797acca02..b3c6d257b 100644 --- a/codex-rs/app-server/src/request_processors/usage_profile_broker.rs +++ b/codex-rs/app-server/src/request_processors/usage_profile_broker.rs @@ -613,8 +613,8 @@ mod tests { chatgpt_profile("third"), ]; let health_by_profile = BTreeMap::from([ - ("second".to_string(), health(20.0)), - ("third".to_string(), health(80.0)), + ("second".to_string(), health(/*remaining_percent*/ 20.0)), + ("third".to_string(), health(/*remaining_percent*/ 80.0)), ]); let now = Instant::now(); @@ -810,7 +810,7 @@ mod tests { &BTreeMap::new(), &BTreeMap::new(), Instant::now(), - 1_000, + /*now_epoch*/ 1_000, ) ); } @@ -902,8 +902,8 @@ mod tests { #[test] fn highest_available_dispatch_selects_healthiest_non_exhausted_profile() { let health_by_profile = BTreeMap::from([ - ("second".to_string(), health(20.0)), - ("third".to_string(), health(80.0)), + ("second".to_string(), health(/*remaining_percent*/ 20.0)), + ("third".to_string(), health(/*remaining_percent*/ 80.0)), ]); assert_eq!( @@ -1009,7 +1009,7 @@ mod tests { /*trigger_window_label*/ None, /*is_fresh*/ true, ), - health(60.0) + health(/*remaining_percent*/ 60.0) ); } } diff --git a/codex-rs/app-server/tests/suite/v2/account.rs b/codex-rs/app-server/tests/suite/v2/account.rs index 0ed469b4f..e3659cd0d 100644 --- a/codex-rs/app-server/tests/suite/v2/account.rs +++ b/codex-rs/app-server/tests/suite/v2/account.rs @@ -1059,7 +1059,10 @@ async fn auth_profile_rpcs_save_list_and_switch_api_key_profiles() -> Result<()> ) .await??; let saved_first: AuthProfileSaveCurrentResponse = to_response(resp)?; - assert_eq!(saved_first.profile, api_key_profile_summary("first", true)); + assert_eq!( + saved_first.profile, + api_key_profile_summary("first", /*active*/ true) + ); assert_account_updated_notification(&mut mcp, Some(AuthMode::ApiKey)).await?; let req_id = mcp @@ -1084,7 +1087,7 @@ async fn auth_profile_rpcs_save_list_and_switch_api_key_profiles() -> Result<()> let saved_second: AuthProfileSaveCurrentResponse = to_response(resp)?; assert_eq!( saved_second.profile, - api_key_profile_summary("second", true) + api_key_profile_summary("second", /*active*/ true) ); assert_account_updated_notification(&mut mcp, Some(AuthMode::ApiKey)).await?; @@ -1100,7 +1103,7 @@ async fn auth_profile_rpcs_save_list_and_switch_api_key_profiles() -> Result<()> assert_eq!( profiles, AuthProfileListResponse { - data: vec![api_key_profile_summary("first", false)], + data: vec![api_key_profile_summary("first", /*active*/ false)], next_cursor: Some("1".to_string()), } ); @@ -1120,7 +1123,7 @@ async fn auth_profile_rpcs_save_list_and_switch_api_key_profiles() -> Result<()> assert_eq!( profiles, AuthProfileListResponse { - data: vec![api_key_profile_summary("second", true)], + data: vec![api_key_profile_summary("second", /*active*/ true)], next_cursor: None, } ); @@ -1134,7 +1137,10 @@ async fn auth_profile_rpcs_save_list_and_switch_api_key_profiles() -> Result<()> ) .await??; let switched: AuthProfileSwitchResponse = to_response(resp)?; - assert_eq!(switched.profile, api_key_profile_summary("first", true)); + assert_eq!( + switched.profile, + api_key_profile_summary("first", /*active*/ true) + ); assert_account_updated_notification(&mut mcp, Some(AuthMode::ApiKey)).await?; let list_id = mcp @@ -1150,8 +1156,8 @@ async fn auth_profile_rpcs_save_list_and_switch_api_key_profiles() -> Result<()> profiles, AuthProfileListResponse { data: vec![ - api_key_profile_summary("first", true), - api_key_profile_summary("second", false), + api_key_profile_summary("first", /*active*/ true), + api_key_profile_summary("second", /*active*/ false), ], next_cursor: None, } @@ -1193,9 +1199,13 @@ async fn auth_profile_switch_accepts_external_subscription_profiles() -> Result< assert_eq!( switched.profile, - external_profile_summary("cursor", AuthProfileSubscriptionProvider::Cursor, true) + external_profile_summary( + "cursor", + AuthProfileSubscriptionProvider::Cursor, + /*active*/ true, + ) ); - assert_account_updated_notification(&mut mcp, None).await?; + assert_account_updated_notification(&mut mcp, /*auth_mode*/ None).await?; let list_id = mcp .send_raw_request("authProfile/list", Some(json!({}))) @@ -1212,7 +1222,7 @@ async fn auth_profile_switch_accepts_external_subscription_profiles() -> Result< data: vec![external_profile_summary( "cursor", AuthProfileSubscriptionProvider::Cursor, - true, + /*active*/ true, )], next_cursor: None, } @@ -1252,8 +1262,8 @@ async fn auth_profile_list_uses_selected_profile_for_active_state() -> Result<() profiles, AuthProfileListResponse { data: vec![ - api_key_profile_summary("personal", false), - api_key_profile_summary("work", true), + api_key_profile_summary("personal", /*active*/ false), + api_key_profile_summary("work", /*active*/ true), ], next_cursor: None, } diff --git a/codex-rs/app-server/tests/suite/v2/background_agent.rs b/codex-rs/app-server/tests/suite/v2/background_agent.rs index 3fc5b06db..e6440d933 100644 --- a/codex-rs/app-server/tests/suite/v2/background_agent.rs +++ b/codex-rs/app-server/tests/suite/v2/background_agent.rs @@ -614,7 +614,7 @@ async fn agent_start_uses_validated_managed_worktree_cwd() -> Result<()> { } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn agent_start_rebinds_workspace_write_permissions_to_managed_worktree() -> Result<()> { +async fn agent_start_rebinds_effective_permissions_to_managed_worktree() -> Result<()> { let codex_home = TempDir::new()?; init_git_repo(codex_home.path())?; let server = create_mock_responses_server_sequence_unchecked(vec![ @@ -664,6 +664,16 @@ exclude_slash_tmp = true )?; let file_system_policy = permission_profile.file_system_sandbox_policy(); let worktree_path = Path::new(created_worktree_path.as_str()); + #[cfg(target_os = "windows")] + { + assert_eq!(permission_profile, PermissionProfile::read_only()); + assert!( + !file_system_policy.can_write_path_with_cwd(worktree_path, worktree_path), + "unsandboxed Windows worker must keep the effective read-only policy: \ + {file_system_policy:?}" + ); + } + #[cfg(not(target_os = "windows"))] assert!( file_system_policy.can_write_path_with_cwd(worktree_path, worktree_path), "managed worktree should be writable, policy: {file_system_policy:?}" @@ -1730,13 +1740,22 @@ async fn worktree_create_reconcile_and_cleanup_use_real_git_worktrees() -> Resul ], )?; let state_db = init_state_db(codex_home.path()).await?; + #[cfg(unix)] + let (_base_repo_alias_root, recorded_base_repo_path) = { + let alias_root = TempDir::new()?; + let alias_path = alias_root.path().join("repo"); + std::os::unix::fs::symlink(codex_home.path(), &alias_path)?; + (alias_root, alias_path) + }; + #[cfg(not(unix))] + let recorded_base_repo_path = codex_home.path().to_path_buf(); state_db .managed_worktrees() .create_managed_worktree(codex_state::ManagedWorktreeCreateParams { worktree_id: Some("outside-root".to_string()), identity: Some("test:outside-root".to_string()), mode: codex_state::ManagedWorktreeMode::IsolatedWorktree, - base_repo_path: codex_home.path().to_path_buf(), + base_repo_path: recorded_base_repo_path, worktree_path: outside_root_path.clone(), branch: Some("codewith/outside-root".to_string()), base_sha: None, diff --git a/codex-rs/app-server/tests/suite/v2/rate_limit_resets.rs b/codex-rs/app-server/tests/suite/v2/rate_limit_resets.rs index 3358e692c..114d3ca8f 100644 --- a/codex-rs/app-server/tests/suite/v2/rate_limit_resets.rs +++ b/codex-rs/app-server/tests/suite/v2/rate_limit_resets.rs @@ -532,7 +532,7 @@ async fn consume_rate_limit_reset_credit_rejects_stale_binding_after_account_a_t let mut mcp = test_app_server(codex_home.path()).await?; timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; - let account_a_fingerprint = read_account_identity(&mut mcp, None).await?; + let account_a_fingerprint = read_account_identity(&mut mcp, /*auth_profile*/ None).await?; let account_b_token = encode_id_token( &ChatGptIdTokenClaims::new() .email("account-b@example.com") @@ -611,7 +611,7 @@ async fn consume_rate_limit_reset_credit_uses_named_auth_profile_and_selected_cr let server = MockServer::start().await; write_chatgpt_base_url(codex_home.path(), &server.uri())?; - mount_usage_response(&server, None).await; + mount_usage_response(&server, /*available_count*/ None).await; Mock::given(method("POST")) .and(path("/api/codex/rate-limit-reset-credits/consume")) @@ -700,7 +700,7 @@ async fn consume_rate_limit_reset_credit_reads_root_auth_profile_when_selected_p let server = MockServer::start().await; write_chatgpt_base_url(codex_home.path(), &server.uri())?; - mount_usage_response(&server, None).await; + mount_usage_response(&server, /*available_count*/ None).await; Mock::given(method("POST")) .and(path("/api/codex/rate-limit-reset-credits/consume")) @@ -725,12 +725,13 @@ async fn consume_rate_limit_reset_credit_reads_root_auth_profile_when_selected_p ) .await?; timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; - let account_identity_fingerprint = read_account_identity(&mut mcp, Some(None)).await?; + let account_identity_fingerprint = + read_account_identity(&mut mcp, Some(/*profile*/ None)).await?; let request_id = mcp .send_consume_account_rate_limit_reset_credit_request( consume_reset_params("root-redeem") - .with_auth_profile(None) + .with_auth_profile(/*profile*/ None) .with_expected_fingerprint(account_identity_fingerprint), ) .await?; @@ -763,7 +764,7 @@ async fn consume_rate_limit_reset_credit_maps_no_credit_outcome() -> Result<()> let server = MockServer::start().await; write_chatgpt_base_url(codex_home.path(), &server.uri())?; - mount_usage_response(&server, None).await; + mount_usage_response(&server, /*available_count*/ None).await; Mock::given(method("POST")) .and(path("/api/codex/rate-limit-reset-credits/consume")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ @@ -774,7 +775,8 @@ async fn consume_rate_limit_reset_credit_maps_no_credit_outcome() -> Result<()> let mut mcp = test_app_server(codex_home.path()).await?; timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; - let account_identity_fingerprint = read_account_identity(&mut mcp, None).await?; + let account_identity_fingerprint = + read_account_identity(&mut mcp, /*auth_profile*/ None).await?; let request_id = mcp .send_consume_account_rate_limit_reset_credit_request( @@ -811,7 +813,7 @@ async fn consume_rate_limit_reset_credit_surfaces_backend_failure() -> Result<() let server = MockServer::start().await; write_chatgpt_base_url(codex_home.path(), &server.uri())?; - mount_usage_response(&server, None).await; + mount_usage_response(&server, /*available_count*/ None).await; Mock::given(method("POST")) .and(path("/api/codex/rate-limit-reset-credits/consume")) .respond_with(ResponseTemplate::new(500).set_body_string("boom")) @@ -820,7 +822,8 @@ async fn consume_rate_limit_reset_credit_surfaces_backend_failure() -> Result<() let mut mcp = test_app_server(codex_home.path()).await?; timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; - let account_identity_fingerprint = read_account_identity(&mut mcp, None).await?; + let account_identity_fingerprint = + read_account_identity(&mut mcp, /*auth_profile*/ None).await?; let request_id = mcp .send_consume_account_rate_limit_reset_credit_request( @@ -858,7 +861,7 @@ async fn consume_rate_limit_reset_credit_timeout_releases_later_request() -> Res let server = MockServer::start().await; write_chatgpt_base_url(codex_home.path(), &server.uri())?; - mount_usage_response(&server, None).await; + mount_usage_response(&server, /*available_count*/ None).await; Mock::given(method("POST")) .and(path("/api/codex/rate-limit-reset-credits/consume")) .and(wiremock::matchers::body_json(json!({ @@ -895,7 +898,8 @@ async fn consume_rate_limit_reset_credit_timeout_releases_later_request() -> Res ) .await?; timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; - let account_identity_fingerprint = read_account_identity(&mut mcp, None).await?; + let account_identity_fingerprint = + read_account_identity(&mut mcp, /*auth_profile*/ None).await?; let request_id = mcp .send_consume_account_rate_limit_reset_credit_request( diff --git a/codex-rs/app-server/tests/suite/v2/remote_dispatch.rs b/codex-rs/app-server/tests/suite/v2/remote_dispatch.rs index 33cea786d..34a17cf50 100644 --- a/codex-rs/app-server/tests/suite/v2/remote_dispatch.rs +++ b/codex-rs/app-server/tests/suite/v2/remote_dispatch.rs @@ -285,8 +285,8 @@ async fn remote_submit_error( Some(remote_submit_params( source_machine_id, target_machine_id, - None, - None, + /*capability_version*/ None, + /*expires_at*/ None, message, )), ) diff --git a/codex-rs/app-server/tests/suite/v2/thread_archive.rs b/codex-rs/app-server/tests/suite/v2/thread_archive.rs index 7bf6fc9bc..3db2bd4e8 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_archive.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_archive.rs @@ -68,6 +68,13 @@ async fn thread_archive_requires_materialized_rollout() -> Result<()> { .is_none(), "thread id should not be discoverable before rollout materialization" ); + let state_db = + StateRuntime::init(codex_home.path().to_path_buf(), "mock_provider".into()).await?; + let stored_thread = state_db + .get_thread(ThreadId::from_string(&thread.id)?) + .await? + .expect("local active-session heartbeat should persist thread metadata"); + assert_eq!(stored_thread.rollout_path, rollout_path); // Archive should fail before the rollout is materialized. let archive_id = mcp diff --git a/codex-rs/app-server/tests/suite/v2/thread_mailbox.rs b/codex-rs/app-server/tests/suite/v2/thread_mailbox.rs index 81c81f5ab..eaa879322 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_mailbox.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_mailbox.rs @@ -410,6 +410,17 @@ async fn thread_mailbox_dispatcher_does_not_steal_live_target_from_dispatch_disa let mut target_mcp = init_mcp(codex_home.path()).await?; let thread_id = start_thread(&mut target_mcp).await?; + let state_db = + StateRuntime::init(codex_home.path().to_path_buf(), "mock_provider".into()).await?; + let target_thread_id = ThreadId::from_string(&thread_id)?; + assert!( + state_db + .local_active_sessions() + .get_session(target_thread_id) + .await? + .is_some(), + "thread/start must advertise the live target before returning" + ); create_config_toml_with_mailbox_dispatcher( codex_home.path(), @@ -533,7 +544,7 @@ async fn thread_mailbox_dispatcher_resume_preserves_persisted_permissions() -> R create_config_toml_with_mailbox_dispatcher_and_sandbox( codex_home.path(), &server.uri(), - /*enabled*/ true, + /*mailbox_dispatcher_enabled*/ true, "danger-full-access", )?; @@ -780,12 +791,8 @@ async fn thread_mailbox_dispatcher_retries_and_poisons_offline_targets() -> Resu /*mailbox_dispatcher_enabled*/ true, )?; - let (retry_thread_id, poison_thread_id) = { - let mut setup_mcp = init_mcp(codex_home.path()).await?; - let retry_thread_id = start_thread(&mut setup_mcp).await?; - let poison_thread_id = start_thread(&mut setup_mcp).await?; - (retry_thread_id, poison_thread_id) - }; + let retry_thread_id = seed_unloaded_thread(codex_home.path()).await?; + let poison_thread_id = seed_unloaded_thread(codex_home.path()).await?; let mut mcp = init_mcp(codex_home.path()).await?; let retry = enqueue_auto_dispatch_message_with_max_attempts(