Skip to content
Merged
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
23 changes: 20 additions & 3 deletions crates/compass-ipc/src/transport.rs
Original file line number Diff line number Diff line change
Expand Up @@ -284,6 +284,13 @@ impl WindowLink {
/// acknowledged by the write succeeding: a successful write only means the
/// bytes reached a kernel buffer, and a window that died between the write
/// and the read would look like a window that showed.
///
/// # Cancellation
///
/// The future may be dropped while it waits, as the engine does when a
/// window takes too long to answer. The window still answers that command
/// in order, before the next one, so the next push skips answers to
/// commands older than its own instead of taking one as its own.
pub async fn push(&mut self, command: WindowCommand) -> Result<WindowOutcome> {
let id = self.next_id;
self.next_id = self.next_id.wrapping_add(1);
Expand All @@ -292,9 +299,19 @@ impl WindowLink {
.send(&ResponseEnvelope::new(id, Response::Window(command)))
.await?;

let envelope = match self.framed.next().await {
Some(frame) => frame?,
None => return Err(Error::ConnectionClosed),
let envelope = loop {
let envelope = match self.framed.next().await {
Some(frame) => frame?,
None => return Err(Error::ConnectionClosed),
};
if envelope.id >= id {
break envelope;
}
tracing::debug!(
late = envelope.id,
waiting_for = id,
"skipping a late window answer"
);
};

if envelope.version != PROTOCOL_VERSION {
Expand Down
64 changes: 64 additions & 0 deletions crates/compass-ipc/tests/window_link.rs
Original file line number Diff line number Diff line change
Expand Up @@ -286,6 +286,70 @@ async fn a_push_to_a_dead_window_fails_rather_than_reporting_success() {
);
}

/// The engine gives up on a slow window by dropping the push. The window
/// still answers that command, late and before the next one, and the next
/// push must not take that late answer for its own.
#[tokio::test]
async fn a_push_given_up_on_does_not_steal_the_next_ones_answer() {
let dir = TempDir::new();
let socket = dir.socket();
let mut links = engine_accepting_one_window(&socket).await;

let (release_tx, release_rx) = oneshot::channel::<()>();
let window = tokio::spawn({
let socket = socket.clone();
async move {
let mut window = WindowClient::attach(socket.as_path())
.await
.expect("attach");
let first = window
.next_command()
.await
.expect("first")
.expect("a command");
assert_eq!(first, WindowCommand::Show);
let _ = release_rx.await;
window
.reply(WindowOutcome::Shown)
.await
.expect("late reply");
let second = window
.next_command()
.await
.expect("second")
.expect("a command");
assert_eq!(second, WindowCommand::Hide);
window.reply(WindowOutcome::Hidden).await.expect("reply");
let _ = window.next_command().await;
}
});

let mut link = tokio::time::timeout(GUARD, links.recv())
.await
.expect("link handed over in time")
.expect("a link");

let given_up =
tokio::time::timeout(Duration::from_millis(100), link.push(WindowCommand::Show)).await;
assert!(
given_up.is_err(),
"the window was holding its answer, got {given_up:?}"
);

let _ = release_tx.send(());
let outcome = tokio::time::timeout(GUARD, link.push(WindowCommand::Hide))
.await
.expect("push answered in time")
.expect("push");
assert_eq!(outcome, WindowOutcome::Hidden);

drop(link);
tokio::time::timeout(GUARD, window)
.await
.expect("window ended in time")
.expect("window task");
}

#[tokio::test]
async fn a_window_whose_engine_went_away_is_told_rather_than_left_waiting() {
let dir = TempDir::new();
Expand Down
41 changes: 40 additions & 1 deletion crates/compass/src/serve.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2544,13 +2544,52 @@ fn no_window(what: &str) -> Response {
/// removed. What it buys is releasing the file descriptor and not paying a
/// doomed write on every subsequent request. See [`WindowSlot`].
pub(crate) async fn forward(slot: &WindowSlot, command: WindowCommand, what: &str) -> Response {
forward_within(slot, command, what, WINDOW_ANSWER_TIMEOUT).await
}

/// How long the engine waits for the window to answer a command.
///
/// The slot is locked while it waits, so every request that needs the window
/// waits too. Every command is answered as soon as the window has acted on it
/// (a dmenu's choice comes later, outside the push), but the first `Show` waits
/// on the new window's renderer, and a machine without a GPU spends seconds
/// probing Vulkan before it settles on GL. So this bounds a stuck window rather
/// than a slow one, and is longer than the CLI's own 10 seconds: a slow first
/// show must not be refused while it is still coming up.
const WINDOW_ANSWER_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30);

/// [`forward`] with the wait for the window's answer bounded by `timeout`.
///
/// A window that does not answer in time keeps its link. It may only be busy,
/// and dropping it would leave the launcher undriven until it restarts; its
/// late answer is skipped by the next push (see [`WindowLink::push`]).
pub(crate) async fn forward_within(
slot: &WindowSlot,
command: WindowCommand,
what: &str,
timeout: std::time::Duration,
) -> Response {
let mut guard = slot.lock().await;

let Some(link) = guard.as_mut() else {
return no_window(what);
};

match link.push(command).await {
let Ok(answer) = tokio::time::timeout(timeout, link.push(command)).await else {
tracing::warn!(
seconds = timeout.as_secs_f32(),
"the launcher window did not answer"
);
return Response::Error(ProtocolError::new(
ErrorKind::Internal,
format!(
"the launcher window did not answer within {} seconds, so it could not {what}",
timeout.as_secs()
),
));
};

match answer {
Ok(WindowOutcome::Shown | WindowOutcome::Hidden) => Response::Ack,
Ok(WindowOutcome::Failed(reason)) => {
// The window is alive and said no. Keeping the link is the point:
Expand Down
58 changes: 58 additions & 0 deletions crates/compass/src/window.rs
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,64 @@ mod tests {
);
}

/// A window that is slow to answer must not hold the engine's slot, and
/// every request queued behind it, for as long as it takes.
#[tokio::test]
async fn the_engine_stops_waiting_on_a_slow_window_and_keeps_it() {
let dir = TempDir::new();
let socket = dir.socket();
let mut links = engine(&socket).await;

let client = tokio::time::timeout(GUARD, WindowClient::attach(socket.as_path()))
.await
.expect("attached in time")
.expect("attach");
let (commands_tx, mut commands_rx) = mpsc::unbounded_channel::<UiCommand>();
let (outcomes_tx, outcomes_rx) = mpsc::unbounded_channel::<UiOutcome>();
tokio::spawn(bridge(client, commands_tx, outcomes_rx));

let link = tokio::time::timeout(GUARD, links.recv())
.await
.expect("link in time")
.expect("a link");
let slot = std::sync::Arc::new(tokio::sync::Mutex::new(Some(link)));

let slow = crate::serve::forward_within(
&slot,
WindowCommand::Show,
"show the launcher",
Duration::from_millis(100),
)
.await;
let Response::Error(err) = slow else {
panic!("a window that has not answered must not be acknowledged, got {slow:?}");
};
assert!(err.message.contains("did not answer"), "{}", err.message);
assert!(slot.lock().await.is_some(), "a slow window keeps its link");

// The window catches up: the late answer is skipped, the next is read.
assert_eq!(commands_rx.recv().await, Some(UiCommand::Show));
outcomes_tx.send(UiOutcome::Shown).expect("late answer");
tokio::spawn(async move {
while let Some(command) = commands_rx.recv().await {
let outcome = match command {
UiCommand::Hide => UiOutcome::Hidden,
_ => UiOutcome::Failed("unexpected".to_owned()),
};
if outcomes_tx.send(outcome).is_err() {
return;
}
}
});
let next = tokio::time::timeout(
GUARD,
crate::serve::forward(&slot, WindowCommand::Hide, "hide the launcher"),
)
.await
.expect("answered in time");
assert!(matches!(next, Response::Ack), "got {next:?}");
}

/// Only used to build a refusing engine in the test below.
fn refusal() -> Response {
Response::Error(ProtocolError::new(
Expand Down
Loading