Skip to content
Closed
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
61 changes: 61 additions & 0 deletions crates/server/src/process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -1990,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<Arc<Record>> {
match self.endpoint.state.lock().unwrap().slots.get(&process_id) {
Some(EndpointSlot::Bound(record)) => Some(record.clone()),
Expand Down
128 changes: 123 additions & 5 deletions crates/server/src/yas_process.rs
Original file line number Diff line number Diff line change
Expand Up @@ -574,16 +574,35 @@ 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 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)
.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
.manager
.wait_native_look(request.process_handle, 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 {
Expand Down Expand Up @@ -978,7 +997,8 @@ fn dispatch_outbound(inner: &Arc<SessionInner>, 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(());
};
Expand Down Expand Up @@ -1183,6 +1203,13 @@ mod tests {
assert_eq!(exits.get(newest).unwrap().code, newest as i32);
}

/// FUTURE, which must end within 5 s.
async fn within_5s<T>(future: impl std::future::Future<Output = T>) -> T {
tokio::time::timeout(Duration::from_secs(5), future)
.await
.expect("within 5 s")
}

fn spawn_request(argv: Vec<Vec<u8>>, env: Vec<EnvEntry>) -> wire::Spawn {
wire::Spawn {
operation_id: [7; 16],
Expand Down Expand Up @@ -1552,10 +1579,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));
Expand Down Expand Up @@ -1597,6 +1627,94 @@ 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;
}

/// 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);
Expand Down
Loading