diff --git a/crates/socket-patch-core/src/api/blob_fetcher.rs b/crates/socket-patch-core/src/api/blob_fetcher.rs index 86a51e19d..836cb67cc 100644 --- a/crates/socket-patch-core/src/api/blob_fetcher.rs +++ b/crates/socket-patch-core/src/api/blob_fetcher.rs @@ -424,10 +424,6 @@ async fn stream_cache_entry_atomic( "cache entry path has no parent directory", )) })?; - let stem = dest - .file_name() - .map(|n| n.to_string_lossy().into_owned()) - .unwrap_or_else(|| "blob".to_string()); // The cache directory (`.socket/blobs/`, `.socket/diffs/`) is created // here, by the first download, and nowhere earlier. A fetch that lands // nothing (all 404, offline, every hash mismatched, every body cut @@ -442,9 +438,7 @@ async fn stream_cache_entry_atomic( } created_dirs.push(dir); } - // Leading dot keeps the stage out of editor/glob views; the uuid suffix - // keeps concurrent writers of the same entry from colliding. - let stage = parent.join(format!(".socket-dl-{}-{}", stem, uuid::Uuid::new_v4())); + let stage = crate::utils::fs::stage_path(dest, ".socket-dl-"); let result = async { // Inside the cleanup scope: a `create_dir_all` that makes some diff --git a/crates/socket-patch-core/src/utils/fs.rs b/crates/socket-patch-core/src/utils/fs.rs index a9e2b2e2d..4f9e67581 100644 --- a/crates/socket-patch-core/src/utils/fs.rs +++ b/crates/socket-patch-core/src/utils/fs.rs @@ -454,10 +454,7 @@ pub(crate) fn normalize_lexically(path: &Path) -> Option { /// Inside an open [`super::group_commit::GroupCommit`] a commit point under /// its root is captured instead of written (see that module). pub(crate) async fn atomic_write_bytes(path: &Path, content: &[u8]) -> std::io::Result<()> { - if super::group_commit::capture_write(path, content, false) { - return Ok(()); - } - atomic_write_bytes_as(path, content, None).await + write_atomic(path, content, WriteOpts::COMMIT_POINT).await } /// [`atomic_write_bytes`], but the new inode keeps the destination's existing @@ -475,65 +472,15 @@ pub async fn atomic_write_bytes_preserving_mode( path: &Path, content: &[u8], ) -> std::io::Result<()> { - if super::group_commit::capture_write(path, content, true) { - return Ok(()); - } - let perms = tokio::fs::metadata(path) - .await - .ok() - .map(|m| m.permissions()); - atomic_write_bytes_as(path, content, perms).await -} - -/// Create the stage file for [`atomic_write_bytes_as`]. On Unix, when the -/// destination's permissions are being preserved, the stage is CREATED with -/// those bits (narrowed further by the umask) rather than the 0666 & ~umask -/// default: the full new content — a `.npmrc` `_authToken`, a 0600 private -/// manifest — is written and fsynced into the stage before the final -/// chmod, so a default-mode stage would expose it to other local users for -/// the whole write, and leave a world-readable copy behind if the process -/// is killed before the rename. -async fn create_stage( - stage: &Path, - perms: Option<&std::fs::Permissions>, -) -> std::io::Result { - let mut options = tokio::fs::OpenOptions::new(); - options.write(true).create_new(true); - #[cfg(unix)] - if let Some(p) = perms { - use std::os::unix::fs::PermissionsExt; - // Permission bits only (never setuid/setgid/sticky on a stage). A - // read-only mode (0400) is fine: O_CREAT still hands back a - // writable descriptor for the file it just created. - options.mode(p.mode() & 0o777); - } - #[cfg(not(unix))] - let _ = perms; - options.open(stage).await -} - -async fn atomic_write_bytes_as( - path: &Path, - content: &[u8], - perms: Option, -) -> std::io::Result<()> { - // A durable commit point: every artifact written without an fsync so - // far is made durable first, so this file never names bytes that could - // still be lost (see `super::durability`). - super::durability::barrier().await?; - let parent = path.parent().unwrap_or_else(|| Path::new(".")); - stage_and_rename(path, parent, content, perms, true).await?; - - // The rename only updated the parent directory entry; fsync the directory - // so the rename itself survives a crash. Best-effort, Unix only. - #[cfg(unix)] - { - if let Ok(dir) = tokio::fs::File::open(parent).await { - let _ = dir.sync_all().await; - } - } - - Ok(()) + write_atomic( + path, + content, + WriteOpts { + preserve_mode: true, + ..WriteOpts::COMMIT_POINT + }, + ) + .await } /// Atomically write a CONTENT-VERIFIED artifact (a vendored `.tgz`, wheel, @@ -543,10 +490,7 @@ async fn atomic_write_bytes_as( /// (see `super::durability` for the crash-safety argument). Readers still /// only ever see the complete old or the complete new bytes. pub(crate) async fn atomic_write_artifact(path: &Path, content: &[u8]) -> std::io::Result<()> { - let parent = path.parent().unwrap_or_else(|| Path::new(".")); - stage_and_rename(path, parent, content, None, false).await?; - super::durability::record(path); - Ok(()) + write_atomic(path, content, WriteOpts::ARTIFACT).await } /// [`atomic_write_artifact`] keeping the destination's permission bits — a @@ -556,14 +500,15 @@ pub(crate) async fn atomic_write_artifact_preserving_mode( path: &Path, content: &[u8], ) -> std::io::Result<()> { - let perms = tokio::fs::metadata(path) - .await - .ok() - .map(|m| m.permissions()); - let parent = path.parent().unwrap_or_else(|| Path::new(".")); - stage_and_rename(path, parent, content, perms, false).await?; - super::durability::record(path); - Ok(()) + write_atomic( + path, + content, + WriteOpts { + preserve_mode: true, + ..WriteOpts::ARTIFACT + }, + ) + .await } /// A group commit's file replacement (see [`super::group_commit`]): stage + @@ -575,50 +520,171 @@ pub(crate) async fn atomic_write_unsynced( content: &[u8], preserve_mode: bool, ) -> std::io::Result<()> { - let perms = if preserve_mode { - tokio::fs::metadata(path) - .await - .ok() - .map(|m| m.permissions()) - } else { - None - }; - let parent = path.parent().unwrap_or_else(|| Path::new(".")); - stage_and_rename(path, parent, content, perms, false).await + write_atomic( + path, + content, + WriteOpts { + preserve_mode, + ..WriteOpts::UNSYNCED + }, + ) + .await } /// The blocking durable writer the group-commit replay runs under the apply /// lock: stage (with the destination's mode when `preserve_mode`) + fsync + -/// rename + directory fsync. +/// rename + directory fsync. No capture and no barrier: the replay is the +/// group commit's own write-out. pub(crate) fn atomic_write_sync( path: &Path, content: &[u8], preserve_mode: bool, ) -> std::io::Result<()> { - use std::io::Write as _; + stage_and_rename_blocking(path, content, preserve_mode, true) +} + +/// What one atomic write does around its stage + rename. Every public +/// writer above is one of these policies; [`write_atomic`] and +/// [`stage_and_rename_blocking`] are the only code that acts on them. +#[derive(Clone, Copy)] +struct WriteOpts { + /// Inside an open group commit, capture the write instead of making it. + capture: bool, + /// A durable commit point: run the durability barrier first, fsync the + /// stage, and fsync the parent directory after the rename. + durable: bool, + /// The new inode keeps the destination's permission bits. + preserve_mode: bool, + /// Record the written path for the next durability barrier. + record: bool, +} + +impl WriteOpts { + const COMMIT_POINT: Self = Self { + capture: true, + durable: true, + preserve_mode: false, + record: false, + }; + const ARTIFACT: Self = Self { + capture: false, + durable: false, + preserve_mode: false, + record: true, + }; + const UNSYNCED: Self = Self { + capture: false, + durable: false, + preserve_mode: false, + record: false, + }; +} + +/// The async writers' shared body: capture or barrier per `opts`, then the +/// one blocking stage-and-rename core on tokio's blocking pool. +async fn write_atomic(path: &Path, content: &[u8], opts: WriteOpts) -> std::io::Result<()> { + if opts.capture && super::group_commit::capture_write(path, content, opts.preserve_mode) { + return Ok(()); + } + if opts.durable { + // A durable commit point: every artifact written without an fsync + // so far is made durable first, so this file never names bytes that + // could still be lost (see `super::durability`). + super::durability::barrier().await?; + } + let (owned_path, owned) = (path.to_path_buf(), content.to_vec()); + run_blocking(move || { + stage_and_rename_blocking(&owned_path, &owned, opts.preserve_mode, opts.durable) + }) + .await?; + if opts.record { + super::durability::record(path); + } + Ok(()) +} + +/// The sibling stage `path` is written to before the rename: a leading dot +/// keeps it out of editor/glob views, the uuid suffix keeps concurrent +/// writers of the same file from colliding. `prefix` is `.socket-stage-` for +/// every [`utils::fs`](self) writer and `.socket-dl-` for the blob cache. +/// The only place a stage name is built. +pub(crate) fn stage_path(path: &Path, prefix: &str) -> PathBuf { let parent = path.parent().unwrap_or_else(|| Path::new(".")); - let perms = preserve_mode - .then(|| std::fs::metadata(path).ok().map(|m| m.permissions())) - .flatten(); let stem = path .file_name() .map(|n| n.to_string_lossy().into_owned()) .unwrap_or_else(|| "file".to_string()); - let stage = parent.join(format!(".socket-stage-{}-{}", stem, uuid::Uuid::new_v4())); + parent.join(format!("{prefix}{stem}-{}", uuid::Uuid::new_v4())) +} + +/// How a stage file is opened. On Unix, when the destination's permissions +/// are being preserved, the stage is CREATED with those bits (narrowed +/// further by the umask) rather than the 0666 & ~umask default: the full +/// new content — a `.npmrc` `_authToken`, a 0600 private manifest — is +/// written and fsynced into the stage before the final chmod, so a +/// default-mode stage would expose it to other local users for the whole +/// write, and leave a world-readable copy behind if the process is killed +/// before the rename. +fn stage_open_options(perms: Option<&std::fs::Permissions>) -> std::fs::OpenOptions { let mut options = std::fs::OpenOptions::new(); options.write(true).create_new(true); #[cfg(unix)] - if let Some(p) = perms.as_ref() { + if let Some(p) = perms { use std::os::unix::fs::{OpenOptionsExt, PermissionsExt}; + // Permission bits only (never setuid/setgid/sticky on a stage). A + // read-only mode (0400) is fine: O_CREAT still hands back a + // writable descriptor for the file it just created. options.mode(p.mode() & 0o777); } - let mut file = options.open(&stage)?; + #[cfg(not(unix))] + let _ = perms; + options +} + +/// Stage `content` next to `path` and rename it over `path` — the one +/// stage-and-rename core every writer above runs. `preserve_mode` keeps the +/// destination's permission bits; `durable` fsyncs the stage before the +/// rename and the parent directory after it. +/// +/// Off Unix every stage is fsynced here, through its own writable handle: +/// Windows' `FlushFileBuffers` needs write access, so the barrier cannot +/// sync the file later by reopening its path (a read-only open is refused +/// with `ERROR_ACCESS_DENIED`, and a write open is refused for a read-only +/// artifact) — see `super::durability::sync_all_blocking`. +fn stage_and_rename_blocking( + path: &Path, + content: &[u8], + preserve_mode: bool, + durable: bool, +) -> std::io::Result<()> { + use std::io::Write as _; + let parent = path.parent().unwrap_or_else(|| Path::new(".")); + let perms = preserve_mode + .then(|| std::fs::metadata(path).ok().map(|m| m.permissions())) + .flatten(); + let sync_file = durable || !super::durability::DEFERS_FILE_SYNC; + let stage = stage_path(path, ".socket-stage-"); + + // `create_new` failing leaves no stage to clean up; every step after it + // does, so they share one error arm. + let mut file = stage_open_options(perms.as_ref()).open(&stage)?; let written = (|| { + // An unbuffered `std::fs::File`: a failed stage write (ENOSPC, EIO, + // quota) surfaces here, before the truncated stage could be renamed + // over the intact target. file.write_all(content)?; - file.sync_all()?; + if sync_file { + file.sync_all()?; + } + // Set the preserved mode on the stage *before* the rename so the + // file never appears at the destination with the wrong bits, even + // briefly. The content is already written through the open handle, + // so a restrictive mode (0400, 0000) cannot fail the write. if let Some(p) = perms { file.set_permissions(p)?; } + // Closed before the rename (Windows refuses to rename an open file) + // and before the error arm's unlink. drop(file); std::fs::rename(&stage, path) })(); @@ -626,76 +692,19 @@ pub(crate) fn atomic_write_sync( let _ = std::fs::remove_file(&stage); return Err(e); } - #[cfg(unix)] - if let Ok(dir) = std::fs::File::open(parent) { - let _ = dir.sync_all(); - } - Ok(()) -} -/// Stage `content` next to `path` and rename it over `path`; `durable` -/// fsyncs the stage before the rename. Off Unix every stage is fsynced here, -/// through its own writable handle: Windows' `FlushFileBuffers` needs write -/// access, so the barrier cannot sync the file later by reopening its path -/// (a read-only open is refused with `ERROR_ACCESS_DENIED`, and a write open -/// is refused for a read-only artifact) — see -/// `super::durability::sync_all_blocking`. -async fn stage_and_rename( - path: &Path, - parent: &Path, - content: &[u8], - perms: Option, - durable: bool, -) -> std::io::Result<()> { - let durable = durable || !super::durability::DEFERS_FILE_SYNC; - let stem = path - .file_name() - .map(|n| n.to_string_lossy().into_owned()) - .unwrap_or_else(|| "file".to_string()); - let stage = parent.join(format!(".socket-stage-{}-{}", stem, uuid::Uuid::new_v4())); - - // `create_new` failing leaves no stage to clean up; every step after it - // does, so they share one error arm. - let file = create_stage(&stage, perms.as_ref()).await?; - if let Err(e) = commit_stage(file, content, perms, &stage, path, durable).await { - let _ = tokio::fs::remove_file(&stage).await; - return Err(e); - } - Ok(()) -} - -/// Write, flush, fsync (when `durable`), (re-mode) and close the stage, then -/// rename it over `path`. Takes the handle by value so it is closed before -/// the rename (Windows refuses to rename an open file) and before the -/// caller's error-path unlink of the stage. -async fn commit_stage( - mut file: tokio::fs::File, - content: &[u8], - perms: Option, - stage: &Path, - path: &Path, - durable: bool, -) -> std::io::Result<()> { - use tokio::io::AsyncWriteExt; - file.write_all(content).await?; - // `write_all` only buffers into tokio's background writer, and - // `sync_all` stores an in-flight write error back into the handle - // instead of returning it — this flush is the only point where a - // failed stage write (ENOSPC, EIO, quota) actually surfaces. Without - // it the truncated stage would be renamed over the intact target. - file.flush().await?; + // The rename only updated the parent directory entry; fsync the + // directory so the rename itself survives a crash. Best-effort, Unix + // only. + #[cfg(unix)] if durable { - file.sync_all().await?; - } - // Set the preserved mode on the stage *before* the rename so the file - // never appears at the destination with the wrong bits, even briefly. - // The content is already written through the open handle, so a - // restrictive mode (0400, 0000) cannot fail the write. - if let Some(p) = perms { - file.set_permissions(p).await?; + if let Ok(dir) = std::fs::File::open(parent) { + let _ = dir.sync_all(); + } } - drop(file); - tokio::fs::rename(stage, path).await + #[cfg(not(unix))] + let _ = parent; + Ok(()) } #[cfg(test)] @@ -933,11 +942,10 @@ mod tests { for mode in [0o600, 0o400, 0o640] { let stage = tmp.path().join(format!(".socket-stage-npmrc-{mode:o}")); let perms = std::fs::Permissions::from_mode(mode); - let mut file = create_stage(&stage, Some(&perms)).await.unwrap(); - use tokio::io::AsyncWriteExt; + let mut file = stage_open_options(Some(&perms)).open(&stage).unwrap(); + use std::io::Write as _; // A read-only preserved mode still yields a writable stage fd. - file.write_all(b"//r/:_authToken=secret\n").await.unwrap(); - file.flush().await.unwrap(); + file.write_all(b"//r/:_authToken=secret\n").unwrap(); let got = std::fs::metadata(&stage).unwrap().permissions().mode() & 0o777; assert_eq!(got & !mode, 0, "stage {got:o} must not exceed {mode:o}"); assert_eq!( @@ -948,10 +956,88 @@ mod tests { } // No preserved mode: the plain umask default, as before. let plain = tmp.path().join(".socket-stage-plain"); - create_stage(&plain, None).await.unwrap(); + stage_open_options(None).open(&plain).unwrap(); assert!(plain.is_file()); } + /// Every former stage-and-rename copy now runs [`stage_and_rename_blocking`]: + /// each of the six writers, on a 0600 destination, writes the exact bytes, + /// leaves no stage behind, and keeps the 0600 bits exactly when its policy + /// preserves mode (the blocking and async writers agree, #728). + #[cfg(unix)] + #[tokio::test] + async fn every_writer_shares_one_stage_and_rename_core() { + use std::os::unix::fs::PermissionsExt; + let tmp = tempfile::tempdir().unwrap(); + let fresh_mode = { + let probe = tmp.path().join("probe"); + std::fs::write(&probe, b"").unwrap(); + let m = std::fs::metadata(&probe).unwrap().permissions().mode() & 0o777; + std::fs::remove_file(&probe).unwrap(); + m + }; + assert_ne!(fresh_mode, 0o600, "umask must make 0600 distinguishable"); + let cases: [(&str, bool); 8] = [ + ("bytes", false), + ("bytes_preserving", true), + ("artifact", false), + ("artifact_preserving", true), + ("unsynced_plain", false), + ("unsynced_preserving", true), + ("sync_plain", false), + ("sync_preserving", true), + ]; + for (name, preserves) in cases { + let dir = tmp.path().join(name); + std::fs::create_dir(&dir).unwrap(); + let path = dir.join("f"); + std::fs::write(&path, b"old").unwrap(); + std::fs::set_permissions(&path, std::fs::Permissions::from_mode(0o600)).unwrap(); + let content = format!("new-{name}"); + let content = content.as_bytes(); + match name { + "bytes" => atomic_write_bytes(&path, content).await, + "bytes_preserving" => atomic_write_bytes_preserving_mode(&path, content).await, + "artifact" => atomic_write_artifact(&path, content).await, + "artifact_preserving" => { + atomic_write_artifact_preserving_mode(&path, content).await + } + "unsynced_plain" => atomic_write_unsynced(&path, content, false).await, + "unsynced_preserving" => atomic_write_unsynced(&path, content, true).await, + "sync_plain" => atomic_write_sync(&path, content, false), + "sync_preserving" => atomic_write_sync(&path, content, true), + _ => unreachable!(), + } + .unwrap_or_else(|e| panic!("{name}: {e}")); + assert_eq!(std::fs::read(&path).unwrap(), content, "{name}"); + let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777; + let want = if preserves { 0o600 } else { fresh_mode }; + assert_eq!(mode, want, "{name}: mode {mode:o}, want {want:o}"); + let names: Vec = std::fs::read_dir(&dir) + .unwrap() + .map(|e| e.unwrap().file_name().to_string_lossy().into_owned()) + .collect(); + assert_eq!(names, ["f"], "{name}: stage litter {names:?}"); + } + } + + /// One stage-name builder for the `utils::fs` writers and the blob + /// cache: a dot-prefixed sibling named after the destination, unique per + /// call. + #[test] + fn stage_path_is_a_unique_hidden_sibling() { + let dest = Path::new("a").join("b").join("pkg.json"); + let one = stage_path(&dest, ".socket-stage-"); + let two = stage_path(&dest, ".socket-stage-"); + assert_ne!(one, two); + assert_eq!(one.parent(), dest.parent()); + let name = one.file_name().unwrap().to_string_lossy().into_owned(); + assert!(name.starts_with(".socket-stage-pkg.json-"), "{name}"); + let dl = stage_path(Path::new("blobs").join("abc").as_path(), ".socket-dl-"); + let dl = dl.file_name().unwrap().to_string_lossy().into_owned(); + assert!(dl.starts_with(".socket-dl-abc-"), "{dl}"); + } + /// The post-rename parent-directory fsync is best-effort: when the /// parent can be traversed and written but not opened for read /// (mode 0o333 — the stage create, the stage write, and the rename @@ -1140,16 +1226,12 @@ mod tests { FsizeGuard { prev, prev_handler } }; - // 1 MiB fits tokio's 2 MiB write buffer in one chunk, so `write_all` - // buffers it and returns Ok before the background write hits EFBIG. + // 1 MiB, then 3 MiB: both past the 256 KiB cap, so the stage's + // `write_all` hits EFBIG (the sizes predate the blocking core, when + // they split tokio's buffered-write and flush error arms). let big = vec![0xABu8; 1024 * 1024]; let res = atomic_write_bytes(&path, &big).await; - // Second case, same capped child: 3 MiB exceeds tokio's single - // 2 MiB write chunk, so `write_all` must wait on the in-flight - // first chunk and returns its EFBIG directly — exercising the - // `write_all` cleanup arm (remove stage, propagate error) instead - // of the flush arm the 1 MiB case pins above. let res_write_all = atomic_write_bytes(&path, &vec![0xCDu8; 3 * 1024 * 1024]).await; drop(guard);