Skip to content

Commit 3d25bd6

Browse files
authored
Merge pull request #5074 from loopx-project/codex/desktop-parallel-service-startup
2 parents 0131464 + 8b000a0 commit 3d25bd6

4 files changed

Lines changed: 318 additions & 103 deletions

File tree

‎apps/desktop/loopx-control-plane/src-tauri/src/lib.rs‎

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -229,6 +229,11 @@ mod tests {
229229
assert!(style.contains("--warning: #f5a623"));
230230
assert!(script.contains("desktop_update_status"));
231231
assert!(script.contains("window.loopxBootRetrying"));
232+
// Services connect concurrently, so the phase names the loopback set
233+
// until one connection outlives its peer and can be named on its own.
234+
assert!(script.contains("正在连接本地服务"));
235+
assert!(script.contains("正在连接状态服务"));
236+
assert!(script.contains("正在连接管家对话服务"));
232237
// The first screen must offer both operator choices, not a repair path
233238
// that silently replaces the CLI runtime.
234239
assert!(html.contains("id=\"pairing-align\""));

‎apps/desktop/loopx-control-plane/src-tauri/src/maintenance.rs‎

Lines changed: 3 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -749,9 +749,10 @@ pub fn start_services(app: &AppHandle) -> Result<Option<crate::services::Service
749749
}
750750
return Err(error);
751751
}
752-
crate::services::ServiceSet::start(|kind| {
752+
crate::services::ServiceSet::start(|pending| {
753+
let service = crate::services::ServiceKind::pending_label(pending);
753754
app.state::<Maintenance>()
754-
.publish("connecting", json!({"service":kind.label()}));
755+
.publish("connecting", json!({"service":service}));
755756
})
756757
.map_err(|e| e.to_string())
757758
})

‎apps/desktop/loopx-control-plane/src-tauri/src/services.rs‎

Lines changed: 191 additions & 101 deletions
Original file line numberDiff line numberDiff line change
@@ -7,10 +7,14 @@ use std::{
77
net::{SocketAddr, TcpStream},
88
path::{Path, PathBuf},
99
process::{Command, Stdio},
10+
sync::Mutex,
1011
thread,
1112
time::{Duration, Instant},
1213
};
1314

15+
/// Every loopback service the App must reach before it opens the workspace.
16+
pub const SERVICE_KINDS: [ServiceKind; 2] = [ServiceKind::Status, ServiceKind::Chat];
17+
1418
const STARTUP_TIMEOUT: Duration = Duration::from_secs(15);
1519
const PROBE_TIMEOUT: Duration = Duration::from_millis(500);
1620
const MAX_PROBE_RESPONSE_BYTES: u64 = 1024 * 1024;
@@ -29,6 +33,17 @@ impl ServiceKind {
2933
}
3034
}
3135

36+
/// Name the services a `connecting` phase is still waiting for. One
37+
/// pending service keeps its own name so a stalled connection stays
38+
/// diagnosable on the boot page; a concurrent connect reports the loopback
39+
/// set, which the boot page renders as "local services".
40+
pub fn pending_label(pending: &[Self]) -> &'static str {
41+
match pending {
42+
[kind] => kind.label(),
43+
_ => "local",
44+
}
45+
}
46+
3247
fn port(self) -> u16 {
3348
match self {
3449
Self::Status => 8766,
@@ -118,115 +133,157 @@ pub struct ServiceSet {
118133
}
119134

120135
impl ServiceSet {
121-
pub fn start(mut progress: impl FnMut(ServiceKind)) -> Result<Self, ServiceError> {
136+
pub fn start(progress: impl Fn(&[ServiceKind]) + Sync) -> Result<Self, ServiceError> {
137+
Self::collect(connect_all(SERVICE_KINDS, connect, progress))
138+
}
139+
140+
/// Fold finished connection attempts into one owned set. Every outcome
141+
/// surrenders its child here, so a set that fails still stops the
142+
/// processes its successful peers started.
143+
fn collect(outcomes: [ServiceOutcome; SERVICE_KINDS.len()]) -> Result<Self, ServiceError> {
122144
let mut services = Self {
123145
owned: Vec::new(),
124146
healed: false,
125147
};
126-
for kind in [ServiceKind::Status, ServiceKind::Chat] {
127-
progress(kind);
128-
if let Err(error) = services.ensure(kind) {
148+
let mut failure = None;
149+
for outcome in outcomes {
150+
services.owned.extend(outcome.owned);
151+
services.healed |= outcome.healed;
152+
if let Err(error) = outcome.result {
153+
failure.get_or_insert(error);
154+
}
155+
}
156+
match failure {
157+
Some(error) => {
129158
services.stop();
130-
return Err(error);
159+
Err(error)
131160
}
161+
None => Ok(services),
132162
}
133-
Ok(services)
134163
}
135164

136-
fn ensure(&mut self, kind: ServiceKind) -> Result<(), ServiceError> {
137-
let executable = loopx_executable();
138-
let expected_runtime_identity = runtime_identity_for_executable(&executable);
139-
let stale_deadline = Instant::now() + STARTUP_TIMEOUT;
140-
loop {
141-
match probe(kind, expected_runtime_identity.as_ref()) {
142-
Probe::Matching => return Ok(()),
143-
Probe::NotReady => return Err(status_readiness_error(kind)),
144-
Probe::Foreign => {
165+
pub fn stop(&mut self) {
166+
for service in self.owned.iter_mut().rev() {
167+
service.stop();
168+
}
169+
self.owned.clear();
170+
}
171+
}
172+
173+
/// One service's connection attempt. The child this App spawned travels with
174+
/// the outcome even when the attempt failed, so `ServiceSet` can stop it
175+
/// instead of leaking a process that no longer has an owner.
176+
struct ServiceOutcome {
177+
owned: Option<OwnedService>,
178+
healed: bool,
179+
result: Result<(), ServiceError>,
180+
}
181+
182+
/// Connect every loopback service at once.
183+
///
184+
/// The services own separate ports, commands and processes, and neither reads
185+
/// the other's readiness, so the window should wait for the slowest one rather
186+
/// than their sum. A start that follows a runtime update pays that difference
187+
/// twice over: each stale listener is replaced and then warms a fresh
188+
/// interpreter before it answers a readiness probe.
189+
///
190+
/// `progress` names the services still being waited on: the whole set while
191+
/// they run together, then whichever connection outlives its peer, so a
192+
/// stalled service is still named on the boot page.
193+
fn connect_all<const N: usize>(
194+
kinds: [ServiceKind; N],
195+
connect: impl Fn(ServiceKind) -> ServiceOutcome + Sync,
196+
progress: impl Fn(&[ServiceKind]) + Sync,
197+
) -> [ServiceOutcome; N] {
198+
let pending = Mutex::new(kinds.to_vec());
199+
progress(&kinds);
200+
thread::scope(|scope| {
201+
kinds
202+
.map(|kind| {
203+
let (connect, progress, pending) = (&connect, &progress, &pending);
204+
scope.spawn(move || {
205+
let outcome = connect(kind);
206+
let remaining = {
207+
let mut pending = pending.lock().expect("pending service lock");
208+
pending.retain(|entry| *entry != kind);
209+
pending.clone()
210+
};
211+
if !remaining.is_empty() {
212+
progress(&remaining);
213+
}
214+
outcome
215+
})
216+
})
217+
.map(|handle| handle.join().expect("service connection thread"))
218+
})
219+
}
220+
221+
fn connect(kind: ServiceKind) -> ServiceOutcome {
222+
let mut owned = None;
223+
let mut healed = false;
224+
let result = connect_service(kind, &mut owned, &mut healed);
225+
ServiceOutcome {
226+
owned,
227+
healed,
228+
result,
229+
}
230+
}
231+
232+
fn connect_service(
233+
kind: ServiceKind,
234+
owned: &mut Option<OwnedService>,
235+
healed: &mut bool,
236+
) -> Result<(), ServiceError> {
237+
let executable = loopx_executable();
238+
let expected_runtime_identity = runtime_identity_for_executable(&executable);
239+
let stale_deadline = Instant::now() + STARTUP_TIMEOUT;
240+
loop {
241+
match probe(kind, expected_runtime_identity.as_ref()) {
242+
Probe::Matching => return Ok(()),
243+
Probe::NotReady => return Err(status_readiness_error(kind)),
244+
Probe::Foreign => {
245+
return Err(ServiceError(format!(
246+
"port {} is occupied by a service that is not LoopX {}",
247+
kind.port(),
248+
kind.label()
249+
)));
250+
}
251+
Probe::Stale => {
252+
// Self-heal: the port is owned by a LoopX service from a
253+
// different installed release (for example after a
254+
// `loopx update`). Terminate that stale listener and keep
255+
// waiting up to the startup timeout so a LaunchAgent-managed
256+
// service (KeepAlive + throttle) has time to restart on the
257+
// current release; unknown (Foreign) processes keep the
258+
// hard error.
259+
terminate_verified_listener(kind, &executable, kind.port())?;
260+
*healed = true;
261+
if Instant::now() >= stale_deadline {
145262
return Err(ServiceError(format!(
146-
"port {} is occupied by a service that is not LoopX {}",
147-
kind.port(),
148-
kind.label()
149-
)));
150-
}
151-
Probe::Stale => {
152-
// Self-heal: the port is owned by a LoopX service from a
153-
// different installed release (for example after a
154-
// `loopx update`). Terminate that stale listener and keep
155-
// waiting up to the startup timeout so a LaunchAgent-managed
156-
// service (KeepAlive + throttle) has time to restart on the
157-
// current release; unknown (Foreign) processes keep the
158-
// hard error.
159-
terminate_verified_listener(kind, &executable, kind.port())?;
160-
self.healed = true;
161-
if Instant::now() >= stale_deadline {
162-
return Err(ServiceError(format!(
163263
"port {} is serving LoopX {} from a different installed runtime and could not be restarted",
164264
kind.port(),
165265
kind.label()
166266
)));
167-
}
168-
thread::sleep(Duration::from_millis(200));
169-
}
170-
Probe::Unresponsive => {
171-
// A bound socket is not HTTP readiness. Give slow startup
172-
// a full grace period, then replace only a verified LoopX
173-
// listener; unknown processes still fail closed.
174-
if Instant::now() < stale_deadline {
175-
thread::sleep(Duration::from_millis(100));
176-
continue;
177-
}
178-
terminate_verified_listener(kind, &executable, kind.port())?;
179-
self.healed = true;
180-
break;
181267
}
182-
Probe::Unavailable => break,
268+
thread::sleep(Duration::from_millis(200));
183269
}
184-
}
185-
186-
if request_platform_managed_start(kind) {
187-
let deadline = Instant::now() + STARTUP_TIMEOUT;
188-
while Instant::now() < deadline {
189-
match probe(kind, expected_runtime_identity.as_ref()) {
190-
Probe::Matching => return Ok(()),
191-
Probe::NotReady => return Err(status_readiness_error(kind)),
192-
Probe::Foreign => {
193-
return Err(ServiceError(format!(
194-
"LoopX {} startup reached an unexpected service on port {}",
195-
kind.label(),
196-
kind.port()
197-
)));
198-
}
199-
Probe::Stale => {
200-
terminate_verified_listener(kind, &executable, kind.port())?;
201-
self.healed = true;
202-
request_platform_managed_start(kind);
203-
}
204-
Probe::Unavailable | Probe::Unresponsive => {}
270+
Probe::Unresponsive => {
271+
// A bound socket is not HTTP readiness. Give slow startup
272+
// a full grace period, then replace only a verified LoopX
273+
// listener; unknown processes still fail closed.
274+
if Instant::now() < stale_deadline {
275+
thread::sleep(Duration::from_millis(100));
276+
continue;
205277
}
206-
thread::sleep(Duration::from_millis(100));
278+
terminate_verified_listener(kind, &executable, kind.port())?;
279+
*healed = true;
280+
break;
207281
}
208-
return Err(ServiceError(format!(
209-
"system-managed LoopX {} did not become ready on port {}",
210-
kind.label(),
211-
kind.port()
212-
)));
282+
Probe::Unavailable => break,
213283
}
284+
}
214285

215-
let mut command = Command::new(&executable);
216-
configure_runtime_environment(&mut command);
217-
command
218-
.args(kind.command_args())
219-
.stdin(Stdio::null())
220-
.stdout(Stdio::null())
221-
.stderr(Stdio::null());
222-
let child = command.group_spawn().map_err(|error| {
223-
ServiceError(format!(
224-
"could not start LoopX {} with `{executable}`: {error}",
225-
kind.label()
226-
))
227-
})?;
228-
self.owned.push(OwnedService { child });
229-
286+
if request_platform_managed_start(kind) {
230287
let deadline = Instant::now() + STARTUP_TIMEOUT;
231288
while Instant::now() < deadline {
232289
match probe(kind, expected_runtime_identity.as_ref()) {
@@ -241,27 +298,60 @@ impl ServiceSet {
241298
}
242299
Probe::Stale => {
243300
terminate_verified_listener(kind, &executable, kind.port())?;
244-
self.healed = true;
245-
thread::sleep(Duration::from_millis(200));
246-
}
247-
Probe::Unavailable | Probe::Unresponsive => {
248-
thread::sleep(Duration::from_millis(100))
301+
*healed = true;
302+
request_platform_managed_start(kind);
249303
}
304+
Probe::Unavailable | Probe::Unresponsive => {}
250305
}
306+
thread::sleep(Duration::from_millis(100));
251307
}
252-
Err(ServiceError(format!(
253-
"LoopX {} did not become ready on port {}",
308+
return Err(ServiceError(format!(
309+
"system-managed LoopX {} did not become ready on port {}",
254310
kind.label(),
255311
kind.port()
256-
)))
312+
)));
257313
}
258314

259-
pub fn stop(&mut self) {
260-
for service in self.owned.iter_mut().rev() {
261-
service.stop();
315+
let mut command = Command::new(&executable);
316+
configure_runtime_environment(&mut command);
317+
command
318+
.args(kind.command_args())
319+
.stdin(Stdio::null())
320+
.stdout(Stdio::null())
321+
.stderr(Stdio::null());
322+
let child = command.group_spawn().map_err(|error| {
323+
ServiceError(format!(
324+
"could not start LoopX {} with `{executable}`: {error}",
325+
kind.label()
326+
))
327+
})?;
328+
*owned = Some(OwnedService { child });
329+
330+
let deadline = Instant::now() + STARTUP_TIMEOUT;
331+
while Instant::now() < deadline {
332+
match probe(kind, expected_runtime_identity.as_ref()) {
333+
Probe::Matching => return Ok(()),
334+
Probe::NotReady => return Err(status_readiness_error(kind)),
335+
Probe::Foreign => {
336+
return Err(ServiceError(format!(
337+
"LoopX {} startup reached an unexpected service on port {}",
338+
kind.label(),
339+
kind.port()
340+
)));
341+
}
342+
Probe::Stale => {
343+
terminate_verified_listener(kind, &executable, kind.port())?;
344+
*healed = true;
345+
thread::sleep(Duration::from_millis(200));
346+
}
347+
Probe::Unavailable | Probe::Unresponsive => thread::sleep(Duration::from_millis(100)),
262348
}
263-
self.owned.clear();
264349
}
350+
Err(ServiceError(format!(
351+
"LoopX {} did not become ready on port {}",
352+
kind.label(),
353+
kind.port()
354+
)))
265355
}
266356

267357
#[cfg(target_os = "macos")]

0 commit comments

Comments
 (0)