diff --git a/crates/bsk-cli/src/cli/doctor.rs b/crates/bsk-cli/src/cli/doctor.rs index f50e63d4..4dbcf3a6 100644 --- a/crates/bsk-cli/src/cli/doctor.rs +++ b/crates/bsk-cli/src/cli/doctor.rs @@ -348,6 +348,10 @@ fn describe_update( } UpdateResult::Skipped => { details.push(match record.skip_reason { + Some(SkipReason::ActiveSessions) => format!( + "bsk {} is available; automatic installation was postponed because agent sessions are active", + record.target_version + ), Some(SkipReason::HostManaged) => format!( "bsk {} is available; this daemon belongs to its terminal or supervisor", record.target_version @@ -365,6 +369,10 @@ fn describe_update( }); if !superseded { hints.push(match record.skip_reason { + Some(SkipReason::ActiveSessions) => { + "the daemon retries on its next update check once sessions are idle" + .to_string() + } Some(SkipReason::HostManaged) => { "run `bsk update`, then restart the daemon in its terminal or supervisor" .to_string() @@ -923,6 +931,27 @@ mod m2_tests { ); } + #[test] + fn active_sessions_postpone_updates_without_installation_advice() { + let record = UpdateRecord::skipped( + UpdateSource::Daemon, + &"0.3.2".parse().unwrap(), + std::path::Path::new("bsk"), + SkipReason::ActiveSessions, + None, + ); + let check = update_check(Some(&record), Some("0.3.1")); + assert!( + check + .detail + .contains("postponed because agent sessions are active") + ); + assert!(!check.detail.contains("cannot write")); + let hint = check.hint.unwrap(); + assert!(hint.contains("next update check")); + assert!(!hint.contains("installer or package manager")); + } + #[test] fn auto_update_check_lists_leftover_files_without_warning() { let leftovers = [update::Leftover { diff --git a/crates/bsk-cli/src/cli/update.rs b/crates/bsk-cli/src/cli/update.rs index b089c608..8f0b93a1 100644 --- a/crates/bsk-cli/src/cli/update.rs +++ b/crates/bsk-cli/src/cli/update.rs @@ -514,20 +514,37 @@ pub(crate) fn install_verified( lock: UpdateLock, cancelled: &dyn Fn() -> bool, ) -> Result { - let binary = match download_candidate_binary(candidate, client) { - Ok(binary) => binary, - Err(err) => { - record.fail(&err, Recovery::Unchanged); - return Err(err); - } - }; + let binary = download_verified(candidate, client, record)?; + install_downloaded(candidate, target, &binary, record, lock, cancelled) +} + +fn download_verified( + candidate: &UpdateCandidate, + client: &reqwest::blocking::Client, + record: &mut UpdateRecord, +) -> Result> { + download_candidate_binary(candidate, client).inspect_err(|err| { + record.fail(err, Recovery::Unchanged); + }) +} + +/// Install an already verified download. The daemon holds session admission +/// closed before entering this step, through self-check and handover or rollback. +pub(crate) fn install_downloaded( + candidate: &UpdateCandidate, + target: &Path, + binary: &[u8], + record: &mut UpdateRecord, + lock: UpdateLock, + cancelled: &dyn Fn() -> bool, +) -> Result { if cancelled() { let err = anyhow::anyhow!("the daemon stopped before the update was installed"); record.fail(&err, Recovery::Unchanged); return Err(err); } record.enter(UpdateStage::Install); - let installed = match install_binary(target, &binary, lock) { + let installed = match install_binary(target, binary, lock) { Ok(installed) => installed, Err(err) => { let recovery = match err.downcast_ref::() { @@ -550,16 +567,13 @@ pub(crate) fn install_verified( Ok(installed) } -/// Daemon-side install with the daemon's own HTTP client. -pub(crate) fn self_install_candidate( +/// Download before closing session admission, using the daemon's HTTP client. +pub(crate) fn self_download_candidate( candidate: &UpdateCandidate, - target: &Path, record: &mut UpdateRecord, - lock: UpdateLock, - cancelled: &dyn Fn() -> bool, -) -> Result { +) -> Result> { let client = update_http_client(ARCHIVE_FETCH_TIMEOUT)?; - install_verified(candidate, target, &client, record, lock, cancelled) + download_verified(candidate, &client, record) } /// Where to turn when bsk cannot write next to its own executable. @@ -603,6 +617,12 @@ pub(crate) enum AutoUpdatePolicy { NotWritable, } +/// Installation may be postponed when a session starts during the download. +pub(crate) enum InstallOutcome { + Installed(T), + PostponedSessions(usize), +} + /// Outcome of one daemon auto-update step (see [`auto_update_step`]). #[derive(Debug, Clone, PartialEq, Eq)] pub(crate) enum AutoUpdateOutcome { @@ -634,7 +654,7 @@ pub(crate) fn auto_update_step( active_sessions: usize, last_attempt: Option<&UpdateRecord>, now_epoch_secs: u64, - install: impl FnOnce(&UpdateCandidate) -> Result, + install: impl FnOnce(&UpdateCandidate) -> Result>, ) -> Result> { let Some(candidate) = candidate else { return Ok(AutoUpdateOutcome::UpToDate); @@ -657,8 +677,12 @@ pub(crate) fn auto_update_step( { return Ok(AutoUpdateOutcome::Deferred { latest, until }); } - let installed = install(candidate)?; - Ok(AutoUpdateOutcome::Installed { latest, installed }) + Ok(match install(candidate)? { + InstallOutcome::Installed(installed) => AutoUpdateOutcome::Installed { latest, installed }, + InstallOutcome::PostponedSessions(sessions) => { + AutoUpdateOutcome::PostponedSessions { latest, sessions } + } + }) } pub fn verify_sha256(bytes: &[u8], expected_hex: &str) -> Result<()> { @@ -908,7 +932,9 @@ fn cached_update_hint( let action = match skipped { Some((SkipReason::NotWritable, record)) => HintAction::Installer(&record.executable), Some((SkipReason::HostManaged, _)) => HintAction::CommandThenHost, - None => HintAction::from_auto_update(cache.auto_update.unwrap_or(auto_update)), + None | Some((SkipReason::ActiveSessions, _)) => { + HintAction::from_auto_update(cache.auto_update.unwrap_or(auto_update)) + } }; Ok(update_hint_for_cache(&cache, current_version, action)) } @@ -2031,7 +2057,7 @@ mod tests { } } - fn no_install(_: &UpdateCandidate) -> Result<()> { + fn no_install(_: &UpdateCandidate) -> Result> { panic!("install must not run") } @@ -2107,7 +2133,7 @@ mod tests { |candidate: &UpdateCandidate| { installs.set(installs.get() + 1); assert_eq!(candidate.latest, LATEST); - Ok("installed") + Ok(InstallOutcome::Installed("installed")) }, ) .unwrap(); @@ -2145,7 +2171,7 @@ mod tests { 0, Some(last), now, - |_| Ok(()), + |_| Ok(InstallOutcome::Installed(())), ) .unwrap() }; @@ -2161,6 +2187,50 @@ mod tests { ); } + #[test] + fn auto_update_step_postpones_a_download_without_failure_backoff() { + let candidate = test_candidate(); + let outcome = auto_update_step( + Some(&candidate), + AutoUpdatePolicy::Install, + 0, + None, + 1_000, + |_| Ok(InstallOutcome::<()>::PostponedSessions(1)), + ) + .unwrap(); + assert_eq!( + outcome, + AutoUpdateOutcome::PostponedSessions { + latest: LATEST, + sessions: 1 + } + ); + let skipped = UpdateRecord::skipped( + UpdateSource::Daemon, + &LATEST, + Path::new("bsk"), + SkipReason::ActiveSessions, + None, + ); + assert_eq!(skipped.retry_blocked_until(&LATEST, 1_001), None); + assert_eq!( + auto_update_step( + Some(&candidate), + AutoUpdatePolicy::Install, + 0, + Some(&skipped), + 1_001, + |_| Ok(InstallOutcome::Installed(())), + ) + .unwrap(), + AutoUpdateOutcome::Installed { + latest: LATEST, + installed: () + } + ); + } + #[test] fn auto_update_step_propagates_install_errors() { let candidate = test_candidate(); @@ -2170,7 +2240,7 @@ mod tests { 0, None, 1_000, - |_| -> Result<()> { bail!("boom") }, + |_| -> Result> { bail!("boom") }, ); assert!(result.is_err()); } diff --git a/crates/bsk-cli/src/cli/update/state.rs b/crates/bsk-cli/src/cli/update/state.rs index b6c06b6a..d79251f2 100644 --- a/crates/bsk-cli/src/cli/update/state.rs +++ b/crates/bsk-cli/src/cli/update/state.rs @@ -55,6 +55,8 @@ pub enum UpdateResult { #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "snake_case")] pub enum SkipReason { + /// A session started while the release was being downloaded. + ActiveSessions, /// The daemon belongs to a terminal or supervisor, which must restart it. HostManaged, /// The directory holding the executable does not accept new files. @@ -148,6 +150,14 @@ impl UpdateRecord { self.save(); } + pub(crate) fn postpone_for_sessions(&mut self) { + self.result = UpdateResult::Skipped; + self.skip_reason = Some(SkipReason::ActiveSessions); + self.retry_after_epoch_secs = None; + self.updated_at_epoch_secs = now_epoch_secs(); + self.save(); + } + pub(crate) fn succeed(&mut self, daemon: Option<(u32, String)>) { self.result = UpdateResult::Succeeded; self.previous_executable = None; diff --git a/crates/bsk-cli/src/daemon/ipc.rs b/crates/bsk-cli/src/daemon/ipc.rs index 7f15fb68..93964fe8 100644 --- a/crates/bsk-cli/src/daemon/ipc.rs +++ b/crates/bsk-cli/src/daemon/ipc.rs @@ -1055,8 +1055,16 @@ pub(super) async fn handle_session_start( } } +#[test] +fn update_admission_error_uses_the_existing_protocol_code() { + let error = map_start_error(StartSessionError::DaemonUpdating); + assert_eq!(error.code, ErrorCode::ProtocolError); + assert!(error.message.contains("retry session start")); +} + fn map_start_error(err: StartSessionError) -> RpcError { let code = match &err { + StartSessionError::DaemonUpdating => ErrorCode::ProtocolError, StartSessionError::NoBrowserConnected => ErrorCode::NoBrowserConnected, StartSessionError::MultipleBrowsersOnline { .. } => ErrorCode::MultipleBrowsersOnline, StartSessionError::BrowserNotFound => ErrorCode::NotFound, diff --git a/crates/bsk-cli/src/daemon/sessions.rs b/crates/bsk-cli/src/daemon/sessions.rs index c4174e01..da1b74de 100644 --- a/crates/bsk-cli/src/daemon/sessions.rs +++ b/crates/bsk-cli/src/daemon/sessions.rs @@ -69,6 +69,8 @@ impl Session { #[derive(Debug, Default)] pub struct SessionRegistry { audit: Option>, + /// Held briefly when reserving a session or claiming an idle installation. + updating: Mutex, inner: Mutex>, /// Operational metadata kept outside the public `Session` wire/domain /// shape so idle enforcement does not break external struct users. @@ -77,6 +79,29 @@ pub struct SessionRegistry { starting: Mutex>, } +#[derive(Debug, PartialEq, Eq, thiserror::Error)] +pub(crate) enum UpdateBlocked { + #[error("{0} agent session(s) active")] + ActiveSessions(usize), + #[error("an automatic update is already installing")] + Updating, +} + +/// Keeps session admission closed without holding a mutex during installation. +#[derive(Debug)] +#[must_use = "hold session admission through installation and handover or rollback"] +pub(crate) struct UpdateGuard(Arc); + +impl Drop for UpdateGuard { + fn drop(&mut self) { + *self + .0 + .updating + .lock() + .unwrap_or_else(|err| err.into_inner()) = false; + } +} + impl SessionRegistry { pub fn with_audit(audit: Arc) -> Self { Self { @@ -88,6 +113,19 @@ impl SessionRegistry { Self::default() } + pub(crate) fn try_begin_update(self: &Arc) -> Result { + let mut updating = self.updating.lock().expect("session admission poisoned"); + if *updating { + return Err(UpdateBlocked::Updating); + } + let active = self.len(); + if active > 0 { + return Err(UpdateBlocked::ActiveSessions(active)); + } + *updating = true; + Ok(UpdateGuard(Arc::clone(self))) + } + pub fn snapshot(&self) -> Vec { self.inner .lock() @@ -126,8 +164,8 @@ impl SessionRegistry { /// Reserve a fresh, collision-free [`SessionId`] under the registry /// lock by inserting a placeholder [`Session`]. Returns `None` - /// after `max_attempts` failed random draws so callers can surface - /// a deterministic error instead of looping forever. + /// while an update is installing, or after `max_attempts` failed random + /// draws so callers can surface a deterministic error instead of looping forever. /// /// The 4-letter id space is `26^4 = 456_976`; with even a small /// number of live sessions the birthday probability of a collision @@ -144,6 +182,20 @@ impl SessionRegistry { max_attempts: u32, now_ms_fn: impl Fn() -> i64, ) -> Option { + self.reserve_id_for_start(browser_id, max_attempts, now_ms_fn) + .ok() + } + + fn reserve_id_for_start( + &self, + browser_id: BrowserId, + max_attempts: u32, + now_ms_fn: impl Fn() -> i64, + ) -> Result { + let updating = self.updating.lock().expect("session admission poisoned"); + if *updating { + return Err(StartSessionError::DaemonUpdating); + } let mut guard = self.inner.lock().expect("session registry poisoned"); for _ in 0..max_attempts { let candidate = SessionId::random(); @@ -164,9 +216,9 @@ impl SessionRegistry { .lock() .expect("session activity registry poisoned") .insert(candidate.clone(), Instant::now()); - return Some(candidate); + return Ok(candidate); } - None + Err(StartSessionError::IdExhausted) } /// Replace a previously [`reserve_id`](Self::reserve_id) placeholder @@ -367,6 +419,8 @@ fn next_rpc_id(prefix: &str) -> RpcId { #[derive(Debug, thiserror::Error)] pub enum StartSessionError { + #[error("the daemon is installing an update; retry session start after the update finishes")] + DaemonUpdating, #[error("no browser is currently connected")] NoBrowserConnected, #[error("more than one browser is online — pass --browser ")] @@ -407,6 +461,7 @@ pub enum StartSessionError { impl StartSessionError { pub fn code(&self) -> &'static str { match self { + StartSessionError::DaemonUpdating => "protocol_error", StartSessionError::NoBrowserConnected => "no_browser_connected", StartSessionError::MultipleBrowsersOnline { .. } => "multiple_browsers_online", StartSessionError::BrowserNotFound => "not_found", @@ -536,9 +591,11 @@ pub(crate) async fn start_session_recoverable( if client.is_unresponsive() { return Err(StartSessionError::ExtensionUnresponsive); } - let session_id = sessions - .reserve_id(client.id.clone(), SESSION_ID_MAX_RESERVE_ATTEMPTS, now_ms) - .ok_or(StartSessionError::IdExhausted)?; + let session_id = sessions.reserve_id_for_start( + client.id.clone(), + SESSION_ID_MAX_RESERVE_ATTEMPTS, + now_ms, + )?; if preserve_cleanup { sessions.starting.lock().unwrap().insert(session_id.clone()); } @@ -1120,6 +1177,67 @@ mod link_tests { } } + #[tokio::test] + async fn update_admission_counts_pending_and_completed_starts() { + let harness = Harness::new(); + let (client, mut rx) = harness.connect(Liveness::default()); + let start = harness.start(Duration::from_secs(30), true); + let request = next_request(&mut rx).await; + assert!(matches!( + harness.sessions.try_begin_update(), + Err(super::UpdateBlocked::ActiveSessions(1)) + )); + answer( + &client, + &request, + ResponseBody::Ok(serde_json::json!({"agent_window_id": 71})), + ); + let session = start.await.unwrap().unwrap(); + assert!(matches!( + harness.sessions.try_begin_update(), + Err(super::UpdateBlocked::ActiveSessions(1)) + )); + harness.sessions.remove(&session.id); + let _admission = harness + .sessions + .try_begin_update() + .expect("stopped sessions do not block updates"); + } + + #[tokio::test] + async fn update_admission_rejects_starts_before_dispatch_and_reopens_on_drop() { + for recoverable in [false, true] { + let harness = Harness::new(); + let (client, mut rx) = harness.connect(Liveness::default()); + let admission = harness.sessions.try_begin_update().unwrap(); + assert!(matches!( + harness.sessions.try_begin_update(), + Err(super::UpdateBlocked::Updating) + )); + let error = start_error(harness.start(Duration::from_secs(30), recoverable)).await; + assert!(matches!(error, StartSessionError::DaemonUpdating)); + assert_eq!(error.code(), "protocol_error"); + assert!(error.to_string().contains("retry session start")); + assert!( + rx.try_recv().is_err(), + "a refused start must not create an Agent Window" + ); + assert!( + harness.sessions.is_empty(), + "a refused start must not reserve an id" + ); + drop(admission); + let start = harness.start(Duration::from_secs(30), recoverable); + let request = next_request(&mut rx).await; + answer( + &client, + &request, + ResponseBody::Ok(serde_json::json!({"agent_window_id": 71})), + ); + start.await.unwrap().unwrap(); + } + } + fn answer(client: &BrowserClient, request: &RequestFrame, body: ResponseBody) { let delivered = client.pending.lock().unwrap().resolve(ResponseFrame { id: request.id.clone(), diff --git a/crates/bsk-cli/src/daemon/start.rs b/crates/bsk-cli/src/daemon/start.rs index 0976091d..1d110bb5 100644 --- a/crates/bsk-cli/src/daemon/start.rs +++ b/crates/bsk-cli/src/daemon/start.rs @@ -498,7 +498,7 @@ enum Transaction { /// Installed and checked; no replacement started yet. Installed(Prepared), /// A replacement has been started and waits for the daemon lock. - HandingOver(handover::Pending), + HandingOver(handover::Pending, super::sessions::UpdateGuard), } impl Transaction { @@ -507,10 +507,11 @@ impl Transaction { use crate::cli::update::state::Recovery; let reason = "the daemon stopped before handing over"; match self { - Transaction::HandingOver(pending) => handover::abandon(pending, reason), + Transaction::HandingOver(pending, _admission) => handover::abandon(pending, reason), Transaction::Installed(Prepared { installed, mut record, + admission: _admission, }) => { let recovery = installed.roll_back(|| Recovery::Restored { daemon_serving: false, @@ -766,7 +767,7 @@ fn serve(cfg: &DaemonConfig, resumed: Option<&mut UpdateRecord>) -> Result { + (Ok(StopReason::Handover), Some(Transaction::HandingOver(pending, _admission))) => { Ok(Stopped::HandOver(Box::new(pending))) } (reason, transaction) => { @@ -1012,8 +1013,8 @@ fn spawn_update_check_task( &cache_path, policy == update::AutoUpdatePolicy::Install, )?; - // The session gate is read after the fetch, as late - // as possible before the binary gets replaced. + // Avoid downloading while sessions are active. The installer + // checks again and closes admission after the download. let active_sessions = state.sessions.len(); let last_attempt = update::state::current(); let outcome = update::auto_update_step( @@ -1028,6 +1029,7 @@ fn spawn_update_check_task( exe_path.as_deref(), &transaction, &stopping, + &state.sessions, ) }, )?; @@ -1143,6 +1145,7 @@ fn spawn_update_check_task( struct Prepared { installed: crate::cli::update::Installed, record: crate::cli::update::state::UpdateRecord, + admission: super::sessions::UpdateGuard, } /// Install and self-check `candidate` under the update lock, leaving the @@ -1153,7 +1156,8 @@ fn prepare_handover( exe: Option<&Path>, transaction: &TransactionSlot, stopping: &std::sync::atomic::AtomicBool, -) -> Result<()> { + sessions: &Arc, +) -> Result> { use crate::cli::update::{ self, state::{UpdateLock, UpdateRecord, UpdateSource}, @@ -1162,10 +1166,24 @@ fn prepare_handover( let lock = UpdateLock::try_acquire(exe)?; let mut record = UpdateRecord::start(UpdateSource::Daemon, &candidate.latest, exe); record.save(); + let binary = update::self_download_candidate(candidate, &mut record)?; + let admission = match sessions.try_begin_update() { + Ok(admission) => admission, + Err(super::sessions::UpdateBlocked::ActiveSessions(active)) => { + record.postpone_for_sessions(); + return Ok(update::InstallOutcome::PostponedSessions(active)); + } + Err(err) => return Err(err.into()), + }; let cancelled = || stopping.load(std::sync::atomic::Ordering::SeqCst); - let installed = update::self_install_candidate(candidate, exe, &mut record, lock, &cancelled)?; - *lock_slot(transaction) = Some(Transaction::Installed(Prepared { installed, record })); - Ok(()) + let installed = + update::install_downloaded(candidate, exe, &binary, &mut record, lock, &cancelled)?; + *lock_slot(transaction) = Some(Transaction::Installed(Prepared { + installed, + record, + admission, + })); + Ok(update::InstallOutcome::Installed(())) } /// Spawn the replacement daemon from the new executable, on the port this @@ -1183,6 +1201,7 @@ fn start_replacement( let Some(Transaction::Installed(Prepared { installed, mut record, + admission, })) = slot.take() else { return false; @@ -1202,13 +1221,16 @@ fn start_replacement( replacement = child_pid, "replacement daemon started; handing over once it is ready" ); - *slot = Some(Transaction::HandingOver(handover::Pending { - child_pid, - child, - port: ws_port, - installed, - record, - })); + *slot = Some(Transaction::HandingOver( + handover::Pending { + child_pid, + child, + port: ws_port, + installed, + record, + }, + admission, + )); true } Err(err) => { diff --git a/crates/bsk-cli/tests/auto_update_handover.rs b/crates/bsk-cli/tests/auto_update_handover.rs index 3accf2c9..a56aae37 100644 --- a/crates/bsk-cli/tests/auto_update_handover.rs +++ b/crates/bsk-cli/tests/auto_update_handover.rs @@ -21,7 +21,9 @@ use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::thread; use std::time::{Duration, Instant}; +use futures_util::{SinkExt, StreamExt}; use sha2::{Digest, Sha256}; +use tokio_tungstenite::tungstenite::{Message, client::IntoClientRequest}; const EXE: &str = if cfg!(windows) { "bsk.exe" } else { "bsk" }; const MARKER: &[u8] = b"auto-update-handover-fixture"; @@ -31,6 +33,7 @@ const MARKER: &[u8] = b"auto-update-handover-fixture"; struct ReleaseServer { url: String, downloads: Arc, + hold_archive: Arc, stop: Arc, worker: Option>, } @@ -58,9 +61,11 @@ impl ReleaseServer { .unwrap(); let stop = Arc::new(AtomicBool::new(false)); let downloads = Arc::new(AtomicUsize::new(0)); + let hold_archive = Arc::new(AtomicBool::new(false)); let worker = { let stop = Arc::clone(&stop); let downloads = Arc::clone(&downloads); + let hold_archive = Arc::clone(&hold_archive); let archive_request = format!("GET /bsk{suffix} "); thread::spawn(move || { while !stop.load(Ordering::SeqCst) { @@ -86,6 +91,9 @@ impl ReleaseServer { } let body = if request.starts_with(archive_request.as_bytes()) { downloads.fetch_add(1, Ordering::SeqCst); + while hold_archive.load(Ordering::SeqCst) && !stop.load(Ordering::SeqCst) { + thread::sleep(Duration::from_millis(10)); + } thread::sleep(archive_delay); &archive } else { @@ -103,6 +111,7 @@ impl ReleaseServer { Self { url: format!("{base}/version.json"), downloads, + hold_archive, stop, worker: Some(worker), } @@ -455,6 +464,146 @@ fn unused_port() -> u16 { .port() } +async fn start_test_session( + fixture: &Fixture, + port: u16, +) -> ( + tokio_tungstenite::WebSocketStream>, + serde_json::Value, +) { + let browser_id = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"; + let mut request = format!("ws://127.0.0.1:{port}/") + .into_client_request() + .unwrap(); + request.headers_mut().insert( + "Origin", + format!("chrome-extension://{browser_id}").parse().unwrap(), + ); + let (mut browser, _) = tokio_tungstenite::connect_async(request).await.unwrap(); + browser + .send(Message::Text( + serde_json::json!({ + "id": "handshake", "method": "system.handshake", "params": { + "client": "browser-skill-extension", "version": env!("CARGO_PKG_VERSION"), + "protocol_version": bsk::daemon::state::PROTOCOL_VERSION, + "instance_id": browser_id, "browser": {"name": "chrome", "version": "149"}, + "min_compatible_protocol": "1.0", "label": "update-race" + } + }) + .to_string(), + )) + .await + .unwrap(); + let handshake = browser.next().await.unwrap().unwrap().into_text().unwrap(); + let handshake: serde_json::Value = serde_json::from_str(&handshake).unwrap(); + assert!(handshake.get("result").is_some(), "{handshake}"); + + let mut start = fixture.command(); + start.args([ + "--json", + "session", + "start", + "--browser", + browser_id, + "--no-focus", + ]); + let started = tokio::task::spawn_blocking(move || start.output().unwrap()); + let request = tokio::time::timeout(Duration::from_secs(10), async { + loop { + let frame = browser.next().await.unwrap().unwrap(); + if let Message::Text(text) = frame { + let request: serde_json::Value = serde_json::from_str(&text).unwrap(); + if request["method"] == "tool.session_start" { + break request; + } + } + } + }) + .await + .unwrap(); + let session_id = request["params"]["session_id"].clone(); + browser + .send(Message::Text( + serde_json::json!({ + "id": request["id"], "result": {"session_id": session_id, "agent_window_id": 71} + }) + .to_string(), + )) + .await + .unwrap(); + let output = started.await.unwrap(); + assert!( + output.status.success(), + "{}", + String::from_utf8_lossy(&output.stderr) + ); + (browser, session_id) +} + +#[test] +fn a_session_started_during_download_postpones_installation() { + let fixture = Fixture::new(mislabelled_bsk); + fixture.server.hold_archive.store(true, Ordering::SeqCst); + let port = unused_port(); + let pid = fixture.start_daemon(port); + fixture.wait_for("the blocked archive request", || { + fixture.server.downloads.load(Ordering::SeqCst) == 1 + }); + + tokio::runtime::Runtime::new().unwrap().block_on(async { + let (_browser, session_id) = start_test_session(&fixture, port).await; + + fixture.server.hold_archive.store(false, Ordering::SeqCst); + fixture.wait_for("the update decision after the download", || { + fixture + .record() + .is_some_and(|record| record["result"] != "in_progress") + }); + let record = fixture.record().unwrap(); + assert_eq!(record["result"], "skipped", "{record}"); + assert_eq!(record["skip_reason"], "active_sessions", "{record}"); + assert_ne!(record["stage"], "install", "{record}"); + assert!(record.get("retry_after_epoch_secs").is_none(), "{record}"); + assert!(fixture.installed() == fixture.original); + assert_eq!(fixture.info().unwrap()["pid"], pid); + let status = fixture + .command() + .args(["--json", "status"]) + .output() + .unwrap(); + assert!(status.status.success()); + let status: serde_json::Value = serde_json::from_slice(&status.stdout).unwrap(); + assert!( + status["sessions"] + .as_array() + .unwrap() + .iter() + .any(|s| s["session_id"] == session_id), + "{status}" + ); + }); +} + +#[test] +fn a_failed_auto_update_reopens_session_admission() { + for release in [mislabelled_bsk, release_whose_daemon_fails] { + let fixture = Fixture::new(release); + let port = unused_port(); + let pid = fixture.start_daemon(port); + fixture.wait_for("rollback after the update failed", || { + fixture + .record() + .is_some_and(|record| record["result"] == "failed") + && fixture.info().is_some_and(|info| info["pid"] == pid) + }); + assert!(fixture.installed() == fixture.original); + tokio::runtime::Runtime::new().unwrap().block_on(async { + let (_browser, _session) = start_test_session(&fixture, port).await; + fixture.status_succeeds(); + }); + } +} + #[test] fn auto_update_exits_only_after_the_new_daemon_serves() { let fixture = Fixture::new(newer_bsk);