Merge current main and make purchase filing atomic under concurrent writes
This commit is contained in:
@@ -313,7 +313,7 @@ async fn image_id(image_ref: &str) -> Option<String> {
|
||||
/// should reference (`localhost/<base>:latest` for build, registry
|
||||
/// URL for pull).
|
||||
async fn ensure_image_present(spec: &CompanionSpec) -> Result<String> {
|
||||
let local_image = format!("localhost/{}:latest", spec.image_base);
|
||||
let mut local_image = format!("localhost/{}:latest", spec.image_base);
|
||||
let local_image_compat = format!("localhost/{}:local", spec.image_base);
|
||||
let registry_image = format!("{}/{}:latest", COMPANION_REGISTRY, spec.image_base);
|
||||
|
||||
@@ -322,11 +322,13 @@ async fn ensure_image_present(spec: &CompanionSpec) -> Result<String> {
|
||||
for dir in spec.build_dir_candidates {
|
||||
let dockerfile = PathBuf::from(dir).join("Dockerfile");
|
||||
if fs::try_exists(&dockerfile).await.unwrap_or(false) {
|
||||
// `:local` is a deliberate manual override — never auto-rebuild it.
|
||||
// Older installers and self-update create :local themselves. It
|
||||
// must receive source updates too; treating it as a permanent
|
||||
// manual override silently kept the old LND UI after an OTA.
|
||||
if image_exists(&local_image_compat).await {
|
||||
return Ok(local_image_compat);
|
||||
local_image = local_image_compat.clone();
|
||||
}
|
||||
// Reuse the auto-built `:latest` only when the build context has NOT
|
||||
// Reuse either local tag only when the build context has NOT
|
||||
// changed since it was built. Without this staleness check an
|
||||
// already-present image is reused forever, so edits to the baked-in
|
||||
// context (Dockerfile, nginx.conf, …) never reach the node — this is
|
||||
@@ -849,20 +851,43 @@ async fn needs_repair(spec: &CompanionSpec) -> Result<bool> {
|
||||
if !matches_known_shape {
|
||||
return Ok(true);
|
||||
}
|
||||
if on_disk.contains(&local_image) && !on_disk.contains(&local_image_compat) {
|
||||
if let Some(image) = managed_local_image(spec, &on_disk) {
|
||||
for dir in spec.build_dir_candidates {
|
||||
let dockerfile = PathBuf::from(dir).join("Dockerfile");
|
||||
if fs::try_exists(&dockerfile).await.unwrap_or(false) {
|
||||
// Conservative on any timeout/error inside: reuse the cache.
|
||||
return Ok(context_is_newer_than_image(dir, &local_image).await);
|
||||
return Ok(context_is_newer_than_image(dir, &image).await);
|
||||
}
|
||||
}
|
||||
}
|
||||
Ok(false)
|
||||
}
|
||||
|
||||
fn managed_local_image(spec: &CompanionSpec, unit: &str) -> Option<String> {
|
||||
["latest", "local"]
|
||||
.iter()
|
||||
.map(|tag| format!("localhost/{}:{tag}", spec.image_base))
|
||||
.find(|image| build_unit(spec, image).render() == unit)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
#[test]
|
||||
fn legacy_installer_local_tag_is_checked_for_source_updates_like_latest() {
|
||||
for spec in ALL_COMPANIONS.iter().flat_map(|group| group.iter()) {
|
||||
for tag in ["local", "latest"] {
|
||||
let image = format!("localhost/{}:{tag}", spec.image_base);
|
||||
let unit = build_unit(spec, &image).render();
|
||||
assert_eq!(managed_local_image(spec, &unit), Some(image));
|
||||
}
|
||||
let registry = format!("{}/{}:latest", COMPANION_REGISTRY, spec.image_base);
|
||||
assert_eq!(
|
||||
managed_local_image(spec, &build_unit(spec, ®istry).render()),
|
||||
None
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
use super::*;
|
||||
|
||||
fn names(specs: &[&'static CompanionSpec]) -> Vec<&'static str> {
|
||||
|
||||
@@ -117,20 +117,24 @@ fn shell_quote(s: &str) -> String {
|
||||
s.replace('\'', "'\\''")
|
||||
}
|
||||
|
||||
/// Save `bytes` into FileBrowser's storage as a new file in `dir`, named
|
||||
/// `name` or, if that's taken, `name (2)`, `name (3)`… Never overwrites.
|
||||
/// Returns the path written.
|
||||
///
|
||||
/// FileBrowser's folders belong to its rootless container range (host uid
|
||||
/// 100000, mode 755), so this service — host uid 1000, outside that range —
|
||||
/// can read them but not write into them, and filing a purchase into Files
|
||||
/// failed with EACCES (2026-09-29). When a direct write is refused, the file
|
||||
/// is written through `podman unshare`, where that range is ours, and given
|
||||
/// the folder's owner so FileBrowser manages it like its own uploads.
|
||||
/// Save a complete purchase without overwriting any existing directory entry.
|
||||
/// Both host and rootless-namespace paths publish with a no-clobber hard link.
|
||||
pub async fn save_new_file(dir: &Path, name: &str, bytes: &[u8]) -> Result<PathBuf> {
|
||||
save_new_file_with(dir, name, bytes, write_via_userns).await
|
||||
}
|
||||
|
||||
fn validate_filename(name: &str) -> Result<()> {
|
||||
anyhow::ensure!(
|
||||
!name.is_empty()
|
||||
&& name != "."
|
||||
&& name != ".."
|
||||
&& !name.contains(['/', '\\', '\0'])
|
||||
&& name.len() <= 255,
|
||||
"Invalid purchased filename"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn save_new_file_with<F, Fut>(
|
||||
dir: &Path,
|
||||
name: &str,
|
||||
@@ -138,100 +142,170 @@ async fn save_new_file_with<F, Fut>(
|
||||
fallback: F,
|
||||
) -> Result<PathBuf>
|
||||
where
|
||||
F: FnOnce(PathBuf, Vec<u8>) -> Fut,
|
||||
Fut: std::future::Future<Output = Result<()>>,
|
||||
F: FnOnce(PathBuf, String, Vec<u8>) -> Fut,
|
||||
Fut: std::future::Future<Output = Result<PathBuf>>,
|
||||
{
|
||||
let target = unused_name(dir, name);
|
||||
match write_direct(dir, &target, bytes).await {
|
||||
Ok(()) => Ok(target),
|
||||
Err(e) if e.kind() == std::io::ErrorKind::PermissionDenied => {
|
||||
fallback(target.clone(), bytes.to_vec())
|
||||
validate_filename(name)?;
|
||||
// Never follow a user-created destination directory symlink.
|
||||
match fs::symlink_metadata(dir).await {
|
||||
Ok(meta) => anyhow::ensure!(meta.is_dir(), "Files destination is not a directory"),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
|
||||
Err(error) => return Err(error.into()),
|
||||
}
|
||||
save_after_direct_result(
|
||||
write_direct(dir, name, bytes).await,
|
||||
dir,
|
||||
name,
|
||||
bytes,
|
||||
fallback,
|
||||
)
|
||||
.await
|
||||
}
|
||||
|
||||
async fn save_after_direct_result<F, Fut>(
|
||||
result: std::io::Result<PathBuf>,
|
||||
dir: &Path,
|
||||
name: &str,
|
||||
bytes: &[u8],
|
||||
fallback: F,
|
||||
) -> Result<PathBuf>
|
||||
where
|
||||
F: FnOnce(PathBuf, String, Vec<u8>) -> Fut,
|
||||
Fut: std::future::Future<Output = Result<PathBuf>>,
|
||||
{
|
||||
match result {
|
||||
Ok(path) => Ok(path),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
|
||||
fallback(dir.to_owned(), name.to_owned(), bytes.to_vec())
|
||||
.await
|
||||
.with_context(|| format!("writing {} via podman unshare", target.display()))?;
|
||||
Ok(target)
|
||||
.context("Saving purchase in Files user namespace")
|
||||
}
|
||||
Err(e) => Err(e).with_context(|| format!("writing {}", target.display())),
|
||||
Err(error) => Err(error).context("Saving purchase in Files"),
|
||||
}
|
||||
}
|
||||
|
||||
/// `dir/name`, or the first free `dir/stem (n).ext` from n = 2.
|
||||
fn unused_name(dir: &Path, name: &str) -> PathBuf {
|
||||
let mut target = dir.join(name);
|
||||
let (stem, ext) = match name.rsplit_once('.') {
|
||||
Some((s, e)) if !s.is_empty() => (s.to_string(), format!(".{e}")),
|
||||
_ => (name.to_string(), String::new()),
|
||||
};
|
||||
let mut n = 2;
|
||||
while target.exists() {
|
||||
target = dir.join(format!("{stem} ({n}){ext}"));
|
||||
n += 1;
|
||||
fn numbered_name(name: &str, attempt: usize) -> String {
|
||||
if attempt == 1 {
|
||||
return name.to_owned();
|
||||
}
|
||||
match name.rsplit_once('.') {
|
||||
Some((stem, extension)) if !stem.is_empty() => format!("{stem} ({attempt}).{extension}"),
|
||||
_ => format!("{name} ({attempt})"),
|
||||
}
|
||||
target
|
||||
}
|
||||
|
||||
async fn write_direct(dir: &Path, target: &Path, bytes: &[u8]) -> std::io::Result<()> {
|
||||
struct PendingFile(PathBuf);
|
||||
impl Drop for PendingFile {
|
||||
fn drop(&mut self) {
|
||||
let _ = std::fs::remove_file(&self.0);
|
||||
}
|
||||
}
|
||||
|
||||
async fn write_direct(dir: &Path, name: &str, bytes: &[u8]) -> std::io::Result<PathBuf> {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
fs::create_dir_all(dir).await?;
|
||||
let mut f = fs::OpenOptions::new()
|
||||
let temp_path = dir.join(format!(".archy-saving-{}", uuid::Uuid::new_v4()));
|
||||
let mut file = fs::OpenOptions::new()
|
||||
.write(true)
|
||||
.create_new(true)
|
||||
.open(target)
|
||||
.mode(0o600)
|
||||
.open(&temp_path)
|
||||
.await?;
|
||||
let written = async {
|
||||
f.write_all(bytes).await?;
|
||||
f.flush().await
|
||||
let temp = PendingFile(temp_path);
|
||||
file.write_all(bytes).await?;
|
||||
file.set_permissions(std::fs::Permissions::from_mode(0o644))
|
||||
.await?;
|
||||
file.sync_all().await?;
|
||||
for attempt in 1..=100 {
|
||||
let target = dir.join(numbered_name(name, attempt));
|
||||
match fs::hard_link(&temp.0, &target).await {
|
||||
Ok(()) => return Ok(target),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => continue,
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
.await;
|
||||
if written.is_err() {
|
||||
let _ = fs::remove_file(target).await;
|
||||
}
|
||||
written
|
||||
Err(std::io::Error::new(
|
||||
std::io::ErrorKind::AlreadyExists,
|
||||
"Too many existing copies; purchase cache retained",
|
||||
))
|
||||
}
|
||||
|
||||
/// Write `bytes` (piped on stdin) to `target` from inside the rootless user
|
||||
/// namespace. It goes to a temp file first and is hard-linked into place, so
|
||||
/// FileBrowser never sees a partial file and an existing file is never
|
||||
/// replaced (`ln` refuses an existing name).
|
||||
async fn write_via_userns(target: PathBuf, bytes: Vec<u8>) -> Result<()> {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
const SCRIPT: &str = r#"set -eu
|
||||
dst=$1
|
||||
dir=$(dirname -- "$dst")
|
||||
// Positional arguments carry all user-controlled text. mktemp prevents temp-name
|
||||
// collisions; ln -T refuses files, symlinks and directories, including races.
|
||||
const WRITE_VIA_USERNS: &str = r#"set -eu
|
||||
dir=$1
|
||||
name=$2
|
||||
expected=$3
|
||||
[ ! -L "$dir" ] || exit 1
|
||||
if [ ! -d "$dir" ]; then
|
||||
mkdir -- "$dir"
|
||||
mkdir -p -- "$dir"
|
||||
chown --reference="$(dirname -- "$dir")" -- "$dir"
|
||||
fi
|
||||
tmp="$dir/.archy-saving.$$"
|
||||
trap 'rm -f -- "$tmp"' EXIT
|
||||
tmp=$(mktemp "$dir/.archy-saving.XXXXXXXXXX")
|
||||
trap 'rm -f -- "$tmp"' EXIT HUP INT TERM
|
||||
cat > "$tmp"
|
||||
[ "$(wc -c < "$tmp")" -eq "$expected" ] || exit 1
|
||||
chown --reference="$dir" -- "$tmp"
|
||||
chmod 0644 -- "$tmp"
|
||||
ln -- "$tmp" "$dst"
|
||||
sync -f -- "$tmp"
|
||||
stem=$name
|
||||
ext=
|
||||
case "$name" in
|
||||
*.*) prefix=${name%.*}; if [ -n "$prefix" ]; then stem=$prefix; ext=.${name##*.}; fi ;;
|
||||
esac
|
||||
n=1
|
||||
while [ "$n" -le 100 ]; do
|
||||
candidate=$name
|
||||
if [ "$n" -gt 1 ]; then candidate="$stem ($n)$ext"; fi
|
||||
dst="$dir/$candidate"
|
||||
if ln -T -- "$tmp" "$dst" 2>/dev/null; then
|
||||
printf '%s' "$candidate"
|
||||
exit 0
|
||||
fi
|
||||
# A conflict may be a dangling symlink; never follow it or overwrite it.
|
||||
if [ ! -e "$dst" ] && [ ! -L "$dst" ]; then exit 1; fi
|
||||
n=$((n + 1))
|
||||
done
|
||||
exit 1
|
||||
"#;
|
||||
|
||||
async fn write_via_userns(dir: PathBuf, name: String, bytes: Vec<u8>) -> Result<PathBuf> {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
let mut child = tokio::process::Command::new("podman")
|
||||
.args(["unshare", "sh", "-c", SCRIPT, "sh"])
|
||||
.arg(&target)
|
||||
.args(["unshare", "sh", "-c", WRITE_VIA_USERNS, "sh"])
|
||||
.arg(&dir)
|
||||
.arg(&name)
|
||||
.arg(bytes.len().to_string())
|
||||
.kill_on_drop(true)
|
||||
.stdin(std::process::Stdio::piped())
|
||||
.stdout(std::process::Stdio::null())
|
||||
.stdout(std::process::Stdio::piped())
|
||||
.stderr(std::process::Stdio::piped())
|
||||
.spawn()
|
||||
.context("Failed to run podman unshare")?;
|
||||
let mut stdin = child.stdin.take().context("podman unshare stdin")?;
|
||||
let fed = stdin.write_all(&bytes).await;
|
||||
drop(stdin);
|
||||
let out = child
|
||||
.wait_with_output()
|
||||
.await
|
||||
.context("Failed to wait for podman unshare")?;
|
||||
if !out.status.success() {
|
||||
anyhow::bail!(
|
||||
"podman unshare exited with {}: {}",
|
||||
out.status,
|
||||
String::from_utf8_lossy(&out.stderr).trim()
|
||||
.context("Starting Files namespace writer")?;
|
||||
let mut stdin = child.stdin.take().context("Files writer stdin missing")?;
|
||||
let operation = async {
|
||||
let fed = stdin.write_all(&bytes).await;
|
||||
drop(stdin);
|
||||
let output = child.wait_with_output().await?;
|
||||
anyhow::ensure!(
|
||||
output.status.success(),
|
||||
"Files namespace writer failed: {}",
|
||||
output.status
|
||||
);
|
||||
}
|
||||
fed.context("Failed to pipe the file to podman unshare")?;
|
||||
Ok(())
|
||||
fed.context("Sending purchase bytes to Files")?;
|
||||
let chosen =
|
||||
String::from_utf8(output.stdout).context("Files writer returned an invalid name")?;
|
||||
validate_filename(&chosen)?;
|
||||
anyhow::ensure!(
|
||||
(1..=100).any(|n| numbered_name(&name, n) == chosen),
|
||||
"Files writer returned an unexpected name"
|
||||
);
|
||||
Ok(dir.join(chosen))
|
||||
};
|
||||
tokio::time::timeout(std::time::Duration::from_secs(120), operation)
|
||||
.await
|
||||
.context("Files namespace writer timed out")?
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -268,95 +342,232 @@ mod tests {
|
||||
let second = ensure_config(&paths).await.unwrap();
|
||||
assert_eq!(second, EnsureOutcome::Unchanged);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unused_name_numbers_duplicates_and_keeps_the_extension() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let d = dir.path();
|
||||
assert_eq!(unused_name(d, "song.mp3"), d.join("song.mp3"));
|
||||
std::fs::write(d.join("song.mp3"), b"").unwrap();
|
||||
assert_eq!(unused_name(d, "song.mp3"), d.join("song (2).mp3"));
|
||||
std::fs::write(d.join("song (2).mp3"), b"").unwrap();
|
||||
assert_eq!(unused_name(d, "song.mp3"), d.join("song (3).mp3"));
|
||||
std::fs::write(d.join("README"), b"").unwrap();
|
||||
assert_eq!(unused_name(d, "README"), d.join("README (2)"));
|
||||
std::fs::write(d.join(".hidden"), b"").unwrap();
|
||||
assert_eq!(unused_name(d, ".hidden"), d.join(".hidden (2)"));
|
||||
#[cfg(test)]
|
||||
mod purchase_write_tests {
|
||||
use super::*;
|
||||
use std::{
|
||||
collections::HashSet,
|
||||
os::unix::fs::{symlink, PermissionsExt},
|
||||
};
|
||||
|
||||
fn no_temps(dir: &Path) {
|
||||
assert!(std::fs::read_dir(dir).unwrap().all(|e| !e
|
||||
.unwrap()
|
||||
.file_name()
|
||||
.to_string_lossy()
|
||||
.starts_with(".archy-saving")));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn save_new_file_writes_directly_into_a_writable_folder() {
|
||||
async fn direct_write_uses_complete_bytes_and_preserves_originals() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let music = dir.path().join("Music");
|
||||
let path = save_new_file_with(&music, "a.mp3", b"abc", |_, _| async {
|
||||
anyhow::bail!("fallback must not run")
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(path, music.join("a.mp3"));
|
||||
assert_eq!(std::fs::read(&path).unwrap(), b"abc");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn save_new_file_never_overwrites_an_existing_file() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
std::fs::write(dir.path().join("a.mp3"), b"original").unwrap();
|
||||
let path = save_new_file_with(dir.path(), "a.mp3", b"new", |_, _| async {
|
||||
anyhow::bail!("fallback must not run")
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(path, dir.path().join("a (2).mp3"));
|
||||
fs::write(dir.path().join("song.mp3"), b"original")
|
||||
.await
|
||||
.unwrap();
|
||||
let target = save_new_file(dir.path(), "song.mp3", b"new").await.unwrap();
|
||||
assert_eq!(target.file_name().unwrap(), "song (2).mp3");
|
||||
assert_eq!(fs::read(target).await.unwrap(), b"new");
|
||||
assert_eq!(
|
||||
std::fs::read(dir.path().join("a.mp3")).unwrap(),
|
||||
fs::read(dir.path().join("song.mp3")).await.unwrap(),
|
||||
b"original"
|
||||
);
|
||||
no_temps(dir.path());
|
||||
}
|
||||
|
||||
/// Regression (2026-09-29): filing a purchase into a FileBrowser folder
|
||||
/// owned by the container's uid range failed with EACCES. A refused
|
||||
/// write must go through the user-namespace fallback, with the same
|
||||
/// target and bytes.
|
||||
#[tokio::test]
|
||||
async fn a_refused_write_goes_through_the_userns_fallback() {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
async fn simultaneous_saves_publish_unique_complete_files() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let music = dir.path().join("Music");
|
||||
std::fs::create_dir(&music).unwrap();
|
||||
std::fs::set_permissions(&music, std::fs::Permissions::from_mode(0o555)).unwrap();
|
||||
if std::fs::File::create(music.join("probe")).is_ok() {
|
||||
return; // running as root: mode bits don't refuse the write
|
||||
let mut tasks = Vec::new();
|
||||
for n in 0..24u8 {
|
||||
let dir = dir.path().to_owned();
|
||||
tasks.push(tokio::spawn(async move {
|
||||
let bytes = vec![n; 32768];
|
||||
let path = save_new_file(&dir, "same.bin", &bytes).await.unwrap();
|
||||
assert_eq!(fs::read(&path).await.unwrap(), bytes);
|
||||
path
|
||||
}));
|
||||
}
|
||||
let mut paths = HashSet::new();
|
||||
for task in tasks {
|
||||
assert!(paths.insert(task.await.unwrap()));
|
||||
}
|
||||
assert_eq!(paths.len(), 24);
|
||||
no_temps(dir.path());
|
||||
}
|
||||
|
||||
let seen = std::sync::Mutex::new(None);
|
||||
let path = save_new_file_with(&music, "a.mp3", b"abc", |target, bytes| {
|
||||
*seen.lock().unwrap() = Some((target, bytes));
|
||||
async { Ok(()) }
|
||||
})
|
||||
#[tokio::test]
|
||||
async fn existing_directories_and_dangling_symlinks_are_conflicts() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
fs::create_dir(dir.path().join("name")).await.unwrap();
|
||||
symlink("missing", dir.path().join("name (2)")).unwrap();
|
||||
let path = save_new_file(dir.path(), "name", b"new").await.unwrap();
|
||||
assert_eq!(path.file_name().unwrap(), "name (3)");
|
||||
assert!(dir.path().join("name").is_dir());
|
||||
assert!(fs::symlink_metadata(dir.path().join("name (2)"))
|
||||
.await
|
||||
.unwrap()
|
||||
.is_symlink());
|
||||
no_temps(dir.path());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn invalid_names_and_symlink_destination_are_refused() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
for name in [
|
||||
"",
|
||||
".",
|
||||
"..",
|
||||
"../escape",
|
||||
"/absolute",
|
||||
"a/b",
|
||||
"a\\b",
|
||||
"a\0b",
|
||||
] {
|
||||
assert!(save_new_file(dir.path(), name, b"bytes").await.is_err());
|
||||
}
|
||||
let outside = tempfile::tempdir().unwrap();
|
||||
symlink(outside.path(), dir.path().join("Music")).unwrap();
|
||||
assert!(save_new_file(&dir.path().join("Music"), "song", b"bytes")
|
||||
.await
|
||||
.is_err());
|
||||
assert_eq!(std::fs::read_dir(outside.path()).unwrap().count(), 0);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn collision_limit_preserves_all_files_and_cleans_temporary_data() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
for n in 1..=100 {
|
||||
fs::write(dir.path().join(numbered_name("a.txt", n)), b"keep")
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
assert!(save_new_file(dir.path(), "a.txt", b"new").await.is_err());
|
||||
for n in 1..=100 {
|
||||
assert_eq!(
|
||||
fs::read(dir.path().join(numbered_name("a.txt", n)))
|
||||
.await
|
||||
.unwrap(),
|
||||
b"keep"
|
||||
);
|
||||
}
|
||||
no_temps(dir.path());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn permission_fallback_is_exercised_without_skipping_as_root() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let result = save_after_direct_result(
|
||||
Err(std::io::ErrorKind::PermissionDenied.into()),
|
||||
dir.path(),
|
||||
"a",
|
||||
b"abc",
|
||||
|dir, name, bytes| async move {
|
||||
assert_eq!(bytes, b"abc");
|
||||
Ok(dir.join(name))
|
||||
},
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(path, music.join("a.mp3"));
|
||||
assert_eq!(
|
||||
seen.into_inner().unwrap(),
|
||||
Some((music.join("a.mp3"), b"abc".to_vec()))
|
||||
);
|
||||
std::fs::set_permissions(&music, std::fs::Permissions::from_mode(0o755)).unwrap();
|
||||
assert_eq!(result, dir.path().join("a"));
|
||||
assert!(save_after_direct_result(
|
||||
Err(std::io::ErrorKind::PermissionDenied.into()),
|
||||
dir.path(),
|
||||
"a",
|
||||
b"abc",
|
||||
|_, _, _| async { anyhow::bail!("namespace unavailable") }
|
||||
)
|
||||
.await
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("namespace"));
|
||||
assert!(save_after_direct_result(
|
||||
Err(std::io::ErrorKind::StorageFull.into()),
|
||||
dir.path(),
|
||||
"a",
|
||||
b"abc",
|
||||
|_, _, _| async { panic!("disk full must not trigger permission fallback") }
|
||||
)
|
||||
.await
|
||||
.is_err());
|
||||
}
|
||||
|
||||
async fn run_script(
|
||||
dir: &Path,
|
||||
name: &str,
|
||||
bytes: &[u8],
|
||||
expected: usize,
|
||||
) -> std::process::Output {
|
||||
use tokio::io::AsyncWriteExt;
|
||||
let mut child = tokio::process::Command::new("sh")
|
||||
.args(["-c", WRITE_VIA_USERNS, "sh"])
|
||||
.arg(dir)
|
||||
.arg(name)
|
||||
.arg(expected.to_string())
|
||||
.stdin(std::process::Stdio::piped())
|
||||
.stdout(std::process::Stdio::piped())
|
||||
.stderr(std::process::Stdio::piped())
|
||||
.spawn()
|
||||
.unwrap();
|
||||
let mut input = child.stdin.take().unwrap();
|
||||
input.write_all(bytes).await.unwrap();
|
||||
drop(input);
|
||||
child.wait_with_output().await.unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_failed_fallback_is_reported() {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
async fn namespace_script_preserves_names_bytes_modes_and_existing_entries() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
std::fs::set_permissions(dir.path(), std::fs::Permissions::from_mode(0o555)).unwrap();
|
||||
if std::fs::File::create(dir.path().join("probe")).is_ok() {
|
||||
return;
|
||||
let folder = dir.path().join("Music");
|
||||
let name = "song ' $() ; #.mp3";
|
||||
for n in 1..=2 {
|
||||
let output = run_script(&folder, name, b"abc", 3).await;
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"{}",
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
let chosen = String::from_utf8(output.stdout).unwrap();
|
||||
assert_eq!(chosen, numbered_name(name, n));
|
||||
let path = folder.join(chosen);
|
||||
assert_eq!(fs::read(&path).await.unwrap(), b"abc");
|
||||
assert_eq!(
|
||||
fs::metadata(path).await.unwrap().permissions().mode() & 0o777,
|
||||
0o644
|
||||
);
|
||||
}
|
||||
let err = save_new_file_with(dir.path(), "a.mp3", b"abc", |_, _| async {
|
||||
anyhow::bail!("no podman")
|
||||
})
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(format!("{err:#}").contains("no podman"));
|
||||
std::fs::set_permissions(dir.path(), std::fs::Permissions::from_mode(0o755)).unwrap();
|
||||
no_temps(&folder);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn namespace_script_refuses_truncated_input_and_cleans_up() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let output = run_script(dir.path(), "never.bin", b"partial", 100).await;
|
||||
assert!(!output.status.success());
|
||||
assert!(!dir.path().join("never.bin").exists());
|
||||
no_temps(dir.path());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn namespace_script_does_not_link_inside_existing_directory() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
fs::create_dir(dir.path().join("name")).await.unwrap();
|
||||
symlink("missing", dir.path().join("name (2)")).unwrap();
|
||||
let output = run_script(dir.path(), "name", b"abc", 3).await;
|
||||
assert!(output.status.success());
|
||||
assert_eq!(output.stdout, b"name (3)");
|
||||
assert_eq!(
|
||||
std::fs::read_dir(dir.path().join("name")).unwrap().count(),
|
||||
0
|
||||
);
|
||||
no_temps(dir.path());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn names_keep_extensions_and_dotfiles() {
|
||||
assert_eq!(numbered_name("a.tar.gz", 2), "a.tar (2).gz");
|
||||
assert_eq!(numbered_name(".hidden", 2), ".hidden (2)");
|
||||
assert_eq!(numbered_name("README", 2), "README (2)");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -89,18 +89,74 @@ bitcoind.estimatemode=ECONOMICAL\n"
|
||||
Ok(EnsureOutcome::Written)
|
||||
}
|
||||
|
||||
/// Bitcoin can accept TCP while returning RPC_IN_WARMUP for many minutes.
|
||||
/// Unlocking LND then triggers its short chain-backend timeout and a restart loop.
|
||||
/// Leave the wallet intact and locked; the next reconciliation retries readiness.
|
||||
async fn bitcoin_rpc_ready() -> bool {
|
||||
let (user, password) = crate::bitcoin_rpc::bitcoin_rpc_credentials().await;
|
||||
let client = match reqwest::Client::builder()
|
||||
.no_proxy()
|
||||
.timeout(std::time::Duration::from_secs(5))
|
||||
.build()
|
||||
{
|
||||
Ok(client) => client,
|
||||
Err(_) => return false,
|
||||
};
|
||||
let response = client.post(crate::constants::BITCOIN_RPC_URL)
|
||||
.basic_auth(user, Some(password))
|
||||
.json(&serde_json::json!({"jsonrpc":"1.0","id":"lnd-readiness","method":"getblockchaininfo","params":[]}))
|
||||
.send().await;
|
||||
match response {
|
||||
Ok(response) if response.status().is_success() => response
|
||||
.json::<serde_json::Value>()
|
||||
.await
|
||||
.is_ok_and(|value| bitcoin_readiness_response(&value)),
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
fn bitcoin_readiness_response(value: &serde_json::Value) -> bool {
|
||||
value.get("error").is_none_or(|e| e.is_null())
|
||||
&& value
|
||||
.pointer("/result/blocks")
|
||||
.and_then(|v| v.as_u64())
|
||||
.is_some()
|
||||
&& value
|
||||
.pointer("/result/initialblockdownload")
|
||||
.and_then(|v| v.as_bool())
|
||||
.is_some()
|
||||
}
|
||||
|
||||
pub async fn ensure_wallet_initialized() -> Result<()> {
|
||||
let admin_macaroon = "/var/lib/archipelago/lnd/data/chain/bitcoin/mainnet/admin.macaroon";
|
||||
let wallet_db = "/var/lib/archipelago/lnd/data/chain/bitcoin/mainnet/wallet.db";
|
||||
if file_exists_as_root(wallet_db).await {
|
||||
// GetInfo can wait for Bitcoin sync even though the wallet is already
|
||||
// unlocked. State RPC stays available during that normal startup phase.
|
||||
let client = reqwest::Client::builder()
|
||||
.no_proxy()
|
||||
.timeout(std::time::Duration::from_secs(5))
|
||||
.danger_accept_invalid_certs(true)
|
||||
.build()?;
|
||||
if wallet_is_unlocked(wallet_state(&client).await.as_deref()) {
|
||||
return Ok(());
|
||||
}
|
||||
if file_exists_as_root(admin_macaroon).await && lnd_getinfo_ready(admin_macaroon).await {
|
||||
return Ok(());
|
||||
}
|
||||
if !bitcoin_rpc_ready().await {
|
||||
tracing::debug!("[lnd] waiting for Bitcoin RPC readiness before wallet unlock");
|
||||
return Ok(());
|
||||
}
|
||||
unlock_existing_wallet_no_wipe().await?;
|
||||
wait_for_admin_macaroon(admin_macaroon).await?;
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
if !bitcoin_rpc_ready().await {
|
||||
tracing::debug!("[lnd] waiting for Bitcoin RPC readiness before wallet initialization");
|
||||
return Ok(());
|
||||
}
|
||||
init_wallet_via_rest().await?;
|
||||
wait_for_admin_macaroon(admin_macaroon).await
|
||||
}
|
||||
@@ -258,6 +314,9 @@ async fn unlock_existing_wallet_via_rest() -> Result<bool> {
|
||||
// exactly the nodes least able to afford it. Waiting longer costs nothing —
|
||||
// a wrong password still exits on the first pass via `all_rejected`.
|
||||
for _ in 0..UNLOCK_NOT_READY_ATTEMPTS {
|
||||
if wallet_is_unlocked(wallet_state(&client).await.as_deref()) {
|
||||
return Ok(true);
|
||||
}
|
||||
let mut all_rejected = true;
|
||||
for pw in &candidates {
|
||||
match try_unlock_once(&client, pw).await {
|
||||
@@ -294,6 +353,10 @@ pub(crate) async fn unlock_existing_wallet_no_wipe() -> Result<()> {
|
||||
}
|
||||
}
|
||||
|
||||
fn wallet_is_unlocked(state: Option<&str>) -> bool {
|
||||
matches!(state, Some("UNLOCKED" | "RPC_ACTIVE" | "SERVER_ACTIVE"))
|
||||
}
|
||||
|
||||
/// Current LND wallet state via the unauthenticated `/v1/state` endpoint
|
||||
/// (NON_EXISTING / LOCKED / UNLOCKED / RPC_ACTIVE / …). None if unreachable.
|
||||
async fn wallet_state(client: &reqwest::Client) -> Option<String> {
|
||||
@@ -1089,3 +1152,44 @@ mod tests {
|
||||
.is_empty());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod bitcoin_readiness_tests {
|
||||
use super::bitcoin_readiness_response;
|
||||
use serde_json::json;
|
||||
#[test]
|
||||
fn only_usable_bitcoin_rpc_allows_wallet_unlock() {
|
||||
for response in [
|
||||
json!({}),
|
||||
json!({"error":{"code":-28,"message":"Loading block index"},"result":null}),
|
||||
json!({"result":{"blocks":null}}),
|
||||
] {
|
||||
assert!(!bitcoin_readiness_response(&response));
|
||||
}
|
||||
// Initial sync is supported by LND. Loading the database is not.
|
||||
for ibd in [true, false] {
|
||||
assert!(bitcoin_readiness_response(
|
||||
&json!({"result":{"blocks":100,"initialblockdownload":ibd},"error":null})
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod syncing_wallet_state_tests {
|
||||
#[test]
|
||||
fn an_unlocked_wallet_waiting_for_chain_sync_is_never_unlocked_again() {
|
||||
for state in ["UNLOCKED", "RPC_ACTIVE", "SERVER_ACTIVE"] {
|
||||
assert!(super::wallet_is_unlocked(Some(state)));
|
||||
}
|
||||
for state in [
|
||||
None,
|
||||
Some("LOCKED"),
|
||||
Some("NON_EXISTING"),
|
||||
Some("WAITING_TO_START"),
|
||||
Some("unknown"),
|
||||
] {
|
||||
assert!(!super::wallet_is_unlocked(state));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -798,6 +798,10 @@ fn host_port_bindings_drifted(
|
||||
}
|
||||
|
||||
async fn ensure_user_podman_socket() -> Result<()> {
|
||||
// Unit tests inject a runtime; they must not restart the host Podman API.
|
||||
if cfg!(test) {
|
||||
return Ok(());
|
||||
}
|
||||
let socket_path = "/run/user/1000/podman/podman.sock";
|
||||
if podman_socket_accepts_connections(socket_path).await {
|
||||
return Ok(());
|
||||
@@ -1170,15 +1174,21 @@ impl ReconcileReport {
|
||||
fn cascade_pairs_for_report<'r>(
|
||||
report: &'r ReconcileReport,
|
||||
user_stopped: &std::collections::HashSet<String>,
|
||||
changed_backends: &HashSet<String>,
|
||||
) -> Vec<(&'r str, &'static str)> {
|
||||
let mut pairs = Vec::new();
|
||||
for (backend, action) in &report.actions {
|
||||
if !matches!(
|
||||
action,
|
||||
ReconcileAction::Installed | ReconcileAction::Started
|
||||
ReconcileAction::NoOp | ReconcileAction::Started | ReconcileAction::Installed
|
||||
) {
|
||||
continue;
|
||||
}
|
||||
// A successful systemctl start can be a no-op after a transient
|
||||
// Podman inspect failure. Require a witnessed lifecycle change.
|
||||
if !changed_backends.contains(backend) {
|
||||
continue;
|
||||
}
|
||||
for dep in crate::app_ops::address_caching_dependents(backend) {
|
||||
let dep_untouched = report
|
||||
.actions
|
||||
@@ -1192,6 +1202,25 @@ fn cascade_pairs_for_report<'r>(
|
||||
pairs
|
||||
}
|
||||
|
||||
/// Only positive runtime evidence permits disrupting an address-caching wallet.
|
||||
/// A known absent/stopped backend becoming running, a new container ID, or a
|
||||
/// changed start timestamp qualifies. A failed observation never does.
|
||||
fn backend_instance_changed(before: Option<&ContainerStatus>, after: &ContainerStatus) -> bool {
|
||||
if after.state != ContainerState::Running || after.id.is_empty() {
|
||||
return false;
|
||||
}
|
||||
let Some(before) = before else {
|
||||
return true;
|
||||
};
|
||||
if before.id.is_empty() {
|
||||
return false;
|
||||
}
|
||||
if before.id != after.id || before.state != ContainerState::Running {
|
||||
return true;
|
||||
}
|
||||
matches!((&before.started_at, &after.started_at), (Some(a), Some(b)) if !a.is_empty() && !b.is_empty() && a != b)
|
||||
}
|
||||
|
||||
#[derive(Debug, Default)]
|
||||
pub struct AdoptionReport {
|
||||
pub adopted: Vec<String>,
|
||||
@@ -1905,14 +1934,40 @@ impl ProdContainerOrchestrator {
|
||||
_ => 2,
|
||||
});
|
||||
// Live container names (any state), for the same recovery check.
|
||||
let present_containers: std::collections::HashSet<String> = self
|
||||
.runtime
|
||||
.list_containers()
|
||||
.await
|
||||
.map(|cs| cs.into_iter().map(|c| c.name).collect())
|
||||
let listed_containers = self.runtime.list_containers().await.ok();
|
||||
let present_containers: HashSet<String> = listed_containers
|
||||
.as_ref()
|
||||
.map(|cs| cs.iter().map(|c| c.name.clone()).collect())
|
||||
.unwrap_or_default();
|
||||
// Keep unknown distinct from confirmed absence. Runtime queries can
|
||||
// fail under load while systemd still has a healthy running backend.
|
||||
let mut backend_before: HashMap<String, Option<ContainerStatus>> = HashMap::new();
|
||||
for lm in &manifests {
|
||||
let id = &lm.manifest.app.id;
|
||||
if crate::app_ops::address_caching_dependents(id).is_empty() {
|
||||
continue;
|
||||
}
|
||||
let name = compute_container_name(&lm.manifest);
|
||||
match self.runtime.get_container_status(&name).await {
|
||||
Ok(status) => {
|
||||
backend_before.insert(id.clone(), Some(status));
|
||||
}
|
||||
Err(_) if listed_containers.is_some() && !present_containers.contains(&name) => {
|
||||
backend_before.insert(id.clone(), None);
|
||||
}
|
||||
Err(err) => {
|
||||
tracing::warn!(backend = %id, error = %err,
|
||||
"cannot observe backend before reconcile; will not infer a dependency restart from an action report");
|
||||
}
|
||||
}
|
||||
}
|
||||
let mut report = ReconcileReport::default();
|
||||
let disk_gb = self.disk_gb().await;
|
||||
let bitcoin_pruned = disk_gb < ARCHIVAL_BITCOIN_DISK_GB
|
||||
|| crate::settings::bitcoin_storage::load(&self.data_dir)
|
||||
.await
|
||||
.map(|settings| settings.prune)
|
||||
.unwrap_or(true);
|
||||
// Register every candidate before the (sequential, possibly slow)
|
||||
// pass so the scanner overlays queued-but-down apps as Restarting
|
||||
// instead of Stopped. Each app is deregistered as its turn finishes,
|
||||
@@ -1952,7 +2007,7 @@ impl ProdContainerOrchestrator {
|
||||
}
|
||||
if mode == ReconcileMode::ExistingOnly
|
||||
&& requires_archival_bitcoin(&app_id)
|
||||
&& disk_gb < ARCHIVAL_BITCOIN_DISK_GB
|
||||
&& bitcoin_pruned
|
||||
{
|
||||
report.record(
|
||||
&app_id,
|
||||
@@ -2087,7 +2142,20 @@ impl ProdContainerOrchestrator {
|
||||
// state recovery, repair recreate, boot InstallMissing) moves the
|
||||
// address behind a running dependent's back — §C "restart lnd after
|
||||
// ANY bitcoin recreate".
|
||||
for (backend, dep) in cascade_pairs_for_report(&report, &user_stopped) {
|
||||
let mut changed_backends = HashSet::new();
|
||||
for (backend, before) in &backend_before {
|
||||
let Some(name) = container_name_by_app_id.get(backend) else {
|
||||
continue;
|
||||
};
|
||||
if let Ok(after) = self.runtime.get_container_status(name).await {
|
||||
if backend_instance_changed(before.as_ref(), &after) {
|
||||
changed_backends.insert(backend.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
// A user stop during a slow reconcile pass still takes precedence.
|
||||
let user_stopped = crate::crash_recovery::load_user_stopped(&self.data_dir).await;
|
||||
for (backend, dep) in cascade_pairs_for_report(&report, &user_stopped, &changed_backends) {
|
||||
// Same rule as the RPC cascade: hold the dependent's op lock
|
||||
// across the restart; skip when a worker is mid-sequence.
|
||||
let lock = crate::app_ops::op_lock(dep);
|
||||
@@ -3226,6 +3294,9 @@ impl ProdContainerOrchestrator {
|
||||
}
|
||||
|
||||
async fn ensure_container_network(&self, manifest: &AppManifest) -> Result<()> {
|
||||
if cfg!(test) {
|
||||
return Ok(());
|
||||
}
|
||||
let Some(network) = manifest.app.container.network.as_deref() else {
|
||||
return Ok(());
|
||||
};
|
||||
@@ -3720,6 +3791,17 @@ impl ProdContainerOrchestrator {
|
||||
}
|
||||
let mut env = manifest.app.environment.clone();
|
||||
env.extend(manifest.app.container.resolve_derived_env(&facts));
|
||||
if matches!(manifest.app.id.as_str(), "bitcoin-core" | "bitcoin-knots") {
|
||||
let storage = crate::settings::bitcoin_storage::load(&self.data_dir).await?;
|
||||
env.retain(|entry| !entry.starts_with("BITCOIN_PRUNE="));
|
||||
if storage.prune {
|
||||
anyhow::ensure!(
|
||||
manifest.app.container.custom_args.iter().any(|arg| arg.contains("BITCOIN_PRUNE")),
|
||||
"This Bitcoin app definition cannot honor the pruning choice. Refresh the app catalog and try again."
|
||||
);
|
||||
env.push("BITCOIN_PRUNE=1".to_string());
|
||||
}
|
||||
}
|
||||
|
||||
// FM_BITCOIND_URL now comes from the manifest's {{BITCOIN_HOST}}
|
||||
// derived_env (works on Knots/Core/any distro). The old hardcoded
|
||||
@@ -6073,6 +6155,48 @@ app:
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn bitcoin_storage_choice_is_applied_and_old_catalog_cannot_silently_ignore_it() {
|
||||
let rt = Arc::new(MockRuntime::default());
|
||||
let mut orch = orch_with(rt).await;
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
orch.set_data_dir(dir.path().to_path_buf());
|
||||
for id in ["bitcoin-core", "bitcoin-knots"] {
|
||||
let mut old = pull_manifest(id, "docker.io/bitcoin/bitcoin:28");
|
||||
// No preference: existing containers need no new environment flag.
|
||||
crate::settings::bitcoin_storage::save(dir.path(), false)
|
||||
.await
|
||||
.unwrap();
|
||||
orch.resolve_dynamic_env(&mut old).await.unwrap();
|
||||
assert!(!old
|
||||
.app
|
||||
.environment
|
||||
.iter()
|
||||
.any(|s| s.starts_with("BITCOIN_PRUNE=")));
|
||||
crate::settings::bitcoin_storage::save(dir.path(), true)
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(orch
|
||||
.resolve_dynamic_env(&mut old)
|
||||
.await
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("cannot honor"));
|
||||
let mut current = pull_manifest(id, "docker.io/bitcoin/bitcoin:28");
|
||||
current
|
||||
.app
|
||||
.container
|
||||
.custom_args
|
||||
.push("if [ ${BITCOIN_PRUNE:-0} = 1 ]; then :; fi".into());
|
||||
orch.resolve_dynamic_env(&mut current).await.unwrap();
|
||||
assert!(current
|
||||
.app
|
||||
.environment
|
||||
.iter()
|
||||
.any(|s| s == "BITCOIN_PRUNE=1"));
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn install_resolves_derived_and_secret_env_before_create() {
|
||||
let rt = Arc::new(MockRuntime::default());
|
||||
@@ -6344,6 +6468,67 @@ app:
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn backend_cascade_requires_observed_instance_change() {
|
||||
let running = ContainerStatus {
|
||||
id: "container-1".into(),
|
||||
name: "bitcoin-core".into(),
|
||||
state: ContainerState::Running,
|
||||
started_at: Some("start-1".into()),
|
||||
health: None,
|
||||
exit_code: None,
|
||||
image: "bitcoin:1".into(),
|
||||
created: "created-1".into(),
|
||||
ports: vec![],
|
||||
lan_address: None,
|
||||
};
|
||||
assert!(!backend_instance_changed(Some(&running), &running));
|
||||
assert!(backend_instance_changed(None, &running));
|
||||
let mut before = running.clone();
|
||||
before.state = ContainerState::Exited;
|
||||
assert!(backend_instance_changed(Some(&before), &running));
|
||||
before = running.clone();
|
||||
before.id = "old-container".into();
|
||||
assert!(backend_instance_changed(Some(&before), &running));
|
||||
before = running.clone();
|
||||
before.started_at = Some("earlier-start".into());
|
||||
assert!(backend_instance_changed(Some(&before), &running));
|
||||
before.started_at = None;
|
||||
assert!(!backend_instance_changed(Some(&before), &running));
|
||||
before.id.clear();
|
||||
assert!(!backend_instance_changed(Some(&before), &running));
|
||||
let mut after = running.clone();
|
||||
after.state = ContainerState::Exited;
|
||||
assert!(!backend_instance_changed(None, &after));
|
||||
after = running.clone();
|
||||
after.id.clear();
|
||||
assert!(!backend_instance_changed(None, &after));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cascade_ignores_false_started_report_but_detects_real_exec_drift() {
|
||||
let none = HashSet::new();
|
||||
let mut report = ReconcileReport {
|
||||
actions: vec![
|
||||
("bitcoin-core".into(), ReconcileAction::Started),
|
||||
("lnd".into(), ReconcileAction::NoOp),
|
||||
],
|
||||
failures: vec![],
|
||||
};
|
||||
// systemctl start of an already active unit does not move its address.
|
||||
assert!(cascade_pairs_for_report(&report, &none, &none).is_empty());
|
||||
// A unit exec rewrite can restart Bitcoin while the outer reconcile
|
||||
// action remains NoOp. Runtime evidence still requires LND to reconnect.
|
||||
let changed = ["bitcoin-core".into()].into();
|
||||
report.actions[0].1 = ReconcileAction::NoOp;
|
||||
assert_eq!(
|
||||
cascade_pairs_for_report(&report, &none, &changed),
|
||||
vec![("bitcoin-core", "lnd")]
|
||||
);
|
||||
report.actions[0].1 = ReconcileAction::Left("lifecycle-op-in-flight".into());
|
||||
assert!(cascade_pairs_for_report(&report, &none, &changed).is_empty());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn cascade_pairs_cover_backend_recreate_with_running_dependent() {
|
||||
use std::collections::HashSet;
|
||||
@@ -6355,6 +6540,7 @@ app:
|
||||
failures: vec![],
|
||||
};
|
||||
let none = HashSet::new();
|
||||
let changed: HashSet<String> = ["bitcoin-core".into(), "bitcoin-knots".into()].into();
|
||||
|
||||
// Backend recreated while lnd sat running (NoOp) → cascade.
|
||||
let r = report(vec![
|
||||
@@ -6362,7 +6548,7 @@ app:
|
||||
("lnd", ReconcileAction::NoOp),
|
||||
]);
|
||||
assert_eq!(
|
||||
cascade_pairs_for_report(&r, &none),
|
||||
cascade_pairs_for_report(&r, &none, &changed),
|
||||
vec![("bitcoin-knots", "lnd")]
|
||||
);
|
||||
|
||||
@@ -6372,7 +6558,7 @@ app:
|
||||
("lnd", ReconcileAction::NoOp),
|
||||
]);
|
||||
assert_eq!(
|
||||
cascade_pairs_for_report(&r, &none),
|
||||
cascade_pairs_for_report(&r, &none, &changed),
|
||||
vec![("bitcoin-core", "lnd")]
|
||||
);
|
||||
|
||||
@@ -6381,7 +6567,7 @@ app:
|
||||
("bitcoin-knots", ReconcileAction::NoOp),
|
||||
("lnd", ReconcileAction::NoOp),
|
||||
]);
|
||||
assert!(cascade_pairs_for_report(&r, &none).is_empty());
|
||||
assert!(cascade_pairs_for_report(&r, &none, &none).is_empty());
|
||||
|
||||
// Dependent itself (re)started this pass → it already resolved the
|
||||
// fresh address; no cascade.
|
||||
@@ -6389,7 +6575,7 @@ app:
|
||||
("bitcoin-knots", ReconcileAction::Installed),
|
||||
("lnd", ReconcileAction::Started),
|
||||
]);
|
||||
assert!(cascade_pairs_for_report(&r, &none).is_empty());
|
||||
assert!(cascade_pairs_for_report(&r, &none, &changed).is_empty());
|
||||
|
||||
// User-stopped dependent is never bounced.
|
||||
let r = report(vec![
|
||||
@@ -6397,14 +6583,14 @@ app:
|
||||
("lnd", ReconcileAction::NoOp),
|
||||
]);
|
||||
let stopped: HashSet<String> = ["lnd".to_string()].into();
|
||||
assert!(cascade_pairs_for_report(&r, &stopped).is_empty());
|
||||
assert!(cascade_pairs_for_report(&r, &stopped, &changed).is_empty());
|
||||
|
||||
// Non-backend recreates don't cascade anything.
|
||||
let r = report(vec![
|
||||
("grafana", ReconcileAction::Installed),
|
||||
("lnd", ReconcileAction::NoOp),
|
||||
]);
|
||||
assert!(cascade_pairs_for_report(&r, &none).is_empty());
|
||||
assert!(cascade_pairs_for_report(&r, &none, &changed).is_empty());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
|
||||
@@ -184,6 +184,7 @@ pub struct QuadletUnit {
|
||||
pub no_new_privileges: bool,
|
||||
pub cpu_quota: Option<u32>,
|
||||
pub restart_policy: RestartPolicy,
|
||||
pub stop_grace_secs: Option<u64>,
|
||||
}
|
||||
|
||||
impl QuadletUnit {
|
||||
@@ -216,6 +217,10 @@ impl QuadletUnit {
|
||||
let _ = writeln!(s, "[Container]");
|
||||
let _ = writeln!(s, "ContainerName={}", self.name);
|
||||
let _ = writeln!(s, "Image={}", self.image);
|
||||
let grace = self
|
||||
.stop_grace_secs
|
||||
.unwrap_or_else(|| archipelago_container::runtime::stop_grace_secs_for(&self.name));
|
||||
let _ = writeln!(s, "StopTimeout={grace}");
|
||||
// Pull=never: companions are pre-pulled or built. A missing image
|
||||
// must surface as a unit start failure, not a silent retry storm.
|
||||
let _ = writeln!(s, "Pull=never");
|
||||
@@ -350,6 +355,15 @@ impl QuadletUnit {
|
||||
// the unit stuck in deactivating. Health/status remains app-level state,
|
||||
// not a systemd start gate.
|
||||
let _ = writeln!(s, "TimeoutStartSec=0");
|
||||
let _ = writeln!(s, "TimeoutStopSec={}", grace.saturating_add(15));
|
||||
// Stop explicitly before Quadlet's generated `podman rm -f`. The
|
||||
// existing container may still carry Podman's old 10-second default;
|
||||
// StopTimeout alone only protects containers created after migration.
|
||||
let _ = writeln!(s, "ExecStop=");
|
||||
let _ = writeln!(
|
||||
s,
|
||||
"ExecStop=/usr/bin/podman stop --ignore --time={grace} --cidfile=%t/%N.cid"
|
||||
);
|
||||
// Restart policy + 10s backoff. RestartSec keeps a crash-loop
|
||||
// from saturating the journal. Companions: Always. Backends:
|
||||
// OnFailure (clean stops stay stopped).
|
||||
@@ -525,6 +539,9 @@ impl QuadletUnit {
|
||||
// Always, not OnFailure: with quadlet's `--rm`, OnFailure left a
|
||||
// cleanly-exited app deleted and unrestarted. See RestartPolicy.
|
||||
restart_policy: RestartPolicy::Always,
|
||||
stop_grace_secs: Some(super::prod_orchestrator::resolve_stop_grace_secs(
|
||||
manifest, name,
|
||||
)),
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -676,6 +693,13 @@ pub async fn unit_exists(name: &str) -> bool {
|
||||
|
||||
/// Resolve the per-user quadlet dir under $HOME. Created if missing.
|
||||
pub async fn unit_dir() -> Result<PathBuf> {
|
||||
#[cfg(test)]
|
||||
{
|
||||
static TEST_UNITS: std::sync::OnceLock<PathBuf> = std::sync::OnceLock::new();
|
||||
return Ok(TEST_UNITS
|
||||
.get_or_init(|| tempfile::tempdir().unwrap().keep())
|
||||
.clone());
|
||||
}
|
||||
let home = std::env::var_os("HOME")
|
||||
.map(PathBuf::from)
|
||||
.ok_or_else(|| anyhow!("HOME not set; cannot locate quadlet unit dir"))?;
|
||||
@@ -785,7 +809,11 @@ pub async fn stop_service(service: &str) -> Result<()> {
|
||||
/// corruption — so the orchestrator passes the per-app grace here. Never waits
|
||||
/// less than `QUADLET_STOP_TIMEOUT`.
|
||||
pub async fn stop_service_with_timeout(service: &str, timeout: Duration) -> Result<()> {
|
||||
let timeout = timeout.max(QUADLET_STOP_TIMEOUT);
|
||||
let name = service.strip_suffix(".service").unwrap_or(service);
|
||||
let body = fs::read_to_string(unit_dir().await?.join(format!("{name}.container")))
|
||||
.await
|
||||
.unwrap_or_default();
|
||||
let timeout = timeout.max(stop_wait_timeout(name, &body));
|
||||
match systemctl_user_status(&["stop", service], timeout).await {
|
||||
Ok(status) if status.success() => Ok(()),
|
||||
Ok(status) => Err(anyhow!("systemctl --user stop {service} exited {status}")),
|
||||
@@ -806,10 +834,29 @@ pub async fn stop_service_with_timeout(service: &str, timeout: Duration) -> Resu
|
||||
}
|
||||
}
|
||||
|
||||
/// The command waiter must outlive both the container grace and systemd's
|
||||
/// stop deadline. Restart/repair callers must not kill Bitcoin at 45 seconds.
|
||||
fn stop_wait_timeout(name: &str, unit_body: &str) -> Duration {
|
||||
Duration::from_secs(stop_grace_from_unit(name, unit_body).saturating_add(30))
|
||||
.max(QUADLET_STOP_TIMEOUT)
|
||||
}
|
||||
|
||||
fn stop_grace_from_unit(name: &str, unit_body: &str) -> u64 {
|
||||
directive_values(unit_body, "StopTimeout=")
|
||||
.last()
|
||||
.and_then(|value| value.parse::<u64>().ok())
|
||||
.unwrap_or_else(|| archipelago_container::runtime::stop_grace_secs_for(name))
|
||||
}
|
||||
|
||||
async fn systemctl_user_status(
|
||||
args: &[&str],
|
||||
timeout: Duration,
|
||||
) -> Result<std::process::ExitStatus> {
|
||||
#[cfg(test)]
|
||||
{
|
||||
use std::os::unix::process::ExitStatusExt;
|
||||
return Ok(std::process::ExitStatus::from_raw(0));
|
||||
}
|
||||
let mut cmd = Command::new("systemctl");
|
||||
cmd.arg("--user").args(args);
|
||||
cmd.kill_on_drop(true);
|
||||
@@ -856,6 +903,10 @@ async fn wait_not_deactivating(service: &str, timeout: Duration) -> bool {
|
||||
}
|
||||
|
||||
async fn systemctl_user_output(args: &[&str], timeout: Duration) -> Result<std::process::Output> {
|
||||
#[cfg(test)]
|
||||
{
|
||||
anyhow::bail!("Unit tests have no real user service manager");
|
||||
}
|
||||
let mut cmd = Command::new("systemctl");
|
||||
cmd.arg("--user").args(args);
|
||||
cmd.kill_on_drop(true);
|
||||
@@ -923,6 +974,10 @@ fn directive_values(unit_body: &str, prefix: &str) -> Vec<String> {
|
||||
/// that systemd no longer knows about.
|
||||
pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> {
|
||||
let svc = format!("{unit_name}.service");
|
||||
let path = dir.join(format!("{unit_name}.container"));
|
||||
let body = fs::read_to_string(&path).await.unwrap_or_default();
|
||||
let timeout = stop_wait_timeout(unit_name, &body);
|
||||
let grace = stop_grace_from_unit(unit_name, &body).to_string();
|
||||
// Stop first; ignore failure (unit may already be down). BOUNDED — on
|
||||
// rootless podman a generated unit can wedge in "deactivating" while
|
||||
// `podman rm -f` hangs underneath it, and an unbounded `systemctl stop`
|
||||
@@ -930,13 +985,12 @@ pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> {
|
||||
// the package entry is stranded in `Removing` (a ghost in My Apps that also
|
||||
// blocks reinstall). If the graceful stop times out, escalate to
|
||||
// SIGKILL + reset-failed so teardown always proceeds.
|
||||
if systemctl_user_status(&["stop", &svc], QUADLET_STOP_TIMEOUT)
|
||||
if systemctl_user_status(&["stop", &svc], timeout)
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
let _ = kill_and_reset_service(&svc).await;
|
||||
}
|
||||
let path = dir.join(format!("{unit_name}.container"));
|
||||
if fs::try_exists(&path).await.unwrap_or(false) {
|
||||
match fs::remove_file(&path).await {
|
||||
Ok(()) => {}
|
||||
@@ -949,9 +1003,9 @@ pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> {
|
||||
// Bounded so a hung podman store can't re-introduce the stall this function
|
||||
// exists to avoid.
|
||||
let _ = tokio::time::timeout(
|
||||
QUADLET_STOP_TIMEOUT,
|
||||
timeout,
|
||||
Command::new("podman")
|
||||
.args(["rm", "-f", unit_name])
|
||||
.args(["rm", "-f", "--ignore", "--time", &grace, unit_name])
|
||||
.status(),
|
||||
)
|
||||
.await;
|
||||
@@ -960,6 +1014,9 @@ pub async fn disable_remove(unit_name: &str, dir: &Path) -> Result<()> {
|
||||
|
||||
/// Is the quadlet-generated service currently active?
|
||||
pub async fn is_active(service: &str) -> bool {
|
||||
if cfg!(test) {
|
||||
return false;
|
||||
}
|
||||
Command::new("systemctl")
|
||||
.args(["--user", "is-active", "--quiet", service])
|
||||
.status()
|
||||
@@ -973,6 +1030,118 @@ mod tests {
|
||||
use super::*;
|
||||
use tempfile::tempdir;
|
||||
|
||||
#[test]
|
||||
fn shutdown_grace_covers_container_systemd_and_caller() {
|
||||
for (name, grace) in [
|
||||
("bitcoin-core", 600),
|
||||
("bitcoin-knots", 600),
|
||||
("lnd", 330),
|
||||
("electrumx", 300),
|
||||
("other", 30),
|
||||
] {
|
||||
let unit = QuadletUnit {
|
||||
name: name.into(),
|
||||
..Default::default()
|
||||
};
|
||||
let body = unit.render();
|
||||
assert!(body.contains(&format!("StopTimeout={grace}\n")));
|
||||
assert!(body.contains(&format!("TimeoutStopSec={}\n", grace + 15)));
|
||||
assert!(body.contains(&format!("podman stop --ignore --time={grace} --cidfile=")));
|
||||
assert_eq!(
|
||||
stop_wait_timeout(name, &body),
|
||||
Duration::from_secs(grace + 30)
|
||||
);
|
||||
// Legacy units have no StopTimeout directive yet.
|
||||
assert_eq!(stop_wait_timeout(name, ""), Duration::from_secs(grace + 30));
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn custom_stop_grace_survives_render_and_restart_budget() {
|
||||
let manifest: AppManifest = serde_yaml::from_str(
|
||||
r#"
|
||||
app:
|
||||
id: custom-db
|
||||
name: Custom database
|
||||
version: 1.0.0
|
||||
stop_grace_secs: 900
|
||||
container:
|
||||
image: example/db:1
|
||||
"#,
|
||||
)
|
||||
.unwrap();
|
||||
let unit = QuadletUnit::from_manifest(&manifest, "custom-db");
|
||||
assert_eq!(unit.stop_grace_secs, Some(900));
|
||||
assert_eq!(
|
||||
stop_wait_timeout("custom-db", &unit.render()),
|
||||
Duration::from_secs(930)
|
||||
);
|
||||
assert_eq!(
|
||||
stop_wait_timeout("lnd", "StopTimeout=invalid"),
|
||||
Duration::from_secs(360)
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn stop_grace_migration_does_not_request_an_execution_restart() {
|
||||
let unit = sample_unit();
|
||||
let new = unit.render();
|
||||
let old = new
|
||||
.lines()
|
||||
.filter(|line| {
|
||||
!line.starts_with("StopTimeout=")
|
||||
&& !line.starts_with("TimeoutStopSec=")
|
||||
&& !line.starts_with("ExecStop=")
|
||||
})
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n");
|
||||
assert!(!exec_changed(&old, &new));
|
||||
assert!(!publish_ports_changed(&old, &new));
|
||||
assert!(!network_aliases_changed(&old, &new));
|
||||
assert!(!health_cmd_changed(&old, &new));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn actual_quadlet_generator_stops_before_forced_removal() {
|
||||
let generator = Path::new("/usr/lib/systemd/system-generators/podman-system-generator");
|
||||
if !generator.exists() {
|
||||
eprintln!(
|
||||
"Quadlet generator unavailable; run this regression on the Linux release host"
|
||||
);
|
||||
return;
|
||||
}
|
||||
let dir = tempdir().unwrap();
|
||||
let unit = QuadletUnit {
|
||||
name: "grace-test".into(),
|
||||
image: "localhost/test:latest".into(),
|
||||
stop_grace_secs: Some(600),
|
||||
..Default::default()
|
||||
};
|
||||
std::fs::write(dir.path().join("grace-test.container"), unit.render()).unwrap();
|
||||
let output = std::process::Command::new(generator)
|
||||
.args(["--user", "--dryrun"])
|
||||
.env("QUADLET_UNIT_DIRS", dir.path())
|
||||
.output()
|
||||
.unwrap();
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"{}",
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
let generated = String::from_utf8_lossy(&output.stdout).to_string()
|
||||
+ &String::from_utf8_lossy(&output.stderr);
|
||||
let stop = generated
|
||||
.find("ExecStop=/usr/bin/podman stop --ignore --time=600")
|
||||
.unwrap();
|
||||
let remove = generated.find("ExecStop=/usr/bin/podman rm ").unwrap();
|
||||
assert!(
|
||||
stop < remove,
|
||||
"Legacy container must stop gracefully before removal"
|
||||
);
|
||||
assert!(generated.contains("--stop-timeout 600"));
|
||||
assert!(generated.contains("TimeoutStopSec=615"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn render_emits_secret_env_by_reference_never_value() {
|
||||
let u = QuadletUnit {
|
||||
|
||||
Reference in New Issue
Block a user