Merge pull request #346 from Tencent/fix/336-windows-in-place-update

fix(update): keep Windows daemons serving across self-update
This commit is contained in:
Zhang GH
2026-09-29 17:14:13 +08:00
committed by GitHub
20 changed files with 4947 additions and 666 deletions
+5 -2
View File
@@ -80,7 +80,10 @@ jobs:
node scripts/check-crate-skill.mjs
- name: Run Windows update tests
run: cargo test -p bsk --lib --locked cli::update::tests
run: |
cargo test -p bsk --lib --locked cli::update::
if ($LASTEXITCODE -ne 0) { exit $LASTEXITCODE }
cargo test -p bsk --test auto_update_policy --locked
- name: Run Windows cancellation process tests
run: cargo test -p bsk --test windows_parent_cancel --locked
@@ -96,7 +99,7 @@ jobs:
windows-daemon-lifecycle:
name: Windows independent daemon and update lifecycle
runs-on: windows-latest
timeout-minutes: 20
timeout-minutes: 45
steps:
- uses: actions/checkout@v6
+44
View File
@@ -15,6 +15,50 @@ Starting from 0.2.0, CLI / Extension / DSH Plugin share the same version number.
at once instead of hanging until their timeout; inputs, transfers and tab borrows keep
their unknown-outcome errors. `bsk doctor` and `bsk browsers` flag a connected
extension that has stopped sending heartbeats.
- A failed auto-update no longer leaves the browser disconnected
([#336](https://github.com/Tencent/BrowserSkill/issues/336)). The daemon
checks that the new executable reports the release's version, starts a
daemon from it, and exits only after a daemon of that version serves the
same port. If the new daemon exits or is not ready within 20 seconds, it is
stopped, the previous executable is put back and the running daemon serves
again on the same port; the release is retried after 6 hours. The update
record says a daemon is serving only once it has published `daemon.json`.
This applies on all platforms.
- Windows self-update no longer depends on a detached script: the running
`bsk.exe` is renamed aside and the new one takes its place.
- `bsk update` installs and checks the new executable before stopping the
daemon. If the restarted daemon is not ready in time, it is stopped, the
previous executable is put back and the previous version restarted. It
restarts the daemon from the installed path, which Linux no longer reports
as the current executable once it is replaced.
- A daemon that stops during an auto-update, for example because it went
idle while the release was downloading or being checked, installs nothing
or puts the previous executable back, and records the attempt as failed.
- `bsk update` leaves a background daemon running when it could not start
one again, such as inside a sandbox whose Windows Job forbids breakaway
(`"daemon": "left_running"`), instead of stopping it.
- A daemon keeps using the executable path it started from, so it can update
again after a rolled-back update on Linux.
- Updates of one executable are serialized through a lock file next to it, so
daemons or `bsk update` runs with different bsk homes cannot overwrite each
other's installation or rollback.
### Changed
- Daemons started with `--foreground` no longer install updates or replace
themselves with a detached process, on any platform. They log a new version
once, and the CLI hint suggests `bsk update` followed by a restart in the
daemon's terminal or supervisor. `bsk update` leaves such a daemon running
instead of stopping it and starting a background daemon in its place
(`"daemon": "left_to_host"` in `--json` output).
- `bsk update` restarts a background daemon on the port it served.
- A daemon that cannot write next to its executable reports new releases
instead of installing them, and the CLI hint points to the installer.
- Each update attempt, including the stage and error of a failure, is kept in
`update-state.json` in the bsk home and shown by `bsk doctor`.
- `bsk update --json` reports `"status": "updated"` on Windows too; the
`"staged"` status is gone.
- README documents `BSK_AUTO_UPDATE=off`.
## [0.3.1] - 2026-09-23
+5 -1
View File
@@ -248,7 +248,11 @@ Finish active browser tasks, then update the CLI:
bsk update --yes
```
For the default local setup, this restarts a running daemon when an update is installed. If Windows reports a staged update, wait for replacement to finish. If you use the installer to replace the binary, restart the daemon afterwards with `bsk daemon restart`.
For the default local setup, this installs the new executable and checks that it runs while the daemon keeps serving, then restarts the daemon from it on the same port. If the restarted daemon does not become ready, the previous executable is put back and the daemon is restarted from it. A daemon started with `--foreground` is left running on the previous version; restart it in its terminal or supervisor. So is a background daemon when `bsk update` runs where it could not start one again, such as inside a sandbox whose Windows Job forbids breakaway; restart it with `bsk daemon restart` outside the sandbox. If you use the installer to replace the binary, restart the daemon afterwards with `bsk daemon restart`.
A daemon that `bsk` started in the background also checks for a new release every 30 minutes and, while no agent session is active, installs it the same way. It then starts a daemon from the new executable and exits only once a daemon of the new version serves the same port. If the new daemon exits or is not ready within 20 seconds, it is stopped, and the running daemon puts the previous executable back, serves again on the same port, and retries that release after 6 hours. Browser connections drop for a moment during a handover and reconnect. A daemon that stops for another reason during an update, such as going idle, puts the previous executable back first. A daemon started with `--foreground` belongs to its terminal or supervisor, so it only reports new releases: run `bsk update`, then restart the daemon there. If bsk cannot write next to its executable, it only reports new releases; update it with the installer or package manager you used. Set `BSK_AUTO_UPDATE=off` to disable daemon-side auto-update while keeping manual `bsk update` available.
`bsk doctor` shows the last update attempt; `update-state.json` in the bsk home keeps its stage, result and error. Updates of one executable run one at a time, even from different bsk homes, through `.bsk.update.lock` (`.bsk.exe.update.lock` on Windows), which stays next to it. Until an update is confirmed, the previous executable stays next to the new one as `.bsk.old-*` (`.bsk.exe.old-*` on Windows); one that is still running remains until its process exits, and a later daemon start removes it. On Windows, self-update renames the running executable, which NTFS supports; where the file system refuses, the update fails and leaves the installation unchanged. bsk supports Windows 10 and Windows Server 2016 or later.
Update the extension through its browser store. Update the DSH plugin separately, then restart its profile:
+5 -1
View File
@@ -248,7 +248,11 @@ BrowserSkill 没有必须使用的云服务,也不收集产品遥测。自动
bsk update --yes
```
默认本地配置下,安装更新后会重启正在运行的 daemon。Windows 如果提示更新已暂存,请等待替换完成。如果使用安装脚本替换了二进制,请随后执行 `bsk daemon restart`。
默认本地配置下,这条命令会在 daemon 继续服务的同时安装新版本并检查它能否运行,然后用新版本在原端口重启 daemon。重启后的 daemon 没有就绪时,会放回旧版本并用它重启 daemon。通过 `--foreground` 启动的 daemon 不会被停止,会继续运行旧版本,请在它所在的终端或进程管理器中重启。如果 `bsk update` 运行在无法重新启动后台 daemon 的环境里(例如 Windows Job 禁止 breakaway 的沙盒),后台 daemon 同样不会被停止,请在沙盒外执行 `bsk daemon restart`。如果使用安装脚本替换了二进制,请随后执行 `bsk daemon restart`。
由 `bsk` 在后台启动的 daemon 还会每 30 分钟检查一次新版本。没有 Agent 会话时,它按同样方式安装新版本,再用新版本启动一个 daemon,确认新版本的 daemon 已在原端口服务后才退出。新 daemon 退出,或 20 秒内没有就绪时,它会被停止,当前 daemon 放回旧版本,在原端口继续服务,6 小时后再尝试这个版本。交接期间浏览器连接会短暂断开并自动重连。更新过程中 daemon 因其他原因(例如空闲超时)停止时,会先放回旧版本再退出。通过 `--foreground` 启动的 daemon 归所在终端或进程管理器管理,只提示新版本:执行 `bsk update`,再在那里重启 daemon。bsk 无法在可执行文件所在目录写入时,也只提示新版本,请用原来的安装脚本或包管理器升级。设置 `BSK_AUTO_UPDATE=off` 可关闭 daemon 自动升级,手动 `bsk update` 仍然可用。
`bsk doctor` 会显示最近一次更新的结果;bsk home 下的 `update-state.json` 记录了它的阶段、结果和错误原因。同一个可执行文件的更新会依次进行,即使来自不同的 bsk home 也是如此,靠的是它旁边的 `.bsk.update.lock`(Windows 上为 `.bsk.exe.update.lock`),这个文件会一直保留。更新确认成功前,旧版本会以 `.bsk.old-*`(Windows 上为 `.bsk.exe.old-*`)的名字留在新版本旁边;仍在运行的旧版本会保留到对应进程退出,之后再启动 daemon 时会清理。Windows 上的自更新依赖重命名正在运行的可执行文件,NTFS 支持这一操作;文件系统不支持时,更新会失败,已安装的版本保持不变。bsk 支持 Windows 10 和 Windows Server 2016 及以上版本。
扩展通过浏览器商店更新。DSH 插件需要单独更新,完成后重启对应 profile:
+358
View File
@@ -12,6 +12,10 @@ use crate::cli::browser_wait::{
};
use crate::cli::ensure_daemon::{AUTO_START_DISABLED_HINT, auto_start_enabled, ensure_daemon};
use crate::cli::status::Output;
use crate::cli::update::{
self,
state::{Recovery, SkipReason, UpdateRecord, UpdateResult, UpdateSource, UpdateStage},
};
use crate::daemon::info::DaemonInfo;
use crate::daemon::paths;
use crate::daemon::probe::{self, Probe};
@@ -166,6 +170,8 @@ enum DaemonState {
Verified {
status: StatusResult,
local_identity_error: Option<String>,
/// Owned by a terminal or supervisor, which must restart it.
host_managed: bool,
},
}
@@ -176,6 +182,16 @@ impl DaemonState {
_ => None,
}
}
fn host_managed(&self) -> bool {
matches!(
self,
DaemonState::Verified {
host_managed: true,
..
}
)
}
}
fn collect_checks(state: DaemonState) -> Vec<CheckResult> {
@@ -185,11 +201,195 @@ fn collect_checks(state: DaemonState) -> Vec<CheckResult> {
check_daemon_running(&state),
check_daemon_management(&state),
check_version_compatible(state.status()),
check_auto_update(state.status(), state.host_managed()),
check_extension_connected(state.status()),
check_browsers_protocol_compatible(state.status()),
]
}
/// An attempt still `in_progress` after this long was interrupted.
const UPDATE_STALLED_AFTER_SECS: u64 = 10 * 60;
fn check_auto_update(status: Option<&StatusResult>, host_managed: bool) -> CheckResult {
let leftovers = std::env::current_exe()
.map(|exe| update::update_leftovers(&exe))
.unwrap_or_default();
auto_update_check(
update::state::current().as_ref(),
status.map(|status| (status.daemon_version.as_str(), host_managed)),
env!("CARGO_PKG_VERSION"),
&leftovers,
update::now_epoch_secs(),
)
}
/// Report the recorded update attempt as it happened, rather than inferring
/// it from which versions exist. `daemon` is the running daemon's version and
/// whether its terminal or supervisor owns it.
fn auto_update_check(
record: Option<&UpdateRecord>,
daemon: Option<(&str, bool)>,
installed_version: &str,
leftovers: &[update::Leftover],
now: u64,
) -> CheckResult {
let name = "auto-update";
let mut details = Vec::new();
let mut hints = Vec::new();
if let Some((daemon_version, host_managed)) =
daemon.filter(|(version, _)| *version != installed_version)
{
details.push(format!(
"the daemon runs bsk {daemon_version}, the installed bsk is {installed_version}"
));
hints.push(if host_managed {
"restart the daemon in its terminal or supervisor to run the installed version"
.to_string()
} else {
"restart the daemon with `bsk daemon restart` to run the installed version".to_string()
});
}
match record {
None => details.push("no update attempt recorded".to_string()),
Some(record) => describe_update(record, installed_version, now, &mut details, &mut hints),
}
if !leftovers.is_empty() {
let names = leftovers
.iter()
.map(|leftover| leftover.path.display().to_string())
.collect::<Vec<_>>()
.join(", ");
details.push(format!(
"kept from earlier updates until nothing runs them: {names}"
));
}
let detail = details.join("; ");
if hints.is_empty() {
CheckResult::ok(name, detail)
} else {
CheckResult::warn(name, detail, hints.join("; "))
}
}
fn describe_update(
record: &UpdateRecord,
installed_version: &str,
now: u64,
details: &mut Vec<String>,
hints: &mut Vec<String>,
) {
let attempt = format!(
"{} {} -> {} ({})",
match record.source {
UpdateSource::Daemon => "auto-update",
UpdateSource::Command => "`bsk update`",
},
record.from_version,
record.target_version,
ago(now, record.updated_at_epoch_secs),
);
let stage = record.stage.map_or("", |stage| match stage {
UpdateStage::Download => " while downloading",
UpdateStage::Install => " while installing",
UpdateStage::Handover => " while handing over to the new daemon",
UpdateStage::Restart => " while restarting the daemon",
});
// An attempt at a version that is installed by now needs no action.
let superseded = semver::Version::parse(&record.target_version)
.ok()
.zip(semver::Version::parse(installed_version).ok())
.is_some_and(|(target, installed)| target <= installed);
let retry = "run `bsk update` to try again now, and check `bsk logs` for details";
match record.result {
UpdateResult::Succeeded => details.push(format!("{attempt} succeeded")),
UpdateResult::InProgress
if now.saturating_sub(record.updated_at_epoch_secs) < UPDATE_STALLED_AFTER_SECS =>
{
details.push(format!("{attempt} in progress{stage}"));
}
UpdateResult::InProgress => {
details.push(format!(
"{attempt} stopped{stage}; the process running it exited"
));
if !superseded {
hints.push(retry.to_string());
}
}
UpdateResult::Failed => {
let error = record.error.as_deref().unwrap_or("unknown error");
let recovery = match &record.recovery {
Some(Recovery::Unchanged) | None => "nothing was changed".to_string(),
Some(Recovery::Restored {
daemon_serving: true,
}) => "the previous version was restored and a daemon kept serving".to_string(),
Some(Recovery::Restored {
daemon_serving: false,
}) => {
hints.push("start the daemon with `bsk daemon start`".to_string());
"the previous version was restored, but no daemon was left serving".to_string()
}
Some(Recovery::RestoreFailed { action }) => {
hints.push(action.clone());
"the previous executable could not be put back".to_string()
}
};
let mut detail = format!("{attempt} failed{stage}: {error}; {recovery}");
if let Some(retry_after) = record.retry_after_epoch_secs.filter(|at| *at > now) {
detail.push_str(&format!(
"; the daemon retries in {}",
duration(retry_after - now)
));
}
details.push(detail);
if !superseded {
hints.push(retry.to_string());
}
}
UpdateResult::Skipped => {
details.push(match record.skip_reason {
Some(SkipReason::HostManaged) => format!(
"bsk {} is available; this daemon belongs to its terminal or supervisor",
record.target_version
),
Some(SkipReason::NotWritable) | None => format!(
"bsk {} is available, but bsk cannot write next to {}{}",
record.target_version,
record.executable.display(),
record
.error
.as_deref()
.map(|error| format!(": {error}"))
.unwrap_or_default()
),
});
if !superseded {
hints.push(match record.skip_reason {
Some(SkipReason::HostManaged) => {
"run `bsk update`, then restart the daemon in its terminal or supervisor"
.to_string()
}
Some(SkipReason::NotWritable) | None => {
update::installer_hint(&record.executable)
}
});
}
}
}
}
fn ago(now: u64, then: u64) -> String {
format!("{} ago", duration(now.saturating_sub(then)))
}
fn duration(secs: u64) -> String {
match secs {
0..60 => format!("{secs}s"),
60..3600 => format!("{}m", secs / 60),
3600..86400 => format!("{}h", secs / 3600),
_ => format!("{}d", secs / 86400),
}
}
fn current_state(browser_wait: Duration) -> DaemonState {
let params = bsk_protocol::StatusParams {
wait_for_browser_ms: wait_for_browser_ms(browser_wait),
@@ -201,6 +401,7 @@ fn current_state(browser_wait: Duration) -> DaemonState {
.require_local_pid()
.err()
.map(|err| format!("{err:#}")),
host_managed: daemon.info.host_managed,
status: daemon.status,
},
Ok(Probe::Absent(Some(info))) => DaemonState::NoListener(info),
@@ -544,11 +745,168 @@ mod m2_tests {
}
}
fn attempt(result: UpdateResult) -> UpdateRecord {
UpdateRecord {
result,
from_version: "0.3.1".into(),
target_version: "0.4.0".into(),
updated_at_epoch_secs: 1_000,
..UpdateRecord::start(
UpdateSource::Daemon,
&semver::Version::new(0, 4, 0),
std::path::Path::new("/home/u/.local/bin/bsk"),
)
}
}
fn update_check(record: Option<&UpdateRecord>, daemon: Option<&str>) -> CheckResult {
auto_update_check(record, daemon.map(|v| (v, false)), "0.3.1", &[], 1_120)
}
#[test]
fn auto_update_check_reports_success_and_version_skew() {
let none = update_check(None, Some("0.3.1"));
assert_eq!(none.status, CheckStatus::Ok);
assert!(none.detail.contains("no update attempt recorded"));
let succeeded = update_check(Some(&attempt(UpdateResult::Succeeded)), Some("0.3.1"));
assert_eq!(succeeded.status, CheckStatus::Ok);
assert!(
succeeded
.detail
.contains("auto-update 0.3.1 -> 0.4.0 (2m ago) succeeded"),
"{}",
succeeded.detail
);
let skewed = update_check(None, Some("0.3.0"));
assert_eq!(skewed.status, CheckStatus::Warning);
assert!(skewed.detail.contains("the daemon runs bsk 0.3.0"));
assert!(skewed.hint.unwrap().contains("bsk daemon restart"));
assert!(!has_failures(&[update_check(None, Some("0.3.0"))]));
// A daemon its terminal or supervisor owns is restarted there.
let host = auto_update_check(None, Some(("0.3.0", true)), "0.3.1", &[], 0);
let hint = host.hint.unwrap();
assert!(hint.contains("its terminal or supervisor"), "{hint}");
assert!(!hint.contains("bsk daemon restart"), "{hint}");
}
#[test]
fn auto_update_check_explains_a_failure_and_its_recovery() {
let mut failed = attempt(UpdateResult::Failed);
failed.stage = Some(UpdateStage::Handover);
failed.error = Some("the replacement daemon (pid 7) exited with exit status: 3".into());
failed.recovery = Some(Recovery::Restored {
daemon_serving: true,
});
failed.retry_after_epoch_secs = Some(1_120 + 2 * 3600);
let check = update_check(Some(&failed), Some("0.3.1"));
assert_eq!(check.status, CheckStatus::Warning);
for text in [
"failed while handing over to the new daemon",
"exited with exit status: 3",
"the previous version was restored and a daemon kept serving",
"the daemon retries in 2h",
] {
assert!(check.detail.contains(text), "{text}: {}", check.detail);
}
assert!(check.hint.unwrap().contains("bsk update"));
failed.recovery = Some(Recovery::RestoreFailed {
action: "stop bsk, then move /a to /b".into(),
});
let check = update_check(Some(&failed), Some("0.3.1"));
assert!(check.hint.unwrap().contains("stop bsk, then move /a to /b"));
// Nothing is left to do once that version is installed.
let superseded = auto_update_check(
Some(&UpdateRecord {
recovery: Some(Recovery::Unchanged),
..failed.clone()
}),
Some(("0.4.0", false)),
"0.4.0",
&[],
1_120,
);
assert_eq!(superseded.status, CheckStatus::Ok, "{}", superseded.detail);
}
#[test]
fn auto_update_check_flags_an_interrupted_attempt() {
let mut interrupted = attempt(UpdateResult::InProgress);
interrupted.stage = Some(UpdateStage::Install);
let running = auto_update_check(Some(&interrupted), None, "0.3.1", &[], 1_060);
assert_eq!(running.status, CheckStatus::Ok);
assert!(running.detail.contains("in progress while installing"));
let stalled = auto_update_check(
Some(&interrupted),
None,
"0.3.1",
&[],
1_000 + UPDATE_STALLED_AFTER_SECS,
);
assert_eq!(stalled.status, CheckStatus::Warning);
assert!(stalled.detail.contains("stopped while installing"));
}
#[test]
fn auto_update_check_names_who_must_install_a_skipped_version() {
let mut skipped = attempt(UpdateResult::Skipped);
skipped.stage = None;
skipped.skip_reason = Some(SkipReason::HostManaged);
let host = update_check(Some(&skipped), Some("0.3.1"));
assert_eq!(host.status, CheckStatus::Warning);
assert!(
host.detail
.contains("belongs to its terminal or supervisor")
);
let hint = host.hint.unwrap();
assert!(
hint.contains(
"run `bsk update`, then restart the daemon in its terminal or supervisor"
),
"{hint}"
);
skipped.skip_reason = Some(SkipReason::NotWritable);
skipped.error = Some("permission denied".into());
let unwritable = update_check(Some(&skipped), Some("0.3.1"));
assert!(
unwritable
.detail
.contains("cannot write next to /home/u/.local/bin/bsk: permission denied")
);
assert!(
unwritable
.hint
.unwrap()
.contains("installer or package manager")
);
}
#[test]
fn auto_update_check_lists_leftover_files_without_warning() {
let leftovers = [update::Leftover {
path: "/bin/.bsk.old-1-1".into(),
owner: Some(1),
expired: true,
}];
let check = auto_update_check(None, None, "0.3.1", &leftovers, 0);
assert_eq!(check.status, CheckStatus::Ok);
assert!(check.detail.contains("/bin/.bsk.old-1-1"));
}
#[test]
fn unverified_local_identity_warns_without_failing_usable_ipc() {
let state = DaemonState::Verified {
status: fake_status(Vec::new(), Vec::new()),
local_identity_error: Some("peer identity is unavailable".into()),
host_managed: false,
};
let running = check_daemon_running(&state);
let management = check_daemon_management(&state);
File diff suppressed because it is too large Load Diff
+466
View File
@@ -0,0 +1,466 @@
//! What the most recent update attempt did, kept in `update-state.json` so a
//! failure stays diagnosable after the process that hit it has exited.
//!
//! Each attempt rewrites the record as it moves through its stages. A record
//! still `in_progress` long after it started therefore names the stage in
//! which that process stopped.
use std::fs::{File, OpenOptions};
use std::path::{Path, PathBuf};
use std::time::Duration;
use anyhow::{Context, Result};
use fs2::FileExt;
use semver::Version;
use serde::{Deserialize, Serialize};
use super::now_epoch_secs;
use crate::daemon::paths;
/// How long a daemon waits before trying a version whose update failed again.
pub(crate) const RETRY_AFTER_FAILURE: Duration = Duration::from_secs(6 * 60 * 60);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum UpdateSource {
/// The daemon's periodic auto-update.
Daemon,
/// `bsk update`.
Command,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum UpdateStage {
/// Fetching and verifying the release archive.
Download,
/// Putting the new executable in place and running it once.
Install,
/// A running daemon handing over to one started from the new executable.
Handover,
/// `bsk update` restarting the daemon it stopped.
Restart,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum UpdateResult {
InProgress,
Succeeded,
Failed,
/// A newer version exists, but this installation cannot apply it itself.
Skipped,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SkipReason {
/// The daemon belongs to a terminal or supervisor, which must restart it.
HostManaged,
/// The directory holding the executable does not accept new files.
NotWritable,
}
/// What a failed attempt left behind.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "state", rename_all = "snake_case")]
pub enum Recovery {
/// The attempt failed before it changed anything.
Unchanged,
/// The previous executable is back in place. `daemon_serving` tells
/// whether a daemon running the previous version kept or resumed serving.
Restored { daemon_serving: bool },
/// The previous executable could not be put back.
RestoreFailed { action: String },
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct UpdateRecord {
pub source: UpdateSource,
pub from_version: String,
pub target_version: String,
pub executable: PathBuf,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub stage: Option<UpdateStage>,
pub result: UpdateResult,
pub started_at_epoch_secs: u64,
pub updated_at_epoch_secs: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub skip_reason: Option<SkipReason>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub recovery: Option<Recovery>,
/// The previous executable, kept until the new version is confirmed.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub previous_executable: Option<PathBuf>,
/// Earliest time a daemon tries this version again.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub retry_after_epoch_secs: Option<u64>,
/// The daemon that answered once a handover or restart finished.
#[serde(default, skip_serializing_if = "Option::is_none")]
pub daemon_pid: Option<u32>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub daemon_version: Option<String>,
}
impl UpdateRecord {
pub(crate) fn start(source: UpdateSource, target_version: &Version, executable: &Path) -> Self {
let now = now_epoch_secs();
Self {
source,
from_version: env!("CARGO_PKG_VERSION").to_string(),
target_version: target_version.to_string(),
executable: executable.to_path_buf(),
stage: Some(UpdateStage::Download),
result: UpdateResult::InProgress,
started_at_epoch_secs: now,
updated_at_epoch_secs: now,
error: None,
skip_reason: None,
recovery: None,
previous_executable: None,
retry_after_epoch_secs: None,
daemon_pid: None,
daemon_version: None,
}
}
pub(crate) fn skipped(
source: UpdateSource,
target_version: &Version,
executable: &Path,
reason: SkipReason,
error: Option<&anyhow::Error>,
) -> Self {
Self {
stage: None,
result: UpdateResult::Skipped,
skip_reason: Some(reason),
error: error.map(|err| format!("{err:#}")),
..Self::start(source, target_version, executable)
}
}
pub(crate) fn enter(&mut self, stage: UpdateStage) {
self.stage = Some(stage);
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;
if let Some((pid, version)) = daemon {
self.daemon_pid = Some(pid);
self.daemon_version = Some(version);
}
self.updated_at_epoch_secs = now_epoch_secs();
self.save();
}
/// Record a failure. Daemon attempts wait [`RETRY_AFTER_FAILURE`] before
/// trying the same version again; `bsk update` retries when asked.
pub(crate) fn fail(&mut self, error: &anyhow::Error, recovery: Recovery) {
let now = now_epoch_secs();
self.result = UpdateResult::Failed;
self.error = Some(format!("{error:#}"));
if !matches!(recovery, Recovery::RestoreFailed { .. }) {
self.previous_executable = None;
}
self.recovery = Some(recovery);
self.retry_after_epoch_secs = (self.source == UpdateSource::Daemon)
.then(|| now.saturating_add(RETRY_AFTER_FAILURE.as_secs()));
self.updated_at_epoch_secs = now;
self.save();
}
/// A failed attempt's previous version serves again: its daemon has bound
/// its endpoints and published `daemon.json`.
pub(crate) fn confirm_serving(&mut self) {
if self.mark_serving() {
self.save();
}
}
fn mark_serving(&mut self) -> bool {
let Some(Recovery::Restored { daemon_serving }) = &mut self.recovery else {
return false;
};
*daemon_serving = true;
self.updated_at_epoch_secs = now_epoch_secs();
true
}
/// A failed attempt's previous version could not serve again either.
pub(crate) fn serving_failed(&mut self, error: &anyhow::Error) {
self.note_serving_failure(error);
self.save();
}
fn note_serving_failure(&mut self, error: &anyhow::Error) {
let handover = self.error.take().unwrap_or_default();
self.error = Some(format!(
"{handover}; the previous version could not serve again: {error:#}"
));
if let Some(Recovery::Restored { daemon_serving }) = &mut self.recovery {
*daemon_serving = false;
}
self.updated_at_epoch_secs = now_epoch_secs();
}
/// When a daemon may next try `target`, if an earlier failure defers it.
pub(crate) fn retry_blocked_until(&self, target: &Version, now: u64) -> Option<u64> {
let retry_after = self.retry_after_epoch_secs?;
(self.result == UpdateResult::Failed
&& self.target_version == target.to_string()
&& now < retry_after)
.then_some(retry_after)
}
/// Whether this already records skipping `target` for `reason`.
pub(crate) fn skips(&self, target: &Version, reason: SkipReason) -> bool {
self.result == UpdateResult::Skipped
&& self.skip_reason == Some(reason)
&& self.target_version == target.to_string()
}
/// Best effort: an attempt must not fail because its diagnostics could
/// not be written.
pub(crate) fn save(&self) {
let result = paths::update_state_path().and_then(|path| write(&path, self));
if let Err(err) = result {
tracing::warn!(error = %format_args!("{err:#}"), "could not record the update state");
}
}
}
pub fn read(path: &Path) -> Result<Option<UpdateRecord>> {
match std::fs::read(path) {
Ok(bytes) => serde_json::from_slice(&bytes)
.map(Some)
.with_context(|| format!("parse {}", path.display())),
Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(err) => Err(anyhow::Error::from(err).context(format!("read {}", path.display()))),
}
}
/// The recorded attempt, if one exists and can be read.
pub(crate) fn current() -> Option<UpdateRecord> {
let path = paths::update_state_path().ok()?;
match read(&path) {
Ok(record) => record,
Err(err) => {
tracing::warn!(error = %format_args!("{err:#}"), "ignoring unreadable update state");
None
}
}
}
pub fn write(path: &Path, record: &UpdateRecord) -> Result<()> {
super::write_json_atomically(path, record)
}
/// Serializes update attempts on one installed executable, whichever bsk
/// home they run for, from installation through confirmation or rollback.
/// The lock file (`.<name>.update.lock`) stays next to the executable.
#[derive(Debug)]
pub(crate) struct UpdateLock(File);
impl UpdateLock {
pub(crate) fn try_acquire(target: &Path) -> Result<Self> {
Self::try_acquire_at(&super::sibling(target, "update.lock")?)
}
fn try_acquire_at(path: &Path) -> Result<Self> {
let file = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(path)
.with_context(|| format!("open {}", path.display()))?;
match file.try_lock_exclusive() {
Ok(()) => Ok(Self(file)),
Err(err) if err.raw_os_error() == fs2::lock_contended_error().raw_os_error() => {
anyhow::bail!("another bsk update is in progress; retry once it finishes")
}
Err(err) => Err(anyhow::Error::new(err).context(format!("lock {}", path.display()))),
}
}
}
impl Drop for UpdateLock {
fn drop(&mut self) {
let _ = <File as FileExt>::unlock(&self.0);
}
}
#[cfg(test)]
mod tests {
use super::*;
const TARGET: Version = Version::new(0, 9, 0);
#[test]
fn records_round_trip_and_omit_empty_fields() {
let tmp = tempfile::TempDir::new().unwrap();
let path = tmp.path().join("update-state.json");
let mut record = UpdateRecord::start(UpdateSource::Daemon, &TARGET, Path::new("/bin/bsk"));
record.previous_executable = Some(PathBuf::from("/bin/.bsk.old-1-1"));
write(&path, &record).unwrap();
assert_eq!(read(&path).unwrap(), Some(record));
let json: serde_json::Value =
serde_json::from_slice(&std::fs::read(&path).unwrap()).unwrap();
assert_eq!(json["result"], "in_progress");
assert_eq!(json["stage"], "download");
assert!(json.get("error").is_none());
assert!(json.get("recovery").is_none());
assert_eq!(read(&tmp.path().join("missing.json")).unwrap(), None);
}
#[test]
fn recovery_serialises_with_a_state_tag() {
for (recovery, expected) in [
(
Recovery::Unchanged,
serde_json::json!({"state": "unchanged"}),
),
(
Recovery::Restored {
daemon_serving: true,
},
serde_json::json!({"state": "restored", "daemon_serving": true}),
),
(
Recovery::RestoreFailed {
action: "rename it".into(),
},
serde_json::json!({"state": "restore_failed", "action": "rename it"}),
),
] {
assert_eq!(serde_json::to_value(&recovery).unwrap(), expected);
}
}
#[test]
fn only_a_failed_daemon_attempt_defers_the_same_version() {
let exe = Path::new("bsk");
let failed = |source| UpdateRecord {
result: UpdateResult::Failed,
retry_after_epoch_secs: (source == UpdateSource::Daemon).then_some(1_000),
..UpdateRecord::start(source, &TARGET, exe)
};
let daemon = failed(UpdateSource::Daemon);
assert_eq!(daemon.retry_blocked_until(&TARGET, 999), Some(1_000));
assert_eq!(daemon.retry_blocked_until(&TARGET, 1_000), None);
assert_eq!(daemon.retry_blocked_until(&Version::new(0, 9, 1), 0), None);
assert_eq!(
failed(UpdateSource::Command).retry_blocked_until(&TARGET, 0),
None
);
let succeeded = UpdateRecord {
result: UpdateResult::Succeeded,
..daemon
};
assert_eq!(succeeded.retry_blocked_until(&TARGET, 0), None);
}
#[test]
fn the_lock_belongs_to_the_installation_not_the_bsk_home() {
let tmp = tempfile::TempDir::new().unwrap();
let shared = tmp.path().join("bin").join("bsk");
let other = tmp.path().join("other").join("bsk");
for exe in [&shared, &other] {
std::fs::create_dir_all(exe.parent().unwrap()).unwrap();
}
// Two daemons with different bsk homes update the same executable.
let first = UpdateLock::try_acquire(&shared).unwrap();
let error = UpdateLock::try_acquire(&shared).unwrap_err();
assert!(format!("{error:#}").contains("another bsk update is in progress"));
let _unrelated = UpdateLock::try_acquire(&other).unwrap();
drop(first);
UpdateLock::try_acquire(&shared).unwrap();
assert!(tmp.path().join("bin").join(".bsk.update.lock").exists());
}
#[test]
fn only_one_update_attempt_holds_the_lock_at_a_time() {
let tmp = tempfile::TempDir::new().unwrap();
let path = tmp.path().join("update.lock");
let barrier = std::sync::Arc::new(std::sync::Barrier::new(8));
let attempts: Vec<_> = (0..8)
.map(|_| {
let path = path.clone();
let barrier = std::sync::Arc::clone(&barrier);
std::thread::spawn(move || {
barrier.wait();
let lock = UpdateLock::try_acquire_at(&path);
// Hold a winning lock until every attempt has run.
std::thread::sleep(Duration::from_millis(200));
lock.map(drop).map_err(|err| format!("{err:#}"))
})
})
.collect();
let results: Vec<_> = attempts.into_iter().map(|t| t.join().unwrap()).collect();
assert_eq!(results.iter().filter(|result| result.is_ok()).count(), 1);
for err in results.iter().filter_map(|result| result.as_ref().err()) {
assert!(err.contains("another bsk update is in progress"), "{err}");
}
UpdateLock::try_acquire_at(&path).expect("released after the winner finishes");
}
#[test]
fn a_resumed_service_is_confirmed_or_its_failure_appended() {
let restored = |serving| UpdateRecord {
result: UpdateResult::Failed,
error: Some("replacement exited".into()),
recovery: Some(Recovery::Restored {
daemon_serving: serving,
}),
..UpdateRecord::start(UpdateSource::Daemon, &TARGET, Path::new("bsk"))
};
let mut confirmed = restored(false);
assert!(confirmed.mark_serving());
assert_eq!(
confirmed.recovery,
Some(Recovery::Restored {
daemon_serving: true
})
);
let mut unchanged = UpdateRecord {
recovery: Some(Recovery::Unchanged),
..restored(false)
};
assert!(!unchanged.mark_serving(), "only a restored daemon resumes");
let mut failed = restored(false);
failed.note_serving_failure(&anyhow::anyhow!("bind WS server: address in use"));
assert_eq!(
failed.recovery,
Some(Recovery::Restored {
daemon_serving: false
})
);
let error = failed.error.unwrap();
assert!(error.starts_with("replacement exited; "), "{error}");
assert!(error.contains("address in use"), "{error}");
}
#[test]
fn skip_records_name_their_version_and_reason() {
let record = UpdateRecord::skipped(
UpdateSource::Daemon,
&TARGET,
Path::new("bsk"),
SkipReason::HostManaged,
None,
);
assert!(record.skips(&TARGET, SkipReason::HostManaged));
assert!(!record.skips(&TARGET, SkipReason::NotWritable));
assert!(!record.skips(&Version::new(1, 0, 0), SkipReason::HostManaged));
assert_eq!(record.stage, None);
}
}
+463 -50
View File
@@ -1,58 +1,471 @@
//! The update helper inherits only its dedicated stdio handles. Keep its
//! existing Job policy: daemon startup separately requires verified breakaway.
//! In-place replacement of a running Windows executable.
//!
//! Windows refuses to overwrite or delete an executable while a process runs
//! it, but the loader opens images with delete sharing, so the file can be
//! renamed. The running image is moved aside and the new binary takes its
//! path; processes keep running the old image, and every later launch uses
//! the new one. The moved-aside image is the previous executable an update
//! restores on failure, and is deleted once nothing runs it.
//!
//! This relies on the file system honouring renames of open images, which
//! NTFS does. Elsewhere the first rename fails and the executable is left as
//! it was.
use std::ffi::OsStr;
use std::fs::File;
use std::fs;
use std::io;
use std::path::Path;
use std::path::{Path, PathBuf};
use std::time::{Duration, Instant};
use windows_sys::Win32::System::Threading::{CREATE_NEW_PROCESS_GROUP, CREATE_NO_WINDOW};
use anyhow::{Context, Result};
use windows_sys::Win32::Foundation::{
ERROR_ACCESS_DENIED, ERROR_LOCK_VIOLATION, ERROR_SHARING_VIOLATION,
};
use crate::windows_process;
pub(super) use crate::windows_process::Process as Helper;
use super::{PreviousNotRestored, sibling, unique_suffix};
pub(super) fn spawn(
script: &Path,
source: &Path,
/// How long to keep retrying a rename that a scanner or indexer briefly blocks.
const RETRY_WINDOW: Duration = Duration::from_secs(5);
const RETRY_DELAY: Duration = Duration::from_millis(50);
/// Files of the former script-based updater are only removed once its helper
/// has certainly finished with them.
const LEGACY_HELPER_GRACE: Duration = Duration::from_secs(60 * 60);
/// Upper bound for a legacy helper report copied into the daemon log.
const LEGACY_REPORT_LIMIT: usize = 4096;
/// Returns where the previous executable now lives. Two concurrent swaps can
/// interleave their renames; callers hold the update lock.
pub(super) fn replace(target: &Path, binary: &[u8]) -> Result<PathBuf> {
super::remove_update_leftovers(target);
let staged = sibling(target, &format!("new-{}", unique_suffix()))?;
super::write_synced(&staged, binary)?;
let installed = install(target, &staged, RETRY_WINDOW, &mut |from, to| {
fs::rename(from, to)
});
if installed.is_err() {
let _ = fs::remove_file(&staged);
}
installed
}
/// Swap `staged` into `target`, putting the original back if that fails.
fn install(
target: &Path,
ready: &Path,
log_path: &Path,
) -> io::Result<Helper> {
let root = std::env::var_os("SystemRoot")
.ok_or_else(|| io::Error::new(io::ErrorKind::NotFound, "SystemRoot is not set"))?;
let application = Path::new(&root).join("System32/cmd.exe");
// /S /C strips one outer pair of quotes. Environment expansion keeps paths
// Unicode and avoids CRT escaping, batch-file encoding, and CALL expansion.
let command = OsStr::new("cmd.exe /D /V:OFF /S /C \"\"%BSK_UPDATE_SCRIPT%\"\"");
let overrides = [
("BSK_UPDATE_SCRIPT", script.as_os_str()),
("BSK_UPDATE_SOURCE", source.as_os_str()),
("BSK_UPDATE_TARGET", target.as_os_str()),
("BSK_UPDATE_READY", ready.as_os_str()),
("BSK_UPDATE_LOG", log_path.as_os_str()),
];
let mut env: Vec<_> = std::env::vars_os()
.filter(|(key, _)| {
let key = key.to_string_lossy();
!overrides
.iter()
.any(|(name, _)| key.eq_ignore_ascii_case(name))
&& !key.eq_ignore_ascii_case(crate::daemon::start::DAEMONIZED_ENV)
&& !key.eq_ignore_ascii_case(crate::daemon::start::DAEMON_REPLACEMENT_WAIT_ENV)
staged: &Path,
window: Duration,
rename: &mut dyn FnMut(&Path, &Path) -> io::Result<()>,
) -> Result<PathBuf> {
let previous = sibling(target, &format!("old-{}", unique_suffix()))?;
retry(window, || rename(target, &previous))
.with_context(|| format!("move {} aside", target.display()))?;
let Err(err) = retry(window, || rename(staged, target)) else {
return Ok(previous);
};
let err = anyhow::Error::new(err).context(format!("install new {}", target.display()));
match retry(window, || rename(&previous, target)) {
Ok(()) => Err(err),
Err(restore) => {
Err(err
.context(format!("restore failed: {restore}"))
.context(PreviousNotRestored {
previous,
target: target.to_path_buf(),
}))
}
}
}
/// Put `previous` back at `target`. The rejected executable is deleted when
/// nothing runs it, and otherwise left for a later cleanup.
pub(super) fn restore(target: &Path, previous: &Path) -> Result<()> {
restore_with(target, previous, RETRY_WINDOW, &mut |from, to| {
fs::rename(from, to)
})
}
fn restore_with(
target: &Path,
previous: &Path,
window: Duration,
rename: &mut dyn FnMut(&Path, &Path) -> io::Result<()>,
) -> Result<()> {
let not_restored = || PreviousNotRestored {
previous: previous.to_path_buf(),
target: target.to_path_buf(),
};
let rejected = sibling(target, &format!("old-{}", unique_suffix()))?;
let moved = match retry(window, || rename(target, &rejected)) {
Ok(()) => true,
Err(err) if err.kind() == io::ErrorKind::NotFound => false,
Err(err) => {
return Err(anyhow::Error::new(err)
.context(format!("move rejected {} aside", target.display()))
.context(not_restored()));
}
};
if let Err(err) = retry(window, || rename(previous, target)) {
// Keep an executable at the path rather than none.
if moved {
let _ = retry(window, || rename(&rejected, target));
}
return Err(anyhow::Error::new(err).context(not_restored()));
}
if moved {
let _ = fs::remove_file(&rejected);
}
Ok(())
}
/// Whether `entry` is a file of the former script-based updater, and if so
/// whether it is old enough to remove.
pub(super) fn legacy_helper_file(entry: &fs::DirEntry, name: &str) -> Option<bool> {
let file_name = entry.file_name();
let file_name = file_name.to_str()?;
file_name
.starts_with(&format!("{name}.update-"))
.then(|| older_than(entry, LEGACY_HELPER_GRACE))
}
/// The former updater only logged next to the executable; surface its last
/// report in the daemon log before the file is removed.
pub(super) fn report_legacy_helper(path: &Path) {
if path.extension().is_none_or(|ext| ext != "log") {
return;
}
let Ok(bytes) = fs::read(path) else {
return;
};
let report = String::from_utf8_lossy(&bytes[..bytes.len().min(LEGACY_REPORT_LIMIT)]);
let report = report.trim();
if !report.is_empty() {
tracing::warn!(path = %path.display(), report, "previous update helper left a report");
}
}
fn older_than(entry: &fs::DirEntry, age: Duration) -> bool {
entry
.metadata()
.and_then(|metadata| metadata.modified())
.is_ok_and(|modified| modified.elapsed().is_ok_and(|elapsed| elapsed >= age))
}
fn retry(window: Duration, mut op: impl FnMut() -> io::Result<()>) -> io::Result<()> {
let deadline = Instant::now() + window;
loop {
match op() {
Err(err) if is_transient(&err) && Instant::now() < deadline => {
std::thread::sleep(RETRY_DELAY);
}
result => return result,
}
}
}
/// Scanners and indexers open fresh files without delete sharing for a moment.
fn is_transient(err: &io::Error) -> bool {
err.raw_os_error().is_some_and(|code| {
[
ERROR_ACCESS_DENIED,
ERROR_SHARING_VIOLATION,
ERROR_LOCK_VIOLATION,
]
.contains(&(code as u32))
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::cli::update::remove_update_leftovers;
use std::os::windows::fs::OpenOptionsExt;
use std::os::windows::process::CommandExt;
use std::time::SystemTime;
use windows_sys::Win32::Storage::FileSystem::FILE_SHARE_READ;
use windows_sys::Win32::System::Threading::CREATE_NO_WINDOW;
fn names(dir: &Path) -> Vec<String> {
let mut names: Vec<_> = fs::read_dir(dir)
.unwrap()
.map(|entry| entry.unwrap().file_name().into_string().unwrap())
.collect();
names.sort();
names
}
/// Renames where exactly the listed calls fail, counting from 1.
fn failing_calls(failing: &'static [usize]) -> impl FnMut(&Path, &Path) -> io::Result<()> {
let mut calls = 0;
move |from, to| {
calls += 1;
if failing.contains(&calls) {
Err(io::Error::from(io::ErrorKind::InvalidInput))
} else {
fs::rename(from, to)
}
}
}
fn fixture() -> (tempfile::TempDir, PathBuf, PathBuf) {
let tmp = tempfile::TempDir::new().unwrap();
let target = tmp.path().join("bsk.exe");
let staged = tmp.path().join(".bsk.exe.new-1-1");
fs::write(&target, b"old binary").unwrap();
fs::write(&staged, b"new binary").unwrap();
(tmp, target, staged)
}
#[test]
#[ignore = "subprocess entry point"]
fn idle_process() {
std::thread::sleep(Duration::from_secs(60));
}
#[test]
fn replaces_a_running_executable() {
let tmp = tempfile::TempDir::new().unwrap();
let dir = tmp.path().join("中文 space %PATH% ! & (update)");
fs::create_dir(&dir).unwrap();
let target = dir.join("bsk.exe");
fs::copy(std::env::current_exe().unwrap(), &target).unwrap();
let mut running = std::process::Command::new(&target)
.args([
"--exact",
"cli::update::windows::tests::idle_process",
"--ignored",
])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.creation_flags(CREATE_NO_WINDOW)
.spawn()
.unwrap();
let replaced = replace(&target, b"new binary");
let still_running = running.try_wait().unwrap().is_none();
let _ = running.kill();
let _ = running.wait();
let previous = replaced.unwrap();
assert!(
still_running,
"replacement must not disturb the running process"
);
assert_eq!(fs::read(&target).unwrap(), b"new binary");
assert!(previous.exists(), "the previous executable is kept");
// The exited image is released asynchronously; scanners may also hold it briefly.
let deadline = Instant::now() + Duration::from_secs(10);
loop {
// This process made the backup, so only an explicit discard removes it.
let _ = fs::remove_file(&previous);
let listed = names(&dir);
if listed == ["bsk.exe"] {
break;
}
assert!(Instant::now() < deadline, "leftovers remain: {listed:?}");
std::thread::sleep(Duration::from_millis(50));
}
}
#[test]
fn restores_the_previous_executable_of_a_running_process() {
let tmp = tempfile::TempDir::new().unwrap();
let target = tmp.path().join("bsk.exe");
fs::copy(std::env::current_exe().unwrap(), &target).unwrap();
let original = fs::read(&target).unwrap();
let mut running = std::process::Command::new(&target)
.args([
"--exact",
"cli::update::windows::tests::idle_process",
"--ignored",
])
.stdin(std::process::Stdio::null())
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.creation_flags(CREATE_NO_WINDOW)
.spawn()
.unwrap();
let previous = replace(&target, b"new binary").unwrap();
let restored = restore(&target, &previous);
let _ = running.kill();
let _ = running.wait();
restored.unwrap();
assert_eq!(fs::read(&target).unwrap(), original);
assert!(!previous.exists());
}
#[test]
fn replaces_in_unicode_space_and_shell_symbol_paths() {
let tmp = tempfile::TempDir::new().unwrap();
for name in [
"ascii",
"with space",
"中文目录",
"literal %PATH% ! & (folder)",
] {
let dir = tmp.path().join(name);
fs::create_dir(&dir).unwrap();
let target = dir.join("bsk.exe");
fs::write(&target, b"old binary").unwrap();
let previous = replace(&target, b"new binary").unwrap();
assert_eq!(fs::read(&target).unwrap(), b"new binary", "{name}");
assert_eq!(fs::read(&previous).unwrap(), b"old binary", "{name}");
let listed = names(&dir);
assert_eq!(listed.len(), 2, "{name}: {listed:?}");
assert!(listed[0].starts_with(".bsk.exe.old-"), "{name}: {listed:?}");
// This process owns the previous executable until it exits.
remove_update_leftovers(&target);
assert_eq!(names(&dir).len(), 2, "{name}");
fs::remove_file(&previous).unwrap();
assert_eq!(names(&dir), ["bsk.exe"], "{name}");
}
}
#[test]
fn keeps_the_original_when_it_cannot_be_moved() {
let tmp = tempfile::TempDir::new().unwrap();
let target = tmp.path().join("bsk.exe");
fs::write(&target, b"old binary").unwrap();
// Without delete sharing the file cannot be renamed.
let _lock = fs::OpenOptions::new()
.read(true)
.share_mode(FILE_SHARE_READ)
.open(&target)
.unwrap();
let staged = tmp.path().join(".bsk.exe.new-1-1");
fs::write(&staged, b"new binary").unwrap();
let error = install(&target, &staged, Duration::from_millis(200), &mut |a, b| {
fs::rename(a, b)
})
.collect();
env.extend(
overrides
.into_iter()
.map(|(key, value)| (key.into(), value.to_owned())),
);
let input = File::open("NUL")?;
let log = File::create(log_path)?;
windows_process::spawn(
application.as_os_str(),
command,
&env,
[&input, &log, &log],
CREATE_NO_WINDOW | CREATE_NEW_PROCESS_GROUP,
)
.unwrap_err();
assert!(format!("{error:#}").contains("aside"), "{error:#}");
assert!(error.downcast_ref::<PreviousNotRestored>().is_none());
assert_eq!(fs::read(&target).unwrap(), b"old binary");
}
#[test]
fn puts_the_original_back_when_the_new_binary_cannot_take_its_place() {
let (_tmp, target, staged) = fixture();
// Call 2 moves the new binary in.
let error =
install(&target, &staged, Duration::ZERO, &mut failing_calls(&[2])).unwrap_err();
assert!(format!("{error:#}").contains("install new"), "{error:#}");
assert!(error.downcast_ref::<PreviousNotRestored>().is_none());
assert_eq!(fs::read(&target).unwrap(), b"old binary");
assert_eq!(fs::read(&staged).unwrap(), b"new binary");
}
#[test]
fn names_the_manual_step_when_the_original_cannot_be_put_back() {
let (_tmp, target, staged) = fixture();
// Call 3 would move the original back.
let error = install(
&target,
&staged,
Duration::ZERO,
&mut failing_calls(&[2, 3]),
)
.unwrap_err();
let missing = error
.downcast_ref::<PreviousNotRestored>()
.unwrap_or_else(|| panic!("{error:#}"));
assert_eq!(missing.target, target);
assert!(!target.exists());
assert_eq!(fs::read(&missing.previous).unwrap(), b"old binary");
let action = missing.action();
assert!(
action.contains(&missing.previous.display().to_string()),
"{action}"
);
assert!(action.contains(&target.display().to_string()), "{action}");
}
#[test]
fn restore_keeps_an_executable_at_the_path_when_it_fails() {
let (tmp, target, _staged) = fixture();
let previous = tmp.path().join(".bsk.exe.old-1-1");
fs::write(&previous, b"previous binary").unwrap();
// Call 2 moves the previous executable back.
let error =
restore_with(&target, &previous, Duration::ZERO, &mut failing_calls(&[2])).unwrap_err();
assert!(error.downcast_ref::<PreviousNotRestored>().is_some());
assert_eq!(fs::read(&target).unwrap(), b"old binary");
assert_eq!(fs::read(&previous).unwrap(), b"previous binary");
}
#[test]
fn retries_while_the_original_is_briefly_locked() {
let tmp = tempfile::TempDir::new().unwrap();
let target = tmp.path().join("bsk.exe");
fs::write(&target, b"old binary").unwrap();
let lock = fs::OpenOptions::new()
.read(true)
.share_mode(FILE_SHARE_READ)
.open(&target)
.unwrap();
let release = std::thread::spawn(move || {
std::thread::sleep(Duration::from_millis(200));
drop(lock);
});
replace(&target, b"new binary").unwrap();
release.join().unwrap();
assert_eq!(fs::read(&target).unwrap(), b"new binary");
}
#[test]
fn removes_only_leftovers_that_are_no_longer_in_use() {
let tmp = tempfile::TempDir::new().unwrap();
let target = tmp.path().join("bsk.exe");
fs::write(&target, b"current").unwrap();
let path = |name: &str| tmp.path().join(name);
let own = std::process::id();
let live_previous = format!(".bsk.exe.old-{own}-1");
let live_staging = format!(".bsk.exe.new-{own}-1");
for name in [
".bsk.exe.old-4294967295-1",
live_previous.as_str(),
".bsk.exe.new-4294967295-1",
live_staging.as_str(),
"bsk.exe.update-7.cmd",
"bsk.exe.update-7.log",
"bsk.exe.update-8.log",
"other.exe.old",
] {
fs::write(path(name), b"leftover").unwrap();
}
let expired = SystemTime::now() - LEGACY_HELPER_GRACE - Duration::from_secs(60);
for name in ["bsk.exe.update-7.cmd", "bsk.exe.update-7.log"] {
fs::File::options()
.write(true)
.open(path(name))
.unwrap()
.set_modified(expired)
.unwrap();
}
remove_update_leftovers(&target);
let mut expected = vec![
live_previous,
live_staging,
"bsk.exe".to_string(),
"bsk.exe.update-8.log".to_string(),
"other.exe.old".to_string(),
];
expected.sort();
assert_eq!(names(tmp.path()), expected);
}
}
+35
View File
@@ -27,6 +27,11 @@ pub struct DaemonInfo {
/// `SystemTime` rendered as RFC 3339-ish seconds-since-epoch for
/// portability across platforms.
pub started_at_epoch_secs: u64,
/// Owned by a terminal or supervisor (`--foreground`, server mode) rather
/// than started by `bsk` in the background. Only its owner restarts it.
/// Omitted when false, so detached daemons write the same file as before.
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub host_managed: bool,
}
impl DaemonInfo {
@@ -40,8 +45,14 @@ impl DaemonInfo {
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0),
host_managed: false,
}
}
pub fn with_host_managed(mut self, host_managed: bool) -> Self {
self.host_managed = host_managed;
self
}
}
/// Atomically write `info` to `daemon.json`. Temp-file-and-rename keeps
@@ -193,4 +204,28 @@ mod tests {
assert_eq!(parsed.pid, 2);
assert_eq!(parsed.version, "b");
}
#[test]
fn only_host_managed_daemons_record_their_owner() {
let detached = DaemonInfo::now(1, "sock".into(), 1, "a");
let json = serde_json::to_value(&detached).unwrap();
assert!(json.get("host_managed").is_none(), "{json}");
let host_managed = detached.clone().with_host_managed(true);
let json = serde_json::to_value(&host_managed).unwrap();
assert_eq!(json["host_managed"], true);
assert_eq!(
serde_json::from_value::<DaemonInfo>(json).unwrap(),
host_managed
);
// Files written before the field existed describe detached daemons.
let mut older = serde_json::to_value(&detached).unwrap();
older.as_object_mut().unwrap().remove("host_managed");
assert!(
!serde_json::from_value::<DaemonInfo>(older)
.unwrap()
.host_managed
);
}
}
+12
View File
@@ -105,6 +105,18 @@ pub fn update_check_path() -> Result<PathBuf> {
Ok(bsk_home()?.join("update-check.json"))
}
/// Outcome of the most recent update attempt (`update-state.json`).
pub fn update_state_path() -> Result<PathBuf> {
Ok(bsk_home()?.join("update-state.json"))
}
/// Why a replacement daemon failed to start, for the daemon handing over to it.
pub fn replacement_failure_path(pid: u32) -> Result<PathBuf> {
Ok(bsk_home()?
.join("run")
.join(format!("replacement-{pid}.error")))
}
pub fn log_dir() -> Result<PathBuf> {
bsk_home()
}
+552 -104
View File
@@ -19,10 +19,11 @@ use std::time::{Duration, Instant};
use anyhow::{Context, Result};
use bsk_protocol::StatusResult;
use tracing::{debug, info, warn};
use tracing::{debug, error, info, warn};
use crate::cli::daemon::StartArgs;
use crate::cli::ensure_daemon::SPAWN_DEADLINE;
use crate::cli::update::state::UpdateRecord;
use crate::daemon::{
browsers::{BROWSER_LIVENESS_TICK, BROWSER_LIVENESS_TIMEOUT, EXTENSION_CONNECT_WAIT},
info as daemon_info, ipc, lockfile, paths,
@@ -32,6 +33,7 @@ use crate::daemon::{
ws,
};
mod handover;
#[cfg(windows)]
mod windows;
@@ -68,6 +70,13 @@ pub struct DaemonConfig {
/// How often the liveness task scans the registry. Defaults to
/// [`BROWSER_LIVENESS_TICK`].
pub browser_liveness_tick: Duration,
/// Started by `bsk` as an independent background process, rather than
/// owned by a terminal or supervisor. Only such a daemon replaces
/// itself after installing an update.
pub detached: bool,
/// The daemon this one takes over from after an auto-update. It waits
/// for that daemon to release the lock before starting to serve.
pub replaces: Option<u32>,
}
impl DaemonConfig {
@@ -84,6 +93,8 @@ impl DaemonConfig {
extension_connect_wait: EXTENSION_CONNECT_WAIT,
browser_liveness_timeout: BROWSER_LIVENESS_TIMEOUT,
browser_liveness_tick: BROWSER_LIVENESS_TICK,
detached: false,
replaces: None,
}
}
@@ -119,6 +130,8 @@ impl From<&StartArgs> for DaemonConfig {
extension_connect_wait: EXTENSION_CONNECT_WAIT,
browser_liveness_timeout: BROWSER_LIVENESS_TIMEOUT,
browser_liveness_tick: BROWSER_LIVENESS_TICK,
detached: false,
replaces: None,
}
}
}
@@ -134,25 +147,69 @@ pub fn run_start(args: StartArgs) -> Result<()> {
// Detached child mode (set by parent before spawn).
if is_daemonized_child() {
wait_for_replaced_daemon();
detach_stdio()?;
cfg.detached = true;
cfg.replaces = replaced_daemon();
return run_foreground(cfg);
}
let exe = std::env::current_exe().context("locate daemon executable")?;
start_detached(&exe, &args).map(drop)
}
/// Reuse the running daemon, or start one from `exe` and wait until it is ready.
pub(crate) fn start_detached(exe: &Path, args: &StartArgs) -> Result<daemon_info::DaemonInfo> {
let deadline = Instant::now() + SPAWN_DEADLINE;
if let Probe::Ready(daemon) = probe::probe(PROBE_TIMEOUT)? {
let status = daemon.status;
validate_existing_start(&args, &status)?;
validate_existing_start(args, &daemon.status)?;
info!(
pid = status.pid,
ws_port = status.ws_port,
pid = daemon.status.pid,
ws_port = daemon.status.ws_port,
"daemon already running"
);
return Ok(());
return Ok(daemon.info);
}
start_background(&args, deadline)?;
Ok(())
start_background_at(exe, args, deadline)
}
/// Whether this process can start a daemon that outlives it, as
/// [`start_owned`] does. On Windows a host Job that forbids breakaway
/// prevents that; the check starts nothing that runs. Elsewhere a daemon can
/// always detach into its own session.
pub(crate) fn check_independent_start(exe: &Path) -> Result<()> {
#[cfg(windows)]
{
windows::check_breakaway(exe)
}
#[cfg(not(windows))]
{
let _ = exe;
Ok(())
}
}
/// Start a daemon from `exe` for `bsk update`, and wait until a daemon of
/// `version` serves `args`' port. Unlike [`start_detached`], the process is
/// this caller's: one that is not ready in time is stopped and reaped, so a
/// restart of the previous version never meets it holding the daemon lock.
pub(crate) fn start_owned(
exe: &Path,
args: &StartArgs,
version: &str,
) -> Result<daemon_info::DaemonInfo> {
let mut child = spawn_detached_at(exe, args, None)?;
let pid = child.id();
let expected = handover::Expected {
port: Some(args.resolved_port()).filter(|port| *port != 0),
version,
};
let daemon = handover::wait_until_serving(&mut child, pid, handover::HANDOVER_TIMEOUT, || {
handover::observe(None, &expected)
})?;
disown_daemon(child);
Ok(daemon)
}
/// Shared explicit/automatic startup, without an intermediate launcher or
@@ -160,13 +217,21 @@ pub fn run_start(args: StartArgs) -> Result<()> {
pub(crate) fn start_background(
args: &StartArgs,
deadline: Instant,
) -> Result<daemon_info::DaemonInfo> {
let exe = std::env::current_exe().context("locate daemon executable")?;
start_background_at(&exe, args, deadline)
}
fn start_background_at(
exe: &Path,
args: &StartArgs,
deadline: Instant,
) -> Result<daemon_info::DaemonInfo> {
anyhow::ensure!(
Instant::now() < deadline,
"daemon startup deadline exceeded"
);
let exe = std::env::current_exe().context("locate daemon executable")?;
let child = spawn_detached_at(&exe, args, None)?;
let child = spawn_detached_at(exe, args, None)?;
wait_for_background(child, args, deadline)
}
@@ -321,24 +386,165 @@ fn wait_for_stopped(expected: &daemon_info::DaemonInfo, timeout: Duration) -> Re
/// Run the daemon in the foreground of the current process: acquire
/// the lock, bind IPC, publish `daemon.json`, and serve until shutdown.
pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
// Before any update can replace the executable (see the function).
let _ = crate::cli::update::installed_executable();
paths::ensure_bsk_home()?;
let _log_guard = init_tracing();
let result = run_daemon(&cfg);
if let Err(err) = &result {
// Detached daemons have no stderr; keep the reason in the log.
error!(error = %format_args!("{err:#}"), "daemon stopped with an error");
}
result
}
/// Serve until shutdown. After a failed auto-update handover this process
/// serves again with the lock it took back.
fn run_daemon(cfg: &DaemonConfig) -> Result<()> {
#[cfg(windows)]
info!(
pid = std::process::id(),
background = is_daemonized_child(),
background = cfg.detached,
in_job = ?crate::windows_process::current_process_in_job(),
"Windows daemon process started"
);
let lock = lockfile::acquire().context("acquire daemon lock")?;
info!(?lock, "daemon lock acquired");
let startup = acquire_daemon_lock(cfg.replaces).and_then(|lock| {
info!(?lock, "daemon lock acquired");
serve(cfg, None).map(|stopped| (lock, stopped))
});
let (mut lock, mut stopped) = match startup {
Ok(started) => started,
Err(err) => {
if cfg.replaces.is_some() {
handover::report_startup_failure(&err);
}
return Err(err);
}
};
// Browsers reconnect to the port served before a failed handover, so a
// `--port 0` daemon resumes on the port it was given, not a new one.
let mut resumed_cfg = cfg.clone();
loop {
let Stopped::HandOver(pending) = stopped else {
return Ok(());
};
drop(lock);
let mut record = match handover::finish(*pending) {
handover::Finished::Exit => return Ok(()),
handover::Finished::Resume {
lock: reclaimed,
record,
port,
} => {
lock = reclaimed;
resumed_cfg.ws_port = port;
record
}
};
stopped = match serve(&resumed_cfg, Some(&mut record)) {
Ok(stopped) => stopped,
Err(err) => {
record.serving_failed(&err);
return Err(err);
}
};
}
}
/// Take the daemon lock. A replacement waits for its predecessor, which
/// keeps serving until the replacement has started and then releases it.
fn acquire_daemon_lock(predecessor: Option<u32>) -> Result<lockfile::DaemonLock> {
let Some(pid) = predecessor else {
return lockfile::acquire().context("acquire daemon lock");
};
info!(
predecessor = pid,
"replacement daemon waiting for its predecessor to release the daemon lock"
);
let deadline = Instant::now() + handover::REPLACEMENT_LOCK_WAIT;
loop {
match lockfile::acquire() {
Ok(lock) => return Ok(lock),
Err(err) if err.is::<lockfile::AlreadyLocked>() && Instant::now() < deadline => {
std::thread::sleep(Duration::from_millis(50));
}
Err(err) => {
return Err(
err.context(format!("acquire daemon lock released by predecessor {pid}"))
);
}
}
}
}
/// Why [`serve`] returned.
enum Stopped {
Shutdown,
/// An auto-update started a replacement that waits for the lock.
HandOver(Box<handover::Pending>),
}
enum StopReason {
Signal,
Idle,
Handover,
}
/// The auto-update a daemon is applying. Its blocking steps store each result
/// here themselves, rather than returning it to the async update task, which
/// shutdown aborts; so [`serve`] always finds an unfinished update once the
/// runtime has waited for those steps, and hands it over or undoes it.
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),
}
impl Transaction {
/// Undo an update this daemon stopped before handing over.
fn abandon(self) {
use crate::cli::update::state::Recovery;
let reason = "the daemon stopped before handing over";
match self {
Transaction::HandingOver(pending) => handover::abandon(pending, reason),
Transaction::Installed(Prepared {
installed,
mut record,
}) => {
let recovery = installed.roll_back(|| Recovery::Restored {
daemon_serving: false,
});
record.fail(&anyhow::anyhow!("handover abandoned: {reason}"), recovery);
}
}
}
}
type TransactionSlot = Arc<Mutex<Option<Transaction>>>;
fn lock_slot(slot: &TransactionSlot) -> std::sync::MutexGuard<'_, Option<Transaction>> {
slot.lock().unwrap_or_else(|poison| poison.into_inner())
}
/// Bind IPC and WS, publish `daemon.json`, and serve until shutdown. Returns
/// with every endpoint released except the daemon lock, which the caller holds.
/// `resumed` records a failed handover whose previous version serves again
/// here; it is confirmed once `daemon.json` is published.
fn serve(cfg: &DaemonConfig, resumed: Option<&mut UpdateRecord>) -> Result<Stopped> {
let runtime = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()
.context("build tokio runtime")?;
let transaction: TransactionSlot = Arc::new(Mutex::new(None));
// Tells an update in progress not to start installing.
let stopping = Arc::new(std::sync::atomic::AtomicBool::new(false));
let cfg = cfg.clone();
runtime.block_on(async move {
let reason = runtime.block_on({
let transaction = Arc::clone(&transaction);
let stopping = Arc::clone(&stopping);
async move {
#[cfg(unix)]
let sock_path = paths::sock_path().context("resolve socket path")?;
#[cfg(windows)]
@@ -355,25 +561,31 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
.context("initialize transfer staging")?;
let session_idle_task = spawn_session_idle_reaper(Arc::clone(&state));
let browser_liveness_task = spawn_browser_liveness_reaper(Arc::clone(&state));
// Fired by the update check task after a successful auto-update:
// the replacement daemon (or Windows update helper) is ready, so
// this process should shut down and let it take over.
let restart_notify = Arc::new(tokio::sync::Notify::new());
let update_check_task =
spawn_update_check_task(Arc::clone(&state), Arc::clone(&restart_notify));
let ws_addr = SocketAddr::new(cfg.listen_ip(), cfg.ws_port);
let ws_handle = ws::WsServer::new(Arc::clone(&state))
.bind(ws_addr)
.await
.with_context(|| format!("bind WS server on {ws_addr}"))?;
let ws_port = ws_handle.local_addr.port();
// Fired by the update check task once it has started a replacement
// daemon, which waits for this process to release the lock. The
// replacement must serve the port bound here.
let restart_notify = Arc::new(tokio::sync::Notify::new());
let update_check_task = spawn_update_check_task(
Arc::clone(&state),
ws_port,
Arc::clone(&restart_notify),
Arc::clone(&transaction),
Arc::clone(&stopping),
);
let info = daemon_info::DaemonInfo::now(
std::process::id(),
sock_path.clone(),
ws_port,
env!("CARGO_PKG_VERSION"),
);
)
.with_host_managed(!cfg.detached);
daemon_info::write(&info).context("write daemon.json")?;
info!(
pid = info.pid,
@@ -381,6 +593,10 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
sock = %sock_path.display(),
"daemon ready"
);
if let Some(record) = resumed {
record.confirm_serving();
}
remove_update_leftovers_after(cfg.replaces);
// Best-effort: keep installed agent skills in step with this
// daemon's bundled SKILL.md. Spawned so a slow/failing fs
@@ -510,19 +726,23 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
let restart_notified = restart_notify.notified();
tokio::pin!(restart_notified);
tokio::select! {
let reason = tokio::select! {
_ = wait_for_shutdown() => {
info!("bsk daemon shutting down (signal)");
StopReason::Signal
}
res = idle_task => {
if matches!(res, Ok(Some(()))) {
info!("bsk daemon shutting down (idle)");
}
StopReason::Idle
}
_ = &mut restart_notified => {
info!("bsk daemon shutting down (auto-update restart)");
info!("bsk daemon releasing its endpoints to the replacement (auto-update)");
StopReason::Handover
}
}
};
stopping.store(true, std::sync::atomic::Ordering::SeqCst);
let _ = ipc_shutdown_tx.send(());
let _ = ipc_task.await;
@@ -537,11 +757,43 @@ pub fn run_foreground(cfg: DaemonConfig) -> Result<()> {
let _ = daemon_info::remove();
let _ = std::fs::remove_file(&sock_path);
Result::<()>::Ok(())
})?;
Result::<StopReason>::Ok(reason)
}
});
// Dropping the runtime waits for blocking update steps still running, so
// the slot now holds whatever they installed or started.
drop(runtime);
let transaction = lock_slot(&transaction).take();
drop(lock);
Ok(())
match (reason, transaction) {
(Ok(StopReason::Handover), Some(Transaction::HandingOver(pending))) => {
Ok(Stopped::HandOver(Box::new(pending)))
}
(reason, transaction) => {
if let Some(transaction) = transaction {
transaction.abandon();
}
reason.map(|_| Stopped::Shutdown)
}
}
}
/// Remove files earlier updates left next to this executable. A predecessor
/// still runs from its previous executable until it exits, so wait for it.
fn remove_update_leftovers_after(predecessor: Option<u32>) {
let Ok(exe) = crate::cli::update::installed_executable() else {
return;
};
// A plain thread, so a slow predecessor never delays this daemon's exit.
std::thread::spawn(move || {
if let Some(pid) = predecessor {
let deadline = Instant::now() + Duration::from_secs(120);
while lockfile::pid_alive(pid) && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(200));
}
}
crate::cli::update::remove_update_leftovers(exe);
});
}
/// Spawn the browser-liveness reaper shared by the production foreground
@@ -663,19 +915,24 @@ pub(crate) fn spawn_session_idle_reaper(state: Arc<DaemonState>) -> tokio::task:
/// it never delays daemon exit (an in-flight fetch is bounded by the
/// update client's own timeout and detached on abort).
///
/// When a tick finds a newer version the daemon also *installs* it
/// When a tick finds a newer version a detached daemon also *installs* it
/// (auto-update, on by default; [`crate::cli::update::AUTO_UPDATE_ENV`]
/// `=off` disables it and keeps the check cache/hint-only). Safety
/// gate: while any agent session is live the tick postpones the
/// install and retries next time. Once the new binary is in place the
/// task spawns the replacement daemon (see
/// [`DAEMON_REPLACEMENT_WAIT_ENV`]) and fires `restart` so this process
/// shuts down and the new version takes over; on Windows the
/// helper waits for this process to exit, replaces the binary, and starts
/// the new daemon with the same configuration.
pub(crate) fn spawn_update_check_task(
/// `=off` disables it and keeps the check cache/hint-only). A daemon owned
/// by a terminal or supervisor only reports the version, since nothing but
/// its owner can restart it. Safety gate: while any agent session is live
/// the tick postpones the install and retries next time. Once the new
/// executable is installed and has passed its self-check, this process
/// spawns the replacement daemon (see [`DAEMON_REPLACEMENT_WAIT_ENV`]),
/// leaves it in `handover_slot` and fires `restart`; [`handover`] then
/// confirms the replacement serves before this process exits, or restores
/// the previous executable and serves again. A failed attempt is retried
/// after [`crate::cli::update::state::RETRY_AFTER_FAILURE`].
fn spawn_update_check_task(
state: Arc<DaemonState>,
ws_port: u16,
restart: Arc<tokio::sync::Notify>,
transaction: TransactionSlot,
stopping: Arc<std::sync::atomic::AtomicBool>,
) -> tokio::task::JoinHandle<()> {
use crate::cli::update;
@@ -692,18 +949,20 @@ pub(crate) fn spawn_update_check_task(
return;
}
};
// Capture our own executable path once, up front: after an
// auto-update replaces the binary, `current_exe` on Linux starts
// returning a " (deleted)"-suffixed path that can neither be
// replaced again nor spawned.
let exe_path = match std::env::current_exe() {
Ok(exe) => Some(exe),
// Captured when the process started: after an update replaced the
// binary, `current_exe` on Linux returns a " (deleted)"-suffixed path
// for good, even once a rollback resumes this process.
let exe_path = match update::installed_executable() {
Ok(exe) => Some(exe.to_path_buf()),
Err(err) => {
warn!(error = %err, "auto-update install disabled: cannot locate current executable");
None
}
};
// The version and reason last reported as someone else's to install,
// so the warning is logged once rather than on every tick.
let mut reported: Option<(semver::Version, update::state::SkipReason)> = None;
let mut ticker = tokio::time::interval(update::UPDATE_CHECK_INTERVAL);
ticker.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
loop {
@@ -725,37 +984,61 @@ pub(crate) fn spawn_update_check_task(
continue;
}
let policy = auto_update_policy(
update::auto_update_enabled() && exe_path.is_some(),
state.config.detached,
);
let result = {
let cache_path = cache_path.clone();
let state = Arc::clone(&state);
let exe_path = exe_path.clone();
let transaction = Arc::clone(&transaction);
let stopping = Arc::clone(&stopping);
tokio::task::spawn_blocking(move || {
let candidate = update::refresh_update_cache(&cache_path)?;
// Checked on every tick: permissions can change while
// the daemon runs.
let not_writable = match (policy, exe_path.as_deref()) {
(update::AutoUpdatePolicy::Install, Some(exe)) => {
update::ensure_replaceable(exe).err()
}
_ => None,
};
let policy = if not_writable.is_some() {
update::AutoUpdatePolicy::NotWritable
} else {
policy
};
let candidate = update::refresh_update_cache(
&cache_path,
policy == update::AutoUpdatePolicy::Install,
)?;
// The session gate is read after the fetch, as late
// as possible before the binary gets replaced.
let active_sessions = state.sessions.len();
let auto_update = update::auto_update_enabled() && exe_path.is_some();
update::auto_update_step(
let last_attempt = update::state::current();
let outcome = update::auto_update_step(
candidate.as_ref(),
auto_update,
policy,
active_sessions,
last_attempt.as_ref(),
update::now_epoch_secs(),
|candidate| {
let target =
exe_path.as_deref().context("current executable unknown")?;
update::self_install_candidate(
prepare_handover(
candidate,
target,
&restart_start_args(&state.config)?,
exe_path.as_deref(),
&transaction,
&stopping,
)
},
)
)?;
anyhow::Ok((outcome, not_writable, last_attempt))
})
.await
};
let outcome = match result {
Ok(Ok(outcome)) => outcome,
let (outcome, not_writable, last_attempt) = match result {
Ok(Ok(checked)) => checked,
Ok(Err(err)) => {
warn!(error = %err, "periodic update check failed");
warn!(error = %format_args!("{err:#}"), "periodic update check failed");
continue;
}
Err(err) => {
@@ -771,41 +1054,83 @@ pub(crate) fn spawn_update_check_task(
%latest,
"periodic update check found a new version; auto-update off, CLI hint only"
),
update::AutoUpdateOutcome::HostManaged { latest } => {
let reason = update::state::SkipReason::HostManaged;
if reported.as_ref() != Some(&(latest.clone(), reason)) {
record_skip(
&latest,
exe_path.as_deref(),
reason,
None,
last_attempt.as_ref(),
);
warn!(
%latest,
"a new bsk version is available; this daemon belongs to its terminal or supervisor, so run `bsk update` and restart it there"
);
reported = Some((latest, reason));
}
}
update::AutoUpdateOutcome::NotWritable { latest } => {
let reason = update::state::SkipReason::NotWritable;
if reported.as_ref() != Some(&(latest.clone(), reason)) {
let error = not_writable
.as_ref()
.map(|err| format!("{err:#}"))
.unwrap_or_default();
let hint = exe_path
.as_deref()
.map(update::installer_hint)
.unwrap_or_default();
record_skip(
&latest,
exe_path.as_deref(),
reason,
not_writable.as_ref(),
last_attempt.as_ref(),
);
warn!(
%latest,
%error,
%hint,
"a new bsk version is available, but bsk cannot install it next to its executable"
);
reported = Some((latest, reason));
}
}
update::AutoUpdateOutcome::PostponedSessions { latest, sessions } => info!(
%latest,
sessions,
"auto-update postponed: agent session(s) active; will retry on the next tick"
),
update::AutoUpdateOutcome::Staged { latest } => {
info!(%latest, "update helper ready; exiting so it can replace and restart the daemon");
restart.notify_one();
return;
}
update::AutoUpdateOutcome::Replaced { latest } => {
update::AutoUpdateOutcome::Deferred { latest, until } => info!(
%latest,
retry_after_epoch_secs = until,
"auto-update of this version failed recently; see `bsk doctor`, or run `bsk update` to retry now"
),
update::AutoUpdateOutcome::Installed {
latest,
installed: (),
} => {
info!(
current = env!("CARGO_PKG_VERSION"),
%latest,
"auto-update installed the new bsk binary; restarting daemon"
"new bsk executable installed and checked; starting the replacement daemon"
);
// `exe_path` is always Some here: the install only
// runs when it was captured.
if let Some(exe) = &exe_path {
match restart_start_args(&state.config).and_then(|args| {
spawn_detached_at(exe, &args, Some(std::process::id())).map(drop)
}) {
Ok(()) => {
info!(
pid = std::process::id(),
"replacement daemon spawned; exiting so it can take over"
);
restart.notify_one();
return;
}
Err(err) => warn!(
error = %err,
"auto-update replaced the binary but failed to spawn the replacement daemon; the next daemon start picks up the new version"
),
let exe = exe_path.clone();
let config = state.config.clone();
let transaction = Arc::clone(&transaction);
let started = tokio::task::spawn_blocking(move || {
start_replacement(&transaction, exe.as_deref(), &config, ws_port)
})
.await;
match started {
Ok(true) => {
restart.notify_one();
return;
}
Ok(false) => {}
Err(err) => warn!(error = %err, "starting the replacement daemon panicked"),
}
}
}
@@ -813,15 +1138,133 @@ pub(crate) fn spawn_update_check_task(
})
}
/// An installed update a daemon has yet to hand over to. Holds the update
/// lock through [`crate::cli::update::Installed`].
struct Prepared {
installed: crate::cli::update::Installed,
record: crate::cli::update::state::UpdateRecord,
}
/// Install and self-check `candidate` under the update lock, leaving the
/// result in `transaction`. Blocking. Installs nothing once the daemon
/// started stopping; any later stop finds the result in `transaction`.
fn prepare_handover(
candidate: &crate::cli::update::UpdateCandidate,
exe: Option<&Path>,
transaction: &TransactionSlot,
stopping: &std::sync::atomic::AtomicBool,
) -> Result<()> {
use crate::cli::update::{
self,
state::{UpdateLock, UpdateRecord, UpdateSource},
};
let exe = exe.context("current executable unknown")?;
let lock = UpdateLock::try_acquire(exe)?;
let mut record = UpdateRecord::start(UpdateSource::Daemon, &candidate.latest, exe);
record.save();
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(())
}
/// Spawn the replacement daemon from the new executable, on the port this
/// daemon serves, and leave it in `transaction`. If it cannot even be
/// started, restore the previous executable and keep serving. Returns
/// whether a replacement waits for the lock. Blocking.
fn start_replacement(
transaction: &TransactionSlot,
exe: Option<&Path>,
config: &DaemonConfig,
ws_port: u16,
) -> bool {
use crate::cli::update::state::{Recovery, UpdateStage};
let mut slot = lock_slot(transaction);
let Some(Transaction::Installed(Prepared {
installed,
mut record,
})) = slot.take()
else {
return false;
};
record.enter(UpdateStage::Handover);
let spawned = exe
.context("current executable unknown")
.and_then(|exe| {
let args = restart_start_args(config, ws_port)?;
spawn_detached_at(exe, &args, Some(std::process::id()))
})
.context("start the replacement daemon");
match spawned {
Ok(child) => {
let child_pid = child.id();
info!(
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,
}));
true
}
Err(err) => {
warn!(
error = %format_args!("{err:#}"),
"could not start the replacement daemon; restoring the previous executable and serving on"
);
let recovery = installed.roll_back(|| Recovery::Restored {
daemon_serving: true,
});
record.fail(&err, recovery);
false
}
}
}
/// Record, once per version, that someone other than this daemon must
/// install `latest`.
fn record_skip(
latest: &semver::Version,
exe: Option<&Path>,
reason: crate::cli::update::state::SkipReason,
error: Option<&anyhow::Error>,
last_attempt: Option<&crate::cli::update::state::UpdateRecord>,
) {
use crate::cli::update::state::{UpdateRecord, UpdateSource};
let (Some(exe), false) = (
exe,
last_attempt.is_some_and(|record| record.skips(latest, reason)),
) else {
return;
};
UpdateRecord::skipped(UpdateSource::Daemon, latest, exe, reason, error).save();
}
/// Only a daemon `bsk` started in the background may replace itself; one
/// owned by a terminal or supervisor must be restarted by that owner.
fn auto_update_policy(enabled: bool, detached: bool) -> crate::cli::update::AutoUpdatePolicy {
use crate::cli::update::AutoUpdatePolicy;
match (enabled, detached) {
(false, _) => AutoUpdatePolicy::Disabled,
(true, false) => AutoUpdatePolicy::HostManaged,
(true, true) => AutoUpdatePolicy::Install,
}
}
/// Rebuild the `StartArgs` for the replacement daemon from the running
/// config so the respawn keeps the same port and idle timeouts.
fn restart_start_args(cfg: &DaemonConfig) -> Result<StartArgs> {
/// config so the respawn keeps the port it serves (`ws_port`, which differs
/// from the configured one for `--port 0`) and its idle timeouts.
fn restart_start_args(cfg: &DaemonConfig, ws_port: u16) -> Result<StartArgs> {
anyhow::ensure!(
cfg.server.is_none(),
"server restart is managed by the deployment supervisor"
);
Ok(StartArgs {
port: Some(cfg.ws_port),
port: Some(ws_port),
foreground: false,
session_idle: Some(cfg.session_idle),
daemon_idle: Some(cfg.daemon_idle),
@@ -945,22 +1388,12 @@ fn is_daemonized_child() -> bool {
std::env::var(DAEMONIZED_ENV).as_deref() == Ok("1")
}
/// Self-update handoff: when the outgoing daemon spawned us as its
/// replacement ([`DAEMON_REPLACEMENT_WAIT_ENV`]), wait for its pid to
/// exit so the daemon lock, IPC socket, and WS port are free before we
/// try to take them over. Bounded on purpose — the lockfile is the real
/// backstop if the predecessor somehow lingers.
fn wait_for_replaced_daemon() {
let pid = std::env::var(DAEMON_REPLACEMENT_WAIT_ENV)
/// The daemon that spawned this one as its replacement
/// ([`DAEMON_REPLACEMENT_WAIT_ENV`]), if any.
fn replaced_daemon() -> Option<u32> {
std::env::var(DAEMON_REPLACEMENT_WAIT_ENV)
.ok()
.and_then(|raw| raw.parse::<u32>().ok());
let Some(pid) = pid else {
return;
};
let deadline = Instant::now() + Duration::from_secs(30);
while lockfile::pid_alive(pid) && Instant::now() < deadline {
std::thread::sleep(Duration::from_millis(50));
}
.and_then(|raw| raw.parse::<u32>().ok())
}
#[cfg(unix)]
@@ -1377,18 +1810,33 @@ mod tests {
#[test]
fn restart_start_args_preserve_the_running_config() {
let cfg = DaemonConfig {
ws_port: 1234,
ws_port: 0,
session_idle: Duration::from_secs(11),
daemon_idle: Duration::from_secs(22),
..DaemonConfig::new(0)
};
let args = restart_start_args(&cfg).unwrap();
// The port actually bound, not the configured `--port 0`.
let args = restart_start_args(&cfg, 1234).unwrap();
assert_eq!(args.port, Some(1234));
assert!(!args.foreground);
assert_eq!(args.session_idle, Some(Duration::from_secs(11)));
assert_eq!(args.daemon_idle, Some(Duration::from_secs(22)));
}
#[test]
fn only_detached_daemons_install_updates_themselves() {
use crate::cli::update::AutoUpdatePolicy;
assert_eq!(auto_update_policy(true, true), AutoUpdatePolicy::Install);
assert_eq!(
auto_update_policy(true, false),
AutoUpdatePolicy::HostManaged
);
assert_eq!(auto_update_policy(false, true), AutoUpdatePolicy::Disabled);
assert_eq!(auto_update_policy(false, false), AutoUpdatePolicy::Disabled);
assert!(!DaemonConfig::new(0).detached);
assert!(!DaemonConfig::from(&StartArgs::default()).detached);
}
#[test]
fn automatic_restart_rejects_server_configuration() {
let mut cfg = DaemonConfig::new(0);
@@ -1399,7 +1847,7 @@ mod tests {
}
.server_config()
.unwrap();
assert!(restart_start_args(&cfg).is_err());
assert!(restart_start_args(&cfg, 52800).is_err());
}
#[test]
+457
View File
@@ -0,0 +1,457 @@
//! Handing a running daemon over to one started from a newly installed
//! executable.
//!
//! The outgoing daemon spawns its replacement while it still serves, then
//! releases the daemon lock, IPC endpoint and port but keeps running. It exits
//! only once a daemon other than itself answers status requests on the IPC
//! endpoint. If the replacement exits or is not ready in time, the outgoing
//! daemon puts the previous executable back, takes the lock again and resumes
//! serving on the same port. Browsers reconnect across the gap either way.
use std::time::{Duration, Instant};
use anyhow::{Context, Result};
use tracing::{error, info, warn};
use super::DaemonChild;
use crate::cli::update::Installed;
use crate::cli::update::state::{Recovery, UpdateRecord};
use crate::daemon::info::DaemonInfo;
use crate::daemon::lockfile::{self, DaemonLock};
use crate::daemon::paths;
use crate::daemon::probe::{self, PROBE_TIMEOUT, Probe};
/// How long the outgoing daemon waits for its replacement to answer.
pub(super) const HANDOVER_TIMEOUT: Duration = Duration::from_secs(20);
/// How long a replacement waits for its predecessor to release the lock. It
/// starts counting before the predecessor has shut down, and is well above
/// [`HANDOVER_TIMEOUT`], so the predecessor decides when a handover failed.
pub(super) const REPLACEMENT_LOCK_WAIT: Duration = Duration::from_secs(60);
/// How long the outgoing daemon tries to take the lock back after a failure.
const RECLAIM_WAIT: Duration = Duration::from_secs(5);
const POLL: Duration = Duration::from_millis(50);
const FAILURE_REPORT_LIMIT: usize = 4096;
/// A replacement that has been spawned and waits for the daemon lock.
pub(crate) struct Pending {
pub(super) child: DaemonChild,
pub(super) child_pid: u32,
/// The WS port this daemon serves, which the replacement must serve too.
pub(super) port: u16,
/// Holds the update lock until the handover is confirmed or undone.
pub(super) installed: Installed,
pub(super) record: UpdateRecord,
}
pub(super) enum Finished {
/// Another daemon serves, or none can; this process exits.
Exit,
/// The handover failed and this process serves again under the lock, on
/// `port`, the one it served before. `record` says so, and is confirmed
/// once serving has resumed.
Resume {
lock: DaemonLock,
record: Box<UpdateRecord>,
port: u16,
},
}
/// The daemon a waiter accepts as serving.
pub(crate) struct Expected<'a> {
/// The WS port browsers connect to; `None` accepts any.
pub(crate) port: Option<u16>,
pub(crate) version: &'a str,
}
/// What one status probe found, judged against [`Expected`].
pub(super) enum Found {
Serving(DaemonInfo),
/// A daemon answers, but not the expected one.
Other(String),
Nothing,
}
/// Complete a handover once this process has stopped serving and released
/// the daemon lock.
pub(super) fn finish(mut pending: Pending) -> Finished {
let own_pid = std::process::id();
let target = pending.record.target_version.clone();
let expected = Expected {
port: Some(pending.port),
version: &target,
};
let waited = wait_until_serving(
&mut pending.child,
pending.child_pid,
HANDOVER_TIMEOUT,
|| observe(Some(own_pid), &expected),
);
match waited {
Ok(daemon) => {
info!(pid = daemon.pid, version = %daemon.version, "replacement daemon is serving; exiting");
pending.record.succeed(Some((daemon.pid, daemon.version)));
pending.installed.discard();
Finished::Exit
}
Err(err) => {
warn!(
error = %format_args!("{err:#}"),
"handover failed; restoring the previous executable"
);
recover(pending, &err)
}
}
}
/// Put the previous executable back and serve again. Restoring comes first,
/// so a daemon another client starts meanwhile runs the previous version too.
/// The record claims a serving daemon only for one that already serves the
/// original port; a resumed one confirms it after publishing `daemon.json`.
fn recover(pending: Pending, err: &anyhow::Error) -> Finished {
let Pending {
installed,
mut record,
port,
..
} = pending;
let restored = installed.restore();
if let Err(restore) = &restored {
error!(error = %format_args!("{restore:#}"), "could not restore the previous bsk executable");
}
let lock = match reclaim_lock() {
Ok(lock) => lock,
Err(lock_err) => {
error!(error = %format_args!("{lock_err:#}"), "could not take the daemon lock back");
None
}
};
let daemon_serving = lock.is_none() && serving_on(port);
record.fail(
err,
match restored {
Ok(()) => Recovery::Restored { daemon_serving },
Err(_) => installed.restore_failed(),
},
);
match lock {
Some(lock) => {
info!("resuming service with the previous version");
Finished::Resume {
lock,
record: Box::new(record),
port,
}
}
None if daemon_serving => {
info!("another daemon serves the port; exiting");
Finished::Exit
}
None => {
error!("no daemon serves the port after the failed handover; run `bsk daemon start`");
Finished::Exit
}
}
}
/// Stop a replacement that has not taken over, and put the previous
/// executable back, when this daemon stops for another reason first.
pub(super) fn abandon(mut pending: Pending, reason: &str) {
stop(&mut pending.child);
let err = anyhow::anyhow!("handover abandoned: {reason}");
let recovery = pending.installed.roll_back(|| Recovery::Restored {
daemon_serving: false,
});
pending.record.fail(&err, recovery);
}
/// Wait until `observe` finds the expected daemon serving, failing as soon as
/// `child` exits or `timeout` passes; a child still running then is stopped
/// and reaped, so it releases whatever it holds. Another client may start
/// the expected daemon meanwhile, which counts; one that answers with another
/// version or port does not, and is named in the error.
pub(super) fn wait_until_serving(
child: &mut DaemonChild,
child_pid: u32,
timeout: Duration,
mut observe: impl FnMut() -> Found,
) -> Result<DaemonInfo> {
let deadline = Instant::now() + timeout;
let mut other = None;
let mut look = |other: &mut Option<String>| match observe() {
Found::Serving(daemon) => Some(daemon),
Found::Other(found) => {
*other = Some(found);
None
}
Found::Nothing => None,
};
loop {
if let Some(daemon) = look(&mut other) {
return Ok(daemon);
}
if let Some(status) = child.try_wait().context("check the new daemon")? {
if let Some(daemon) = look(&mut other) {
return Ok(daemon);
}
anyhow::bail!(
"the new daemon (pid {child_pid}) exited with {status} before it was ready{}{}",
startup_failure(child_pid),
describe_other(other.as_deref())
);
}
if Instant::now() >= deadline {
stop(child);
anyhow::bail!(
"the new daemon (pid {child_pid}) was not ready within {timeout:?} and was stopped{}{}",
startup_failure(child_pid),
describe_other(other.as_deref())
);
}
std::thread::sleep(POLL);
}
}
fn describe_other(other: Option<&str>) -> String {
other.map_or_else(String::new, |other| {
format!("; meanwhile a different daemon answered: {other}")
})
}
/// Probe the IPC endpoint once. `exclude` is a daemon that never counts,
/// such as the one handing over.
pub(super) fn observe(exclude: Option<u32>, expected: &Expected<'_>) -> Found {
match probe::probe(PROBE_TIMEOUT) {
Ok(Probe::Ready(daemon)) if Some(daemon.status.pid) != exclude => judge(
daemon.status.pid,
&daemon.status.daemon_version,
daemon.status.ws_port,
expected,
)
.map_or_else(Found::Other, |()| Found::Serving(daemon.info)),
_ => Found::Nothing,
}
}
/// Whether a daemon answering as `pid`, `version` on `port` is the expected
/// one; if not, a description of what answered instead.
fn judge(pid: u32, version: &str, port: u16, expected: &Expected<'_>) -> Result<(), String> {
let port_matches = expected.port.is_none_or(|expected| expected == port);
if version == expected.version && port_matches {
return Ok(());
}
Err(format!(
"pid {pid}, bsk {version} on port {port}, expected bsk {}{}",
expected.version,
expected
.port
.map(|port| format!(" on port {port}"))
.unwrap_or_default()
))
}
/// A replacement that fails to start leaves the reason for its predecessor,
/// which reports it in the update record.
pub(super) fn report_startup_failure(err: &anyhow::Error) {
if let Ok(path) = paths::replacement_failure_path(std::process::id()) {
let _ = std::fs::write(path, format!("{err:#}"));
}
}
/// `": <reason>"` from a replacement's startup failure report, if it left one.
fn startup_failure(pid: u32) -> String {
let Ok(path) = paths::replacement_failure_path(pid) else {
return String::new();
};
let Ok(bytes) = std::fs::read(&path) else {
return String::new();
};
let _ = std::fs::remove_file(&path);
let reason = String::from_utf8_lossy(&bytes[..bytes.len().min(FAILURE_REPORT_LIMIT)]);
match reason.trim() {
"" => String::new(),
reason => format!(": {reason}"),
}
}
/// Whether some daemon other than this process serves `port`.
fn serving_on(port: u16) -> bool {
matches!(
probe::probe(PROBE_TIMEOUT),
Ok(Probe::Ready(daemon))
if daemon.status.pid != std::process::id() && daemon.status.ws_port == port
)
}
/// `None` when another process keeps the lock: either it serves already, or
/// it stays stuck and this process cannot serve either.
fn reclaim_lock() -> Result<Option<DaemonLock>> {
let deadline = Instant::now() + RECLAIM_WAIT;
loop {
match lockfile::acquire() {
Ok(lock) => return Ok(Some(lock)),
Err(err) if err.is::<lockfile::AlreadyLocked>() => {
if Instant::now() >= deadline {
return Ok(None);
}
std::thread::sleep(POLL);
}
Err(err) => return Err(err),
}
}
}
fn stop(child: &mut DaemonChild) {
let _ = child.kill();
let _ = child.wait();
}
#[cfg(test)]
mod judge_tests {
use super::*;
#[test]
fn only_the_expected_version_on_the_expected_port_counts() {
let expected = Expected {
port: Some(52719),
version: "999.0.0",
};
assert!(judge(7, "999.0.0", 52719, &expected).is_ok());
let other = judge(7, "0.3.1", 52720, &expected).unwrap_err();
assert_eq!(
other,
"pid 7, bsk 0.3.1 on port 52720, expected bsk 999.0.0 on port 52719"
);
assert!(
judge(7, "0.3.1", 52719, &expected).is_err(),
"wrong version"
);
assert!(judge(7, "999.0.0", 52800, &expected).is_err(), "wrong port");
let any_port = Expected {
port: None,
version: "999.0.0",
};
assert!(judge(7, "999.0.0", 40000, &any_port).is_ok());
}
}
#[cfg(all(test, unix))]
mod tests {
use super::*;
use crate::daemon::test_support::isolated;
fn spawn(script: &str) -> DaemonChild {
std::process::Command::new("/bin/sh")
.args(["-c", script])
.spawn()
.unwrap()
}
fn serving(pid: u32) -> DaemonInfo {
DaemonInfo::now(pid, "sock".into(), 52719, "999.0.0")
}
#[test]
fn the_expected_daemon_counts_even_when_another_client_started_it() {
let mut child = spawn("sleep 30");
let pid = child.id();
let mut probes = 0;
let daemon = wait_until_serving(&mut child, pid, Duration::from_secs(10), || {
probes += 1;
if probes < 3 {
Found::Nothing
} else {
Found::Serving(serving(4242))
}
})
.unwrap();
assert_eq!(daemon.pid, 4242);
stop(&mut child);
}
#[test]
fn a_different_daemon_answering_is_not_a_successful_handover() {
isolated(
concat!(
module_path!(),
"::a_different_daemon_answering_is_not_a_successful_handover"
),
|| {
let mut child = spawn("sleep 0.3; exit 1");
let pid = child.id();
let error = wait_until_serving(&mut child, pid, Duration::from_secs(10), || {
Found::Other(
"pid 9, bsk 0.3.1 on port 52720, expected bsk 999.0.0 on port 52719".into(),
)
})
.unwrap_err();
let error = format!("{error:#}");
assert!(error.contains("exited with"), "{error}");
assert!(
error.contains("a different daemon answered: pid 9, bsk 0.3.1 on port 52720"),
"{error}"
);
},
);
}
#[test]
fn a_replacement_that_exits_early_fails_with_its_reported_reason() {
isolated(
concat!(
module_path!(),
"::a_replacement_that_exits_early_fails_with_its_reported_reason"
),
|| {
paths::ensure_bsk_home().unwrap();
let mut child = spawn("sleep 0.2; exit 3");
let pid = child.id();
let report = paths::replacement_failure_path(pid).unwrap();
std::fs::write(&report, "bind WS server: address in use").unwrap();
let error =
wait_until_serving(&mut child, pid, Duration::from_secs(10), || Found::Nothing)
.unwrap_err();
let error = format!("{error:#}");
assert!(error.contains(&format!("pid {pid}")), "{error}");
assert!(error.contains("before it was ready"), "{error}");
assert!(error.contains("address in use"), "{error}");
assert!(!report.exists(), "the report is consumed");
},
);
}
#[test]
fn a_replacement_that_never_becomes_ready_is_stopped() {
isolated(
concat!(
module_path!(),
"::a_replacement_that_never_becomes_ready_is_stopped"
),
|| {
paths::ensure_bsk_home().unwrap();
let mut child = spawn("sleep 30");
let pid = child.id();
let started = Instant::now();
let error = wait_until_serving(&mut child, pid, Duration::from_millis(300), || {
Found::Nothing
})
.unwrap_err();
assert!(started.elapsed() < Duration::from_secs(10));
let error = format!("{error:#}");
assert!(error.contains("was not ready within"), "{error}");
assert!(
child.try_wait().unwrap().is_some(),
"the replacement must not linger and take over later"
);
},
);
}
}
@@ -56,6 +56,35 @@ pub(super) fn spawn(exe: &Path, args: &StartArgs, predecessor_pid: Option<u32>)
Ok(child)
}
/// Whether [`spawn`] can start a daemon outside this process's Jobs. Creates
/// `exe` suspended with the same flags, checks it left every Job, and
/// terminates it without ever letting it run.
pub(super) fn check_breakaway(exe: &Path) -> Result<()> {
let command_line = command_line([exe.as_os_str(), OsStr::new("--version")].into_iter());
let env: Vec<_> = std::env::vars_os().collect();
let input = File::open("NUL").context("open probe stdin")?;
let output = File::options()
.write(true)
.open("NUL")
.context("open probe output")?;
let mut probe = windows_process::spawn(
exe.as_os_str(),
&command_line,
&env,
[&input, &output, &output],
DETACHED_PROCESS | CREATE_NEW_PROCESS_GROUP | CREATE_BREAKAWAY_FROM_JOB | CREATE_SUSPENDED,
)
.context(DETACH_HINT)?;
let outside = probe.outside_job();
let _ = probe.kill();
let _ = probe.wait();
anyhow::ensure!(
outside.context("check the probe's Job membership")?,
"{DETACH_HINT}: breakaway from an outer Job is not allowed"
);
Ok(())
}
/// Quote argv directly for the Windows CRT, without invoking a shell. Preserve
/// UTF-16 paths, embedded quotes and backslashes before a closing quote.
pub(super) fn command_line<'a>(args: impl Iterator<Item = &'a OsStr>) -> OsString {
+11 -1
View File
@@ -17,7 +17,7 @@ use windows_sys::Win32::Foundation::{
use windows_sys::Win32::System::JobObjects::IsProcessInJob;
use windows_sys::Win32::System::Threading::{
CREATE_UNICODE_ENVIRONMENT, CreateProcessW, DeleteProcThreadAttributeList,
EXTENDED_STARTUPINFO_PRESENT, GetCurrentProcess, GetExitCodeProcess,
EXTENDED_STARTUPINFO_PRESENT, GetCurrentProcess, GetExitCodeProcess, GetProcessId,
InitializeProcThreadAttributeList, PROC_THREAD_ATTRIBUTE_HANDLE_LIST, PROCESS_INFORMATION,
ResumeThread, STARTF_USESTDHANDLES, STARTUPINFOEXW, TerminateProcess,
UpdateProcThreadAttribute, WaitForSingleObject,
@@ -29,6 +29,11 @@ pub(crate) struct Process {
}
impl Process {
pub(crate) fn id(&self) -> u32 {
// SAFETY: the owned process handle remains valid for the call.
unsafe { GetProcessId(self.handle.as_raw_handle()) }
}
pub(crate) fn try_wait(&mut self) -> io::Result<Option<ExitStatus>> {
// SAFETY: the owned process handle remains valid for both calls.
match unsafe { WaitForSingleObject(self.handle.as_raw_handle(), 0) } {
@@ -63,6 +68,11 @@ impl Process {
Ok(())
}
/// Whether the process belongs to no Job, including outer ones.
pub(crate) fn outside_job(&self) -> io::Result<bool> {
Ok(!process_in_job(self.handle.as_raw_handle())?)
}
/// Called only for a child created suspended. A nested Job may allow only
/// partial breakaway, so successful CreateProcessW is not sufficient.
pub(crate) fn resume_outside_job(&self) -> io::Result<()> {
@@ -0,0 +1,804 @@
//! A detached daemon that auto-updates hands over to a daemon started from
//! the new executable, and exits only once that daemon serves the release's
//! version on its port. When the new executable fails its self-check, or its
//! daemon cannot start, the previous executable is put back and the running
//! daemon keeps serving on its port. `bsk update` restarts the daemon with the
//! same guarantees.
//!
//! On Windows these tests need a host that permits Job breakaway; CI runs
//! them from `scripts/test-windows-daemon.ps1`.
mod release_fixture;
use std::cell::RefCell;
use std::fs;
use std::io::{Cursor, Read, Write};
use std::net::TcpListener;
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::thread;
use std::time::{Duration, Instant};
use sha2::{Digest, Sha256};
const EXE: &str = if cfg!(windows) { "bsk.exe" } else { "bsk" };
const MARKER: &[u8] = b"auto-update-handover-fixture";
/// Serves a manifest naming [`release_fixture::newer_version`] and an
/// archive holding `binary`.
struct ReleaseServer {
url: String,
downloads: Arc<AtomicUsize>,
stop: Arc<AtomicBool>,
worker: Option<thread::JoinHandle<()>>,
}
impl ReleaseServer {
/// `archive_delay` holds back the archive, keeping a download in flight.
fn new(binary: &[u8], archive_delay: Duration) -> Self {
let (archive, suffix) = archive(binary);
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let base = format!("http://{}", listener.local_addr().unwrap());
let mut assets = serde_json::Map::new();
assets.insert(
bsk::cli::update::current_platform_key()
.unwrap()
.to_string(),
serde_json::json!({
"url": format!("{base}/bsk{suffix}"),
"sha256": Sha256::digest(&archive).iter().map(|byte| format!("{byte:02x}")).collect::<String>(),
}),
);
let manifest = serde_json::to_vec(
&serde_json::json!({"version": release_fixture::newer_version(), "assets": assets}),
)
.unwrap();
let stop = Arc::new(AtomicBool::new(false));
let downloads = Arc::new(AtomicUsize::new(0));
let worker = {
let stop = Arc::clone(&stop);
let downloads = Arc::clone(&downloads);
let archive_request = format!("GET /bsk{suffix} ");
thread::spawn(move || {
while !stop.load(Ordering::SeqCst) {
let mut stream = match listener.accept() {
Ok((stream, _)) => stream,
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(10));
continue;
}
Err(err) => panic!("accept release request: {err}"),
};
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
let mut request = Vec::new();
let mut chunk = [0; 1024];
while !request.windows(4).any(|window| window == b"\r\n\r\n") {
match stream.read(&mut chunk) {
Ok(0) | Err(_) => break,
Ok(len) => request.extend_from_slice(&chunk[..len]),
}
}
let body = if request.starts_with(archive_request.as_bytes()) {
downloads.fetch_add(1, Ordering::SeqCst);
thread::sleep(archive_delay);
&archive
} else {
&manifest
};
let _ = write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
);
let _ = stream.write_all(body);
}
})
};
Self {
url: format!("{base}/version.json"),
downloads,
stop,
worker: Some(worker),
}
}
}
impl Drop for ReleaseServer {
fn drop(&mut self) {
self.stop.store(true, Ordering::SeqCst);
let _ = self.worker.take().unwrap().join();
}
}
#[cfg(windows)]
fn archive(binary: &[u8]) -> (Vec<u8>, &'static str) {
let mut archive = zip::ZipWriter::new(Cursor::new(Vec::new()));
archive
.start_file(
"bsk.exe",
zip::write::SimpleFileOptions::default()
.compression_method(zip::CompressionMethod::Stored),
)
.unwrap();
archive.write_all(binary).unwrap();
(archive.finish().unwrap().into_inner(), ".zip")
}
#[cfg(not(windows))]
fn archive(binary: &[u8]) -> (Vec<u8>, &'static str) {
let mut tar = tar::Builder::new(Vec::new());
let mut header = tar::Header::new_gnu();
header.set_size(binary.len() as u64);
header.set_mode(0o755);
header.set_cksum();
tar.append_data(&mut header, "bsk", Cursor::new(binary))
.unwrap();
let mut gzip = flate2::write::GzEncoder::new(Vec::new(), flate2::Compression::fast());
gzip.write_all(&tar.into_inner().unwrap()).unwrap();
(gzip.finish().unwrap(), ".tar.gz")
}
struct Fixture {
_tmp: tempfile::TempDir,
exe: PathBuf,
home: PathBuf,
original: Vec<u8>,
release: Vec<u8>,
server: ReleaseServer,
daemon: RefCell<Option<Child>>,
}
impl Fixture {
/// An installation of the current bsk whose next release is `release`.
fn new(release: impl FnOnce(&Path) -> Vec<u8>) -> Self {
Self::with_archive_delay(release, Duration::ZERO)
}
fn with_archive_delay(release: impl FnOnce(&Path) -> Vec<u8>, delay: Duration) -> Self {
let tmp = tempfile::TempDir::new().unwrap();
let dir = tmp.path().join("bin dir");
fs::create_dir(&dir).unwrap();
let exe = dir.join(EXE);
fs::copy(env!("CARGO_BIN_EXE_bsk"), &exe).unwrap();
let home = tmp.path().join("home");
fs::create_dir(&home).unwrap();
let release = release(tmp.path());
let server = ReleaseServer::new(&release, delay);
Self {
original: fs::read(&exe).unwrap(),
_tmp: tmp,
exe,
home,
release,
server,
daemon: RefCell::new(None),
}
}
fn command(&self) -> Command {
let mut command = Command::new(&self.exe);
command
.env("BSK_HOME", &self.home)
.env("BSK_UPDATE_MANIFEST_URL", &self.server.url)
.env("BSK_AUTO_UPDATE", "off")
.env("RUST_LOG", "info")
.env("NO_PROXY", "127.0.0.1,localhost")
.env_remove("BSK_DAEMONIZED")
.env_remove("BSK_DAEMON_REPLACES_PID")
.stdin(Stdio::null());
#[cfg(windows)]
{
use std::os::windows::process::CommandExt;
command.creation_flags(0x0800_0000); // CREATE_NO_WINDOW
}
command
}
/// Start a daemon as `bsk` starts one in the background, so it installs
/// updates itself, and return its pid. It checks for updates at once.
fn start_daemon(&self, port: u16) -> u32 {
self.start_daemon_idling_after(port, "60s")
}
fn start_daemon_idling_after(&self, port: u16, idle: &str) -> u32 {
let child = self
.command()
.env("BSK_AUTO_UPDATE", "on")
.env("BSK_DAEMONIZED", "1")
.args([
"daemon",
"start",
"--port",
&port.to_string(),
"--daemon-idle",
idle,
])
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap();
let pid = child.id();
*self.daemon.borrow_mut() = Some(child);
pid
}
/// Start a daemon owned by this test, as a terminal or supervisor would.
fn start_foreground_daemon(&self, port: u16) -> u32 {
let child = self
.command()
.args([
"daemon",
"start",
"--foreground",
"--port",
&port.to_string(),
])
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.unwrap();
let pid = child.id();
*self.daemon.borrow_mut() = Some(child);
pid
}
/// `bsk --json update --yes`, returning its report.
fn update(&self) -> serde_json::Value {
let out = self
.command()
.args(["--json", "update", "--yes"])
.output()
.unwrap();
assert!(
out.status.success(),
"{}",
String::from_utf8_lossy(&out.stderr)
);
serde_json::from_slice(&out.stdout).unwrap()
}
fn daemon_exited(&self) -> bool {
self.daemon
.borrow_mut()
.as_mut()
.is_some_and(|daemon| daemon.try_wait().unwrap().is_some())
}
fn json(&self, name: &str) -> Option<serde_json::Value> {
serde_json::from_slice(&fs::read(self.home.join(name)).ok()?).ok()
}
fn info(&self) -> Option<serde_json::Value> {
self.json("daemon.json")
}
fn record(&self) -> Option<serde_json::Value> {
self.json("update-state.json")
}
fn installed(&self) -> Vec<u8> {
// A scanner may briefly hold a freshly renamed file.
let deadline = Instant::now() + Duration::from_secs(5);
loop {
match fs::read(&self.exe) {
Ok(bytes) => return bytes,
Err(err) if Instant::now() < deadline => {
let _ = err;
thread::sleep(Duration::from_millis(50));
}
Err(err) => panic!("read {}: {err}", self.exe.display()),
}
}
}
/// Files next to the executable other than the executable itself and
/// the update lock, which stays.
fn leftovers(&self) -> Vec<String> {
let lock = format!(".{EXE}.update.lock");
fs::read_dir(self.exe.parent().unwrap())
.unwrap()
.flatten()
.map(|entry| entry.file_name().to_string_lossy().into_owned())
.filter(|name| name != EXE && *name != lock)
.collect()
}
fn wait_for(&self, description: &str, mut check: impl FnMut() -> bool) {
let deadline = Instant::now() + Duration::from_secs(60);
while Instant::now() < deadline {
if check() {
return;
}
thread::sleep(Duration::from_millis(50));
}
let logs: Vec<_> = fs::read_dir(&self.home)
.unwrap()
.flatten()
.filter(|entry| entry.file_name().to_string_lossy().contains("daemon.log"))
.map(|entry| fs::read_to_string(entry.path()).unwrap_or_default())
.collect();
panic!(
"timed out waiting for {description}; record: {:?}; daemon.json: {:?}; files next to {EXE}: {:?}; logs: {logs:?}",
self.record(),
self.info(),
self.leftovers()
);
}
fn status_succeeds(&self) {
let status = self.command().args(["--json", "status"]).output().unwrap();
assert!(
status.status.success(),
"{}",
String::from_utf8_lossy(&status.stderr)
);
}
}
impl Drop for Fixture {
fn drop(&mut self) {
let _ = self.command().args(["daemon", "stop"]).output();
if let Some(daemon) = self.daemon.get_mut().as_mut() {
let _ = daemon.kill();
let _ = daemon.wait();
}
}
}
/// The release: the bsk under test, reporting the newer version.
fn newer_bsk(dir: &Path) -> Vec<u8> {
release_fixture::newer_bsk(dir)
}
/// The bsk under test with a trailing marker: a different file that still
/// reports the current version, not the one its manifest names.
fn mislabelled_bsk(_: &Path) -> Vec<u8> {
let mut binary = fs::read(env!("CARGO_BIN_EXE_bsk")).unwrap();
binary.extend_from_slice(MARKER);
binary
}
/// Answers `--version` like the release; its daemon then runs `daemon`.
fn stand_in(dir: &Path, name: &str, daemon: &str) -> Vec<u8> {
release_fixture::compiled(
dir,
name,
&format!(
r#"
if std::env::args().nth(1).as_deref() == Some("--version") {{
println!("bsk {}");
return;
}}
{daemon}
"#,
release_fixture::newer_version()
),
)
}
/// Its daemon fails to start.
fn release_whose_daemon_fails(dir: &Path) -> Vec<u8> {
stand_in(dir, "daemon_fails", "std::process::exit(3);")
}
/// Its daemon takes the daemon lock, then never serves.
fn release_whose_daemon_hangs(dir: &Path) -> Vec<u8> {
stand_in(
dir,
"daemon_hangs",
r#"
let home = std::env::var_os("BSK_HOME").unwrap();
let lock = std::fs::OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(std::path::Path::new(&home).join("daemon.lock"))
.unwrap();
lock.lock().unwrap();
std::thread::sleep(std::time::Duration::from_secs(600));
"#,
)
}
/// Cannot even report its version.
fn release_that_cannot_run(dir: &Path) -> Vec<u8> {
release_fixture::compiled(dir, "cannot_run", "std::process::exit(1);")
}
/// Passes the self-check only after 3 seconds; its daemon fails to start.
fn release_with_a_slow_self_check(dir: &Path) -> Vec<u8> {
release_fixture::compiled(
dir,
"slow_self_check",
&format!(
r#"
if std::env::args().nth(1).as_deref() == Some("--version") {{
std::thread::sleep(std::time::Duration::from_secs(3));
println!("bsk {}");
return;
}}
std::process::exit(3);
"#,
release_fixture::newer_version()
),
)
}
/// Start a daemon from the installed executable, as the next command would.
fn daemon_starts_again(fixture: &Fixture) {
let port = unused_port();
let start = fixture
.command()
.args(["daemon", "start", "--port", &port.to_string()])
.output()
.unwrap();
assert!(
start.status.success(),
"{}",
String::from_utf8_lossy(&start.stderr)
);
fixture.status_succeeds();
}
fn unused_port() -> u16 {
TcpListener::bind("127.0.0.1:0")
.unwrap()
.local_addr()
.unwrap()
.port()
}
#[test]
fn auto_update_exits_only_after_the_new_daemon_serves() {
let fixture = Fixture::new(newer_bsk);
let port = unused_port();
let old_pid = fixture.start_daemon(port);
fixture.wait_for("the replacement daemon", || {
fixture
.info()
.is_some_and(|info| info["pid"] != old_pid && info["ws_port"] == port)
});
fixture.wait_for("the previous daemon to exit", || fixture.daemon_exited());
let record = fixture.record().expect("the update is recorded");
let new_pid = fixture.info().unwrap()["pid"].clone();
assert_eq!(record["source"], "daemon", "{record}");
assert_eq!(record["result"], "succeeded", "{record}");
assert_eq!(record["stage"], "handover", "{record}");
assert_eq!(
record["target_version"],
release_fixture::newer_version(),
"{record}"
);
assert_eq!(
record["daemon_version"],
release_fixture::newer_version(),
"{record}"
);
assert_eq!(record["daemon_pid"], new_pid, "{record}");
assert!(record.get("previous_executable").is_none(), "{record}");
assert!(fixture.installed() == fixture.release);
// The replacement removes the executable its predecessor ran from.
fixture.wait_for("leftover cleanup", || fixture.leftovers().is_empty());
fixture.status_succeeds();
assert_eq!(
fixture.server.downloads.load(Ordering::SeqCst),
1,
"the replacement must not download the release again"
);
}
#[test]
fn a_failed_handover_restores_the_previous_executable_and_keeps_serving() {
let fixture = Fixture::new(release_whose_daemon_fails);
let port = unused_port();
let old_pid = fixture.start_daemon(port);
fixture.wait_for("the failed handover to be recorded", || {
fixture
.record()
.is_some_and(|record| record["result"] == "failed")
});
fixture.wait_for("the previous daemon to serve again", || {
fixture
.info()
.is_some_and(|info| info["pid"] == old_pid && info["ws_port"] == port)
});
// Confirmed only once the resumed daemon has published daemon.json.
fixture.wait_for("the resumed service to be confirmed", || {
fixture.record().is_some_and(|record| {
record["recovery"] == serde_json::json!({"state": "restored", "daemon_serving": true})
})
});
let record = fixture.record().unwrap();
assert_eq!(record["stage"], "handover", "{record}");
let error = record["error"].as_str().unwrap();
assert!(error.contains("before it was ready"), "{error}");
assert!(record["retry_after_epoch_secs"].is_u64(), "{record}");
assert!(
!fixture.daemon_exited(),
"the previous daemon keeps running"
);
assert!(
fixture.installed() == fixture.original,
"the previous executable is back in place"
);
fixture.status_succeeds();
fixture.wait_for("leftover cleanup", || {
fixture
.leftovers()
.iter()
.all(|name| !name.contains(".new-"))
});
}
#[test]
fn an_update_installed_after_the_daemon_went_idle_is_undone() {
// The daemon idles out while the new executable's self-check runs, so the
// install finishes after shutdown began and must be rolled back.
let fixture = Fixture::new(release_with_a_slow_self_check);
fixture.start_daemon_idling_after(unused_port(), "1s");
fixture.wait_for("the daemon to idle out", || fixture.daemon_exited());
let record = fixture.record().unwrap();
assert_eq!(record["result"], "failed", "{record}");
assert_eq!(
record["recovery"],
serde_json::json!({"state": "restored", "daemon_serving": false}),
"{record}"
);
let error = record["error"].as_str().unwrap();
assert!(error.contains("handover abandoned"), "{error}");
assert!(
fixture.installed() == fixture.original,
"the previous executable is back in place"
);
daemon_starts_again(&fixture);
}
#[test]
fn a_download_still_in_flight_when_the_daemon_idles_out_installs_nothing() {
let fixture = Fixture::with_archive_delay(release_whose_daemon_fails, Duration::from_secs(4));
fixture.start_daemon_idling_after(unused_port(), "1s");
fixture.wait_for("the daemon to idle out", || fixture.daemon_exited());
let record = fixture.record().unwrap();
assert_eq!(record["result"], "failed", "{record}");
assert_eq!(
record["recovery"],
serde_json::json!({"state": "unchanged"}),
"{record}"
);
let error = record["error"].as_str().unwrap();
assert!(
error.contains("stopped before the update was installed"),
"{error}"
);
assert!(fixture.installed() == fixture.original);
daemon_starts_again(&fixture);
}
#[test]
fn a_dynamic_port_daemon_resumes_on_the_port_it_was_given() {
let fixture = Fixture::new(release_whose_daemon_fails);
let old_pid = fixture.start_daemon(0);
let mut port = None;
fixture.wait_for("the daemon to serve", || {
port = fixture
.info()
.filter(|info| info["pid"] == old_pid)
.and_then(|info| info["ws_port"].as_u64());
port.is_some()
});
let port = port.unwrap();
fixture.wait_for("the resumed service to be confirmed", || {
fixture.record().is_some_and(|record| {
record["recovery"] == serde_json::json!({"state": "restored", "daemon_serving": true})
})
});
let info = fixture.info().unwrap();
assert_eq!(info["pid"], old_pid, "the previous daemon serves again");
assert_eq!(info["ws_port"], port, "{info}");
let port = u16::try_from(port).unwrap();
std::net::TcpStream::connect(("127.0.0.1", port))
.expect("browsers can reconnect to the original port");
fixture.status_succeeds();
}
#[test]
fn manual_update_leaves_a_host_managed_daemon_to_its_owner() {
let fixture = Fixture::new(newer_bsk);
let port = unused_port();
let pid = fixture.start_foreground_daemon(port);
fixture.wait_for("the foreground daemon", || {
fixture
.info()
.is_some_and(|info| info["pid"] == pid && info["host_managed"] == true)
});
let report = fixture.update();
assert_eq!(report["status"], "updated", "{report}");
assert_eq!(report["daemon"], "left_to_host", "{report}");
assert!(
report["message"]
.as_str()
.unwrap()
.contains("restart it there"),
"{report}"
);
assert!(fixture.installed() == fixture.release);
assert!(!fixture.daemon_exited(), "the owner's daemon keeps running");
assert_eq!(fixture.info().unwrap()["pid"], pid);
fixture.status_succeeds();
assert_eq!(fixture.record().unwrap()["result"], "succeeded");
}
#[test]
fn manual_update_restarts_a_background_daemon_on_its_port() {
let fixture = Fixture::new(newer_bsk);
let port = unused_port();
let start = fixture
.command()
.args(["daemon", "start", "--port", &port.to_string()])
.output()
.unwrap();
assert!(
start.status.success(),
"{}",
String::from_utf8_lossy(&start.stderr)
);
let old_pid = fixture.info().unwrap()["pid"].clone();
let report = fixture.update();
assert_eq!(report["daemon"], "restarted", "{report}");
assert!(fixture.installed() == fixture.release);
let info = fixture.info().unwrap();
assert_ne!(info["pid"], old_pid);
assert_eq!(info["ws_port"], port, "restarted on the port it served");
assert!(info.get("host_managed").is_none(), "{info}");
let record = fixture.record().unwrap();
assert_eq!(record["result"], "succeeded", "{record}");
assert_eq!(record["stage"], "restart", "{record}");
assert_eq!(record["daemon_pid"], info["pid"], "{record}");
fixture.status_succeeds();
}
#[test]
fn a_release_that_cannot_run_or_reports_another_version_is_never_handed_over_to() {
for (release, expected_error) in [
(
release_that_cannot_run as fn(&Path) -> Vec<u8>,
"exited with",
),
(
mislabelled_bsk,
concat!("printed \"bsk ", env!("CARGO_PKG_VERSION"), "\""),
),
] {
let fixture = Fixture::new(release);
let port = unused_port();
let old_pid = fixture.start_daemon(port);
fixture.wait_for("the daemon to serve", || {
fixture.info().is_some_and(|info| info["pid"] == old_pid)
});
fixture.wait_for("the failed update to be recorded", || {
fixture
.record()
.is_some_and(|record| record["result"] == "failed")
});
let record = fixture.record().unwrap();
assert_eq!(record["stage"], "install", "{record}");
assert_eq!(
record["recovery"],
serde_json::json!({"state": "unchanged"}),
"{record}"
);
let error = record["error"].as_str().unwrap();
assert!(error.contains("self-check"), "{error}");
assert!(error.contains(expected_error), "{error}");
assert!(fixture.installed() == fixture.original);
assert_eq!(fixture.info().unwrap()["pid"], old_pid, "never stopped");
assert!(!fixture.daemon_exited());
fixture.status_succeeds();
}
}
#[test]
fn manual_update_stops_a_new_daemon_stuck_on_the_lock_and_restores_the_service() {
let fixture = Fixture::new(release_whose_daemon_hangs);
let port = unused_port();
let start = fixture
.command()
.args(["daemon", "start", "--port", &port.to_string()])
.output()
.unwrap();
assert!(
start.status.success(),
"{}",
String::from_utf8_lossy(&start.stderr)
);
let out = fixture
.command()
.args(["--json", "update", "--yes"])
.output()
.unwrap();
// `--json` reports the error on stdout.
let output = format!(
"{}{}",
String::from_utf8_lossy(&out.stdout),
String::from_utf8_lossy(&out.stderr)
);
assert!(!out.status.success(), "the update must fail: {output}");
assert!(output.contains("rolled back"), "{output}");
assert!(fixture.installed() == fixture.original);
// The stuck daemon was stopped, so the previous version got the lock and
// serves the original port again.
let info = fixture.info().expect("a daemon serves again");
assert_eq!(info["ws_port"], port, "{info}");
assert_eq!(info["version"], env!("CARGO_PKG_VERSION"), "{info}");
fixture.status_succeeds();
let record = fixture.record().unwrap();
assert_eq!(record["stage"], "restart", "{record}");
assert_eq!(
record["recovery"],
serde_json::json!({"state": "restored", "daemon_serving": true}),
"{record}"
);
let error = record["error"].as_str().unwrap();
assert!(error.contains("was not ready within"), "{error}");
}
#[test]
fn a_resumed_daemon_that_cannot_serve_again_is_recorded_as_not_serving() {
let fixture = Fixture::new(release_whose_daemon_fails);
let port = unused_port();
fixture.start_daemon(port);
fixture.wait_for("the daemon to serve", || {
fixture.info().is_some_and(|info| info["ws_port"] == port)
});
// Take the port the moment the daemon releases it for the handover, so
// the previous version cannot bind it again after the replacement fails.
let listener = loop {
if let Ok(listener) = TcpListener::bind(("127.0.0.1", port)) {
break listener;
}
assert!(!fixture.daemon_exited(), "the daemon exited early");
};
fixture.wait_for("the daemon to give up", || fixture.daemon_exited());
let record = fixture.record().unwrap();
assert_eq!(record["result"], "failed", "{record}");
assert_eq!(
record["recovery"],
serde_json::json!({"state": "restored", "daemon_serving": false}),
"{record}"
);
let error = record["error"].as_str().unwrap();
assert!(error.contains("before it was ready"), "{error}");
assert!(
error.contains("the previous version could not serve again"),
"{error}"
);
assert!(fixture.installed() == fixture.original);
drop(listener);
}
+181
View File
@@ -0,0 +1,181 @@
//! A daemon owned by a terminal or supervisor (`--foreground`) reports new
//! versions but never replaces itself: only its owner can restart it.
use std::fs;
use std::io::{Read, Write};
use std::net::TcpListener;
use std::path::Path;
use std::process::{Child, Command, Stdio};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::thread;
use std::time::{Duration, Instant};
/// Serves a manifest that names a newer release and counts archive downloads.
struct ReleaseServer {
url: String,
downloads: Arc<AtomicUsize>,
stop: Arc<AtomicBool>,
worker: Option<thread::JoinHandle<()>>,
}
impl ReleaseServer {
fn start() -> Self {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let base = format!("http://{}", listener.local_addr().unwrap());
let mut assets = serde_json::Map::new();
assets.insert(
bsk::cli::update::current_platform_key()
.unwrap()
.to_string(),
serde_json::json!({"url": format!("{base}/bsk.archive"), "sha256": "00"}),
);
let manifest =
serde_json::to_vec(&serde_json::json!({"version": "999.0.0", "assets": assets}))
.unwrap();
let stop = Arc::new(AtomicBool::new(false));
let downloads = Arc::new(AtomicUsize::new(0));
let worker = {
let stop = Arc::clone(&stop);
let downloads = Arc::clone(&downloads);
thread::spawn(move || {
while !stop.load(Ordering::SeqCst) {
let mut stream = match listener.accept() {
Ok((stream, _)) => stream,
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(10));
continue;
}
Err(err) => panic!("accept release request: {err}"),
};
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
let mut request = Vec::new();
let mut chunk = [0; 1024];
while !request.windows(4).any(|window| window == b"\r\n\r\n") {
match stream.read(&mut chunk) {
Ok(0) | Err(_) => break,
Ok(len) => request.extend_from_slice(&chunk[..len]),
}
}
if request.starts_with(b"GET /bsk.archive ") {
downloads.fetch_add(1, Ordering::SeqCst);
}
let _ = write!(
stream,
"HTTP/1.1 200 OK\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
manifest.len()
);
let _ = stream.write_all(&manifest);
}
})
};
Self {
url: format!("{base}/version.json"),
downloads,
stop,
worker: Some(worker),
}
}
}
impl Drop for ReleaseServer {
fn drop(&mut self) {
self.stop.store(true, Ordering::SeqCst);
let _ = self.worker.take().unwrap().join();
}
}
struct Daemon(Child);
impl Drop for Daemon {
fn drop(&mut self) {
let _ = self.0.kill();
let _ = self.0.wait();
}
}
fn bsk(home: &Path, server: &ReleaseServer) -> Command {
let mut command = Command::new(env!("CARGO_BIN_EXE_bsk"));
command
.env("BSK_HOME", home)
.env("BSK_UPDATE_MANIFEST_URL", &server.url)
.env("BSK_AUTO_UPDATE", "on")
.env("RUST_LOG", "info")
.env("NO_PROXY", "127.0.0.1,localhost")
.env_remove("BSK_DAEMONIZED")
.env_remove("BSK_DAEMON_REPLACES_PID")
.stdin(Stdio::null());
command
}
#[test]
fn foreground_daemon_reports_updates_without_replacing_itself() {
let server = ReleaseServer::start();
let tmp = tempfile::TempDir::new().unwrap();
let home = tmp.path().join("home");
fs::create_dir(&home).unwrap();
let log_path = tmp.path().join("foreground.log");
let log = fs::File::create(&log_path).unwrap();
let mut daemon = Daemon(
bsk(&home, &server)
.args(["daemon", "start", "--foreground", "--port", "0"])
.stdout(log.try_clone().unwrap())
.stderr(log)
.spawn()
.unwrap(),
);
let deadline = Instant::now() + Duration::from_secs(30);
while !fs::read_to_string(&log_path)
.unwrap_or_default()
.contains("belongs to its terminal or supervisor")
{
assert!(
daemon.0.try_wait().unwrap().is_none(),
"daemon exited: {}",
fs::read_to_string(&log_path).unwrap_or_default()
);
assert!(
Instant::now() < deadline,
"no update report: {}",
fs::read_to_string(&log_path).unwrap_or_default()
);
thread::sleep(Duration::from_millis(50));
}
assert!(
daemon.0.try_wait().unwrap().is_none(),
"the daemon keeps serving"
);
assert_eq!(
server.downloads.load(Ordering::SeqCst),
0,
"a host-managed daemon must not download the release"
);
let cache: serde_json::Value =
serde_json::from_slice(&fs::read(home.join("update-check.json")).unwrap()).unwrap();
assert_eq!(cache["latest_version"], "999.0.0");
assert_eq!(cache["auto_update"], false);
let record: serde_json::Value =
serde_json::from_slice(&fs::read(home.join("update-state.json")).unwrap()).unwrap();
assert_eq!(record["result"], "skipped", "{record}");
assert_eq!(record["skip_reason"], "host_managed", "{record}");
assert_eq!(record["target_version"], "999.0.0", "{record}");
// The CLI hint follows the daemon's policy, not the CLI's own switch.
let status = bsk(&home, &server).arg("status").output().unwrap();
let stderr = String::from_utf8_lossy(&status.stderr);
assert!(
stderr.contains(
"999.0.0. Run `bsk update`, then restart the daemon in its terminal or supervisor."
),
"{stderr}"
);
let info: serde_json::Value =
serde_json::from_slice(&fs::read(home.join("daemon.json")).unwrap()).unwrap();
assert_eq!(info["host_managed"], true, "{info}");
}
+101
View File
@@ -0,0 +1,101 @@
//! Releases of the bsk under test for update tests. A release must report the
//! version its manifest names, so tests serve a copy whose version string is
//! replaced by a newer one of the same length.
#![allow(dead_code)]
use std::collections::HashMap;
use std::fs;
use std::path::Path;
use std::sync::{Mutex, OnceLock};
/// Builds once per test process: each fixture takes seconds on Windows.
fn cached(key: &str, build: impl FnOnce() -> Vec<u8>) -> Vec<u8> {
static BUILT: OnceLock<Mutex<HashMap<String, Vec<u8>>>> = OnceLock::new();
let built = BUILT.get_or_init(Default::default);
if let Some(binary) = built.lock().unwrap().get(key) {
return binary.clone();
}
let binary = build();
built
.lock()
.unwrap()
.insert(key.to_string(), binary.clone());
binary
}
/// Newer than the bsk under test and of the same length: every digit becomes
/// 9, so `0.3.1` becomes `9.9.9`.
pub fn newer_version() -> String {
let current = env!("CARGO_PKG_VERSION");
let newer: String = current
.chars()
.map(|c| if c.is_ascii_digit() { '9' } else { c })
.collect();
assert_ne!(newer, current, "the bsk under test is already {current}");
newer
}
/// The bsk under test, reporting [`newer_version`]. On macOS the copy is
/// signed again (ad hoc), since the kernel refuses modified signed code.
pub fn newer_bsk(dir: &Path) -> Vec<u8> {
cached("newer bsk", || build_newer_bsk(dir))
}
fn build_newer_bsk(dir: &Path) -> Vec<u8> {
let current = env!("CARGO_PKG_VERSION").as_bytes();
let newer = newer_version();
let original = fs::read(env!("CARGO_BIN_EXE_bsk")).unwrap();
let mut binary = Vec::with_capacity(original.len());
let mut rest = original.as_slice();
while let Some(at) = rest
.windows(current.len())
.position(|window| window == current)
{
binary.extend_from_slice(&rest[..at]);
binary.extend_from_slice(newer.as_bytes());
rest = &rest[at + current.len()..];
}
binary.extend_from_slice(rest);
assert_ne!(binary, original, "the version string was not found");
resign(dir, binary)
}
#[cfg(target_os = "macos")]
fn resign(dir: &Path, binary: Vec<u8>) -> Vec<u8> {
let path = dir.join("newer-bsk");
fs::write(&path, binary).unwrap();
let status = std::process::Command::new("codesign")
.args(["--force", "--sign", "-"])
.arg(&path)
.stderr(std::process::Stdio::null())
.status()
.unwrap();
assert!(status.success(), "codesign {}", path.display());
fs::read(path).unwrap()
}
#[cfg(not(target_os = "macos"))]
fn resign(_dir: &Path, binary: Vec<u8>) -> Vec<u8> {
binary
}
/// Compile a stand-in for a broken release from the body of its `main`.
pub fn compiled(dir: &Path, name: &str, main: &str) -> Vec<u8> {
cached(&format!("compiled {name}"), || compile(dir, name, main))
}
fn compile(dir: &Path, name: &str, main: &str) -> Vec<u8> {
let source = dir.join(format!("{name}.rs"));
fs::write(&source, format!("fn main() {{ {main} }}")).unwrap();
let output = dir.join(format!("{name}{}", std::env::consts::EXE_SUFFIX));
let rustc = std::env::var_os("RUSTC").unwrap_or_else(|| "rustc".into());
let status = std::process::Command::new(rustc)
.args(["--edition", "2021", "-o"])
.arg(&output)
.arg(&source)
.status()
.unwrap();
assert!(status.success(), "compile {name}");
fs::read(output).unwrap()
}
+199 -40
View File
@@ -1,6 +1,8 @@
//! Exercise self-update with a real executable and a local release server.
#![cfg(windows)]
mod release_fixture;
use std::fs;
use std::io::{Cursor, Read, Write};
use std::net::{TcpListener, TcpStream};
@@ -15,8 +17,6 @@ use std::time::{Duration, Instant};
use bsk::daemon::info::DaemonInfo;
use sha2::{Digest, Sha256};
const MARKER: &[u8] = b"windows-update-regression-fixture";
struct ReleaseServer {
url: String,
requests: Arc<AtomicUsize>,
@@ -40,7 +40,7 @@ impl ReleaseServer {
listener.set_nonblocking(true).unwrap();
let url = format!("http://{}", listener.local_addr().unwrap());
let manifest = serde_json::to_vec(&serde_json::json!({
"version": "999.0.0",
"version": release_fixture::newer_version(),
"assets": {"windows-x64": {
"url": format!("{url}/bsk.zip"),
"sha256": Sha256::digest(&archive).iter().map(|byte| format!("{byte:02x}")).collect::<String>(),
@@ -138,7 +138,7 @@ fn release_server_waits_for_delayed_and_fragmented_request_headers() {
assert!(response.starts_with("HTTP/1.1 200 OK\r\n"), "{response}");
let (_, body) = response.split_once("\r\n\r\n").unwrap();
let manifest: serde_json::Value = serde_json::from_str(body).unwrap();
assert_eq!(manifest["version"], "999.0.0");
assert_eq!(manifest["version"], release_fixture::newer_version());
}
struct Fixture {
@@ -159,10 +159,8 @@ impl Fixture {
let home = tmp.path().join("home");
fs::create_dir(&home).unwrap();
fs::copy(env!("CARGO_BIN_EXE_bsk"), &exe).unwrap();
// A PE overlay distinguishes the replacement without requiring a
// second build or changing the executable's behavior/version.
let mut binary = fs::read(&exe).unwrap();
binary.extend_from_slice(MARKER);
// The same program reporting a newer version, as a release must.
let binary = release_fixture::newer_bsk(tmp.path());
let server = ReleaseServer::new(&binary);
Self {
_tmp: tmp,
@@ -196,18 +194,7 @@ impl Fixture {
}
thread::sleep(Duration::from_millis(50));
}
let diagnostics: Vec<_> = fs::read_dir(self.exe.parent().unwrap())
.unwrap()
.flatten()
.filter(|entry| entry.path().extension().is_some_and(|ext| ext == "log"))
.map(|entry| {
(
entry.path(),
String::from_utf8_lossy(&fs::read(entry.path()).unwrap_or_default())
.into_owned(),
)
})
.collect();
let leftovers = self.leftovers();
let home_logs: Vec<_> = fs::read_dir(&self.home)
.unwrap()
.flatten()
@@ -221,7 +208,7 @@ impl Fixture {
})
.collect();
panic!(
"timed out waiting for {description}; helper logs: {diagnostics:?}; daemon logs: {home_logs:?}"
"timed out waiting for {description}; files next to bsk.exe: {leftovers:?}; daemon logs: {home_logs:?}"
);
}
@@ -233,6 +220,17 @@ impl Fixture {
// A locked executable may briefly reject reads during replacement.
fs::read(&self.exe).is_ok_and(|binary| binary == self.binary)
}
/// Files next to the executable other than the executable itself and
/// the update lock, which stays.
fn leftovers(&self) -> Vec<String> {
fs::read_dir(self.exe.parent().unwrap())
.unwrap()
.flatten()
.map(|entry| entry.file_name().to_string_lossy().into_owned())
.filter(|name| name != "bsk.exe" && name != ".bsk.exe.update.lock")
.collect()
}
}
impl Drop for Fixture {
@@ -247,8 +245,148 @@ impl Drop for Fixture {
}
}
/// A Job that forbids breakaway, like the one a sandboxed agent runs
/// commands in.
struct RestrictiveJob(std::os::windows::io::OwnedHandle);
impl RestrictiveJob {
fn new() -> Self {
use std::os::windows::io::{AsRawHandle, FromRawHandle, OwnedHandle};
use windows_sys::Win32::System::JobObjects::{
CreateJobObjectW, JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE,
JOBOBJECT_EXTENDED_LIMIT_INFORMATION, JobObjectExtendedLimitInformation,
SetInformationJobObject,
};
// SAFETY: no name or inheritable security descriptor; this test owns it.
let handle = unsafe { CreateJobObjectW(std::ptr::null(), std::ptr::null()) };
assert!(!handle.is_null(), "{}", std::io::Error::last_os_error());
let job = unsafe { OwnedHandle::from_raw_handle(handle) };
let mut limits: JOBOBJECT_EXTENDED_LIMIT_INFORMATION = unsafe { std::mem::zeroed() };
limits.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE;
assert_ne!(
unsafe {
SetInformationJobObject(
job.as_raw_handle(),
JobObjectExtendedLimitInformation,
(&limits as *const JOBOBJECT_EXTENDED_LIMIT_INFORMATION).cast(),
std::mem::size_of_val(&limits) as u32,
)
},
0,
"{}",
std::io::Error::last_os_error()
);
Self(job)
}
fn assign(&self, child: &Child) {
use std::os::windows::io::AsRawHandle;
use windows_sys::Win32::System::JobObjects::AssignProcessToJobObject;
assert_ne!(
unsafe { AssignProcessToJobObject(self.0.as_raw_handle(), child.as_raw_handle()) },
0,
"{}",
std::io::Error::last_os_error()
);
}
}
/// Runs `bsk --json update --yes` once its parent has put it in a Job, and
/// writes the outcome to `BSK_GATED_OUTPUT`.
#[test]
fn manual_update_replaces_the_running_executable_after_cli_exit() {
#[ignore = "subprocess entry point"]
fn gated_update_process() {
let mut go = String::new();
std::io::stdin().read_line(&mut go).unwrap();
let output = Command::new(std::env::var_os("BSK_GATED_EXE").unwrap())
.args(["--json", "update", "--yes"])
.stdin(Stdio::null())
.creation_flags(0x0800_0000)
.output()
.unwrap();
let result = serde_json::json!({
"code": output.status.code(),
"stdout": String::from_utf8_lossy(&output.stdout),
"stderr": String::from_utf8_lossy(&output.stderr),
});
fs::write(
std::env::var_os("BSK_GATED_OUTPUT").unwrap(),
serde_json::to_vec(&result).unwrap(),
)
.unwrap();
}
#[test]
fn update_from_a_restrictive_job_leaves_a_background_daemon_running() {
let fixture = Fixture::new();
let port = unused_port();
let start = fixture
.command()
.args(["daemon", "start", "--port", &port.to_string()])
.output()
.unwrap();
assert!(
start.status.success(),
"{}",
String::from_utf8_lossy(&start.stderr)
);
let before = fixture.info().unwrap();
let job = RestrictiveJob::new();
let result_path = fixture._tmp.path().join("gated-update.json");
let mut gated = Command::new(std::env::current_exe().unwrap())
.args(["--exact", "gated_update_process", "--ignored", "--quiet"])
.env("BSK_GATED_EXE", &fixture.exe)
.env("BSK_GATED_OUTPUT", &result_path)
.env("BSK_HOME", &fixture.home)
.env("BSK_UPDATE_MANIFEST_URL", &fixture.server.url)
.env("BSK_AUTO_UPDATE", "off")
.env("RUST_LOG", "info")
.env("NO_PROXY", "127.0.0.1,localhost")
.env_remove("BSK_DAEMONIZED")
.env_remove("BSK_DAEMON_REPLACES_PID")
.stdin(Stdio::piped())
.stdout(Stdio::null())
.stderr(Stdio::null())
.creation_flags(0x0800_0000)
.spawn()
.unwrap();
job.assign(&gated);
gated.stdin.take().unwrap().write_all(b"go\n").unwrap();
assert!(gated.wait().unwrap().success());
let result: serde_json::Value =
serde_json::from_slice(&fs::read(&result_path).unwrap()).unwrap();
assert_eq!(result["code"], 0, "{result}");
let report: serde_json::Value =
serde_json::from_str(result["stdout"].as_str().unwrap()).unwrap();
assert_eq!(report["daemon"], "left_running", "{report}");
assert!(
report["message"]
.as_str()
.unwrap()
.contains("cannot start an independent daemon"),
"{report}"
);
assert!(fixture.updated(), "the release is installed");
// The daemon it could not have started again keeps serving.
let after = fixture.info().expect("the daemon keeps serving");
assert_eq!(after.pid, before.pid);
assert_eq!(after.ws_port, port);
let status = fixture
.command()
.args(["--json", "status"])
.output()
.unwrap();
assert!(
status.status.success(),
"{}",
String::from_utf8_lossy(&status.stderr)
);
}
#[test]
fn manual_update_replaces_the_running_executable_in_place() {
let fixture = Fixture::new();
let out = fixture
.command()
@@ -261,15 +399,34 @@ fn manual_update_replaces_the_running_executable_after_cli_exit() {
String::from_utf8_lossy(&out.stderr)
);
let report: serde_json::Value = serde_json::from_slice(&out.stdout).unwrap();
assert_eq!(report["status"], "staged");
fixture.wait_for("CLI replacement and helper cleanup", || {
fixture.updated() && fs::read_dir(fixture.exe.parent().unwrap()).unwrap().count() == 1
});
assert_eq!(report["status"], "updated");
assert!(
fixture.updated(),
"the new binary must be in place when the command returns"
);
assert!(
fixture.info().is_none(),
"an update must not start a previously absent daemon"
);
assert_eq!(fixture.server.requests.load(Ordering::SeqCst), 1);
// The CLI ran from the image it moved aside; the next daemon removes it.
let leftovers = fixture.leftovers();
assert!(
leftovers.len() == 1 && leftovers[0].starts_with(".bsk.exe.old-"),
"{leftovers:?}"
);
let start = fixture
.command()
.args(["daemon", "start", "--port", &unused_port().to_string()])
.output()
.unwrap();
assert!(
start.status.success(),
"{}",
String::from_utf8_lossy(&start.stderr)
);
fixture.wait_for("leftover cleanup", || fixture.leftovers().is_empty());
}
#[test]
@@ -304,19 +461,21 @@ fn automatic_update_exits_old_daemon_and_restarts_on_the_same_port() {
.info()
.is_some_and(|info| info.pid != old_pid && info.ws_port == port)
});
assert!(
fixture
.daemon
.as_mut()
.unwrap()
.try_wait()
.unwrap()
.is_some(),
"old daemon must exit"
);
fixture.wait_for("helper cleanup", || {
fs::read_dir(fixture.exe.parent().unwrap()).unwrap().count() == 1
});
// The old daemon exits only after it has seen the replacement serve.
let deadline = Instant::now() + Duration::from_secs(30);
while fixture
.daemon
.as_mut()
.unwrap()
.try_wait()
.unwrap()
.is_none()
{
assert!(Instant::now() < deadline, "old daemon must exit");
thread::sleep(Duration::from_millis(50));
}
// The replacement removes the image its predecessor ran from.
fixture.wait_for("leftover cleanup", || fixture.leftovers().is_empty());
let status = fixture
.command()
.args(["--json", "status"])
@@ -330,7 +489,7 @@ fn automatic_update_exits_old_daemon_and_restarts_on_the_same_port() {
assert_eq!(
fixture.server.requests.load(Ordering::SeqCst),
1,
"the replacement must not stage another update immediately"
"the replacement must not download another update immediately"
);
}
+4
View File
@@ -143,6 +143,10 @@ daemon; bsk cannot make a process outlive the environment that owns it.
Start and stop the shared daemon in this owning environment. Browser task
cleanup is `bsk session stop`, which leaves other sessions and the daemon alone.
A foreground daemon never replaces itself: when a new release is available it
logs the version and the CLI suggests `bsk update`. `bsk update` installs the
release but leaves a foreground daemon running on the previous version, so
restart the host task afterwards to use it.
The daemon's existing idle-exit behavior is unchanged: with no connected
browsers, active sessions or IPC clients, its default idle timeout is 10 minutes.
If the host workflow needs a longer idle window, pass the existing `--daemon-idle`
+46 -8
View File
@@ -37,6 +37,8 @@ public static class DaemonTestHost {
Set-Location -LiteralPath $manifest.workspace
$env:CI = 'true' # Default-port cold startup must execute, never skip.
if ($manifest.runnerTrackingId) { $env:RUNNER_TRACKING_ID = $manifest.runnerTrackingId }
# Update tests compile stand-ins for broken releases.
$env:RUSTC = $manifest.rustc
for ($round = 1; $round -le $manifest.repeat; $round++) {
foreach ($suite in $manifest.suites) {
$name = "$($suite.name)-$round"
@@ -47,18 +49,29 @@ public static class DaemonTestHost {
$process.StartInfo.CreateNoWindow = $true
$process.StartInfo.RedirectStandardOutput = $true
$process.StartInfo.RedirectStandardError = $true
# The suite's copies of bsk and the daemons they start live in
# its own temporary directory, so a timeout stops exactly those.
$suiteTemp = Join-Path $runDir "temp/$name"
$null = New-Item -ItemType Directory -Path $suiteTemp -Force
$suiteTemp = (Resolve-Path -LiteralPath $suiteTemp).Path.TrimEnd('\') + '\'
$process.StartInfo.EnvironmentVariables['TEMP'] = $suiteTemp
$process.StartInfo.EnvironmentVariables['TMP'] = $suiteTemp
$exitCode = -1
$errorText = $null
$stdout = $null
$stderr = $null
$timedOut = $false
$watch = [Diagnostics.Stopwatch]::StartNew()
try {
if (-not $process.Start()) { throw 'Test process did not start' }
$stdout = $process.StandardOutput.ReadToEndAsync()
$stderr = $process.StandardError.ReadToEndAsync()
if (-not $process.WaitForExit(180000)) { throw 'Test process exceeded 180 seconds' }
if (-not $process.WaitForExit($manifest.suiteSeconds * 1000)) {
$timedOut = $true
throw "Test process exceeded $($manifest.suiteSeconds) seconds"
}
$exitCode = $process.ExitCode
if (-not $stdout.Wait(5000) -or -not $stderr.Wait(5000)) { throw 'Test process exited without pipe EOF' }
$stdout.Result | Set-Content -LiteralPath "$runDir/$name.stdout.log" -Encoding utf8
$stderr.Result | Set-Content -LiteralPath "$runDir/$name.stderr.log" -Encoding utf8
} catch {
$exitCode = -1
$errorText = $_.ToString()
@@ -66,6 +79,22 @@ public static class DaemonTestHost {
try {
if (-not $process.HasExited) { $process.Kill(); $null = $process.WaitForExit(5000) }
} catch { }
if ($timedOut) {
# Daemons the tests started run outside the test's process
# tree; stop those running from this suite's copies.
Get-CimInstance Win32_Process | Where-Object {
$_.ExecutablePath -and $_.ExecutablePath.StartsWith($suiteTemp, [StringComparison]::OrdinalIgnoreCase)
} | ForEach-Object { Stop-Process -Id $_.ProcessId -Force -ErrorAction SilentlyContinue }
}
# Keep the logs, not the copies of bsk, for the artifact upload.
Get-ChildItem -LiteralPath $suiteTemp -Recurse -Filter '*.exe' -ErrorAction SilentlyContinue |
Remove-Item -Force -ErrorAction SilentlyContinue
# Keep what the suite printed, also after a timeout.
foreach ($stream in @(@{ task=$stdout; file="$runDir/$name.stdout.log" }, @{ task=$stderr; file="$runDir/$name.stderr.log" })) {
try {
if ($stream.task -and $stream.task.Wait(5000)) { $stream.task.Result | Set-Content -LiteralPath $stream.file -Encoding utf8 }
} catch { }
}
$process.Dispose()
}
$results += [pscustomobject]@{ name=$name; exitCode=$exitCode; seconds=$watch.Elapsed.TotalSeconds; error=$errorText }
@@ -88,7 +117,7 @@ Push-Location -LiteralPath $workspace
try {
# Build in the normal runner environment; the WMI host needs no inherited
# Cargo PATH, credentials or toolchain environment. Use Cargo's exact paths.
$artifacts = @(& cargo test -p bsk --lib --test windows_daemon_start --test windows_update --no-run --locked --message-format=json-render-diagnostics)
$artifacts = @(& cargo test -p bsk --lib --test windows_daemon_start --test windows_update --test auto_update_handover --no-run --locked --message-format=json-render-diagnostics)
if ($LASTEXITCODE -ne 0) { throw 'Could not build Windows lifecycle tests' }
$tests = @{}
foreach ($line in $artifacts) {
@@ -101,13 +130,21 @@ try {
$suites = @(
@{ name='startup'; executable=$tests['windows_daemon_start']; arguments=@('--test-threads=1','--nocapture') },
@{ name='launcher-lifetime'; executable=$tests['bsk']; arguments=@('daemon::start::tests','--test-threads=1','--nocapture') },
@{ name='update'; executable=$tests['windows_update']; arguments=@('--test-threads=1','--nocapture') }
@{ name='update'; executable=$tests['windows_update']; arguments=@('--test-threads=1','--nocapture') },
@{ name='update-handover'; executable=$tests['auto_update_handover']; arguments=@('--test-threads=1','--nocapture') }
)
foreach ($suite in $suites) {
if (-not $suite.executable -or -not (Test-Path -LiteralPath $suite.executable)) { throw "Missing executable for $($suite.name)" }
}
$sysroot = (& rustc --print sysroot).Trim()
if ($LASTEXITCODE -ne 0) { throw 'Could not locate rustc' }
$rustc = Join-Path $sysroot 'bin/rustc.exe'
if (-not (Test-Path -LiteralPath $rustc)) { throw "Missing $rustc" }
# Update suites wait on 20-second handover deadlines and take several
# minutes on slower machines.
$suiteSeconds = 480
$manifest = @{
workspace=$workspace; repeat=$Repeat; suites=$suites
workspace=$workspace; repeat=$Repeat; suites=$suites; rustc=$rustc; suiteSeconds=$suiteSeconds
userSid=[Security.Principal.WindowsIdentity]::GetCurrent().User.Value
runnerTrackingId=$env:RUNNER_TRACKING_ID
}
@@ -122,7 +159,8 @@ try {
try {
$workerProcess = [Diagnostics.Process]::GetProcessById($created.ProcessId)
$null = $workerProcess.Handle # Retain identity for timeout cleanup.
$deadline = [DateTime]::UtcNow.AddSeconds(600 * $Repeat)
# Every suite may use its whole budget, plus time to start the host.
$deadline = [DateTime]::UtcNow.AddSeconds(($suites.Count * $suiteSeconds + 120) * $Repeat)
while (-not $workerProcess.WaitForExit(1000)) {
if ([DateTime]::UtcNow -ge $deadline) { throw 'Independent test host timed out' }
}
@@ -132,7 +170,7 @@ try {
$report = Get-Content -LiteralPath "$runDir/result.json" -Raw | ConvertFrom-Json
$report.results | Format-Table -AutoSize | Out-Host
if ($report.error) { throw $report.error }
if (@($report.results).Count -ne (3 * $Repeat)) { throw 'Not all lifecycle suites executed' }
if (@($report.results).Count -ne ($suites.Count * $Repeat)) { throw 'Not all lifecycle suites executed' }
if (@($report.results | Where-Object exitCode -ne 0).Count) { throw "Lifecycle regression failed; logs: $runDir" }
Write-Output "All lifecycle suites passed; logs: $runDir"
} finally {