diff --git a/src/ui/remote_workspace.rs b/src/ui/remote_workspace.rs index 4e59d511..e64b9a79 100644 --- a/src/ui/remote_workspace.rs +++ b/src/ui/remote_workspace.rs @@ -956,6 +956,12 @@ fn pump_tick(cx: &mut gpui::App) -> bool { changed = true; log::info!("link to {target} is attached"); crate::ui::machine_mirror::MachineMirrors::refresh(cx, host); + // A link this machine's windows never asked for — the switcher + // connected it, or `finish_connect` installed it — comes up + // without any reconnect attempt finishing, so nothing else + // tells the windows on it that their machine can be reached + // now. One of them may be sitting empty owing a pull. + crate::ui::tree_sync::on_link_up(cx, host); } continue; } diff --git a/src/ui/tree_sync.rs b/src/ui/tree_sync.rs index 05c3a2e6..bef802f0 100644 --- a/src/ui/tree_sync.rs +++ b/src/ui/tree_sync.rs @@ -714,6 +714,19 @@ struct WsState { /// from it until the pull is retried — an empty window diffs into /// "close every tab" and would wipe the layout off the machine. rehydrate: Option, + /// How many pulls in a row this window has owed, which paces the retry. + /// + /// Counts consecutive failures, so it is cleared by anything that ends the + /// run: a pull that lands (`finish_hydration`), a prime that lands + /// (`finish_prime` — the machine answered, which is the whole question), + /// and a debt abandoned rather than paid (`take_rehydrate` dropping a + /// `Replace` the user has overtaken). A machine that hiccups once is then + /// asked again promptly, and one that is really gone is not asked in a + /// loop. + /// + /// Leaving it standing after the run ends is what makes a *first* failure + /// wait the cap: the count would still be carrying an outage that is over. + rehydrate_attempts: u32, /// Whether this window has already been told why it opened empty. /// /// The retry is as quiet as the failure was, so a window whose machine @@ -736,6 +749,7 @@ impl Default for WsState { informed: false, epoch: 0, rehydrate: None, + rehydrate_attempts: 0, said_why_empty: false, } } @@ -819,7 +833,14 @@ fn take_rehydrate(cx: &mut App, client_ws: WorkspaceId, window_is_empty: bool) - .windows .get_mut(&client_ws)?; let adopt = state.rehydrate.take()?; - (window_is_empty || adopt == Adopt::IfEmpty).then_some(adopt) + if !window_is_empty && adopt == Adopt::Replace { + // Abandoned, not paid — but the run of failures is over either way, and + // a count left standing would make the next window's first failure wait + // the cap on an outage that has nothing to do with it. + state.rehydrate_attempts = 0; + return None; + } + Some(adopt) } /// Whether a window with no tabs may delete `client_ws` outright — from the @@ -1070,6 +1091,9 @@ fn finish_prime(cx: &mut App, client_ws: WorkspaceId, epoch: u64, outcome: io::R let landed = match outcome { Ok(mirror) => { state.informed |= mirror.tabs.is_empty(); + // The machine answered, which is the only thing the retry was + // waiting to find out, so the next failure starts its backoff over. + state.rehydrate_attempts = 0; let landed = (mirror.tabs.clone(), mirror.active); state.sync = SyncPhase::Primed(mirror); landed @@ -1252,7 +1276,7 @@ enum Adopt { fn hydrate(cx: &mut App, client_ws: WorkspaceId, adopt: Adopt) { let host = WorkspaceStore::host_of(cx, client_ws); let machine_ws = tree_workspace_id(cx, client_ws); - let epoch = { + let (epoch, failures) = { let state = cx .default_global::() .windows @@ -1266,22 +1290,29 @@ fn hydrate(cx: &mut App, client_ws: WorkspaceId, adopt: Adopt) { state.epoch += 1; // This attempt takes over the debt; it re-records it if it fails too. state.rehydrate = None; - state.epoch + (state.epoch, state.rehydrate_attempts) }; + // How many times in a row this window has already failed, which is what + // decides whether another failure is news or the same news again. + let level = hydration_log_level(failures, log::Level::Warn); cx.spawn(async move |cx| { let deadline = std::time::Instant::now() + HYDRATE_LINK_DEADLINE; let client = loop { match cx.update(|cx| tree_control_for(cx, host)) { TreeLink::Ready(client) => break Some(client), TreeLink::Unserved => { - log::warn!( + log::log!( + level, "workspace {client_ws}: its machine's server does not serve the \ machine tree; opening empty" ); break None; } TreeLink::Down if std::time::Instant::now() > deadline => { - log::warn!("workspace {client_ws}: no link to its machine; opening empty"); + log::log!( + level, + "workspace {client_ws}: no link to its machine; opening empty" + ); break None; } TreeLink::Down => cx.background_executor().timer(HYDRATE_LINK_POLL).await, @@ -1334,12 +1365,23 @@ fn say_why_the_window_is_empty(cx: &mut App, client_ws: WorkspaceId) { }); } -/// Records that a hydration failed and still owes `client_ws` its layout. +/// Records that a hydration failed and still owes `client_ws` its layout, and +/// arms the retry that pays it back. /// /// Nothing else recovers on its own: the window stays empty, and without this /// the next `sync_window` would push that emptiness to the machine as "close -/// every tab". Instead the pull is retried the next time the window syncs — -/// which is what a reconnect does through `on_link_up`. +/// every tab". The debt is settled by the next sync of this window — a +/// reconnect drives one through `on_link_up`, an edit in the window drives one +/// through `save_session`, and [`arm_rehydrate_retry`] drives one when neither +/// happens. +/// +/// That last driver is the load-bearing one. A pull can fail with the link +/// perfectly healthy — a `MachineGet` that overran its ten seconds on a slow +/// link, or a create that lost its race with `start_prime` — and then no link +/// ever comes back up to notice, and an empty window has nothing to edit. The +/// window sat empty until the app was restarted, with every tab and every +/// shell still on the machine: "tty7 lost my session" for a request that +/// needed asking twice. /// /// Returns whether the debt was taken on. A superseded attempt gets `false`: /// a newer hydration owns the window now, and this one speaks for nothing. @@ -1354,45 +1396,162 @@ fn owe_rehydration(cx: &mut App, client_ws: WorkspaceId, epoch: u64, adopt: Adop *priming = false; } state.rehydrate = Some(adopt); - log::info!("workspace {client_ws}: will pull its layout again once its machine answers"); + state.rehydrate_attempts = state.rehydrate_attempts.saturating_add(1); + let attempts = state.rehydrate_attempts; + log::log!( + // Once settled this line says the same thing every thirty seconds until + // the window closes, which is a fact about the machine and not an event. + hydration_log_level(attempts, log::Level::Info), + "workspace {client_ws}: will pull its layout again once its machine answers \ + (attempt {attempts})" + ); + arm_rehydrate_retry(cx, client_ws, epoch, attempts); true } +/// Whether the debt this retry was armed for is still the window's own. +/// +/// A newer epoch means another hydration took the window over while the +/// backoff ran, and this retry speaks for nothing. +fn still_owed(cx: &App, client_ws: WorkspaceId, epoch: u64) -> bool { + cx.try_global::() + .and_then(|t| t.windows.get(&client_ws)) + .is_some_and(|s| s.rehydrate.is_some() && s.epoch == epoch) +} + +/// The attempt from which the backoff no longer grows. +/// +/// Also the point where a window stops being a fresh failure and becomes a +/// standing one, which is what [`hydration_log_level`] keys off. +const REHYDRATE_SETTLED: u32 = 5; +const REHYDRATE_BACKOFF_CAP: std::time::Duration = std::time::Duration::from_secs(30); + +/// The first retry is soon enough to look instant to someone watching an empty +/// window; the backoff is what keeps a machine that is really unreachable from +/// being asked on a loop for as long as its window stays open. +fn rehydrate_backoff(attempts: u32) -> std::time::Duration { + std::time::Duration::from_secs(2u64.saturating_pow(attempts.min(REHYDRATE_SETTLED))) + .min(REHYDRATE_BACKOFF_CAP) +} + +/// Steps `fresh` down to `debug` once this window's failures have stopped being +/// events and become a standing condition. +/// +/// The first few are news: something that was working stopped. Once the backoff +/// has settled at its cap the window is in a steady state — a machine that is +/// simply not there — and the retry will go on failing every thirty seconds for +/// as long as the window stays open. Repeating that at full volume buries +/// whatever else is in the log. The retry stays exactly as persistent either +/// way; only the volume drops. +fn hydration_log_level(attempts: u32, fresh: log::Level) -> log::Level { + if attempts >= REHYDRATE_SETTLED { + log::Level::Debug + } else { + fresh + } +} + +/// Asks `client_ws` to sync once the backoff is up, if it still owes a pull. +/// +/// Deliberately routed through `sync_window` rather than straight into +/// `hydrate`: that is where the rules about *whether* a window may still adopt +/// the machine's layout live — a preempted workspace stays out of it, and a +/// `Replace` is dropped once the user has filled the window in themselves. +fn arm_rehydrate_retry(cx: &mut App, client_ws: WorkspaceId, epoch: u64, attempts: u32) { + let delay = rehydrate_backoff(attempts); + cx.spawn(async move |cx| { + cx.background_executor().timer(delay).await; + let _ = cx.update(|cx| { + if !still_owed(cx, client_ws, epoch) { + return; + } + // No window left to fill, so asking its machine now would be work + // for nobody. Closing a window drops its whole `WsState` through + // `forget`, debt and all, so `still_owed` above normally answers + // first; this covers the window that is on its way out and has + // already dropped its app. + let Some(app) = crate::ui::windows::WindowRegistry::app_for(cx, client_ws) + .and_then(|app| app.upgrade()) + else { + return; + }; + app.update(cx, |app, cx| sync_window(app, cx)); + }); + }) + .detach(); +} + fn pull_workspace( client: &ControlClient, machine_ws: WorkspaceId, ) -> io::Result<(Machine, WsMirror, Session)> { - let machine: Machine = match client.call(ControlRequest::MachineGet)? { - ReplyOk::MachineTree(m) => *m, - other => return Err(io::Error::other(format!("MachineGet answered {other:?}"))), + let machine = match layout_of(machine_get(client)?, machine_ws) { + Ok(pulled) => return Ok(pulled), + Err(machine) => machine, }; - match machine.workspaces.iter().find(|w| w.id == machine_ws) { - Some(ws) => { - let mirror = WsMirror { - tabs: ws.tabs.clone(), - active: ws.active_tab, - }; - let session = session_from_tree(ws, &machine.panes); - Ok((machine, mirror, session)) - } - None => { - // The whole tree is already in hand, so the taken names can be read - // straight off it rather than passed down from the main thread. - let taken: Vec<&str> = machine - .workspaces - .iter() - .filter_map(|w| w.name.as_deref()) - .collect(); - let name = tty7_core::core::codename::unique(|n| taken.contains(&n)); - client.call(ControlRequest::WorkspaceCreate { - name: Some(name), - workspace: Some(machine_ws), - })?; - Ok((machine, WsMirror::default(), Session::default())) + // The whole tree is already in hand, so the taken names can be read + // straight off it rather than passed down from the main thread. + let taken: Vec<&str> = machine + .workspaces + .iter() + .filter_map(|w| w.name.as_deref()) + .collect(); + let name = tty7_core::core::codename::unique(|n| taken.contains(&n)); + match client.call(ControlRequest::WorkspaceCreate { + name: Some(name), + workspace: Some(machine_ws), + }) { + Ok(_) => Ok((machine, WsMirror::default(), Session::default())), + // Losing this create is not a failed hydration. Opening a remote + // workspace runs two pulls at once — this one and `start_prime`'s — + // and both create when the tree they read did not hold it yet, so the + // loser is told it already exists. The workspace the create was for is + // on the machine either way, and it may already hold tabs: read the + // tree again and hydrate from what is really there. Treating this as a + // failure left the window empty over a workspace that was fine. + // + // Any refusal is worth the second look, not just "already exists": what + // matters is whether the workspace is there now, and the tree answers + // that better than the error text does. If it still is not there, the + // create's own refusal is the honest error to report — the reread + // happened on its behalf and has nothing of its own to say. + Err(refused) => { + log::debug!( + "workspace {machine_ws} could not be created ({refused}); reading the tree \ + again in case something else created it first" + ); + match machine_get(client) { + Ok(machine) => layout_of(machine, machine_ws).map_err(|_| refused), + Err(_) => Err(refused), + } } } } +fn machine_get(client: &ControlClient) -> io::Result { + match client.call(ControlRequest::MachineGet)? { + ReplyOk::MachineTree(m) => Ok(*m), + other => Err(io::Error::other(format!("MachineGet answered {other:?}"))), + } +} + +/// This workspace's layout as `machine` has it, or the tree handed back +/// untouched when the machine does not hold the workspace at all. +fn layout_of( + machine: Machine, + machine_ws: WorkspaceId, +) -> Result<(Machine, WsMirror, Session), Machine> { + let Some(ws) = machine.workspaces.iter().find(|w| w.id == machine_ws) else { + return Err(machine); + }; + let mirror = WsMirror { + tabs: ws.tabs.clone(), + active: ws.active_tab, + }; + let session = session_from_tree(ws, &machine.panes); + Ok((machine, mirror, session)) +} + fn finish_hydration( cx: &mut App, client_ws: WorkspaceId, @@ -1412,7 +1571,15 @@ fn finish_hydration( let (machine, mirror, session) = match outcome { Ok(pulled) => pulled, Err(e) => { - log::warn!("could not hydrate workspace {client_ws} from its machine: {e}"); + let failures = cx + .default_global::() + .windows + .get(&client_ws) + .map_or(0, |s| s.rehydrate_attempts); + log::log!( + hydration_log_level(failures, log::Level::Warn), + "could not hydrate workspace {client_ws} from its machine: {e}" + ); let _ = owe_rehydration(cx, client_ws, epoch, adopt); return; } @@ -1427,6 +1594,8 @@ fn finish_hydration( let dirty = matches!(state.sync, SyncPhase::Unprimed { dirty: true, .. }); state.informed |= machine_was_empty; state.sync = SyncPhase::Primed(mirror); + // The machine answered, so the next failure starts its backoff over. + state.rehydrate_attempts = 0; // The machine answered, so the explanation has been overtaken by events // and a later outage deserves its own. state.said_why_empty = false; @@ -2150,6 +2319,193 @@ mod tests { }); } + /// The debt an owed pull records is worth nothing without something that + /// pays it. A pull can fail with the link up and healthy — a `MachineGet` + /// past its deadline on a slow link, a create that lost its race — and + /// then no reconnect ever happens to notice, and an empty window has no + /// edit in it to drive a sync. The window sat there empty, with every tab + /// still on the machine, until the app was restarted. + #[gpui::test] + fn an_owed_pull_is_retried_until_it_is_paid_or_superseded(cx: &mut gpui::TestAppContext) { + cx.update(|cx| { + let ws = WorkspaceId::new(); + let epoch = cx + .default_global::() + .windows + .entry(ws) + .or_default() + .epoch; + owe_rehydration(cx, ws, epoch, Adopt::IfEmpty); + assert!( + still_owed(cx, ws, epoch), + "the retry armed for this debt must still recognise it" + ); + + // Paid: the pull landed, so the retry that is still in flight has + // to stand down rather than replay the machine over the window. + cx.default_global::() + .windows + .get_mut(&ws) + .expect("owed above") + .rehydrate = None; + assert!(!still_owed(cx, ws, epoch)); + + // Superseded: a newer hydration owns the window now. + let state = cx + .default_global::() + .windows + .get_mut(&ws) + .expect("owed above"); + state.rehydrate = Some(Adopt::IfEmpty); + state.epoch += 1; + assert!(!still_owed(cx, ws, epoch)); + assert!(still_owed(cx, ws, epoch + 1)); + }); + } + + #[test] + fn the_retry_backs_off_and_settles_at_a_cap() { + let secs = |n| rehydrate_backoff(n).as_secs(); + assert_eq!(secs(1), 2, "the first retry is prompt: a window is empty"); + assert!( + secs(1) < secs(2) && secs(2) < secs(3), + "a machine that keeps refusing must be asked less often, not more" + ); + assert_eq!(secs(REHYDRATE_SETTLED), 30); + assert_eq!( + secs(50), + 30, + "a window left open on an unreachable machine settles at the cap" + ); + } + + /// Once the backoff stops growing the same failure repeats every thirty + /// seconds for as long as the window stays open. Reporting each one at full + /// volume turns one unreachable machine into a log nobody can read past. + #[test] + fn a_standing_failure_stops_shouting_once_the_backoff_settles() { + assert_eq!( + hydration_log_level(1, log::Level::Warn), + log::Level::Warn, + "the first failures are news and must stay news" + ); + assert_eq!( + hydration_log_level(REHYDRATE_SETTLED, log::Level::Warn), + log::Level::Debug + ); + assert_eq!( + hydration_log_level(REHYDRATE_SETTLED, log::Level::Info), + log::Level::Debug, + "the step down is to debug from wherever it started, not to warn" + ); + } + + /// The count paces the retry, so it has to mean "failures in a row". Left + /// standing after the run ends, it makes the next *first* failure wait the + /// cap on an outage that was already over. + #[gpui::test] + fn the_backoff_count_ends_with_the_run_of_failures(cx: &mut gpui::TestAppContext) { + cx.update(|cx| { + let _ = tty7_core::core::config::set_config_dir( + std::env::temp_dir().join(format!("tty7-backoff-count-{}", std::process::id())), + ); + let view = crate::core::session::WindowView::default(); + let ws = view.id; + WorkspaceStore::install_for_test( + cx, + crate::core::session::WindowViews { + views: vec![view], + active: Some(ws), + }, + ); + let unprimed = |cx: &mut App| { + cx.default_global::() + .windows + .entry(ws) + .or_default() + .sync = SyncPhase::Unprimed { + dirty: false, + priming: true, + }; + }; + let attempts = + |cx: &mut App| cx.default_global::().windows[&ws].rehydrate_attempts; + + unprimed(cx); + let epoch = cx.default_global::().windows[&ws].epoch; + for expected in 1..=3 { + unprimed(cx); + owe_rehydration(cx, ws, epoch, Adopt::IfEmpty); + assert_eq!( + attempts(cx), + expected, + "each failure in the run paces the next" + ); + } + + // The machine answered. Whatever it was, it is over. + unprimed(cx); + finish_prime(cx, ws, epoch, Ok(WsMirror::default())); + assert_eq!( + attempts(cx), + 0, + "a prime landing is the machine answering, which is the whole question" + ); + + // Abandoned rather than paid: the user filled the window in + // themselves, so the `Replace` is dropped — and the run is over too. + { + let state = cx + .default_global::() + .windows + .get_mut(&ws) + .unwrap(); + state.rehydrate = Some(Adopt::Replace); + state.rehydrate_attempts = 4; + } + assert!(take_rehydrate(cx, ws, false).is_none()); + assert_eq!( + attempts(cx), + 0, + "a debt nobody owes any more cannot go on pacing the next one" + ); + }); + } + + /// The retry fires on a timer, so the window it was armed for can be gone + /// by the time it runs. It has to notice and stand down — and leave the + /// debt where it is, because a window that is not there is not one that + /// has been paid. + #[gpui::test] + async fn a_retry_that_finds_no_window_stands_down(cx: &mut gpui::TestAppContext) { + let ws = cx.update(|cx| { + crate::ui::windows::WindowRegistry::init(cx); + let ws = WorkspaceId::new(); + let epoch = cx + .default_global::() + .windows + .entry(ws) + .or_default() + .epoch; + owe_rehydration(cx, ws, epoch, Adopt::IfEmpty); + ws + }); + + // Well past the first backoff: the armed retry really runs, rather than + // the test ending while it is still asleep. + cx.executor().advance_clock(rehydrate_backoff(1) * 2); + cx.executor().run_until_parked(); + + cx.update(|cx| { + assert!( + cx.default_global::().windows[&ws] + .rehydrate + .is_some(), + "the debt outlives a retry that found nothing to pay it into" + ); + }); + } + #[gpui::test] fn a_window_that_filled_up_while_owed_keeps_what_it_has(cx: &mut gpui::TestAppContext) { cx.update(|cx| {