From 3ec30d58fa5c3a2b4df79a90bd15c8204d3a57cb Mon Sep 17 00:00:00 2001 From: Pierre Carrier Date: Wed, 30 Sep 2026 11:52:10 +0000 Subject: [PATCH 1/4] 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/4] 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); From 45ff11ecd801e69b3fd8f328b46d4b7e78546f38 Mon Sep 17 00:00:00 2001 From: Pierre Carrier Date: Wed, 30 Sep 2026 12:31:49 +0000 Subject: [PATCH 3/4] Keep an ordinary process's exit for its owner when it missed it An ordinary process belongs to its spawning session. yas-client WAITs for its exit when the attachment that would report it goes first: a stream reset or dropped, a DETACH, or a connection that drops and comes back just as the process exits. That WAIT answered NOT_FOUND whenever the exit came in between. With no binding of the owner's to take the exit, the record was released and nothing was kept, since only a detachable process's final is retained. The same happened in two other cases. When the owner's route failed as the exit was queued, the adapter dropped the exit because no route was there to take it. When the exit found the owner's event queue full, it was dropped, and the attachment waited for it forever. The owner now finds the exit in each case: - The adapter records the replay of every exit it is handed, whether its route is there or not. NativeEvent::Exit now carries the process handle for that. - For the owner's endpoint alone, the native side keeps the final of an ordinary process whose exit that endpoint missed: it had no binding when the exit was queued, or its exit event was dropped before dispatch. WriterGuard now tells its action whether its event was dispatched. The endpoint keeps the newest exit_replays() of these finals, as many as the adapter keeps replays, until it shuts down. WATCH, and so WAIT and ATTACH, finds them after the public finals. They are stored before the release moves the catalogue, which is what a WAIT that found the exit on its way waits for. - A binding whose exit can't be queued is evicted, as a binding that falls behind is. Its attachment fails and its client WAITs, instead of waiting for an exit that won't come. New tests, each of which failed with NotFound before this change: - an_owner_whose_look_went_before_the_exit_still_waits_for_it - an_owner_whose_route_failed_as_the_exit_came_still_waits_for_it - an_owner_whose_queue_was_full_as_the_exit_came_still_finds_it an_owner_that_drops_its_stream_after_the_exit_still_waits_for_it no longer depends on whether its detach comes before or after the exit. a_watcher_whose_queue_fills_after_the_exit_fails_alone assumed that 80 lines make 80 frames, which breaks when a loaded reader takes several lines in one read. Its child now writes each line only after the test has received the last one, as a FIFO tells it, and runs without stdin, so there are no stdin events. The test checks each frame, and checks that only the watcher's attachment failed, evicted. --- crates/server/src/process.rs | 128 ++++++++++++++++--- crates/server/src/yas_process.rs | 207 +++++++++++++++++++++++++------ docs/design/processes.md | 7 +- docs/design/yas.md | 8 +- 4 files changed, 286 insertions(+), 64 deletions(-) diff --git a/crates/server/src/process.rs b/crates/server/src/process.rs index b66217ce..4bf4a8b5 100644 --- a/crates/server/src/process.rs +++ b/crates/server/src/process.rs @@ -451,22 +451,30 @@ impl Policy { } } +/// Runs its action once, told whether the event it rides was dispatched: true after its +/// dispatch, false when it is dropped without (its queue was full, or its reader gone). struct WriterGuard { - action: Option>, + action: Option>, } impl WriterGuard { - fn new(f: impl FnOnce() + Send + 'static) -> Self { + fn new(f: impl FnOnce(bool) + Send + 'static) -> Self { Self { action: Some(Box::new(f)), } } + + fn dispatched(mut self) { + if let Some(f) = self.action.take() { + f(true); + } + } } impl Drop for WriterGuard { fn drop(&mut self) { if let Some(f) = self.action.take() { - f(); + f(false); } } } @@ -619,13 +627,14 @@ pub(crate) enum NativeEvent { }, Exit { process_id: u32, + process_handle: u64, exit: NativeExit, }, } pub(crate) struct NativeEventEnvelope { pub(crate) event: NativeEvent, - _guard: Option, + guard: Option, } impl NativeEventEnvelope { @@ -636,7 +645,9 @@ impl NativeEventEnvelope { /// before a concurrent final output acknowledgement observes retirement. pub(crate) fn dispatch(self, dispatch: impl FnOnce(NativeEvent) -> T) -> T { let result = dispatch(self.event); - drop(self._guard); + if let Some(guard) = self.guard { + guard.dispatched(); + } result } } @@ -697,10 +708,7 @@ impl EndpointOutput { fn send_native(&self, event: NativeEvent, guard: Option) -> bool { self.events - .try_send(NativeEventEnvelope { - event, - _guard: guard, - }) + .try_send(NativeEventEnvelope { event, guard }) .is_ok() } @@ -727,8 +735,21 @@ impl EndpointOutput { ) } - fn send_exit(&self, process_id: u32, exit: NativeExit, guard: WriterGuard) -> bool { - self.send_native(NativeEvent::Exit { process_id, exit }, Some(guard)) + fn send_exit( + &self, + process_id: u32, + process_handle: u64, + exit: NativeExit, + guard: WriterGuard, + ) -> bool { + self.send_native( + NativeEvent::Exit { + process_id, + process_handle, + exit, + }, + Some(guard), + ) } } @@ -1093,6 +1114,36 @@ struct EndpointState { slots: FxHashMap, /// Ordinary processes remain owned after their creator unsubscribes. owned: FxHashMap>, + /// The finals of this endpoint's ordinary processes whose exits it missed: a WAIT of its + /// own finds them, as it finds the exits it took. + missed: MissedExits, +} + +/// The newest finals of an endpoint's ordinary processes whose exits it missed: it had no +/// binding to take the exit (its attachment went first), or the exit was dropped on its way +/// (its queue was full). Nobody else sees them. +#[derive(Default)] +struct MissedExits { + values: FxHashMap>, + order: VecDeque, +} + +impl MissedExits { + fn insert(&mut self, final_record: Arc, capacity: usize) { + let generation = final_record.generation; + if self.values.insert(generation, final_record).is_none() { + self.order.push_back(generation); + } + while self.order.len() > capacity.max(1) { + if let Some(retired) = self.order.pop_front() { + self.values.remove(&retired); + } + } + } + + fn get(&self, generation: u64) -> Option<&Arc> { + self.values.get(&generation) + } } enum EndpointSlot { @@ -1988,7 +2039,11 @@ impl Manager { drop(server); record.changed.notify_waiters(); Ok(watched) - } else if let Some(record) = server.finals.get(&process_handle) { + } else if let Some(record) = server + .finals + .get(&process_handle) + .or_else(|| endpoint.missed.get(process_handle)) + { if endpoint_usage(&endpoint) >= self.server.0.policy.max_per_endpoint { return Err(NativeError::ResourceExhausted); } @@ -2062,6 +2117,7 @@ impl Manager { let (slots, owned) = { let mut endpoint = self.endpoint.state.lock().unwrap(); endpoint.accepting = false; + endpoint.missed = MissedExits::default(); ( std::mem::take(&mut endpoint.slots), std::mem::take(&mut endpoint.owned), @@ -3200,7 +3256,7 @@ fn deliver( offset, data: data.to_vec(), }, - _guard: None, + guard: None, }); let binding = &mut inner.bindings[index]; let credit = if stream == PROCESS_STREAM_STDOUT { @@ -3806,27 +3862,49 @@ fn try_queue_terminal(record: &Arc) { record.terminal_notify.notify_waiters(); let (bindings, final_record) = terminal; if bindings.is_empty() { - finish_terminal(record.clone(), final_record); + // Nobody takes the exit, its owner included. + finish_terminal(record.clone(), final_record, true); return; } + // Whether the owner misses the exit: it has no binding to take it (its attachment went + // first), or the exit is dropped on its way to it. + let owner = record.owner.upgrade().map(|owner| owner.id); + let owner_missed = Arc::new(AtomicBool::new( + !bindings + .iter() + .any(|binding| Some(binding.endpoint_id) == owner), + )); let remaining = Arc::new(AtomicUsize::new(bindings.len())); for binding in bindings { let endpoint = binding.endpoint.upgrade(); let record_for_guard = record.clone(); let final_for_guard = final_record.clone(); let remaining_for_guard = remaining.clone(); + let owner_missed = owner_missed.clone(); + let owners = Some(binding.endpoint_id) == owner; let process_id = binding.process_id; - let guard = WriterGuard::new(move || { + let guard = WriterGuard::new(move |dispatched| { + if owners && !dispatched { + owner_missed.store(true, Ordering::Release); + } if let Some(endpoint) = endpoint { remove_bound_slot(&endpoint, process_id, &record_for_guard); } if remaining_for_guard.fetch_sub(1, Ordering::AcqRel) == 1 { - finish_terminal(record_for_guard, final_for_guard); + let owner_missed = owner_missed.load(Ordering::Acquire); + finish_terminal(record_for_guard, final_for_guard, owner_missed); } }); - let _ = binding - .out - .send_exit(binding.process_id, final_record.exit(), guard); + if !binding.out.send_exit( + binding.process_id, + record.generation, + final_record.exit(), + guard, + ) { + // Its queue is full, and the exit lost to it: its attachment fails, as one that + // falls behind does, and its client WAITs. + binding.out.evict(binding.process_id); + } } } @@ -3898,11 +3976,21 @@ fn portable_signal_reason(signal: u32) -> u8 { } } -fn finish_terminal(record: Arc, final_record: Arc) { +/// `owner_missed`: the exit did not reach the owner. An ordinary process's final is then kept +/// for it, before the release moves the catalogue: a WAIT of its that found the exit on its way +/// looks again at that change. +fn finish_terminal(record: Arc, final_record: Arc, owner_missed: bool) { if record.detachable { let server = record.server.clone(); server.finish_detached(record, final_record); } else { + if owner_missed && let Some(owner) = record.owner.upgrade() { + let mut state = owner.state.lock().unwrap(); + if state.accepting { + let capacity = record.server.0.maxima.exit_replays(); + state.missed.insert(final_record, capacity); + } + } record.server.release_record(&record); } } diff --git a/crates/server/src/yas_process.rs b/crates/server/src/yas_process.rs index 0aae1353..e1340449 100644 --- a/crates/server/src/yas_process.rs +++ b/crates/server/src/yas_process.rs @@ -991,27 +991,25 @@ fn dispatch_outbound(inner: &Arc, event: process::NativeEvent) -> } Ok(()) } - process::NativeEvent::Exit { process_id, exit } => { + process::NativeEvent::Exit { + process_id, + process_handle, + exit, + } => { let exit = native_exit_info(exit); + // The replay goes first, and whether the route is there or not: a WAIT in this + // session that no longer finds the route (it left, or failed as the exit was queued) + // must find the exit. + inner + .exits + .lock() + .unwrap() + .insert(process_handle, exit.clone()); let route = { 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 - // 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(()); }; - let process_handle = route.process_handle.load(Ordering::Acquire); - if process_handle != 0 { - // Record the replay before the route disappears: a WAIT - // that no longer finds the route must find the exit. - inner - .exits - .lock() - .unwrap() - .insert(process_handle, exit.clone()); - } routes.remove(&process_id); route.exit.send_replace(Some(exit.clone())); route @@ -1536,50 +1534,66 @@ mod tests { #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn a_watcher_whose_queue_fills_after_the_exit_fails_alone() { + use std::io::Write; + use std::os::unix::ffi::OsStrExt; let server = Server::new(false, true); let runtime = Runtime::new(server.clone()); let owner = runtime.session([6; 16], None).unwrap(); let watcher = runtime.session([7; 16], None).unwrap(); - let mut attachment = owner - .spawn( - &spawn_request( - sh( - "{sleep} 0.3; i=0; while [ $i -lt 80 ]; do echo x; {sleep} 0.02; \ - i=$((i+1)); done; printf z; exit 0", - ), - Vec::new(), - ), - None, - ) - .await - .unwrap(); - // Never read: its route queue (80 events) fills with the 80 lines, and the last - // frame comes after the child exits. - let _watched = watcher + // The child writes its next line once the test has the last one, as a FIFO says: one + // line a frame, however slow the machine. + let dir = tempfile::tempdir().unwrap(); + let fifo = dir.path().join("go"); + let path = std::ffi::CString::new(fifo.as_os_str().as_bytes()).unwrap(); + // SAFETY: a NUL-terminated path. + assert_eq!(unsafe { libc::mkfifo(path.as_ptr(), 0o600) }, 0); + let mut request = spawn_request( + sh(&format!( + "exec 3<'{}'; i=0; while [ $i -lt 80 ]; do echo x; read -r go <&3; \ + i=$((i+1)); done; printf z; exit 0", + fifo.display() + )), + Vec::new(), + ); + // No stdin, so no stdin events: the watcher's queue takes output alone. + request.flags = schema::process::SPAWN_STDIN_NULL as u16; + let mut attachment = owner.spawn(&request, None).await.unwrap(); + // Never read: its route queue (80 events) fills with the 80 lines, and the next frame + // comes after the child exits. + let (_watch, watched) = watcher .attach(&watch_request(attachment.process_handle)) .await - .unwrap(); + .unwrap() + .split(); + // The child's first line waits for this: the watcher sees every line. + let mut go = within_5s(tokio::task::spawn_blocking(move || { + std::fs::OpenOptions::new().write(true).open(fifo) + })) + .await + .unwrap() + .unwrap(); // The owner acknowledges its first 48 frames only: with 32 unacknowledged, the reader // waits for it, and `z` stays in the pipe as the child exits. let (mut frames, mut end) = (0, 0u64); while frames < 80 { - match tokio::time::timeout(Duration::from_secs(5), attachment.next()) - .await - .unwrap() - .unwrap() - { + match within_5s(attachment.next()).await.unwrap() { Event::Output { stream, lifetime_offset, data, } => { frames += 1; - end = lifetime_offset + data.len() as u64; + assert_eq!( + (lifetime_offset, data.as_slice()), + (end, &b"x\n"[..]), + "frame {frames}" + ); + end += data.len() as u64; if frames <= 48 { attachment.acknowledge_output(stream, end).unwrap(); } + go.write_all(b"\n").unwrap(); } - // Its stdin reports when the child lets go of it, among the frames. Event::StdinProgress { .. } => {} other => panic!("{other:?}"), } @@ -1589,7 +1603,10 @@ mod tests { 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)); - tokio::time::sleep(Duration::from_millis(300)).await; + // `z` found the watcher's queue full: its attachment alone failed. + let mut failed = watched.failed.clone(); + within_5s(failed.wait_for(Option::is_some)).await.unwrap(); + assert_eq!(watched.failure().as_deref(), Some(ROUTE_EVICTED)); assert!(watcher.inner.closed.borrow().is_none()); assert_eq!(runs_a_command(&watcher).await, (b"after\n".to_vec(), 0)); owner.shutdown().await; @@ -1603,7 +1620,10 @@ mod tests { let runtime = Runtime::new(server.clone()); let owner = runtime.session([10; 16], None).unwrap(); // The reader takes a window (1 MiB) and waits for the owner; the last 32 KiB stay in - // the pipe as the child exits. + // the pipe as the child exits. (Or a loaded reader's small reads reach the frames the + // owner may leave unacknowledged first, and the child exits after the detach, with no + // look of its owner's to take the exit. Its WAIT gets it all the same: the server keeps + // it for the owner.) let attachment = owner .spawn( &spawn_request(sh("{head} -c 1081344 /dev/zero; exit 4"), Vec::new()), @@ -1627,6 +1647,111 @@ mod tests { server.shutdown().await; } + fn wait_request(process_handle: u64) -> wire::Wait { + wire::Wait { + process_handle, + timeout_ns: 5_000_000_000, + extensions: Extensions::default(), + } + } + + /// Until `handle` has left the catalogue: its exit is final and its record released. + async fn left_the_catalogue(server: &Server, handle: u64) { + loop { + let revision = server.native_catalogue_revision(); + let snapshot = server.native_snapshot(); + if !snapshot + .records + .iter() + .any(|record| record.process_handle == handle) + { + return; + } + server.wait_native_catalogue_change(revision).await; + } + } + + /// The owner's look at its process went before the exit (its client dropped the streams), + /// so no route took the exit, and the process left the catalogue: the owner's WAIT still + /// answers with the exit. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn an_owner_whose_look_went_before_the_exit_still_waits_for_it() { + let server = Server::new(false, true); + let runtime = Runtime::new(server.clone()); + let owner = runtime.session([16; 16], None).unwrap(); + let attachment = owner + .spawn(&spawn_request(sh("{sleep} 0.3; exit 4"), Vec::new()), None) + .await + .unwrap(); + let handle = attachment.process_handle; + let (control, _events) = attachment.split(); + control.detach().await.unwrap(); + within_5s(left_the_catalogue(&server, handle)).await; + let exit = owner.wait(&wait_request(handle)).await.unwrap(); + assert_eq!((exit.code, exit.detail.as_slice()), (4, &b""[..])); + owner.shutdown().await; + server.shutdown().await; + } + + /// The owner's route failed as its process's exit came (an eviction fails the route before + /// it detaches the binding): the exit reached the session with no route to take it. The + /// owner's WAIT still answers with it. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn an_owner_whose_route_failed_as_the_exit_came_still_waits_for_it() { + let server = Server::new(false, true); + let runtime = Runtime::new(server.clone()); + let owner = runtime.session([17; 16], None).unwrap(); + let attachment = owner + .spawn(&spawn_request(sh("{sleep} 0.3; exit 4"), Vec::new()), None) + .await + .unwrap(); + let handle = attachment.process_handle; + fail_route(&owner.inner, attachment.route.process_id, ROUTE_EVICTED); + within_5s(left_the_catalogue(&server, handle)).await; + let exit = owner.wait(&wait_request(handle)).await.unwrap(); + assert_eq!((exit.code, exit.detail.as_slice()), (4, &b""[..])); + drop(attachment); + owner.shutdown().await; + server.shutdown().await; + } + + /// The owner's queue was full as its process's exit came, so the exit never reached the + /// session: its binding is evicted (its attachment fails, and its client WAITs), and a + /// look at the process still finds the exit. + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn an_owner_whose_queue_was_full_as_the_exit_came_still_finds_it() { + let server = Server::new(false, true); + // Room for one event, never taken: the output takes it, and the exit finds it full. + let (manager, _events, evictions) = server.native_endpoint_with_session([18; 16], 1); + let started = manager + .spawn_native( + process::NativeSpawnRequest { + process_id: 1, + flags: schema::process::SPAWN_STDIN_NULL as u8, + preserve_residual: false, + residue_grace: None, + keep_output: None, + cwd: None, + argv: sh("printf x; exit 4"), + env: Vec::new(), + clear_environment: true, + }, + None, + ) + .await + .unwrap(); + let handle = started.process_handle; + within_5s(left_the_catalogue(&server, handle)).await; + let watched = manager + .watch_native(2, handle, false) + .expect("the exit is kept for its owner"); + assert!(!watched.running); + assert_eq!(watched.exit.map(|exit| exit.code), Some(4)); + assert_eq!(evictions.take(), vec![1]); + manager.shutdown().await; + 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. diff --git a/docs/design/processes.md b/docs/design/processes.md index f77f9d5e..a19393fe 100644 --- a/docs/design/processes.md +++ b/docs/design/processes.md @@ -107,7 +107,12 @@ yas-client's `Command::keep_output` asks for it where offered, and `Process::elided` reads the counts. Closing stdin half-closes the child stream. An ordinary child belongs to its spawning session and is terminated when it disappears; a detachable -child remains discoverable until its retained final record expires. +child remains discoverable until its retained final record expires. The spawning +session can still WAIT for an ordinary child that has gone when its exit did not +reach it: its attachment went first (a stream reset or dropped, a DETACH, a +route failed as the exit came), or the exit found its queue full. The server +keeps such exits for that session alone, as many as the process generations +maximum (64 by default, never fewer), until the session ends. Arguments and environment values preserve arbitrary bytes on Unix and use exact UTF-8-to-native conversion on Windows. The cwd union is server default, native diff --git a/docs/design/yas.md b/docs/design/yas.md index 4d0738bf..aaf732e2 100644 --- a/docs/design/yas.md +++ b/docs/design/yas.md @@ -3317,8 +3317,12 @@ and otherwise WAITs, so either side may be older. Catalog records contain argv0, native PID for diagnostics, lifecycle, owner session, detachable flag, stream offsets, exit record, and retention deadline. An ordinary process is owned by its spawning session and terminated when that -session disappears. A detachable process survives without watchers and remains -discoverable until its retained exit result expires. +session disappears. That session's WAIT answers with its exit even after it has +gone, when no attachment of the session's took the exit (one went first, or the +exit came as it failed): the server keeps such exits for the session alone, as +many as `MAX_PROCESSES` (never fewer than 64), until the session ends. A +detachable process survives without watchers and remains discoverable until its +retained exit result expires. ATTACH returns new stdout/stderr Transfers beginning at the process's current lifetime offsets; earlier output is explicitly reported as a gap and is not From 29350d9a6e5a8237f45c73dbdd93d058347e52ce Mon Sep 17 00:00:00 2001 From: Pierre Carrier Date: Wed, 30 Sep 2026 14:37:45 +0000 Subject: [PATCH 4/4] Wait for the queue-full test's eviction before counting it The failed try_send drops the exit's envelope inside send_native, whose guard keeps the final and releases the record before try_queue_terminal evicts the binding: the test could see the process gone and the final kept with no eviction pushed yet. It waits for the eviction's notice (a permit is kept when it came first). Under load, it failed 7 of 24 parallel runs of yas_process::tests::a, and 3 of 40 alone, in review; with the wait, none of 24 and none of 40. --- crates/server/src/yas_process.rs | 3 +++ 1 file changed, 3 insertions(+) diff --git a/crates/server/src/yas_process.rs b/crates/server/src/yas_process.rs index e1340449..a0afd569 100644 --- a/crates/server/src/yas_process.rs +++ b/crates/server/src/yas_process.rs @@ -1747,6 +1747,9 @@ mod tests { .expect("the exit is kept for its owner"); assert!(!watched.running); assert_eq!(watched.exit.map(|exit| exit.code), Some(4)); + // The eviction follows the release that `left_the_catalogue` saw: wait for it (a permit + // is kept when it came first). + within_5s(evictions.notified()).await; assert_eq!(evictions.take(), vec![1]); manager.shutdown().await; server.shutdown().await;