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
221 changes: 221 additions & 0 deletions apps/codex-plus-launcher/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -161,6 +161,21 @@ async fn activate_existing_codex_app(options: &LaunchOptions) -> anyhow::Result<
let helper_port = hooks.select_helper_port(options.helper_port);
let settings = hooks.load_settings().await?;
let app_dir = hooks.resolve_app_dir(options.app_dir.as_deref(), &settings)?;
let has_pending_recovery = hooks.has_pending_remote_control_session_recoveries();
let blocking_process_ids = if has_pending_recovery {
codex_plus_core::watcher::find_session_index_cleanup_blocking_processes()
} else {
Vec::new()
};
if should_finalize_pending_remote_control_recovery(has_pending_recovery, &blocking_process_ids)
{
hooks.run_remote_control_session_recovery().await?;
} else if has_pending_recovery {
let _ = codex_plus_core::diagnostic_log::append_diagnostic_log(
"launcher.remote_control_session_finalization_deferred_existing_app",
json!({"blocking_process_ids": blocking_process_ids}),
);
}
let launch_result = hooks
.launch_codex(
&app_dir,
Expand Down Expand Up @@ -215,6 +230,13 @@ async fn activate_existing_codex_app(options: &LaunchOptions) -> anyhow::Result<
launch_result.map(|_| ())
}

fn should_finalize_pending_remote_control_recovery(
has_pending_recovery: bool,
blocking_process_ids: &[u32],
) -> bool {
has_pending_recovery && blocking_process_ids.is_empty()
}

fn log_launcher_already_running(debug_port: u16) {
let _ = codex_plus_core::diagnostic_log::append_diagnostic_log(
"launcher.already_running",
Expand Down Expand Up @@ -310,6 +332,103 @@ impl LaunchHooks for LauncherHooks {
Ok(())
}

fn has_pending_remote_control_session_recoveries(&self) -> bool {
codex_plus_core::paths::default_pending_remote_control_recovery_path().exists()
}

fn remote_control_session_recovery_is_safe_to_run(&self) -> bool {
codex_plus_core::watcher::find_session_index_cleanup_blocking_processes().is_empty()
}

async fn run_remote_control_session_recovery(&self) -> anyhow::Result<()> {
let outcomes = tokio::task::spawn_blocking(|| {
let requests = codex_plus_core::remote_control_recovery::load_pending_remote_control_recoveries(None)?;
let settings = codex_plus_core::settings::SettingsStore::default()
.load()?;
let mut outcomes = Vec::with_capacity(requests.len());
for request in requests {
let current_profile = settings
.relay_profiles
.iter()
.find(|profile| profile.id == request.profile_id);
let request_is_current = settings.active_relay_id == request.profile_id
&& current_profile.is_some_and(|profile| {
codex_plus_core::remote_control_recovery::config_generation(
profile,
&request.target_provider,
) == request.config_generation
});
if !request_is_current {
outcomes.push((
request,
codex_plus_data::ProviderSyncResult {
status: codex_plus_data::ProviderSyncStatus::Skipped,
message: "Remote Control session finalization deferred after relay profile changed".to_string(),
target_provider: String::new(),
backup_dir: None,
changed_session_files: 0,
sqlite_rows_updated: 0,
sqlite_provider_rows_updated: 0,
sqlite_user_event_rows_updated: 0,
sqlite_cwd_rows_updated: 0,
sqlite_catalog_rows_inserted: 0,
updated_workspace_roots: 0,
skipped_locked_rollout_files: Vec::new(),
encrypted_content_warning: None,
},
None,
));
continue;
}
let result = codex_plus_data::run_remote_control_session_finalization_for_thread_with_target(
None,
&request.thread_id,
&request.target_provider,
);
let completed = result.status == codex_plus_data::ProviderSyncStatus::Synced;
let completion_error = if completed {
codex_plus_core::remote_control_recovery::complete_pending_remote_control_recovery(
None,
&request.thread_id,
)
.err()
.map(|error| error.to_string())
} else {
None
};
outcomes.push((request, result, completion_error));
}
Ok::<_, anyhow::Error>(outcomes)
})
.await
.map_err(|error| anyhow::anyhow!("Remote Control session recovery task failed: {error}"))?;
match outcomes {
Ok(outcomes) => {
for (request, result, completion_error) in outcomes {
let _ = codex_plus_core::diagnostic_log::append_diagnostic_log(
"launcher.remote_control_session_finalization",
json!({
"thread_id": request.thread_id,
"profile_id": request.profile_id,
"target_provider": request.target_provider,
"config_generation": request.config_generation,
"status": result.status,
"message": result.message,
"completion_error": completion_error
}),
);
}
}
Err(error) => {
let _ = codex_plus_core::diagnostic_log::append_diagnostic_log(
"launcher.remote_control_session_finalization_failed_nonfatal",
json!({"message": error.to_string()}),
);
}
}
Ok(())
}

async fn apply_active_relay_profile(
&self,
settings: &codex_plus_core::settings::BackendSettings,
Expand Down Expand Up @@ -514,6 +633,74 @@ impl BridgeDataService for LauncherDataService {
.await
.map_err(|error| anyhow::anyhow!("thread sort keys task failed: {error}"))
}

async fn recover_remote_control_session(&self, thread_id: String) -> anyhow::Result<Value> {
let settings = codex_plus_core::settings::SettingsStore::default()
.load()
.unwrap_or_default();
let profile = settings.active_relay_profile();
if !settings.relay_profiles_enabled
|| profile.relay_mode != codex_plus_core::settings::RelayMode::Official
|| !profile.official_mix_api_key
{
return Ok(json!({
"status": "skipped",
"message": "Remote Control session recovery is disabled for the active profile"
}));
}
let home = codex_plus_core::codex_sqlite::default_codex_home_dir();
let target_provider =
codex_plus_core::model_catalog::codex_model_provider_for_relay_profile(&home, &profile);
if target_provider.trim().is_empty() || target_provider == "openai" {
return Ok(json!({
"status": "skipped",
"message": "Remote Control session recovery requires a non-openai target provider"
}));
}
let candidate_thread_id = thread_id.clone();
let candidate = tokio::task::spawn_blocking(move || {
codex_plus_data::remote_control_session_recovery_candidate_exists(
None,
&candidate_thread_id,
)
})
.await
.map_err(|error| anyhow::anyhow!("Remote Control candidate check failed: {error}"))??;
if !candidate {
return Ok(json!({
"status": "skipped",
"message": "Remote Control session recovery is waiting for a recent openai thread"
}));
}
let request = codex_plus_core::remote_control_recovery::PendingRemoteControlRecovery {
thread_id: thread_id.clone(),
profile_id: profile.id.clone(),
target_provider: target_provider.clone(),
config_generation: codex_plus_core::remote_control_recovery::config_generation(
&profile,
&target_provider,
),
created_at: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs() as i64,
};
codex_plus_core::remote_control_recovery::enqueue_pending_remote_control_recovery(
None, request,
)?;
tokio::task::spawn_blocking(move || {
serde_json::to_value(
codex_plus_data::run_remote_control_session_catalog_recovery_for_thread_with_target(
None,
&thread_id,
&target_provider,
),
)
.map_err(anyhow::Error::from)
})
.await
.map_err(|error| anyhow::anyhow!("Remote Control session recovery task failed: {error}"))?
}
}

impl LauncherDataService {
Expand Down Expand Up @@ -864,6 +1051,40 @@ mod tests {
assert!(source.contains("launcher.already_running"));
}

#[test]
fn existing_launcher_path_drains_pending_remote_control_recovery_before_activation() {
let source = include_str!("main.rs");
let start = source
.find("async fn activate_existing_codex_app")
.expect("existing launcher activation function");
let body = &source[start..];
let recovery = body
.find(
"let has_pending_recovery = hooks.has_pending_remote_control_session_recoveries()",
)
.expect("pending recovery guard");
let launch = body
.find("let launch_result = hooks")
.expect("Codex activation");

assert!(recovery < launch);
assert!(body[recovery..launch].contains("find_session_index_cleanup_blocking_processes"));
assert!(body[recovery..launch].contains("should_finalize_pending_remote_control_recovery"));
assert!(
body[recovery..launch].contains("hooks.run_remote_control_session_recovery().await?")
);
}

#[test]
fn pending_remote_control_finalization_requires_an_idle_desktop() {
assert!(should_finalize_pending_remote_control_recovery(true, &[]));
assert!(!should_finalize_pending_remote_control_recovery(false, &[]));
assert!(!should_finalize_pending_remote_control_recovery(
true,
&[42]
));
}

#[test]
fn launcher_hooks_forward_runtime_watchdogs_and_computer_use_guard_methods() {
let source = include_str!("main.rs");
Expand Down
Loading
Loading