From 3ec30d58fa5c3a2b4df79a90bd15c8204d3a57cb Mon Sep 17 00:00:00 2001 From: Pierre Carrier Date: Wed, 30 Sep 2026 11:52:10 +0000 Subject: [PATCH 1/2] Answer a WAIT that looks too soon with the exit, not CONFLICT A WAIT on a process its session holds no route to looks at it with a temporary WATCH. That look is refused as CONFLICT while the process's exit is on its way to its watchers (the record is final only once they all have it), or while the session's own previous look, such as the one a KILL's CONTROL just made, is still leaving. So a WAIT right after a KILL of a detached process could answer CONFLICT instead of the exit: client_host's a_process_whose_attachment_went_is_waited_for_and_held_by_nothing failed so on a loaded machine. Such a WAIT now waits for the process catalogue's next change, with which either settles, and looks again, within its timeout. a_wait_right_after_a_kill_of_a_detached_process_gets_its_exit (64 rounds of spawn detachable, detach, KILL, WAIT) failed 1 run in 10 without the change, and none in 20 with it. a_watcher_whose_queue_fills_after_the_exit_fails_alone took only output in its frame loop and panicked on the stdin's progress event, which a slow machine delivers among the frames; it lets it pass now, and waits for the child's reaping (Server::wait_reaped, tests only) rather than 500 ms. --- crates/server/src/process.rs | 22 +++++++++ crates/server/src/yas_process.rs | 77 ++++++++++++++++++++++++++++++-- 2 files changed, 95 insertions(+), 4 deletions(-) diff --git a/crates/server/src/process.rs b/crates/server/src/process.rs index a0b1b455..6bb89c71 100644 --- a/crates/server/src/process.rs +++ b/crates/server/src/process.rs @@ -889,6 +889,28 @@ impl Server { .map(|record| record.cwd.clone()) } + /// The catalogue's revision, which every change of it moves: what + /// [`Self::wait_native_catalogue_change`] waits past. + pub(crate) fn native_catalogue_revision(&self) -> u64 { + self.0.state.lock().unwrap().catalog_revision + } + + /// Tests: until the child of `process_handle` is reaped (at once when it is not live). + #[cfg(all(test, unix))] + pub(crate) async fn wait_reaped(&self, process_handle: u64) { + let record = self + .0 + .state + .lock() + .unwrap() + .live + .get(&process_handle) + .and_then(Weak::upgrade); + if let Some(record) = record { + record.wait_reaped().await; + } + } + pub(crate) async fn wait_native_catalogue_change(&self, revision: u64) { loop { let notified = self.0.catalog_changed.notified(); diff --git a/crates/server/src/yas_process.rs b/crates/server/src/yas_process.rs index 42cbbfa7..fea053e7 100644 --- a/crates/server/src/yas_process.rs +++ b/crates/server/src/yas_process.rs @@ -574,16 +574,32 @@ impl Session { // replay is recorded before the route is removed. return Ok(Some(exit)); } else { + // A look refused as CONFLICT is one that comes too soon: the process's exit is + // on its way to its watchers (it is final once they have it), or this session's + // own last look at it, a CONTROL's, is still leaving. Either settles with the + // catalogue's next change, after which the WAIT looks again. + let revision = self.inner.server.native_catalogue_revision(); match self .watch_process(request.process_handle, false, true) - .await? + .await { - WatchOutcome::Exited(exit) => return Ok(Some(exit)), - WatchOutcome::Running(attachment) => { + Ok(WatchOutcome::Exited(exit)) => return Ok(Some(exit)), + Ok(WatchOutcome::Running(attachment)) => { let exit = attachment.route.exit.subscribe(); let failed = attachment.route.failed.subscribe(); (exit, failed, Some(attachment)) } + Err(Error::Conflict) => { + let settled = self.inner.server.wait_native_catalogue_change(revision); + match deadline { + None => settled.await, + Some(deadline) => tokio::time::timeout_at(deadline, settled) + .await + .map_err(|_| Error::Timeout)?, + } + return Ok(None); + } + Err(error) => return Err(error), } }; let wait = async { @@ -1183,6 +1199,13 @@ mod tests { assert_eq!(exits.get(newest).unwrap().code, newest as i32); } + /// FUTURE, which must end within 5 s. + async fn within_5s(future: impl std::future::Future) -> T { + tokio::time::timeout(Duration::from_secs(5), future) + .await + .expect("within 5 s") + } + fn spawn_request(argv: Vec>, env: Vec) -> wire::Spawn { wire::Spawn { operation_id: [7; 16], @@ -1552,10 +1575,13 @@ mod tests { attachment.acknowledge_output(stream, end).unwrap(); } } + // Its stdin reports when the child lets go of it, among the frames. + Event::StdinProgress { .. } => {} other => panic!("{other:?}"), } } - tokio::time::sleep(Duration::from_millis(500)).await; + // `z` is in the pipe once the child is gone. + within_5s(server.wait_reaped(attachment.process_handle)).await; attachment.acknowledge_output(Stream::Stdout, end).unwrap(); let (rest, exit) = output_and_exit(&mut attachment, Duration::from_secs(5), end).await; assert_eq!((rest, exit.code), (b"z".to_vec(), 0)); @@ -1597,6 +1623,49 @@ mod tests { server.shutdown().await; } + /// WAIT on a detached process right after this session KILLed it: the CONTROL's own look at + /// the process may still be leaving, or its exit still on its way to its watchers, when the + /// WAIT looks. The WAIT waits for that to settle and answers with the exit, never CONFLICT. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn a_wait_right_after_a_kill_of_a_detached_process_gets_its_exit() { + let server = Server::new(false, true); + let runtime = Runtime::new(server.clone()); + let session = runtime.session([12; 16], None).unwrap(); + for round in 1..=64u8 { + let mut request = spawn_request(sh("exec {sleep} 30"), Vec::new()); + request.flags = schema::process::SPAWN_DETACHABLE as u16; + let attachment = session.spawn(&request, None).await.unwrap(); + let handle = attachment.process_handle; + let (control_half, _events) = attachment.split(); + control_half.detach().await.unwrap(); + session + .control(&wire::Control { + process_handle: handle, + operation_id: [round; 16], + action: wire::ControlAction::Kill, + value: 0, + extensions: Extensions::default(), + }) + .await + .unwrap(); + let wait = wire::Wait { + process_handle: handle, + timeout_ns: 5_000_000_000, + extensions: Extensions::default(), + }; + let exit = tokio::time::timeout(Duration::from_secs(10), session.wait(&wait)) + .await + .expect("the WAIT answers") + .expect("with the exit"); + assert!( + matches!(exit.kind, wire::ExitKind::Killed | wire::ExitKind::Signal), + "{exit:?}" + ); + } + session.shutdown().await; + server.shutdown().await; + } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn an_owner_slow_to_take_its_window_after_the_exit_gets_all_its_output() { let server = Server::new(false, true); From d7217561d77ba0375b9055cdc97e12546e24bc60 Mon Sep 17 00:00:00 2001 From: Pierre Carrier Date: Wed, 30 Sep 2026 12:18:36 +0000 Subject: [PATCH 2/2] Let a WAIT look again once this session's own look has gone A WAIT's temporary look can be refused as CONFLICT because this session's own look at the process is still bound: a concurrent CONTROL's, or a route that an eviction failed and has yet to detach (dispatch_outbound fails the route, then sends the Detach). Such a WAIT waited for the catalogue's next change, which a Detach doesn't make, so it slept until the process ended and then found nothing of an ordinary process: NOT_FOUND. It now waits with Manager::wait_native_look, which also wakes when the record changes and returns once nothing refuses a look: no exit in flight and no binding of this endpoint. a_wait_that_finds_its_own_look_still_bound_answers_once_it_goes fails a watcher's route, WAITs, and detaches the look 200 ms later. Without the change the WAIT answered NOT_FOUND once the child's second had passed; with it, the exit. The Exit arm's comment still said such a WAIT answers CONFLICT; it says what it does now. --- crates/server/src/process.rs | 39 ++++++++++++++++++++++ crates/server/src/yas_process.rs | 57 +++++++++++++++++++++++++++++--- 2 files changed, 92 insertions(+), 4 deletions(-) diff --git a/crates/server/src/process.rs b/crates/server/src/process.rs index 6bb89c71..b66217ce 100644 --- a/crates/server/src/process.rs +++ b/crates/server/src/process.rs @@ -2012,6 +2012,45 @@ impl Manager { } } + /// Until a look at `process_handle` that WATCH refused as CONFLICT, with the catalogue at + /// `revision`, is worth another: the catalogue has moved on (the process's exit reached its + /// watchers and is final, or it left), or nothing refuses the look now (this endpoint's + /// own binding on it, a concurrent CONTROL's or a failed route's, has gone). + pub(crate) async fn wait_native_look(&self, process_handle: u64, revision: u64) { + let record = { + let state = self.server.0.state.lock().unwrap(); + if state.catalog_revision != revision { + return; + } + state.live.get(&process_handle).and_then(Weak::upgrade) + }; + let Some(record) = record else { + return; + }; + loop { + // Both before the looks: notify_waiters reaches futures made before it. + let catalogue = self.server.0.catalog_changed.notified(); + let changed = record.changed.notified(); + if self.server.0.state.lock().unwrap().catalog_revision != revision { + return; + } + { + let inner = record.inner.lock().unwrap(); + let own = inner + .bindings + .iter() + .any(|binding| binding.endpoint_id == self.endpoint.id); + if !inner.terminal_queued && !own { + return; + } + } + tokio::select! { + () = catalogue => {} + () = changed => {} + } + } + } + fn get(&self, process_id: u32) -> Option> { match self.endpoint.state.lock().unwrap().slots.get(&process_id) { Some(EndpointSlot::Bound(record)) => Some(record.clone()), diff --git a/crates/server/src/yas_process.rs b/crates/server/src/yas_process.rs index fea053e7..0aae1353 100644 --- a/crates/server/src/yas_process.rs +++ b/crates/server/src/yas_process.rs @@ -576,8 +576,8 @@ impl Session { } else { // A look refused as CONFLICT is one that comes too soon: the process's exit is // on its way to its watchers (it is final once they have it), or this session's - // own last look at it, a CONTROL's, is still leaving. Either settles with the - // catalogue's next change, after which the WAIT looks again. + // own look at it is still bound (a concurrent CONTROL's, or a route that failed + // and has yet to detach). The WAIT waits for that to settle, then looks again. let revision = self.inner.server.native_catalogue_revision(); match self .watch_process(request.process_handle, false, true) @@ -590,7 +590,10 @@ impl Session { (exit, failed, Some(attachment)) } Err(Error::Conflict) => { - let settled = self.inner.server.wait_native_catalogue_change(revision); + let settled = self + .inner + .manager + .wait_native_look(request.process_handle, revision); match deadline { None => settled.await, Some(deadline) => tokio::time::timeout_at(deadline, settled) @@ -994,7 +997,8 @@ fn dispatch_outbound(inner: &Arc, event: process::NativeEvent) -> let mut routes = inner.routes.lock().unwrap(); // Its route failed as the exit was queued: that attachment is gone already. No // replay is recorded (the handle left with the route), so a WAIT in this session - // answers CONFLICT while the exit is in flight, then NOT_FOUND. + // waits for the exit to be final, then finds the final record of a detachable + // process and nothing (NOT_FOUND) of an ordinary one. let Some(route) = routes.get(&process_id).cloned() else { return Ok(()); }; @@ -1666,6 +1670,51 @@ mod tests { server.shutdown().await; } + /// WAIT while this session's own look at the process is still bound: an eviction fails the + /// route first and detaches its binding next, and the WAIT looks between the two. The + /// Detach settles it: the WAIT looks again, attaches and answers with the exit. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn a_wait_that_finds_its_own_look_still_bound_answers_once_it_goes() { + let server = Server::new(false, true); + let runtime = Runtime::new(server.clone()); + let owner = runtime.session([14; 16], None).unwrap(); + let watcher = runtime.session([15; 16], None).unwrap(); + let attachment = owner + .spawn(&spawn_request(sh("{sleep} 1; exit 4"), Vec::new()), None) + .await + .unwrap(); + let handle = attachment.process_handle; + let look = watcher.attach(&watch_request(handle)).await.unwrap(); + let process_id = look.route.process_id; + // What an eviction does first. + fail_route(&watcher.inner, process_id, ROUTE_EVICTED); + let wait = tokio::spawn({ + let watcher = watcher.clone(); + async move { + watcher + .wait(&wire::Wait { + process_handle: handle, + timeout_ns: 5_000_000_000, + extensions: Extensions::default(), + }) + .await + } + }); + tokio::time::sleep(Duration::from_millis(200)).await; + // And next. + watcher + .inner + .manager + .control_native(process_id, process::NativeControl::Detach) + .unwrap(); + let exit = within_5s(wait).await.unwrap().expect("the exit"); + assert_eq!(exit.code, 4); + drop((look, attachment)); + owner.shutdown().await; + watcher.shutdown().await; + server.shutdown().await; + } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn an_owner_slow_to_take_its_window_after_the_exit_gets_all_its_output() { let server = Server::new(false, true);