diff --git a/apps/desktop/loopx-control-plane/src-tauri/src/lib.rs b/apps/desktop/loopx-control-plane/src-tauri/src/lib.rs index c0e1c09f1a..beea153f33 100644 --- a/apps/desktop/loopx-control-plane/src-tauri/src/lib.rs +++ b/apps/desktop/loopx-control-plane/src-tauri/src/lib.rs @@ -229,6 +229,11 @@ mod tests { assert!(style.contains("--warning: #f5a623")); assert!(script.contains("desktop_update_status")); assert!(script.contains("window.loopxBootRetrying")); + // Services connect concurrently, so the phase names the loopback set + // until one connection outlives its peer and can be named on its own. + assert!(script.contains("正在连接本地服务")); + assert!(script.contains("正在连接状态服务")); + assert!(script.contains("正在连接管家对话服务")); // The first screen must offer both operator choices, not a repair path // that silently replaces the CLI runtime. assert!(html.contains("id=\"pairing-align\"")); diff --git a/apps/desktop/loopx-control-plane/src-tauri/src/maintenance.rs b/apps/desktop/loopx-control-plane/src-tauri/src/maintenance.rs index 827124bab3..34828b8f36 100644 --- a/apps/desktop/loopx-control-plane/src-tauri/src/maintenance.rs +++ b/apps/desktop/loopx-control-plane/src-tauri/src/maintenance.rs @@ -749,9 +749,10 @@ pub fn start_services(app: &AppHandle) -> Result() - .publish("connecting", json!({"service":kind.label()})); + .publish("connecting", json!({"service":service})); }) .map_err(|e| e.to_string()) }) diff --git a/apps/desktop/loopx-control-plane/src-tauri/src/services.rs b/apps/desktop/loopx-control-plane/src-tauri/src/services.rs index d3c26d3552..a187b027c5 100644 --- a/apps/desktop/loopx-control-plane/src-tauri/src/services.rs +++ b/apps/desktop/loopx-control-plane/src-tauri/src/services.rs @@ -7,10 +7,14 @@ use std::{ net::{SocketAddr, TcpStream}, path::{Path, PathBuf}, process::{Command, Stdio}, + sync::Mutex, thread, time::{Duration, Instant}, }; +/// Every loopback service the App must reach before it opens the workspace. +pub const SERVICE_KINDS: [ServiceKind; 2] = [ServiceKind::Status, ServiceKind::Chat]; + const STARTUP_TIMEOUT: Duration = Duration::from_secs(15); const PROBE_TIMEOUT: Duration = Duration::from_millis(500); const MAX_PROBE_RESPONSE_BYTES: u64 = 1024 * 1024; @@ -29,6 +33,17 @@ impl ServiceKind { } } + /// Name the services a `connecting` phase is still waiting for. One + /// pending service keeps its own name so a stalled connection stays + /// diagnosable on the boot page; a concurrent connect reports the loopback + /// set, which the boot page renders as "local services". + pub fn pending_label(pending: &[Self]) -> &'static str { + match pending { + [kind] => kind.label(), + _ => "local", + } + } + fn port(self) -> u16 { match self { Self::Status => 8766, @@ -118,115 +133,157 @@ pub struct ServiceSet { } impl ServiceSet { - pub fn start(mut progress: impl FnMut(ServiceKind)) -> Result { + pub fn start(progress: impl Fn(&[ServiceKind]) + Sync) -> Result { + Self::collect(connect_all(SERVICE_KINDS, connect, progress)) + } + + /// Fold finished connection attempts into one owned set. Every outcome + /// surrenders its child here, so a set that fails still stops the + /// processes its successful peers started. + fn collect(outcomes: [ServiceOutcome; SERVICE_KINDS.len()]) -> Result { let mut services = Self { owned: Vec::new(), healed: false, }; - for kind in [ServiceKind::Status, ServiceKind::Chat] { - progress(kind); - if let Err(error) = services.ensure(kind) { + let mut failure = None; + for outcome in outcomes { + services.owned.extend(outcome.owned); + services.healed |= outcome.healed; + if let Err(error) = outcome.result { + failure.get_or_insert(error); + } + } + match failure { + Some(error) => { services.stop(); - return Err(error); + Err(error) } + None => Ok(services), } - Ok(services) } - fn ensure(&mut self, kind: ServiceKind) -> Result<(), ServiceError> { - let executable = loopx_executable(); - let expected_runtime_identity = runtime_identity_for_executable(&executable); - let stale_deadline = Instant::now() + STARTUP_TIMEOUT; - loop { - match probe(kind, expected_runtime_identity.as_ref()) { - Probe::Matching => return Ok(()), - Probe::NotReady => return Err(status_readiness_error(kind)), - Probe::Foreign => { + pub fn stop(&mut self) { + for service in self.owned.iter_mut().rev() { + service.stop(); + } + self.owned.clear(); + } +} + +/// One service's connection attempt. The child this App spawned travels with +/// the outcome even when the attempt failed, so `ServiceSet` can stop it +/// instead of leaking a process that no longer has an owner. +struct ServiceOutcome { + owned: Option, + healed: bool, + result: Result<(), ServiceError>, +} + +/// Connect every loopback service at once. +/// +/// The services own separate ports, commands and processes, and neither reads +/// the other's readiness, so the window should wait for the slowest one rather +/// than their sum. A start that follows a runtime update pays that difference +/// twice over: each stale listener is replaced and then warms a fresh +/// interpreter before it answers a readiness probe. +/// +/// `progress` names the services still being waited on: the whole set while +/// they run together, then whichever connection outlives its peer, so a +/// stalled service is still named on the boot page. +fn connect_all( + kinds: [ServiceKind; N], + connect: impl Fn(ServiceKind) -> ServiceOutcome + Sync, + progress: impl Fn(&[ServiceKind]) + Sync, +) -> [ServiceOutcome; N] { + let pending = Mutex::new(kinds.to_vec()); + progress(&kinds); + thread::scope(|scope| { + kinds + .map(|kind| { + let (connect, progress, pending) = (&connect, &progress, &pending); + scope.spawn(move || { + let outcome = connect(kind); + let remaining = { + let mut pending = pending.lock().expect("pending service lock"); + pending.retain(|entry| *entry != kind); + pending.clone() + }; + if !remaining.is_empty() { + progress(&remaining); + } + outcome + }) + }) + .map(|handle| handle.join().expect("service connection thread")) + }) +} + +fn connect(kind: ServiceKind) -> ServiceOutcome { + let mut owned = None; + let mut healed = false; + let result = connect_service(kind, &mut owned, &mut healed); + ServiceOutcome { + owned, + healed, + result, + } +} + +fn connect_service( + kind: ServiceKind, + owned: &mut Option, + healed: &mut bool, +) -> Result<(), ServiceError> { + let executable = loopx_executable(); + let expected_runtime_identity = runtime_identity_for_executable(&executable); + let stale_deadline = Instant::now() + STARTUP_TIMEOUT; + loop { + match probe(kind, expected_runtime_identity.as_ref()) { + Probe::Matching => return Ok(()), + Probe::NotReady => return Err(status_readiness_error(kind)), + Probe::Foreign => { + return Err(ServiceError(format!( + "port {} is occupied by a service that is not LoopX {}", + kind.port(), + kind.label() + ))); + } + Probe::Stale => { + // Self-heal: the port is owned by a LoopX service from a + // different installed release (for example after a + // `loopx update`). Terminate that stale listener and keep + // waiting up to the startup timeout so a LaunchAgent-managed + // service (KeepAlive + throttle) has time to restart on the + // current release; unknown (Foreign) processes keep the + // hard error. + terminate_verified_listener(kind, &executable, kind.port())?; + *healed = true; + if Instant::now() >= stale_deadline { return Err(ServiceError(format!( - "port {} is occupied by a service that is not LoopX {}", - kind.port(), - kind.label() - ))); - } - Probe::Stale => { - // Self-heal: the port is owned by a LoopX service from a - // different installed release (for example after a - // `loopx update`). Terminate that stale listener and keep - // waiting up to the startup timeout so a LaunchAgent-managed - // service (KeepAlive + throttle) has time to restart on the - // current release; unknown (Foreign) processes keep the - // hard error. - terminate_verified_listener(kind, &executable, kind.port())?; - self.healed = true; - if Instant::now() >= stale_deadline { - return Err(ServiceError(format!( "port {} is serving LoopX {} from a different installed runtime and could not be restarted", kind.port(), kind.label() ))); - } - thread::sleep(Duration::from_millis(200)); - } - Probe::Unresponsive => { - // A bound socket is not HTTP readiness. Give slow startup - // a full grace period, then replace only a verified LoopX - // listener; unknown processes still fail closed. - if Instant::now() < stale_deadline { - thread::sleep(Duration::from_millis(100)); - continue; - } - terminate_verified_listener(kind, &executable, kind.port())?; - self.healed = true; - break; } - Probe::Unavailable => break, + thread::sleep(Duration::from_millis(200)); } - } - - if request_platform_managed_start(kind) { - let deadline = Instant::now() + STARTUP_TIMEOUT; - while Instant::now() < deadline { - match probe(kind, expected_runtime_identity.as_ref()) { - Probe::Matching => return Ok(()), - Probe::NotReady => return Err(status_readiness_error(kind)), - Probe::Foreign => { - return Err(ServiceError(format!( - "LoopX {} startup reached an unexpected service on port {}", - kind.label(), - kind.port() - ))); - } - Probe::Stale => { - terminate_verified_listener(kind, &executable, kind.port())?; - self.healed = true; - request_platform_managed_start(kind); - } - Probe::Unavailable | Probe::Unresponsive => {} + Probe::Unresponsive => { + // A bound socket is not HTTP readiness. Give slow startup + // a full grace period, then replace only a verified LoopX + // listener; unknown processes still fail closed. + if Instant::now() < stale_deadline { + thread::sleep(Duration::from_millis(100)); + continue; } - thread::sleep(Duration::from_millis(100)); + terminate_verified_listener(kind, &executable, kind.port())?; + *healed = true; + break; } - return Err(ServiceError(format!( - "system-managed LoopX {} did not become ready on port {}", - kind.label(), - kind.port() - ))); + Probe::Unavailable => break, } + } - let mut command = Command::new(&executable); - configure_runtime_environment(&mut command); - command - .args(kind.command_args()) - .stdin(Stdio::null()) - .stdout(Stdio::null()) - .stderr(Stdio::null()); - let child = command.group_spawn().map_err(|error| { - ServiceError(format!( - "could not start LoopX {} with `{executable}`: {error}", - kind.label() - )) - })?; - self.owned.push(OwnedService { child }); - + if request_platform_managed_start(kind) { let deadline = Instant::now() + STARTUP_TIMEOUT; while Instant::now() < deadline { match probe(kind, expected_runtime_identity.as_ref()) { @@ -241,27 +298,60 @@ impl ServiceSet { } Probe::Stale => { terminate_verified_listener(kind, &executable, kind.port())?; - self.healed = true; - thread::sleep(Duration::from_millis(200)); - } - Probe::Unavailable | Probe::Unresponsive => { - thread::sleep(Duration::from_millis(100)) + *healed = true; + request_platform_managed_start(kind); } + Probe::Unavailable | Probe::Unresponsive => {} } + thread::sleep(Duration::from_millis(100)); } - Err(ServiceError(format!( - "LoopX {} did not become ready on port {}", + return Err(ServiceError(format!( + "system-managed LoopX {} did not become ready on port {}", kind.label(), kind.port() - ))) + ))); } - pub fn stop(&mut self) { - for service in self.owned.iter_mut().rev() { - service.stop(); + let mut command = Command::new(&executable); + configure_runtime_environment(&mut command); + command + .args(kind.command_args()) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + let child = command.group_spawn().map_err(|error| { + ServiceError(format!( + "could not start LoopX {} with `{executable}`: {error}", + kind.label() + )) + })?; + *owned = Some(OwnedService { child }); + + let deadline = Instant::now() + STARTUP_TIMEOUT; + while Instant::now() < deadline { + match probe(kind, expected_runtime_identity.as_ref()) { + Probe::Matching => return Ok(()), + Probe::NotReady => return Err(status_readiness_error(kind)), + Probe::Foreign => { + return Err(ServiceError(format!( + "LoopX {} startup reached an unexpected service on port {}", + kind.label(), + kind.port() + ))); + } + Probe::Stale => { + terminate_verified_listener(kind, &executable, kind.port())?; + *healed = true; + thread::sleep(Duration::from_millis(200)); + } + Probe::Unavailable | Probe::Unresponsive => thread::sleep(Duration::from_millis(100)), } - self.owned.clear(); } + Err(ServiceError(format!( + "LoopX {} did not become ready on port {}", + kind.label(), + kind.port() + ))) } #[cfg(target_os = "macos")] diff --git a/apps/desktop/loopx-control-plane/src-tauri/src/services_tests.rs b/apps/desktop/loopx-control-plane/src-tauri/src/services_tests.rs index e2186e8416..8bc670752a 100644 --- a/apps/desktop/loopx-control-plane/src-tauri/src/services_tests.rs +++ b/apps/desktop/loopx-control-plane/src-tauri/src/services_tests.rs @@ -424,3 +424,122 @@ fn service_supervisor_reuses_matching_replaces_stale_and_rejects_foreign() { fs::remove_dir_all(&fixture_root).expect("remove service supervisor fixture"); } + +#[test] +fn pending_service_label_names_one_service_and_the_concurrent_set() { + // A single pending service keeps its own name so a stalled connection is + // still diagnosable; a set that is connecting together has no single name. + assert_eq!(ServiceKind::pending_label(&[ServiceKind::Status]), "status"); + assert_eq!(ServiceKind::pending_label(&[ServiceKind::Chat]), "chat"); + assert_eq!(ServiceKind::pending_label(&SERVICE_KINDS), "local"); + assert_eq!(ServiceKind::pending_label(&[]), "local"); +} + +#[test] +fn service_connections_run_concurrently_and_name_the_remaining_service() { + // Concurrency is the contract, not a timing coincidence: each connection + // must be able to observe its peer in flight. A sequential implementation + // can never satisfy the peer wait, and fails on the bounded timeout + // instead of hanging the suite. + let in_flight = std::sync::Mutex::new(0usize); + let peer_arrived = std::sync::Condvar::new(); + let saw_peer = Mutex::new(Vec::new()); + let published = Mutex::new(Vec::new()); + + let outcomes = connect_all( + SERVICE_KINDS, + |kind| { + let mut count = in_flight.lock().expect("in-flight lock"); + *count += 1; + peer_arrived.notify_all(); + let mut timed_out = false; + while *count < SERVICE_KINDS.len() && !timed_out { + let (waited, timeout) = peer_arrived + .wait_timeout(count, Duration::from_secs(10)) + .expect("peer wait"); + count = waited; + timed_out = timeout.timed_out(); + } + let observed = *count >= SERVICE_KINDS.len(); + drop(count); + saw_peer + .lock() + .expect("observation lock") + .push((kind, observed)); + ServiceOutcome { + owned: None, + healed: false, + result: Ok(()), + } + }, + |pending| { + published + .lock() + .expect("published lock") + .push(pending.to_vec()) + }, + ); + + assert!(outcomes.iter().all(|outcome| outcome.result.is_ok())); + let saw_peer = saw_peer.into_inner().expect("observation lock"); + assert_eq!(saw_peer.len(), SERVICE_KINDS.len()); + assert!( + saw_peer.iter().all(|(_, observed)| *observed), + "every service must connect while its peer is still in flight: {saw_peer:?}" + ); + + // The boot page first sees the set connecting together, then the single + // service whose connection outlived its peer. A finished set publishes + // nothing, because there is no remaining service to name. + let published = published.into_inner().expect("published lock"); + assert_eq!(published.len(), 2, "{published:?}"); + assert_eq!(published[0], SERVICE_KINDS.to_vec()); + assert_eq!(published[1].len(), 1, "{published:?}"); +} + +#[cfg(unix)] +#[test] +fn a_failed_service_set_stops_the_child_its_peer_started() { + // Ownership must travel with every outcome: a peer that failed still + // leaves this App responsible for the process it already spawned. + let mut command = Command::new("sh"); + command + .args(["-c", "exec sleep 30"]) + .stdin(Stdio::null()) + .stdout(Stdio::null()) + .stderr(Stdio::null()); + let child = command.group_spawn().expect("spawn owned service fixture"); + let pid = child.id(); + + let Err(error) = ServiceSet::collect([ + ServiceOutcome { + owned: Some(OwnedService { child }), + healed: false, + result: Ok(()), + }, + ServiceOutcome { + owned: None, + healed: false, + result: Err(ServiceError("LoopX chat did not become ready".into())), + }, + ]) else { + panic!("a failed peer must fail the whole set"); + }; + assert!(error.to_string().contains("did not become ready")); + + let mut alive = true; + for _ in 0..50 { + if !Command::new("kill") + .args(["-0", &pid.to_string()]) + .stderr(Stdio::null()) + .status() + .expect("kill -0") + .success() + { + alive = false; + break; + } + thread::sleep(Duration::from_millis(20)); + } + assert!(!alive, "owned service {pid} outlived the failed set"); +}