From bd9d645f5e7a14f9448f3f9ad07dd33138b46aa1 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Fri, 25 Sep 2026 03:03:06 +0800 Subject: [PATCH 1/8] feat(cache): restore task results from the remote cache After a local miss in `read` or `read-write` mode, fetch the entry from the remote cache. An exact entry that passes validation is restored: its output archive is downloaded into the cache directory, checked by decoding it, and the entry is recorded locally before it replays as a hit. Hits never upload. A fallback entry, a failed validation, or a failed read is a miss. When the local cache has no entry for the task, the miss reason comes from the remote cache. Read failures produce no warning. Co-authored-by: Claude Opus 5.5 --- CHANGELOG.md | 2 +- crates/vt/src/session/cache/archive.rs | 42 +++ crates/vt/src/session/cache/display.rs | 3 + crates/vt/src/session/cache/mod.rs | 157 +++++++-- crates/vt/src/session/cache/remote.rs | 201 +++++++++++- crates/vt/src/session/execute/mod.rs | 8 +- crates/vt/src/session/reporter/summary.rs | 7 + .../fixtures/remote_cache/snapshots.toml | 188 ++++++++++- .../remote_cache/snapshots/corrupt_archive.md | 41 +++ .../remote_cache/snapshots/fallback.md | 30 ++ .../snapshots/invalid_endpoint.md | 6 +- .../fixtures/remote_cache/snapshots/read.md | 3 + .../snapshots/read_invalid_endpoint.md | 27 ++ .../snapshots/read_unreachable_endpoint.md | 9 + .../remote_cache/snapshots/read_write.md | 3 +- .../remote_cache/snapshots/restore.md | 72 +++++ .../snapshots/unreachable_endpoint.md | 6 +- crates/vt_remote_cache/Cargo.toml | 2 +- crates/vt_remote_cache/README.md | 10 +- crates/vt_remote_cache/src/lib.rs | 300 ++++++++++++++++-- 20 files changed, 1033 insertions(+), 84 deletions(-) create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fallback.md create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore.md diff --git a/CHANGELOG.md b/CHANGELOG.md index d2a0ecfb3..c968a1a86 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,7 +2,7 @@ - **Fixed** An invalid glob in `--filter` no longer shows its error message twice ([#763](https://github.com/voidzero-dev/vite-task/pull/763)). - **Changed** The detailed summary from `vp run --verbose` and `vp run --last-details` now shows each underlying cause of an error on its own line ([#761](https://github.com/voidzero-dev/vite-task/pull/761)). -- **Added** Remote caching. Configure an endpoint with the workspace's `cache: { remote: { url } }` or `VP_REMOTE_CACHE_URL`, and choose access with `--remote-cache=off|read|read-write` or `VP_REMOTE_CACHE`. The default is `read` with an endpoint and `off` without one. In `read-write` mode, `vp run` uploads the results of successful, cacheable tasks after caching them locally. A failed upload doesn't fail the task; the run summary shows a warning instead. Tasks can opt out with `cache: { remote: false }` ([#727](https://github.com/voidzero-dev/vite-task/pull/727), [#755](https://github.com/voidzero-dev/vite-task/pull/755)). +- **Added** Remote caching. Configure an endpoint with the workspace's `cache: { remote: { url } }` or `VP_REMOTE_CACHE_URL`, and choose access with `--remote-cache=off|read|read-write` or `VP_REMOTE_CACHE`. The default is `read` with an endpoint and `off` without one. After a local cache miss, `vp run` looks the task up in the remote cache and, on a hit, restores its outputs and caches it locally. A failed read is just a cache miss, with the failure as its reason. In `read-write` mode, `vp run` also uploads the results of successful, cacheable tasks after caching them locally. A failed upload doesn't fail the task; the run summary shows a warning instead. Tasks can opt out with `cache: { remote: false }` ([#727](https://github.com/voidzero-dev/vite-task/pull/727), [#755](https://github.com/voidzero-dev/vite-task/pull/755), [#756](https://github.com/voidzero-dev/vite-task/pull/756)). - **Fixed** On Windows, environment variable names used by `vp run` now match regardless of ASCII letter case. Assignments in task commands override earlier assignments and inherited variables spelled differently, and `FORCE_COLOR`, `VP_RUN_CONCURRENCY_LIMIT`, and variables requested through `@voidzero-dev/vite-task-client` are found under any spelling ([#747](https://github.com/voidzero-dev/vite-task/pull/747)). - **Changed** A task's cache settings now go inside `cache`, e.g. `cache: { env: ["NODE_ENV"], input: ["src/**"] }`; `cache: true` is the same as `cache: {}`. `env`, `untrackedEnv`, `input`, and `output` are no longer supported at the top level of a task ([#749](https://github.com/voidzero-dev/vite-task/pull/749)). - **Fixed** Cached tasks on macOS no longer intermittently fail with exit 2 and `oils I/O error (main): No such process` when a fast command finishes before the shell gets scheduled. The bundled shell that runs task commands is updated to Oils 0.38.0, which fixes this race ([#702](https://github.com/voidzero-dev/vite-task/issues/702), [#703](https://github.com/voidzero-dev/vite-task/pull/703)). diff --git a/crates/vt/src/session/cache/archive.rs b/crates/vt/src/session/cache/archive.rs index 787b72100..b600ec19c 100644 --- a/crates/vt/src/session/cache/archive.rs +++ b/crates/vt/src/session/cache/archive.rs @@ -61,3 +61,45 @@ pub fn extract_output_archive( archive.unpack(workspace_root.as_path())?; Ok(()) } + +/// Read a tar.zst archive to the end without writing any files, to check that +/// it decodes. +/// +/// # Errors +/// +/// Returns an error if opening the archive fails or it doesn't decode. +pub fn check_output_archive(archive_path: &AbsolutePath) -> io::Result<()> { + let file = File::open(archive_path.as_path())?; + let mut archive = tar::Archive::new(zstd::Decoder::new(file)?); + for entry in archive.entries()? { + io::copy(&mut entry?, &mut io::sink())?; + } + // Decode the rest of the stream too, so a truncated zstd frame after the + // end of the tar data is caught. + io::copy(&mut archive.into_inner(), &mut io::sink())?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use vt_path::AbsolutePathBuf; + + use super::*; + + #[test] + fn check_accepts_only_complete_archives() { + let tmp = tempfile::TempDir::new().unwrap(); + let dir = AbsolutePathBuf::new(tmp.path().to_path_buf()).unwrap(); + let output = RelativePathBuf::new("output.txt").unwrap(); + std::fs::write(dir.join(&output).as_path(), "built\n".repeat(1000)).unwrap(); + let archive_path = dir.join("output.tar.zst"); + create_output_archive(&dir, &[output], &archive_path).unwrap(); + check_output_archive(&archive_path).unwrap(); + + let archive = std::fs::read(archive_path.as_path()).unwrap(); + for corrupt in [&archive[..archive.len() - 1], b"corrupt"] { + std::fs::write(archive_path.as_path(), corrupt).unwrap(); + assert!(check_output_archive(&archive_path).is_err()); + } + } +} diff --git a/crates/vt/src/session/cache/display.rs b/crates/vt/src/session/cache/display.rs index 87d86958a..31e00752e 100644 --- a/crates/vt/src/session/cache/display.rs +++ b/crates/vt/src/session/cache/display.rs @@ -193,6 +193,9 @@ pub fn format_cache_status_inline(cache_status: &CacheStatus) -> Option { }; Some(vt_str::format!("○ cache miss: {reason}, executing")) } + CacheStatus::Miss(CacheMiss::RemoteReadFailed(reason)) => { + Some(vt_str::format!("○ cache miss: {reason}, executing")) + } CacheStatus::Disabled(_) => Some(Str::from("⊘ cache disabled")), } } diff --git a/crates/vt/src/session/cache/mod.rs b/crates/vt/src/session/cache/mod.rs index b72530b25..a0e8222f9 100644 --- a/crates/vt/src/session/cache/mod.rs +++ b/crates/vt/src/session/cache/mod.rs @@ -30,7 +30,7 @@ use wincode::{ io::{Reader, Writer}, }; -use self::remote::{RemoteClients, UploadError}; +use self::remote::{ReadError, RemoteClients, RemoteEntry, UploadError}; use super::execute::{ fingerprint::{PostRunFingerprint, TrackedEnvQuery}, pipe::StdOutput, @@ -149,13 +149,12 @@ pub struct ExecutionCache { } #[derive(Debug, Clone, Serialize)] -#[expect( - clippy::large_enum_variant, - reason = "FingerprintMismatch contains SpawnFingerprint which is intentionally large; boxing would add unnecessary indirection for a short-lived enum" -)] pub enum CacheMiss { NotFound, FingerprintMismatch(FingerprintMismatch), + /// Reading the remote cache failed, and the local cache has no entry for + /// the task. The message names the cause. + RemoteReadFailed(Str), } #[derive(Debug, Clone, Copy, Serialize, Deserialize)] @@ -327,19 +326,63 @@ impl ExecutionCache { /// Try to hit cache by looking up the cache entry key and validating inputs. /// Returns `Ok(Ok(cache_value))` on cache hit, `Ok(Err(cache_miss))` on miss. + /// + /// After a local miss, the remote cache is queried if the task has one. A + /// remote hit is recorded locally, with its output archive downloaded into + /// `cache_dir`, and is never uploaded. If the local cache has an entry for + /// the task, its miss reason is kept. Otherwise the reason comes from the + /// remote cache. #[tracing::instrument(level = "debug", skip_all)] pub async fn try_hit( &self, cache_metadata: &CacheMetadata, globbed_inputs: &BTreeMap, workspace_root: &AbsolutePath, + cache_dir: &AbsolutePath, ) -> anyhow::Result> { - let execution_cache_key = &cache_metadata.execution_cache_key; - let cache_key = CacheEntryKey::from_metadata(cache_metadata); + let local_miss = match self + .try_hit_local(cache_metadata, &cache_key, globbed_inputs, workspace_root) + .await? + { + Ok(cache_value) => return Ok(Ok(cache_value)), + Err(miss) => miss, + }; + let Some(ResolvedRemoteCacheConfig { url, .. }) = &cache_metadata.remote_cache else { + return Ok(Err(local_miss)); + }; + let remote_miss = match self + .try_hit_remote( + url, + cache_metadata, + &cache_key, + globbed_inputs, + workspace_root, + cache_dir, + ) + .await? + { + Ok(cache_value) => return Ok(Ok(cache_value)), + Err(miss) => miss, + }; + Ok(Err(match local_miss { + CacheMiss::NotFound => remote_miss, + local_miss => local_miss, + })) + } + + async fn try_hit_local( + &self, + cache_metadata: &CacheMetadata, + cache_key: &CacheEntryKey, + globbed_inputs: &BTreeMap, + workspace_root: &AbsolutePath, + ) -> anyhow::Result> { + let execution_cache_key = &cache_metadata.execution_cache_key; + // Try to find the cache entry by key (spawn fingerprint + input config) - if let Some(cache_value) = self.get_by_cache_key(&cache_key).await? { + if let Some(cache_value) = self.get_by_cache_key(cache_key).await? { if let Some(mismatch) = cache_value.validate(cache_metadata, globbed_inputs, workspace_root)? { @@ -347,7 +390,7 @@ impl ExecutionCache { } // Associate the execution key to the cache entry key if not already, // so that next time we can find it and report what changed - self.upsert_task_fingerprint(execution_cache_key, &cache_key).await?; + self.upsert_task_fingerprint(execution_cache_key, cache_key).await?; return Ok(Ok(cache_value)); } @@ -358,18 +401,95 @@ impl ExecutionCache { { // `get_by_cache_key` above returned None for the *current* cache key, // so the associated key must differ. - let mismatch = old_cache_key.into_mismatch(&cache_key); + let mismatch = old_cache_key.into_mismatch(cache_key); return Ok(Err(CacheMiss::FingerprintMismatch(mismatch))); } Ok(Err(CacheMiss::NotFound)) } - /// Update cache after successful execution. + /// Fetch the entry from the remote cache at `endpoint`. An exact entry + /// that passes validation is a hit once its output archive is downloaded + /// and the entry is recorded locally. A fallback entry, a failed + /// validation, or a failed read is a miss. + async fn try_hit_remote( + &self, + endpoint: &Arc, + cache_metadata: &CacheMetadata, + cache_key: &CacheEntryKey, + globbed_inputs: &BTreeMap, + workspace_root: &AbsolutePath, + cache_dir: &AbsolutePath, + ) -> anyhow::Result> { + let read_failed = |err: ReadError| { + tracing::debug!(?err, "remote cache read failed"); + CacheMiss::from(err) + }; + + let fetched = self + .remote_clients + .fetch(endpoint, cache_key, &cache_metadata.execution_cache_key) + .await; + let (cache_value, blob_id) = match fetched { + Ok(RemoteEntry::Exact { value, blob_id }) => (value, blob_id), + Ok(RemoteEntry::Fallback { key }) => { + return Ok(Err(CacheMiss::FingerprintMismatch(key.into_mismatch(cache_key)))); + } + Ok(RemoteEntry::NotFound) => return Ok(Err(CacheMiss::NotFound)), + Err(err) => return Ok(Err(read_failed(err))), + }; + if let Some(mismatch) = + cache_value.validate(cache_metadata, globbed_inputs, workspace_root)? + { + return Ok(Err(CacheMiss::FingerprintMismatch(mismatch))); + } + + let output_archive = match blob_id { + Some(blob_id) => { + match self.remote_clients.download_archive(endpoint, &blob_id, cache_dir).await { + Ok(archive_name) => Some(archive_name), + Err(err) => return Ok(Err(read_failed(err))), + } + } + None => None, + }; + let cache_value = CacheEntryValue { output_archive, ..cache_value }; + self.record(cache_key, &cache_metadata.execution_cache_key, &cache_value, cache_dir) + .await?; + Ok(Ok(cache_value)) + } + + /// Record an entry locally. /// /// If a previous entry exists for the same cache key with a different /// `output_archive`, the stale archive file in `cache_dir` is removed /// (best-effort) so it doesn't accumulate on disk. + async fn record( + &self, + cache_key: &CacheEntryKey, + execution_cache_key: &ExecutionCacheKey, + cache_value: &CacheEntryValue, + cache_dir: &AbsolutePath, + ) -> anyhow::Result<()> { + // If a previous entry exists with a stale output archive, delete the + // old file so the cache directory doesn't accumulate orphaned archives. + if let Some(old_value) = self.get_by_cache_key(cache_key).await? + && let Some(old_archive) = old_value.output_archive + && cache_value.output_archive.as_ref() != Some(&old_archive) + { + let old_archive_path = cache_dir.join(old_archive.as_str()); + // Best-effort cleanup: a missing file (e.g. after a crash or manual + // cache clear) is fine, so we ignore the error. + let _ = std::fs::remove_file(old_archive_path.as_path()); + } + + self.upsert_cache_entry(cache_key, cache_value).await?; + self.upsert_task_fingerprint(execution_cache_key, cache_key).await?; + Ok(()) + } + + /// Update cache after successful execution, recording the entry locally + /// as [`Self::record`] does. /// /// In `read-write` remote mode, the entry is then uploaded to the remote /// cache. Returns `Ok(Err(_))` if the local update succeeded but the @@ -385,20 +505,7 @@ impl ExecutionCache { let cache_key = CacheEntryKey::from_metadata(cache_metadata); - // If a previous entry exists with a stale output archive, delete the - // old file so the cache directory doesn't accumulate orphaned archives. - if let Some(old_value) = self.get_by_cache_key(&cache_key).await? - && let Some(old_archive) = old_value.output_archive - && cache_value.output_archive.as_ref() != Some(&old_archive) - { - let old_archive_path = cache_dir.join(old_archive.as_str()); - // Best-effort cleanup: a missing file (e.g. after a crash or manual - // cache clear) is fine, so we ignore the error. - let _ = std::fs::remove_file(old_archive_path.as_path()); - } - - self.upsert_cache_entry(&cache_key, &cache_value).await?; - self.upsert_task_fingerprint(execution_cache_key, &cache_key).await?; + self.record(&cache_key, execution_cache_key, &cache_value, cache_dir).await?; let Some(ResolvedRemoteCacheConfig { access: RemoteCacheAccess::ReadWrite, url }) = &cache_metadata.remote_cache diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index 5f76c6dfd..7d95c980f 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -1,5 +1,6 @@ -//! The remote cache tier. After a local update in `read-write` mode, the entry -//! is uploaded to the remote cache as opaque bytes: +//! The remote cache tier. After a local miss, the entry is fetched from the +//! remote cache, and after a local update in `read-write` mode, the entry is +//! uploaded to it. Entries are opaque bytes to the remote cache: //! //! | Field | Contents | //! | --------------- | ------------------------------------- | @@ -17,7 +18,7 @@ use std::sync::{Arc, Mutex, PoisonError}; use rustc_hash::FxHashMap; use vt_path::AbsolutePath; use vt_plan::cache_metadata::ExecutionCacheKey; -use vt_remote_cache::Client; +use vt_remote_cache::{Client, Fetched}; use vt_str::Str; use wincode::{ SchemaWrite, @@ -25,7 +26,8 @@ use wincode::{ }; use super::{ - CACHE_SCHEMA_VERSION, CacheEntryKey, CacheEntryValue, TaskCacheConfig, serialize_cache, + CACHE_SCHEMA_VERSION, CacheEntryKey, CacheEntryValue, CacheMiss, TaskCacheConfig, archive, + deserialize_cache, serialize_cache, }; /// Why an entry wasn't uploaded. @@ -37,6 +39,47 @@ pub enum UploadError { Encode(#[from] WriteError), } +/// Why no entry could be read from the remote cache. It's a cache miss, and +/// the message is its reason. The message names only the kind of failure, so +/// it's the same on every platform. +#[derive(Debug, thiserror::Error)] +pub enum ReadError { + #[error("remote cache fetch failed ({0})")] + Fetch(#[source] vt_remote_cache::Error), + #[error("remote cache download failed ({0})")] + Download(#[source] vt_remote_cache::Error), + #[error("remote cache value is corrupt")] + CorruptValue(#[source] wincode::error::ReadError), + #[error("remote cache key is corrupt")] + CorruptKey(#[source] Option), + #[error("downloaded archive is corrupt")] + CorruptArchive(#[source] std::io::Error), + #[error("failed to encode the cache key")] + Encode(#[from] WriteError), +} + +impl From for CacheMiss { + fn from(err: ReadError) -> Self { + Self::RemoteReadFailed(vt_str::format!("{err}")) + } +} + +/// A decoded fetch response. +#[derive(Debug)] +pub(super) enum RemoteEntry { + /// The entry stored under the current key. Its blob, if any, isn't + /// downloaded yet. + Exact { + value: CacheEntryValue, + blob_id: Option, + }, + /// The entry last stored for this execution, under a different key. + Fallback { + key: CacheEntryKey, + }, + NotFound, +} + /// Remote cache clients, each created when its endpoint is first used. #[derive(Debug, Default)] pub struct RemoteClients { @@ -55,6 +98,46 @@ impl RemoteClients { Ok(client) } + /// Fetch the entry stored under `cache_key`, falling back to the entry + /// last stored for `execution_cache_key`. + pub(super) async fn fetch( + &self, + endpoint: &Arc, + cache_key: &CacheEntryKey, + execution_cache_key: &ExecutionCacheKey, + ) -> Result { + let client = self.client(endpoint).map_err(ReadError::Fetch)?; + let key = encode_key(cache_key)?; + let secondary_key = encode_key(execution_cache_key)?; + let fetched = client.fetch(&key, &secondary_key).await.map_err(ReadError::Fetch)?; + decode_fetched(fetched) + } + + /// Download the blob `blob_id` into `cache_dir` and check that it decodes + /// as an output archive. Returns the archive's file name. If either step + /// fails, the file is removed. + pub(super) async fn download_archive( + &self, + endpoint: &Arc, + blob_id: &str, + cache_dir: &AbsolutePath, + ) -> Result { + let client = self.client(endpoint).map_err(ReadError::Download)?; + let archive_name = vt_str::format!("{}.tar.zst", uuid::Uuid::new_v4()); + let archive_path = cache_dir.join(archive_name.as_str()); + let result = match client.download(blob_id, &archive_path).await { + Ok(()) => { + archive::check_output_archive(&archive_path).map_err(ReadError::CorruptArchive) + } + Err(err) => Err(ReadError::Download(err)), + }; + if result.is_err() { + // Best-effort cleanup: the file may not have been created. + let _ = std::fs::remove_file(archive_path.as_path()); + } + result.map(|()| archive_name) + } + /// Upload an entry that was just recorded locally, along with its output /// archive in `cache_dir`. pub(super) async fn upload( @@ -75,6 +158,18 @@ impl RemoteClients { } } +/// Decode the value of an exact response, or the key of a fallback response. +fn decode_fetched(fetched: Fetched) -> Result { + Ok(match fetched { + Fetched::Exact { value, blob_id } => RemoteEntry::Exact { + value: deserialize_cache(&value).map_err(ReadError::CorruptValue)?, + blob_id, + }, + Fetched::Fallback { key } => RemoteEntry::Fallback { key: decode_key(&key)? }, + Fetched::NotFound => RemoteEntry::NotFound, + }) +} + #[derive(SchemaWrite)] struct KeyHeader { cache_schema_version: u32, @@ -82,14 +177,102 @@ struct KeyHeader { arch: Str, } -/// Encode a key with the header for this build. -fn encode_key>(key: &K) -> WriteResult> { - let header = KeyHeader { +/// The key header for this build. +fn encode_header() -> WriteResult> { + serialize_cache(&KeyHeader { cache_schema_version: CACHE_SCHEMA_VERSION, os: Str::from(std::env::consts::OS), arch: Str::from(std::env::consts::ARCH), - }; - let mut bytes = serialize_cache(&header)?; + }) +} + +/// Encode a key with the header for this build. +fn encode_key>(key: &K) -> WriteResult> { + let mut bytes = encode_header()?; bytes.extend(serialize_cache(key)?); Ok(bytes) } + +/// Decode a key stored with the header for this build. +fn decode_key(bytes: &[u8]) -> Result { + let header = encode_header()?; + let key = bytes.strip_prefix(header.as_slice()).ok_or(ReadError::CorruptKey(None))?; + deserialize_cache(key).map_err(|err| ReadError::CorruptKey(Some(err))) +} + +#[cfg(test)] +mod tests { + use std::{collections::BTreeMap, time::Duration}; + + use super::*; + use crate::session::execute::{ + fingerprint::PostRunFingerprint, + pipe::{OutputKind, StdOutput}, + }; + + fn cache_value() -> CacheEntryValue { + CacheEntryValue { + post_run_fingerprint: PostRunFingerprint::default(), + std_outputs: Arc::new([StdOutput { + kind: OutputKind::StdOut, + content: b"built\n".to_vec(), + }]), + duration: Duration::from_millis(5), + globbed_inputs: BTreeMap::new(), + output_archive: Some(Str::from("uploader.tar.zst")), + } + } + + #[test] + fn exact_response_gives_the_entry_to_restore() { + let value = serialize_cache(&cache_value()).unwrap(); + let fetched = Fetched::Exact { value, blob_id: Some(Str::from("1")) }; + let RemoteEntry::Exact { value, blob_id } = decode_fetched(fetched).unwrap() else { + panic!("expected an exact entry"); + }; + assert_eq!(blob_id.as_deref(), Some("1")); + assert_eq!(value.std_outputs[0].content, b"built\n"); + assert_eq!(value.duration, Duration::from_millis(5)); + } + + #[test] + fn not_found_response_has_no_entry() { + assert!(matches!(decode_fetched(Fetched::NotFound), Ok(RemoteEntry::NotFound))); + } + + fn miss_reason(error: ReadError) -> Str { + match CacheMiss::from(error) { + CacheMiss::RemoteReadFailed(reason) => reason, + miss => panic!("expected a read failure, got {miss:?}"), + } + } + + #[test] + fn value_that_does_not_decode_is_a_corrupt_entry() { + let mut value = serialize_cache(&cache_value()).unwrap(); + value.push(0); + for value in [b"not a cache value".to_vec(), value] { + let error = decode_fetched(Fetched::Exact { value, blob_id: None }).unwrap_err(); + assert!(matches!(error, ReadError::CorruptValue(_)), "{error:?}"); + assert_eq!(miss_reason(error), "remote cache value is corrupt"); + } + } + + #[test] + fn fallback_key_that_does_not_decode_is_a_corrupt_entry() { + let mut other_header = serialize_cache(&KeyHeader { + cache_schema_version: CACHE_SCHEMA_VERSION + 1, + os: Str::from(std::env::consts::OS), + arch: Str::from(std::env::consts::ARCH), + }) + .unwrap(); + other_header.extend(b"key"); + let mut garbage = encode_header().unwrap(); + garbage.extend(b"not a cache key"); + for key in [b"not a cache key".to_vec(), other_header, garbage] { + let error = decode_fetched(Fetched::Fallback { key }).unwrap_err(); + assert!(matches!(error, ReadError::CorruptKey(_)), "{error:?}"); + assert_eq!(miss_reason(error), "remote cache key is corrupt"); + } + } +} diff --git a/crates/vt/src/session/execute/mod.rs b/crates/vt/src/session/execute/mod.rs index ae31c485e..2c26587c9 100644 --- a/crates/vt/src/session/execute/mod.rs +++ b/crates/vt/src/session/execute/mod.rs @@ -378,7 +378,7 @@ async fn run( // 1. Determine cache status FIRST by trying cache hit, so the reporter can // display cache status immediately when execution begins. On a lookup // error, `start()` is never called — there is no valid status to show. - let lookup = lookup_cache(cache_metadata, cache, workspace_root).await?; + let lookup = lookup_cache(cache_metadata, cache, workspace_root, cache_dir).await?; // 2. Report execution start with the looked-up cache status (`start()` // runs exactly once on every arm) and either replay the hit — no need @@ -509,11 +509,13 @@ enum CacheLookup { Disabled, } -/// Phase 1: compute the globbed inputs and try to hit the cache. +/// Phase 1: compute the globbed inputs and try to hit the cache. A remote hit +/// downloads its output archive into `cache_dir`. async fn lookup_cache( cache_metadata: Option<&CacheMetadata>, cache: &ExecutionCache, workspace_root: &Arc, + cache_dir: &AbsolutePath, ) -> Result { let Some(cache_metadata) = cache_metadata else { return Ok(CacheLookup::Disabled); @@ -530,7 +532,7 @@ async fn lookup_cache( Report::failed(ExecutionError::Cache { kind: CacheErrorKind::Lookup, source: err }) })?; - match cache.try_hit(cache_metadata, &globbed_inputs, workspace_root).await { + match cache.try_hit(cache_metadata, &globbed_inputs, workspace_root, cache_dir).await { Ok(Ok(cached)) => Ok(CacheLookup::Hit(cached)), Ok(Err(miss)) => Ok(CacheLookup::Miss { miss, globbed_inputs }), Err(err) => { diff --git a/crates/vt/src/session/reporter/summary.rs b/crates/vt/src/session/reporter/summary.rs index 284b1646b..2edf2e80f 100644 --- a/crates/vt/src/session/reporter/summary.rs +++ b/crates/vt/src/session/reporter/summary.rs @@ -150,6 +150,9 @@ pub enum SavedCacheMissReason { /// A runner-aware tool reported a tracked bulk env query whose match-set changed /// between runs. Carries the first differing entry. TrackedEnvQueryChanged { query: TrackedEnvQuery, mismatch: EnvMismatch }, + /// Reading the remote cache failed, and the local cache had no entry. + /// Carries the failure's message. + RemoteReadFailed(Str), } /// An error's message and the messages of its causes, outermost first. @@ -274,6 +277,7 @@ impl SavedCacheMissReason { } } }, + CacheMiss::RemoteReadFailed(reason) => Self::RemoteReadFailed(reason.clone()), } } } @@ -590,6 +594,9 @@ impl TaskResult { | SavedCacheMissReason::TrackedEnvQueryChanged { mismatch, .. } => { vt_str::format!("→ Cache miss: {mismatch}") } + SavedCacheMissReason::RemoteReadFailed(reason) => { + vt_str::format!("→ Cache miss: {reason}") + } }, }, } diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml index 02e279e93..34d8c265f 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml @@ -16,7 +16,7 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], - ], comment = "A new execution is uploaded with one store request." }, + ], comment = "The fetch finds no entry. The new execution is uploaded with one store request." }, { argv = [ "remote-cache-server", "vt", @@ -60,6 +60,153 @@ steps = [ ], comment = "An endpoint without a mode selects read. The task reruns without uploading." }, ] +[[e2e]] +name = "restore" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ] }, + [ + "vt", + "cache", + "clean", + ], + [ + "vtt", + "rm", + "-rf", + "dist", + ], + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ], comment = "A remote hit downloads the output archive. Hits never upload." }, + { argv = [ + "vtt", + "print-file", + "dist/output.txt", + ], comment = "The outputs are restored." }, + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], comment = "The remote hit was recorded locally, so this is a local hit with no requests." }, + [ + "vt", + "cache", + "clean", + ], + [ + "vtt", + "write-file", + "src/a.txt", + "changed", + ], + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], comment = "The exact entry fails validation, so the archive isn't downloaded." }, +] + +[[e2e]] +name = "fallback" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ] }, + [ + "vt", + "cache", + "clean", + ], + [ + "vtt", + "replace-file-content", + "vite-task.json", + "built", + "rebuilt", + ], + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], comment = "The entry stored for this task has a different key. The miss reason compares it with the current key." }, +] + +[[e2e]] +name = "corrupt_archive" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ] }, + { argv = [ + "vtt", + "write-file", + "remote-cache/blobs/1", + "corrupt", + ], comment = "Overwrite the stored archive." }, + [ + "vt", + "cache", + "clean", + ], + { argv = [ + "remote-cache-server", + "vt", + "run", + "build", + ], comment = "The downloaded archive doesn't decode, so the task reruns." }, + { argv = [ + "vtt", + "list-dir", + "node_modules/.vite/task-cache", + "--ext", + ".tar.zst", + "--recursive", + ], comment = "Only the rerun's archive is on disk. The corrupt download was removed." }, +] + [[e2e]] name = "invalid_endpoint" steps = [ @@ -76,7 +223,7 @@ steps = [ "VP_REMOTE_CACHE_URL", "cache.example/projects/test", ], - ], comment = "The failed upload is a warning. The task succeeds." }, + ], comment = "The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds." }, { argv = [ "vt", "run", @@ -114,10 +261,45 @@ steps = [ "VP_REMOTE_CACHE_URL", "http://127.0.0.1:0/projects/test", ], - ], comment = "Nothing can listen on port 0. The failed upload is a warning. The task succeeds." }, + ], comment = "Nothing can listen on port 0. The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds." }, { argv = [ "vt", "run", "--last-details", ], comment = "The details include the underlying error." }, ] + +[[e2e]] +name = "read_invalid_endpoint" +steps = [ + { argv = [ + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE_URL", + "cache.example/projects/test", + ], + ], comment = "The failed fetch is the miss reason. Read failures aren't warnings." }, + [ + "vt", + "run", + "--last-details", + ], +] + +[[e2e]] +name = "read_unreachable_endpoint" +steps = [ + { argv = [ + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE_URL", + "http://127.0.0.1:1/projects/test", + ], + ], comment = "Nothing listens on port 1. The failed fetch is the miss reason. Read failures aren't warnings." }, +] diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md new file mode 100644 index 000000000..d49aca15b --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md @@ -0,0 +1,41 @@ +# corrupt_archive + +## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` + +``` +$ vtt write-file dist/output.txt built + +[remote-cache] POST /fetch 200 not_found +[remote-cache] POST /store 200 +``` + +## `vtt write-file remote-cache/blobs/1 corrupt` + +Overwrite the stored archive. + +``` +``` + +## `vt cache clean` + +``` +``` + +## `remote-cache-server vt run build` + +The downloaded archive doesn't decode, so the task reruns. + +``` +$ vtt write-file dist/output.txt built ○ cache miss: downloaded archive is corrupt, executing + +[remote-cache] POST /fetch 200 exact +[remote-cache] GET /blob/1 200 +``` + +## `vtt list-dir node_modules/.vite/task-cache --ext .tar.zst --recursive` + +Only the rerun's archive is on disk. The corrupt download was removed. + +``` +.tar.zst +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fallback.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fallback.md new file mode 100644 index 000000000..25d43608e --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fallback.md @@ -0,0 +1,30 @@ +# fallback + +## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` + +``` +$ vtt write-file dist/output.txt built + +[remote-cache] POST /fetch 200 not_found +[remote-cache] POST /store 200 +``` + +## `vt cache clean` + +``` +``` + +## `vtt replace-file-content vite-task.json built rebuilt` + +``` +``` + +## `remote-cache-server vt run build` + +The entry stored for this task has a different key. The miss reason compares it with the current key. + +``` +$ vtt write-file dist/output.txt rebuilt ○ cache miss: args changed, executing + +[remote-cache] POST /fetch 200 fallback +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md index 10b158c81..01d3f1414 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md @@ -2,10 +2,10 @@ ## `VP_REMOTE_CACHE=read-write VP_REMOTE_CACHE_URL=cache.example/projects/test vt run build` -The failed upload is a warning. The task succeeds. +The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds. ``` -$ vtt write-file dist/output.txt built +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (invalid endpoint), executing --- vt run: remote-cache#build not uploaded to the remote cache: invalid endpoint. (Run `vt run --last-details` for full details) @@ -27,7 +27,7 @@ Performance: 0% cache hit rate Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ - → Cache miss: no previous cache entry found + → Cache miss: remote cache fetch failed (invalid endpoint) ⚠ Not uploaded to the remote cache: invalid endpoint ↳ relative URL without a base ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read.md index 5ade066f1..3f11f5be7 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read.md @@ -5,6 +5,7 @@ ``` $ vtt write-file dist/output.txt built +[remote-cache] POST /fetch 200 not_found [remote-cache] POST /store 200 ``` @@ -19,4 +20,6 @@ An endpoint without a mode selects read. The task reruns without uploading. ``` $ vtt write-file dist/output.txt built ○ cache miss: 'src/a.txt' modified, executing + +[remote-cache] POST /fetch 200 exact ``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md new file mode 100644 index 000000000..a4bc8e07b --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md @@ -0,0 +1,27 @@ +# read_invalid_endpoint + +## `VP_REMOTE_CACHE_URL=cache.example/projects/test vt run build` + +The failed fetch is the miss reason. Read failures aren't warnings. + +``` +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (invalid endpoint), executing +``` + +## `vt run --last-details` + +``` + +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + Vite+ Task Runner • Execution Summary +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + +Statistics: 1 tasks • 0 cache hits • 1 cache misses +Performance: 0% cache hit rate + +Task Details: +──────────────────────────────────────────────── + [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ + → Cache miss: remote cache fetch failed (invalid endpoint) +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md new file mode 100644 index 000000000..aec1de48b --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md @@ -0,0 +1,9 @@ +# read_unreachable_endpoint + +## `VP_REMOTE_CACHE_URL=http://127.0.0.1:1/projects/test vt run build` + +Nothing listens on port 1. The failed fetch is the miss reason. Read failures aren't warnings. + +``` +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (network error), executing +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_write.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_write.md index 10fc112fd..70d75c261 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_write.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_write.md @@ -2,11 +2,12 @@ ## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` -A new execution is uploaded with one store request. +The fetch finds no entry. The new execution is uploaded with one store request. ``` $ vtt write-file dist/output.txt built +[remote-cache] POST /fetch 200 not_found [remote-cache] POST /store 200 ``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore.md new file mode 100644 index 000000000..b3434525f --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore.md @@ -0,0 +1,72 @@ +# restore + +## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` + +``` +$ vtt write-file dist/output.txt built + +[remote-cache] POST /fetch 200 not_found +[remote-cache] POST /store 200 +``` + +## `vt cache clean` + +``` +``` + +## `vtt rm -rf dist` + +``` +``` + +## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` + +A remote hit downloads the output archive. Hits never upload. + +``` +$ vtt write-file dist/output.txt built ◉ cache hit, replaying + +--- +vt run: cache hit. +[remote-cache] POST /fetch 200 exact +[remote-cache] GET /blob/1 200 +``` + +## `vtt print-file dist/output.txt` + +The outputs are restored. + +``` +built +``` + +## `remote-cache-server vt run build` + +The remote hit was recorded locally, so this is a local hit with no requests. + +``` +$ vtt write-file dist/output.txt built ◉ cache hit, replaying + +--- +vt run: cache hit. +``` + +## `vt cache clean` + +``` +``` + +## `vtt write-file src/a.txt changed` + +``` +``` + +## `remote-cache-server vt run build` + +The exact entry fails validation, so the archive isn't downloaded. + +``` +$ vtt write-file dist/output.txt built ○ cache miss: 'src/a.txt' modified, executing + +[remote-cache] POST /fetch 200 exact +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md index f705f47ee..047d22762 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md @@ -2,10 +2,10 @@ ## `VP_REMOTE_CACHE=read-write VP_REMOTE_CACHE_URL=http://127.0.0.1:0/projects/test vt run build` -Nothing can listen on port 0. The failed upload is a warning. The task succeeds. +Nothing can listen on port 0. The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds. ``` -$ vtt write-file dist/output.txt built +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (network error), executing --- vt run: remote-cache#build not uploaded to the remote cache: network error. (Run `vt run --last-details` for full details) @@ -27,7 +27,7 @@ Performance: 0% cache hit rate Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ - → Cache miss: no previous cache entry found + → Cache miss: remote cache fetch failed (network error) ⚠ Not uploaded to the remote cache: network error ↳ error sending request for url (http://127.0.0.1:0/projects/test/store) ↳ client error (Connect) diff --git a/crates/vt_remote_cache/Cargo.toml b/crates/vt_remote_cache/Cargo.toml index 4351ca4c7..0f1ca24ce 100644 --- a/crates/vt_remote_cache/Cargo.toml +++ b/crates/vt_remote_cache/Cargo.toml @@ -14,7 +14,7 @@ rustls = { workspace = true, features = ["ring", "std"] } serde = { workspace = true, features = ["derive"] } serde_bytes = { workspace = true } thiserror = { workspace = true } -tokio = { workspace = true, features = ["fs"] } +tokio = { workspace = true, features = ["fs", "io-util"] } url = { workspace = true } vt_path = { workspace = true } vt_str = { workspace = true } diff --git a/crates/vt_remote_cache/README.md b/crates/vt_remote_cache/README.md index e54d4651e..96d356981 100644 --- a/crates/vt_remote_cache/README.md +++ b/crates/vt_remote_cache/README.md @@ -4,8 +4,14 @@ Client for the [remote cache server API](https://github.com/voidzero-dev/vite-ta `Client::new` takes the configured endpoint, which can include a namespace path, such as `https://cache.example.com/projects/my-project`. Each operation appends its route to that path, so a store goes to `https://cache.example.com/projects/my-project/store`. Endpoints that aren't HTTP or HTTPS URLs are rejected. -`Client::store` sends a multipart request: a CBOR `metadata` part with the key, secondary key, and value as byte strings, and an optional `blob` part streamed from a file. Only HTTP 200 counts as success. The response body isn't decoded. +`Client::fetch` sends a key and a secondary key as a CBOR map of byte strings. It decodes the response into `Fetched`: an exact match with its value and optional blob ID, a fallback match with the key it's stored under, or no match. The fallback's value and blob ID aren't decoded. + +`Client::download` gets a blob by its ID and streams it into a file. The file is created only after a 200 response. A download that fails after that can leave the file incomplete, so the caller removes it. + +`Client::store` sends a multipart request: a CBOR `metadata` part with the key, secondary key, and value as byte strings, and an optional `blob` part streamed from a file. The response body isn't decoded. + +Only HTTP 200 counts as success for every operation. reqwest configures TLS. It uses the process's default rustls crypto provider, which the client installs as ring unless one is already installed, and verifies certificates with the operating system's verifier. Connections time out after 10 seconds. Reads time out after 60 seconds, and until the response headers arrive, that limit also covers sending the request. -`Error` names the kind of failure: an invalid endpoint, a client that couldn't be created, a blob file that couldn't be read, a network error (including timeouts), or a status other than 200. Its messages contain no OS-specific details, so they can be shown to users as is. The details are in the source: the underlying error, the parse error for an endpoint that isn't a URL, or the message in an error response's body. +`Error` names the kind of failure: an invalid endpoint, a client that couldn't be created, a blob file that couldn't be read or written, a network error (including timeouts and responses that end early), a status other than 200, or a malformed fetch response. Its messages contain no OS-specific details, so they can be shown to users as is. The details are in the source: the underlying error, the parse error for an endpoint that isn't a URL, or the message in an error response's body. diff --git a/crates/vt_remote_cache/src/lib.rs b/crates/vt_remote_cache/src/lib.rs index f7396ed74..702d9f9b3 100644 --- a/crates/vt_remote_cache/src/lib.rs +++ b/crates/vt_remote_cache/src/lib.rs @@ -5,9 +5,11 @@ use std::time::Duration; use reqwest::{ Response, StatusCode, + header::CONTENT_TYPE, multipart::{Form, Part}, }; -use serde::Serialize; +use serde::{Deserialize, Serialize}; +use tokio::io::AsyncWriteExt as _; use url::{ParseError, Url}; use vt_path::AbsolutePath; use vt_str::Str; @@ -35,14 +37,20 @@ pub enum Error { /// The blob file couldn't be opened. #[error("failed to read the blob")] ReadBlob(#[source] std::io::Error), + /// The downloaded blob couldn't be written to its file. + #[error("failed to write the blob")] + WriteBlob(#[source] std::io::Error), /// No complete response arrived, for example because the connection - /// failed or timed out. + /// failed, timed out, or closed before the whole body arrived. #[error("network error")] Network(#[source] reqwest::Error), /// The server responded with a status other than 200. The source is the /// message in the response body, if any. #[error("HTTP status {}", .0.as_u16())] Status(StatusCode, #[source] Option), + /// The response body isn't a fetch response. + #[error("malformed response")] + MalformedResponse(#[source] ciborium::de::Error), } /// The message in the body of an error response. @@ -61,11 +69,47 @@ struct StoreMetadata<'a> { value: &'a [u8], } +/// The keys a fetch looks up. +#[derive(Serialize)] +struct FetchRequest<'a> { + #[serde(with = "serde_bytes")] + key: &'a [u8], + #[serde(with = "serde_bytes")] + secondary_key: &'a [u8], +} + +/// The result of a fetch. +#[derive(Debug, PartialEq, Eq, Deserialize)] +#[serde(tag = "kind", rename_all = "snake_case")] +pub enum Fetched { + /// The entry stored under the requested key. + Exact { + /// The stored value. + #[serde(with = "serde_bytes")] + value: Vec, + /// The ID to download the entry's blob with, if it has one. + blob_id: Option, + }, + /// No entry is stored under the requested key, but the secondary key is + /// associated with the entry stored under `key`. The fallback entry's + /// value and blob ID aren't decoded. + Fallback { + /// The key the fallback entry is stored under. + #[serde(with = "serde_bytes")] + key: Vec, + }, + /// Neither key matched an entry. + NotFound, +} + /// A client for one remote cache endpoint. #[derive(Debug)] pub struct Client { http: reqwest::Client, + fetch_url: Url, store_url: Url, + /// `{endpoint}/blob`, to which each download appends a blob ID. + blob_url: Url, } impl Client { @@ -78,7 +122,9 @@ impl Client { /// [`Error::HttpClient`] if the HTTP client can't be created. pub fn new(endpoint: &str) -> Result { let endpoint = parse_endpoint(endpoint)?; + let fetch_url = route_url(&endpoint, "fetch")?; let store_url = route_url(&endpoint, "store")?; + let blob_url = route_url(&endpoint, "blob")?; // reqwest configures TLS with the process's default crypto provider. // Installing fails if one is already installed; vite-plus installs // ring too. @@ -88,7 +134,50 @@ impl Client { .read_timeout(READ_TIMEOUT) .build() .map_err(Error::HttpClient)?; - Ok(Self { http, store_url }) + Ok(Self { http, fetch_url, store_url, blob_url }) + } + + /// Fetch the entry stored under `key` with `POST {endpoint}/fetch`, + /// falling back to the entry associated with `secondary_key`. + /// + /// # Errors + /// + /// Returns an error if the request fails, the server responds with a + /// status other than 200, or the response isn't a fetch response. + pub async fn fetch(&self, key: &[u8], secondary_key: &[u8]) -> Result { + let body = encode_cbor(&FetchRequest { key, secondary_key }); + let response = self + .http + .post(self.fetch_url.clone()) + .header(CONTENT_TYPE, "application/cbor") + .body(body) + .send() + .await + .map_err(Error::Network)?; + let body = check_status(response).await?.bytes().await.map_err(Error::Network)?; + decode_fetched(&body) + } + + /// Download the blob `blob_id` with `GET {endpoint}/blob/{blob_id}`, + /// writing it to the file at `path`. The file is created after a 200 + /// response. If the download fails after that, it may be left incomplete. + /// + /// # Errors + /// + /// Returns an error if the request fails, the server responds with a + /// status other than 200, the body ends early, or the file can't be + /// written. + pub async fn download(&self, blob_id: &str, path: &AbsolutePath) -> Result<(), Error> { + let mut url = self.blob_url.clone(); + url.path_segments_mut().map_err(|()| Error::InvalidEndpoint(None))?.push(blob_id); + let response = self.http.get(url).send().await.map_err(Error::Network)?; + let mut response = check_status(response).await?; + let mut file = tokio::fs::File::create(path).await.map_err(Error::WriteBlob)?; + while let Some(chunk) = response.chunk().await.map_err(Error::Network)? { + file.write_all(&chunk).await.map_err(Error::WriteBlob)?; + } + file.flush().await.map_err(Error::WriteBlob)?; + Ok(()) } /// Store `value` under `key` with `POST {endpoint}/store`, uploading the @@ -118,16 +207,31 @@ impl Client { .send() .await .map_err(Error::Network)?; - if response.status() != StatusCode::OK { - return Err(status_error(response).await); - } // The response's blob ID isn't needed. Read the body anyway, so the // connection can be reused. - response.bytes().await.map_err(Error::Network)?; + check_status(response).await?.bytes().await.map_err(Error::Network)?; Ok(()) } } +/// Only HTTP 200 counts as success. For other statuses, the error includes the +/// message in the response body. +async fn check_status(response: Response) -> Result { + let status = response.status(); + if status == StatusCode::OK { + return Ok(response); + } + let message = response.text().await.ok().and_then(|text| { + let text = text.trim(); + (!text.is_empty()).then(|| ServerMessage(Str::from(text))) + }); + Err(Error::Status(status, message)) +} + +fn decode_fetched(body: &[u8]) -> Result { + ciborium::from_reader(body).map_err(Error::MalformedResponse) +} + fn parse_endpoint(endpoint: &str) -> Result { let url = Url::parse(endpoint).map_err(|err| Error::InvalidEndpoint(Some(err)))?; if !matches!(url.scheme(), "http" | "https") { @@ -143,26 +247,14 @@ fn route_url(endpoint: &Url, route: &str) -> Result { Ok(url) } -/// The error for a response with a status other than 200, with the message in -/// its body. -async fn status_error(response: Response) -> Error { - let status = response.status(); - let message = response.text().await.ok().and_then(|text| { - let text = text.trim(); - (!text.is_empty()).then(|| ServerMessage(Str::from(text))) - }); - Error::Status(status, message) -} - -fn encode_metadata(metadata: &StoreMetadata<'_>) -> Vec { +fn encode_cbor(value: &impl Serialize) -> Vec { let mut bytes = Vec::new(); - ciborium::into_writer(metadata, &mut bytes) - .expect("encoding byte strings into a Vec can't fail"); + ciborium::into_writer(value, &mut bytes).expect("encoding byte strings into a Vec can't fail"); bytes } fn metadata_part(metadata: &StoreMetadata<'_>) -> Part { - Part::bytes(encode_metadata(metadata)).mime_str("application/cbor").expect("valid MIME type") + Part::bytes(encode_cbor(metadata)).mime_str("application/cbor").expect("valid MIME type") } async fn blob_part(path: &AbsolutePath) -> Result { @@ -215,7 +307,64 @@ mod tests { expected.extend(b"\x63key\x41k"); expected.extend(b"\x6dsecondary_key\x40"); expected.extend(b"\x65value\x42\x00\xff"); - assert_eq!(encode_metadata(&metadata), expected); + assert_eq!(encode_cbor(&metadata), expected); + } + + fn cbor_map(fields: Vec<(&str, ciborium::Value)>) -> Vec { + let map = fields.into_iter().map(|(name, value)| (name.into(), value)).collect(); + encode_cbor(&ciborium::Value::Map(map)) + } + + #[test] + fn decodes_each_kind_of_fetch_response() { + let bytes = |bytes: &[u8]| ciborium::Value::Bytes(bytes.to_vec()); + let exact = cbor_map(vec![ + ("kind", "exact".into()), + ("value", bytes(b"\x00value")), + ("blob_id", "1".into()), + ]); + assert_eq!( + decode_fetched(&exact).unwrap(), + Fetched::Exact { value: b"\x00value".to_vec(), blob_id: Some(Str::from("1")) } + ); + + let without_blob = cbor_map(vec![ + ("kind", "exact".into()), + ("value", bytes(b"value")), + ("blob_id", ciborium::Value::Null), + ]); + assert_eq!( + decode_fetched(&without_blob).unwrap(), + Fetched::Exact { value: b"value".to_vec(), blob_id: None } + ); + + let fallback = cbor_map(vec![ + ("kind", "fallback".into()), + ("key", bytes(b"stored key")), + ("value", bytes(b"value")), + ("blob_id", "2".into()), + ]); + assert_eq!( + decode_fetched(&fallback).unwrap(), + Fetched::Fallback { key: b"stored key".to_vec() } + ); + + let not_found = cbor_map(vec![("kind", "not_found".into())]); + assert_eq!(decode_fetched(¬_found).unwrap(), Fetched::NotFound); + } + + #[test] + fn rejects_malformed_fetch_responses() { + for body in [ + b"\xffnot cbor".to_vec(), + cbor_map(vec![("kind", "unknown".into())]), + // The value must be a byte string. + cbor_map(vec![("kind", "exact".into()), ("value", 1.into())]), + ] { + let error = decode_fetched(&body).unwrap_err(); + assert!(matches!(error, Error::MalformedResponse(_)), "{error:?}"); + assert_eq!(error.to_string(), "malformed response"); + } } fn contains(haystack: &[u8], needle: &[u8]) -> bool { @@ -225,6 +374,13 @@ mod tests { /// Accept one HTTP request, respond with `status_line` and `body`, and /// return the raw request. fn serve_once(listener: &TcpListener, status_line: &str, body: &[u8]) -> Vec { + let headers = vt_str::format!("{status_line}\r\ncontent-length: {}\r\n\r\n", body.len()); + serve_raw_once(listener, &[headers.as_bytes(), body].concat()) + } + + /// Accept one HTTP request, write `response`, close the connection, and + /// return the raw request. + fn serve_raw_once(listener: &TcpListener, response: &[u8]) -> Vec { let (mut stream, _) = listener.accept().unwrap(); let mut request = Vec::new(); let mut buf = [0; 4096]; @@ -243,42 +399,122 @@ mod tests { let (name, value) = line.split_once(':')?; name.eq_ignore_ascii_case("content-length").then(|| value.trim().parse().unwrap()) }) - .expect("request has a content length"); + .unwrap_or(0); while request.len() < header_end + content_length { let n = stream.read(&mut buf).unwrap(); assert_ne!(n, 0, "connection closed before the request body ended"); request.extend_from_slice(&buf[..n]); } - let headers = vt_str::format!("{status_line}\r\ncontent-length: {}\r\n\r\n", body.len()); - stream.write_all(&[headers.as_bytes(), body].concat()).unwrap(); + stream.write_all(response).unwrap(); request } + fn client_for(listener: &TcpListener) -> Client { + let port = listener.local_addr().unwrap().port(); + Client::new(&vt_str::format!("http://127.0.0.1:{port}/projects/test")).unwrap() + } + /// The `metadata` part of a store request for key `k`, secondary key `s`, /// and value `v`. fn metadata_part_bytes() -> Vec { let metadata = StoreMetadata { key: b"k", secondary_key: b"s", value: b"v" }; [ b"name=\"metadata\"\r\nContent-Type: application/cbor\r\n\r\n".as_slice(), - &encode_metadata(&metadata), + &encode_cbor(&metadata), b"\r\n", ] .concat() } + #[tokio::test] + async fn fetch_posts_the_keys_as_cbor() { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let client = client_for(&listener); + let body = cbor_map(vec![("kind", "not_found".into())]); + let server = std::thread::spawn(move || serve_once(&listener, "HTTP/1.1 200 OK", &body)); + + assert_eq!(client.fetch(b"k", b"s").await.unwrap(), Fetched::NotFound); + + let request = server.join().unwrap(); + assert!(request.starts_with(b"POST /projects/test/fetch HTTP/1.1\r\n")); + assert!(contains(&request, b"content-type: application/cbor\r\n")); + let keys = encode_cbor(&FetchRequest { key: b"k", secondary_key: b"s" }); + assert!(request.ends_with(&[b"\r\n\r\n".as_slice(), &keys].concat())); + } + + #[tokio::test] + async fn fetch_fails_on_an_error_status() { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let client = client_for(&listener); + let server = std::thread::spawn(move || { + serve_once(&listener, "HTTP/1.1 503 Service Unavailable", b"try later") + }); + + let error = client.fetch(b"k", b"s").await.unwrap_err(); + assert!(matches!(error, Error::Status(StatusCode::SERVICE_UNAVAILABLE, _)), "{error:?}"); + assert_eq!(std::error::Error::source(&error).unwrap().to_string(), "try later"); + server.join().unwrap(); + } + + #[tokio::test] + async fn download_writes_the_blob_to_a_file() { + let dir = tempfile::tempdir().unwrap(); + let path = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let client = client_for(&listener); + let server = + std::thread::spawn(move || serve_once(&listener, "HTTP/1.1 200 OK", b"archive bytes")); + + client.download("7", &path).await.unwrap(); + + let request = server.join().unwrap(); + assert!(request.starts_with(b"GET /projects/test/blob/7 HTTP/1.1\r\n")); + assert_eq!(std::fs::read(path.as_path()).unwrap(), b"archive bytes"); + } + + #[tokio::test] + async fn download_of_a_missing_blob_creates_no_file() { + let dir = tempfile::tempdir().unwrap(); + let path = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let client = client_for(&listener); + let server = std::thread::spawn(move || { + serve_once(&listener, "HTTP/1.1 404 Not Found", b"Blob not found") + }); + + let error = client.download("7", &path).await.unwrap_err(); + assert_eq!(error.to_string(), "HTTP status 404"); + assert!(!path.as_path().exists()); + server.join().unwrap(); + } + + #[tokio::test] + async fn incomplete_download_is_a_network_error() { + let dir = tempfile::tempdir().unwrap(); + let path = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let client = client_for(&listener); + // The connection closes before the announced length arrives. + let server = std::thread::spawn(move || { + serve_raw_once(&listener, b"HTTP/1.1 200 OK\r\ncontent-length: 100\r\n\r\npartial") + }); + + let error = client.download("7", &path).await.unwrap_err(); + assert!(matches!(error, Error::Network(_)), "{error:?}"); + server.join().unwrap(); + } + #[tokio::test] async fn store_posts_metadata_and_blob_parts() { let dir = tempfile::tempdir().unwrap(); let blob = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); std::fs::write(blob.as_path(), b"archive bytes").unwrap(); let listener = TcpListener::bind("127.0.0.1:0").unwrap(); - let port = listener.local_addr().unwrap().port(); + let client = client_for(&listener); let server = std::thread::spawn(move || { serve_once(&listener, "HTTP/1.1 500 Internal Server Error", b"storage failed\n") }); - let client = - Client::new(&vt_str::format!("http://127.0.0.1:{port}/projects/test")).unwrap(); let err = client.store(b"k", b"s", b"v", Some(&blob)).await.unwrap_err(); assert!(matches!(err, Error::Status(StatusCode::INTERNAL_SERVER_ERROR, _))); assert_eq!(err.to_string(), "HTTP status 500"); @@ -296,13 +532,11 @@ mod tests { #[tokio::test] async fn store_without_a_blob_sends_only_metadata() { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); - let port = listener.local_addr().unwrap().port(); + let client = client_for(&listener); // `0xff` can't start a CBOR item, so this body doesn't decode. let server = std::thread::spawn(move || serve_once(&listener, "HTTP/1.1 200 OK", b"\xffnot cbor")); - let client = - Client::new(&vt_str::format!("http://127.0.0.1:{port}/projects/test")).unwrap(); client.store(b"k", b"s", b"v", None).await.unwrap(); let request = server.join().unwrap(); From e4539f2b6695badcbd74be0e2e8b0475230c8910 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Fri, 25 Sep 2026 03:30:22 +0800 Subject: [PATCH 2/8] fix(cache): treat remote validation errors as misses An error while validating a remote entry now makes the lookup a miss with the reason "remote cache entry couldn't be validated", instead of failing the task. Local validation is unchanged. Exact fetch responses must include `blob_id`, using null for no blob, so a response without it is a malformed response rather than a hit that restores no outputs. The decision from a fetch result to a hit or miss moves into `remote::resolve`, with unit tests for each outcome. Co-authored-by: Claude Opus 5.5 --- crates/vt/src/session/cache/mod.rs | 39 ++-- crates/vt/src/session/cache/remote.rs | 265 +++++++++++++++++----- crates/vt/src/session/cache/validation.rs | 15 +- crates/vt_remote_cache/README.md | 2 +- crates/vt_remote_cache/src/lib.rs | 15 +- 5 files changed, 242 insertions(+), 94 deletions(-) diff --git a/crates/vt/src/session/cache/mod.rs b/crates/vt/src/session/cache/mod.rs index a0e8222f9..1fba45d88 100644 --- a/crates/vt/src/session/cache/mod.rs +++ b/crates/vt/src/session/cache/mod.rs @@ -30,7 +30,7 @@ use wincode::{ io::{Reader, Writer}, }; -use self::remote::{ReadError, RemoteClients, RemoteEntry, UploadError}; +use self::remote::{RemoteClients, Restore, UploadError}; use super::execute::{ fingerprint::{PostRunFingerprint, TrackedEnvQuery}, pipe::StdOutput, @@ -383,9 +383,11 @@ impl ExecutionCache { // Try to find the cache entry by key (spawn fingerprint + input config) if let Some(cache_value) = self.get_by_cache_key(cache_key).await? { - if let Some(mismatch) = - cache_value.validate(cache_metadata, globbed_inputs, workspace_root)? - { + if let Some(mismatch) = cache_value.validate( + &cache_metadata.unfiltered_envs, + globbed_inputs, + workspace_root, + )? { return Ok(Err(CacheMiss::FingerprintMismatch(mismatch))); } // Associate the execution key to the cache entry key if not already, @@ -411,7 +413,8 @@ impl ExecutionCache { /// Fetch the entry from the remote cache at `endpoint`. An exact entry /// that passes validation is a hit once its output archive is downloaded /// and the entry is recorded locally. A fallback entry, a failed - /// validation, or a failed read is a miss. + /// validation, or a failed read is a miss. An error while validating + /// counts as a failed read, so the remote entry never fails the task. async fn try_hit_remote( &self, endpoint: &Arc, @@ -421,34 +424,24 @@ impl ExecutionCache { workspace_root: &AbsolutePath, cache_dir: &AbsolutePath, ) -> anyhow::Result> { - let read_failed = |err: ReadError| { - tracing::debug!(?err, "remote cache read failed"); - CacheMiss::from(err) - }; - let fetched = self .remote_clients .fetch(endpoint, cache_key, &cache_metadata.execution_cache_key) .await; - let (cache_value, blob_id) = match fetched { - Ok(RemoteEntry::Exact { value, blob_id }) => (value, blob_id), - Ok(RemoteEntry::Fallback { key }) => { - return Ok(Err(CacheMiss::FingerprintMismatch(key.into_mismatch(cache_key)))); - } - Ok(RemoteEntry::NotFound) => return Ok(Err(CacheMiss::NotFound)), - Err(err) => return Ok(Err(read_failed(err))), + let validate = |cache_value: &CacheEntryValue| { + cache_value.validate(&cache_metadata.unfiltered_envs, globbed_inputs, workspace_root) }; - if let Some(mismatch) = - cache_value.validate(cache_metadata, globbed_inputs, workspace_root)? - { - return Ok(Err(CacheMiss::FingerprintMismatch(mismatch))); - } + let Restore { value: cache_value, blob_id } = + match remote::resolve(fetched, cache_key, validate) { + Ok(restore) => restore, + Err(miss) => return Ok(Err(miss)), + }; let output_archive = match blob_id { Some(blob_id) => { match self.remote_clients.download_archive(endpoint, &blob_id, cache_dir).await { Ok(archive_name) => Some(archive_name), - Err(err) => return Ok(Err(read_failed(err))), + Err(err) => return Ok(Err(err.into_miss())), } } None => None, diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index 7d95c980f..4ba5a5127 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -26,8 +26,8 @@ use wincode::{ }; use super::{ - CACHE_SCHEMA_VERSION, CacheEntryKey, CacheEntryValue, CacheMiss, TaskCacheConfig, archive, - deserialize_cache, serialize_cache, + CACHE_SCHEMA_VERSION, CacheEntryKey, CacheEntryValue, CacheMiss, FingerprintMismatch, + TaskCacheConfig, archive, deserialize_cache, serialize_cache, }; /// Why an entry wasn't uploaded. @@ -54,30 +54,58 @@ pub enum ReadError { CorruptKey(#[source] Option), #[error("downloaded archive is corrupt")] CorruptArchive(#[source] std::io::Error), + #[error("remote cache entry couldn't be validated")] + Validate(#[source] anyhow::Error), #[error("failed to encode the cache key")] Encode(#[from] WriteError), } -impl From for CacheMiss { - fn from(err: ReadError) -> Self { - Self::RemoteReadFailed(vt_str::format!("{err}")) +impl ReadError { + /// The miss for this failure. The full error is only logged. + pub(super) fn into_miss(self) -> CacheMiss { + tracing::debug!(err = ?self, "remote cache read failed"); + CacheMiss::RemoteReadFailed(vt_str::format!("{self}")) } } -/// A decoded fetch response. +/// An exact entry that passed validation. It's a hit once its blob, if any, +/// is downloaded. #[derive(Debug)] -pub(super) enum RemoteEntry { - /// The entry stored under the current key. Its blob, if any, isn't - /// downloaded yet. - Exact { - value: CacheEntryValue, - blob_id: Option, - }, - /// The entry last stored for this execution, under a different key. - Fallback { - key: CacheEntryKey, - }, - NotFound, +pub(super) struct Restore { + pub value: CacheEntryValue, + pub blob_id: Option, +} + +/// Turn the result of a fetch into an entry to restore or a miss. `validate` +/// checks an exact entry against the current execution. A fallback's miss +/// reason compares its key with `cache_key`. A failed fetch, an entry that +/// doesn't decode, or a validation error is a read failure. +#[expect( + clippy::result_large_err, + reason = "`CacheMiss` is intentionally large, and a lookup returns it once" +)] +pub(super) fn resolve( + fetched: Result, + cache_key: &CacheEntryKey, + validate: impl FnOnce(&CacheEntryValue) -> anyhow::Result>, +) -> Result { + let (value, blob_id) = match fetched.map_err(ReadError::into_miss)? { + Fetched::Exact { value, blob_id } => { + let value = deserialize_cache(&value) + .map_err(|err| ReadError::CorruptValue(err).into_miss())?; + (value, blob_id) + } + Fetched::Fallback { key } => { + let key = decode_key(&key).map_err(ReadError::into_miss)?; + return Err(CacheMiss::FingerprintMismatch(key.into_mismatch(cache_key))); + } + Fetched::NotFound => return Err(CacheMiss::NotFound), + }; + match validate(&value) { + Ok(None) => Ok(Restore { value, blob_id }), + Ok(Some(mismatch)) => Err(CacheMiss::FingerprintMismatch(mismatch)), + Err(err) => Err(ReadError::Validate(err).into_miss()), + } } /// Remote cache clients, each created when its endpoint is first used. @@ -105,12 +133,11 @@ impl RemoteClients { endpoint: &Arc, cache_key: &CacheEntryKey, execution_cache_key: &ExecutionCacheKey, - ) -> Result { + ) -> Result { let client = self.client(endpoint).map_err(ReadError::Fetch)?; let key = encode_key(cache_key)?; let secondary_key = encode_key(execution_cache_key)?; - let fetched = client.fetch(&key, &secondary_key).await.map_err(ReadError::Fetch)?; - decode_fetched(fetched) + client.fetch(&key, &secondary_key).await.map_err(ReadError::Fetch) } /// Download the blob `blob_id` into `cache_dir` and check that it decodes @@ -158,18 +185,6 @@ impl RemoteClients { } } -/// Decode the value of an exact response, or the key of a fallback response. -fn decode_fetched(fetched: Fetched) -> Result { - Ok(match fetched { - Fetched::Exact { value, blob_id } => RemoteEntry::Exact { - value: deserialize_cache(&value).map_err(ReadError::CorruptValue)?, - blob_id, - }, - Fetched::Fallback { key } => RemoteEntry::Fallback { key: decode_key(&key)? }, - Fetched::NotFound => RemoteEntry::NotFound, - }) -} - #[derive(SchemaWrite)] struct KeyHeader { cache_schema_version: u32, @@ -204,12 +219,54 @@ fn decode_key(bytes: &[u8]) -> Result { mod tests { use std::{collections::BTreeMap, time::Duration}; + use vt_graph::config::ResolvedGlobConfig; + use vt_path::RelativePathBuf; + use vt_plan::cache_metadata::{EnvValueHash, SpawnFingerprint}; + use super::*; - use crate::session::execute::{ - fingerprint::PostRunFingerprint, - pipe::{OutputKind, StdOutput}, + use crate::session::{ + cache::InputChangeKind, + execute::{ + fingerprint::{PostRunFingerprint, TrackedEnvQuery}, + pipe::{OutputKind, StdOutput}, + }, }; + /// A key whose spawn fingerprint runs `vtt build`. `SpawnFingerprint`'s + /// fields are private to `vt_plan`, so it's decoded from types with the + /// same encoding. Update them if `SpawnFingerprint` changes. + fn cache_key(input_config: ResolvedGlobConfig) -> CacheEntryKey { + #[derive(SchemaWrite)] + enum ProgramFingerprintLayout { + OutsideWorkspace { program_name: Str }, + } + #[derive(SchemaWrite)] + struct SpawnFingerprintLayout { + cwd: RelativePathBuf, + program_fingerprint: ProgramFingerprintLayout, + args: Arc<[Str]>, + fingerprinted_envs: BTreeMap, + untracked_env_config: Arc<[Str]>, + } + + let layout = SpawnFingerprintLayout { + cwd: RelativePathBuf::default(), + program_fingerprint: ProgramFingerprintLayout::OutsideWorkspace { + program_name: Str::from("vtt"), + }, + args: Arc::from([Str::from("build")]), + fingerprinted_envs: BTreeMap::new(), + untracked_env_config: Arc::from([]), + }; + let spawn_fingerprint: SpawnFingerprint = + deserialize_cache(&serialize_cache(&layout).unwrap()).unwrap(); + CacheEntryKey { + spawn_fingerprint, + input_config, + output_config: ResolvedGlobConfig::default_auto(), + } + } + fn cache_value() -> CacheEntryValue { CacheEntryValue { post_run_fingerprint: PostRunFingerprint::default(), @@ -223,56 +280,140 @@ mod tests { } } - #[test] - fn exact_response_gives_the_entry_to_restore() { - let value = serialize_cache(&cache_value()).unwrap(); - let fetched = Fetched::Exact { value, blob_id: Some(Str::from("1")) }; - let RemoteEntry::Exact { value, blob_id } = decode_fetched(fetched).unwrap() else { - panic!("expected an exact entry"); - }; - assert_eq!(blob_id.as_deref(), Some("1")); - assert_eq!(value.std_outputs[0].content, b"built\n"); - assert_eq!(value.duration, Duration::from_millis(5)); + fn exact(value: &CacheEntryValue) -> Fetched { + Fetched::Exact { value: serialize_cache(value).unwrap(), blob_id: Some(Str::from("1")) } } - #[test] - fn not_found_response_has_no_entry() { - assert!(matches!(decode_fetched(Fetched::NotFound), Ok(RemoteEntry::NotFound))); + /// Validate as `try_hit_remote` does, with no envs and `globbed_inputs` as + /// the current inputs. + fn validate_against( + globbed_inputs: BTreeMap, + ) -> impl FnOnce(&CacheEntryValue) -> anyhow::Result> { + move |value| { + let workspace_root = vt_path::current_dir().unwrap(); + value.validate(&FxHashMap::default(), &globbed_inputs, &workspace_root) + } + } + + fn not_validated(_: &CacheEntryValue) -> anyhow::Result> { + panic!("only exact entries that decode are validated") } - fn miss_reason(error: ReadError) -> Str { - match CacheMiss::from(error) { + fn read_failure(miss: CacheMiss) -> Str { + match miss { CacheMiss::RemoteReadFailed(reason) => reason, miss => panic!("expected a read failure, got {miss:?}"), } } + #[test] + fn exact_entry_that_validates_is_restored() { + let key = cache_key(ResolvedGlobConfig::default_auto()); + let restore = + resolve(Ok(exact(&cache_value())), &key, validate_against(BTreeMap::new())).unwrap(); + assert_eq!(restore.blob_id.as_deref(), Some("1")); + assert_eq!(restore.value.std_outputs[0].content, b"built\n"); + assert_eq!(restore.value.duration, Duration::from_millis(5)); + } + + #[test] + fn exact_entry_that_fails_validation_is_a_mismatch() { + let current_inputs = BTreeMap::from([(RelativePathBuf::new("src/a.txt").unwrap(), 1)]); + let miss = resolve( + Ok(exact(&cache_value())), + &cache_key(ResolvedGlobConfig::default_auto()), + validate_against(current_inputs), + ) + .unwrap_err(); + assert!( + matches!( + &miss, + CacheMiss::FingerprintMismatch(FingerprintMismatch::InputChanged { + kind: InputChangeKind::Added, + path, + }) if path.as_str() == "src/a.txt" + ), + "{miss:?}" + ); + } + + #[test] + fn exact_entry_that_cannot_be_validated_is_a_read_failure() { + // The stored env query isn't a valid glob, so validating it errors. + let mut value = cache_value(); + value + .post_run_fingerprint + .tracked_env_queries + .insert(TrackedEnvQuery::Glob(Str::from("PROBE_[")), BTreeMap::new()); + let miss = resolve( + Ok(exact(&value)), + &cache_key(ResolvedGlobConfig::default_auto()), + validate_against(BTreeMap::new()), + ) + .unwrap_err(); + assert_eq!(read_failure(miss), "remote cache entry couldn't be validated"); + } + + #[test] + fn fallback_is_a_mismatch_with_the_stored_key() { + let stored_key = encode_key(&cache_key(ResolvedGlobConfig::default_auto())).unwrap(); + let mut input_config = ResolvedGlobConfig::default_auto(); + input_config.positive_globs.insert(Str::from("src/**")); + let miss = resolve( + Ok(Fetched::Fallback { key: stored_key }), + &cache_key(input_config), + not_validated, + ) + .unwrap_err(); + assert!( + matches!(miss, CacheMiss::FingerprintMismatch(FingerprintMismatch::InputConfig)), + "{miss:?}" + ); + } + + #[test] + fn not_found_is_a_miss_without_an_entry() { + let key = cache_key(ResolvedGlobConfig::default_auto()); + let miss = resolve(Ok(Fetched::NotFound), &key, not_validated).unwrap_err(); + assert!(matches!(miss, CacheMiss::NotFound), "{miss:?}"); + } + + #[test] + fn failed_fetch_is_a_read_failure() { + let key = cache_key(ResolvedGlobConfig::default_auto()); + let fetched = Err(ReadError::Fetch(vt_remote_cache::Error::InvalidEndpoint(None))); + let miss = resolve(fetched, &key, not_validated).unwrap_err(); + assert_eq!(read_failure(miss), "remote cache fetch failed (invalid endpoint)"); + } + #[test] fn value_that_does_not_decode_is_a_corrupt_entry() { - let mut value = serialize_cache(&cache_value()).unwrap(); - value.push(0); - for value in [b"not a cache value".to_vec(), value] { - let error = decode_fetched(Fetched::Exact { value, blob_id: None }).unwrap_err(); - assert!(matches!(error, ReadError::CorruptValue(_)), "{error:?}"); - assert_eq!(miss_reason(error), "remote cache value is corrupt"); + let key = cache_key(ResolvedGlobConfig::default_auto()); + let mut trailing = serialize_cache(&cache_value()).unwrap(); + trailing.push(0); + for value in [b"not a cache value".to_vec(), trailing] { + let fetched = Ok(Fetched::Exact { value, blob_id: None }); + let miss = resolve(fetched, &key, not_validated).unwrap_err(); + assert_eq!(read_failure(miss), "remote cache value is corrupt"); } } #[test] fn fallback_key_that_does_not_decode_is_a_corrupt_entry() { + let key = cache_key(ResolvedGlobConfig::default_auto()); let mut other_header = serialize_cache(&KeyHeader { cache_schema_version: CACHE_SCHEMA_VERSION + 1, os: Str::from(std::env::consts::OS), arch: Str::from(std::env::consts::ARCH), }) .unwrap(); - other_header.extend(b"key"); + other_header.extend(serialize_cache(&key).unwrap()); let mut garbage = encode_header().unwrap(); garbage.extend(b"not a cache key"); - for key in [b"not a cache key".to_vec(), other_header, garbage] { - let error = decode_fetched(Fetched::Fallback { key }).unwrap_err(); - assert!(matches!(error, ReadError::CorruptKey(_)), "{error:?}"); - assert_eq!(miss_reason(error), "remote cache key is corrupt"); + for stored_key in [b"not a cache key".to_vec(), other_header, garbage] { + let fetched = Ok(Fetched::Fallback { key: stored_key }); + let miss = resolve(fetched, &key, not_validated).unwrap_err(); + assert_eq!(read_failure(miss), "remote cache key is corrupt"); } } } diff --git a/crates/vt/src/session/cache/validation.rs b/crates/vt/src/session/cache/validation.rs index 1ab853227..859afec7d 100644 --- a/crates/vt/src/session/cache/validation.rs +++ b/crates/vt/src/session/cache/validation.rs @@ -1,22 +1,25 @@ //! Cache entry validation and fallback diagnostics, independent of storage. -use std::collections::BTreeMap; +use std::{collections::BTreeMap, ffi::OsStr, sync::Arc}; +use rustc_hash::FxHashMap; +use vt_casefold::EnvName; use vt_path::{AbsolutePath, RelativePathBuf}; -use vt_plan::cache_metadata::CacheMetadata; use super::{CacheEntryKey, CacheEntryValue, FingerprintMismatch, InputChangeKind}; impl CacheEntryValue { - /// Validate explicit inputs, then inferred inputs and tracked environment values. - /// Returns the first mismatch, or `None` when the entry is valid. + /// Validate explicit inputs, then inferred inputs and tracked environment + /// values against `unfiltered_envs`, the execution's + /// `CacheMetadata::unfiltered_envs`. Returns the first mismatch, or `None` + /// when the entry is valid. /// /// # Errors /// /// Propagates errors from post-run fingerprint validation. pub(crate) fn validate( &self, - cache_metadata: &CacheMetadata, + unfiltered_envs: &FxHashMap>, Arc>, globbed_inputs: &BTreeMap, workspace_root: &AbsolutePath, ) -> anyhow::Result> { @@ -25,7 +28,7 @@ impl CacheEntryValue { } self.post_run_fingerprint - .validate(workspace_root, &cache_metadata.unfiltered_envs) + .validate(workspace_root, unfiltered_envs) .map(|mismatch| mismatch.map(FingerprintMismatch::from)) } } diff --git a/crates/vt_remote_cache/README.md b/crates/vt_remote_cache/README.md index 96d356981..b24cd7a27 100644 --- a/crates/vt_remote_cache/README.md +++ b/crates/vt_remote_cache/README.md @@ -4,7 +4,7 @@ Client for the [remote cache server API](https://github.com/voidzero-dev/vite-ta `Client::new` takes the configured endpoint, which can include a namespace path, such as `https://cache.example.com/projects/my-project`. Each operation appends its route to that path, so a store goes to `https://cache.example.com/projects/my-project/store`. Endpoints that aren't HTTP or HTTPS URLs are rejected. -`Client::fetch` sends a key and a secondary key as a CBOR map of byte strings. It decodes the response into `Fetched`: an exact match with its value and optional blob ID, a fallback match with the key it's stored under, or no match. The fallback's value and blob ID aren't decoded. +`Client::fetch` sends a key and a secondary key as a CBOR map of byte strings. It decodes the response into `Fetched`: an exact match with its value and blob ID, a fallback match with the key it's stored under, or no match. An exact match must include `blob_id`, which is null when there's no blob. The fallback's value and blob ID aren't decoded. `Client::download` gets a blob by its ID and streams it into a file. The file is created only after a 200 response. A download that fails after that can leave the file incomplete, so the caller removes it. diff --git a/crates/vt_remote_cache/src/lib.rs b/crates/vt_remote_cache/src/lib.rs index 702d9f9b3..62c4b3538 100644 --- a/crates/vt_remote_cache/src/lib.rs +++ b/crates/vt_remote_cache/src/lib.rs @@ -87,7 +87,9 @@ pub enum Fetched { /// The stored value. #[serde(with = "serde_bytes")] value: Vec, - /// The ID to download the entry's blob with, if it has one. + /// The ID to download the entry's blob with, if it has one. The field + /// must be present, with null for no blob. + #[serde(deserialize_with = "Option::deserialize")] blob_id: Option, }, /// No entry is stored under the requested key, but the secondary key is @@ -359,7 +361,16 @@ mod tests { b"\xffnot cbor".to_vec(), cbor_map(vec![("kind", "unknown".into())]), // The value must be a byte string. - cbor_map(vec![("kind", "exact".into()), ("value", 1.into())]), + cbor_map(vec![ + ("kind", "exact".into()), + ("value", 1.into()), + ("blob_id", ciborium::Value::Null), + ]), + // `blob_id` is nullable, but it must be present. + cbor_map(vec![ + ("kind", "exact".into()), + ("value", ciborium::Value::Bytes(b"value".to_vec())), + ]), ] { let error = decode_fetched(&body).unwrap_err(); assert!(matches!(error, Error::MalformedResponse(_)), "{error:?}"); From c65d98d58ed40df903771253b0b86db428858ec1 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Fri, 25 Sep 2026 12:09:27 +0800 Subject: [PATCH 3/8] test(cache): use port 0 for the unreachable endpoint in `read` mode Co-authored-by: Claude Opus 5.5 --- .../tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml | 4 ++-- .../remote_cache/snapshots/read_unreachable_endpoint.md | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml index 34d8c265f..32dbb5e54 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml @@ -299,7 +299,7 @@ steps = [ ], envs = [ [ "VP_REMOTE_CACHE_URL", - "http://127.0.0.1:1/projects/test", + "http://127.0.0.1:0/projects/test", ], - ], comment = "Nothing listens on port 1. The failed fetch is the miss reason. Read failures aren't warnings." }, + ], comment = "Nothing can listen on port 0. The failed fetch is the miss reason. Read failures aren't warnings." }, ] diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md index aec1de48b..0b5138060 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md @@ -1,8 +1,8 @@ # read_unreachable_endpoint -## `VP_REMOTE_CACHE_URL=http://127.0.0.1:1/projects/test vt run build` +## `VP_REMOTE_CACHE_URL=http://127.0.0.1:0/projects/test vt run build` -Nothing listens on port 1. The failed fetch is the miss reason. Read failures aren't warnings. +Nothing can listen on port 0. The failed fetch is the miss reason. Read failures aren't warnings. ``` $ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (network error), executing From 25c34fdd38364593f5b5752bc795066bf56387df Mon Sep 17 00:00:00 2001 From: wan9chi Date: Sun, 27 Sep 2026 20:23:21 +0800 Subject: [PATCH 4/8] test(cache): show `--last-details` for remote cache read failures Co-authored-by: Claude Opus 5.5 --- .../fixtures/remote_cache/snapshots.toml | 10 ++++++++++ .../remote_cache/snapshots/corrupt_archive.md | 20 +++++++++++++++++++ .../snapshots/read_unreachable_endpoint.md | 20 +++++++++++++++++++ 3 files changed, 50 insertions(+) diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml index 32dbb5e54..44b1c1182 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml @@ -197,6 +197,11 @@ steps = [ "run", "build", ], comment = "The downloaded archive doesn't decode, so the task reruns." }, + { argv = [ + "vt", + "run", + "--last-details", + ], comment = "The details include the underlying error." }, { argv = [ "vtt", "list-dir", @@ -302,4 +307,9 @@ steps = [ "http://127.0.0.1:0/projects/test", ], ], comment = "Nothing can listen on port 0. The failed fetch is the miss reason. Read failures aren't warnings." }, + { argv = [ + "vt", + "run", + "--last-details", + ], comment = "The details include the underlying error." }, ] diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md index d49aca15b..d091f409a 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md @@ -32,6 +32,26 @@ $ vtt write-file dist/output.txt built ○ cache miss: downloaded archive is cor [remote-cache] GET /blob/1 200 ``` +## `vt run --last-details` + +The details include the underlying error. + +``` + +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + Vite+ Task Runner • Execution Summary +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + +Statistics: 1 tasks • 0 cache hits • 1 cache misses +Performance: 0% cache hit rate + +Task Details: +──────────────────────────────────────────────── + [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ + → Cache miss: downloaded archive is corrupt +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ +``` + ## `vtt list-dir node_modules/.vite/task-cache --ext .tar.zst --recursive` Only the rerun's archive is on disk. The corrupt download was removed. diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md index 0b5138060..a941b5728 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md @@ -7,3 +7,23 @@ Nothing can listen on port 0. The failed fetch is the miss reason. Read failures ``` $ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (network error), executing ``` + +## `vt run --last-details` + +The details include the underlying error. + +``` + +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + Vite+ Task Runner • Execution Summary +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + +Statistics: 1 tasks • 0 cache hits • 1 cache misses +Performance: 0% cache hit rate + +Task Details: +──────────────────────────────────────────────── + [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ + → Cache miss: remote cache fetch failed (network error) +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ +``` From 7f606231c45838f5a5117e87cf1b841c047b4650 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Sun, 27 Sep 2026 21:07:09 +0800 Subject: [PATCH 5/8] refactor(cache): check the remote cache archive while it downloads The client returns the blob's chunks instead of writing them to a file. Each chunk is written to the archive file and passed to the archive check, which reads the chunks on a blocking thread as they arrive, so the archive isn't read back from disk. A failed check stops the download. Co-authored-by: Claude Opus 5.5 --- Cargo.lock | 2 + Cargo.toml | 1 + crates/vt/Cargo.toml | 2 + crates/vt/src/session/cache/archive.rs | 23 ++++---- crates/vt/src/session/cache/remote.rs | 78 ++++++++++++++++++++++---- crates/vt_remote_cache/Cargo.toml | 3 +- crates/vt_remote_cache/README.md | 4 +- crates/vt_remote_cache/src/lib.rs | 72 ++++++++++++++---------- 8 files changed, 129 insertions(+), 56 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 867cf2389..4c4838b5b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4973,6 +4973,7 @@ dependencies = [ "anstream 1.0.0", "anyhow", "async-trait", + "bytes", "clap", "ctrlc", "derive_more", @@ -5216,6 +5217,7 @@ dependencies = [ name = "vt_remote_cache" version = "0.0.0" dependencies = [ + "bytes", "ciborium", "reqwest", "rustls", diff --git a/Cargo.toml b/Cargo.toml index fb0603fc7..cfb40df9a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -46,6 +46,7 @@ assert2 = "0.4.0" assertables = "10.0.0" async-trait = "0.1.89" base64 = "0.22.1" +bytes = "1.11.1" url = "2.5.8" wincode = "0.6.0" bindgen = "0.72.1" diff --git a/crates/vt/Cargo.toml b/crates/vt/Cargo.toml index e99b614e1..57bcabc7a 100644 --- a/crates/vt/Cargo.toml +++ b/crates/vt/Cargo.toml @@ -15,6 +15,7 @@ workspace = true anstream = { workspace = true } anyhow = { workspace = true } async-trait = { workspace = true } +bytes = { workspace = true } wincode = { workspace = true, features = ["derive"] } clap = { workspace = true, features = ["derive"] } ctrlc = { workspace = true } @@ -34,6 +35,7 @@ thiserror = { workspace = true } tar = { workspace = true } tokio = { workspace = true, features = [ "rt-multi-thread", + "fs", "io-std", "io-util", "macros", diff --git a/crates/vt/src/session/cache/archive.rs b/crates/vt/src/session/cache/archive.rs index b600ec19c..dc110e29c 100644 --- a/crates/vt/src/session/cache/archive.rs +++ b/crates/vt/src/session/cache/archive.rs @@ -1,6 +1,9 @@ //! Output archive creation and extraction using tar + zstd compression. -use std::{fs::File, io}; +use std::{ + fs::File, + io::{self, Read}, +}; use vt_path::{AbsolutePath, RelativePathBuf}; @@ -62,15 +65,14 @@ pub fn extract_output_archive( Ok(()) } -/// Read a tar.zst archive to the end without writing any files, to check that -/// it decodes. +/// Read a tar.zst archive from `reader` to the end without writing any files, +/// to check that it decodes. /// /// # Errors /// -/// Returns an error if opening the archive fails or it doesn't decode. -pub fn check_output_archive(archive_path: &AbsolutePath) -> io::Result<()> { - let file = File::open(archive_path.as_path())?; - let mut archive = tar::Archive::new(zstd::Decoder::new(file)?); +/// Returns an error if reading fails or the archive doesn't decode. +pub fn check_output_archive(reader: impl Read) -> io::Result<()> { + let mut archive = tar::Archive::new(zstd::Decoder::new(reader)?); for entry in archive.entries()? { io::copy(&mut entry?, &mut io::sink())?; } @@ -94,12 +96,11 @@ mod tests { std::fs::write(dir.join(&output).as_path(), "built\n".repeat(1000)).unwrap(); let archive_path = dir.join("output.tar.zst"); create_output_archive(&dir, &[output], &archive_path).unwrap(); - check_output_archive(&archive_path).unwrap(); - let archive = std::fs::read(archive_path.as_path()).unwrap(); + check_output_archive(archive.as_slice()).unwrap(); + for corrupt in [&archive[..archive.len() - 1], b"corrupt"] { - std::fs::write(archive_path.as_path(), corrupt).unwrap(); - assert!(check_output_archive(&archive_path).is_err()); + assert!(check_output_archive(corrupt).is_err()); } } } diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index 4ba5a5127..16dbc60ae 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -13,9 +13,14 @@ //! target OS and architecture. Platforms share an endpoint's namespace, but //! their keys differ. -use std::sync::{Arc, Mutex, PoisonError}; +use std::{ + io, + sync::{Arc, Mutex, PoisonError}, +}; +use bytes::Bytes; use rustc_hash::FxHashMap; +use tokio::{io::AsyncWriteExt as _, sync::mpsc}; use vt_path::AbsolutePath; use vt_plan::cache_metadata::ExecutionCacheKey; use vt_remote_cache::{Client, Fetched}; @@ -53,7 +58,9 @@ pub enum ReadError { #[error("remote cache key is corrupt")] CorruptKey(#[source] Option), #[error("downloaded archive is corrupt")] - CorruptArchive(#[source] std::io::Error), + CorruptArchive(#[source] io::Error), + #[error("failed to write the downloaded archive")] + WriteArchive(#[source] io::Error), #[error("remote cache entry couldn't be validated")] Validate(#[source] anyhow::Error), #[error("failed to encode the cache key")] @@ -140,9 +147,9 @@ impl RemoteClients { client.fetch(&key, &secondary_key).await.map_err(ReadError::Fetch) } - /// Download the blob `blob_id` into `cache_dir` and check that it decodes - /// as an output archive. Returns the archive's file name. If either step - /// fails, the file is removed. + /// Download the blob `blob_id` into `cache_dir`, checking that it decodes + /// as an output archive as it arrives. Returns the archive's file name. If + /// the download or the check fails, the file is removed. pub(super) async fn download_archive( &self, endpoint: &Arc, @@ -152,12 +159,7 @@ impl RemoteClients { let client = self.client(endpoint).map_err(ReadError::Download)?; let archive_name = vt_str::format!("{}.tar.zst", uuid::Uuid::new_v4()); let archive_path = cache_dir.join(archive_name.as_str()); - let result = match client.download(blob_id, &archive_path).await { - Ok(()) => { - archive::check_output_archive(&archive_path).map_err(ReadError::CorruptArchive) - } - Err(err) => Err(ReadError::Download(err)), - }; + let result = download_checked(&client, blob_id, &archive_path).await; if result.is_err() { // Best-effort cleanup: the file may not have been created. let _ = std::fs::remove_file(archive_path.as_path()); @@ -185,6 +187,60 @@ impl RemoteClients { } } +/// Chunks buffered between the download and the archive check. +const CHECK_BUFFER_CHUNKS: usize = 16; + +/// Download the blob `blob_id` to the file at `path`. Each chunk is written to +/// the file and passed to the archive check, which runs on a blocking thread. +async fn download_checked( + client: &Client, + blob_id: &str, + path: &AbsolutePath, +) -> Result<(), ReadError> { + let mut download = client.download(blob_id).await.map_err(ReadError::Download)?; + let mut file = + tokio::fs::File::create(path.as_path()).await.map_err(ReadError::WriteArchive)?; + let (sender, receiver) = mpsc::channel(CHECK_BUFFER_CHUNKS); + let check = tokio::task::spawn_blocking(move || { + archive::check_output_archive(ChunkReader { receiver, chunk: Bytes::new() }) + }); + while let Some(chunk) = download.chunk().await.map_err(ReadError::Download)? { + file.write_all(&chunk).await.map_err(ReadError::WriteArchive)?; + // The check stops reading when it fails. + if sender.send(chunk).await.is_err() { + break; + } + } + drop(sender); + check + .await + .map_err(io::Error::from) + .and_then(|checked| checked) + .map_err(ReadError::CorruptArchive)?; + file.flush().await.map_err(ReadError::WriteArchive) +} + +/// Reads the chunks sent through a channel, in order, until the sender is +/// dropped. Reading blocks, so it's only for blocking threads. +struct ChunkReader { + receiver: mpsc::Receiver, + chunk: Bytes, +} + +impl io::Read for ChunkReader { + fn read(&mut self, buf: &mut [u8]) -> io::Result { + while self.chunk.is_empty() { + match self.receiver.blocking_recv() { + Some(chunk) => self.chunk = chunk, + None => return Ok(0), + } + } + let len = buf.len().min(self.chunk.len()); + buf[..len].copy_from_slice(&self.chunk.split_to(len)); + Ok(len) + } +} + #[derive(SchemaWrite)] struct KeyHeader { cache_schema_version: u32, diff --git a/crates/vt_remote_cache/Cargo.toml b/crates/vt_remote_cache/Cargo.toml index 0f1ca24ce..0ab0131a8 100644 --- a/crates/vt_remote_cache/Cargo.toml +++ b/crates/vt_remote_cache/Cargo.toml @@ -8,13 +8,14 @@ publish = false rust-version.workspace = true [dependencies] +bytes = { workspace = true } ciborium = { workspace = true } reqwest = { workspace = true, features = ["multipart", "rustls-no-provider", "stream"] } rustls = { workspace = true, features = ["ring", "std"] } serde = { workspace = true, features = ["derive"] } serde_bytes = { workspace = true } thiserror = { workspace = true } -tokio = { workspace = true, features = ["fs", "io-util"] } +tokio = { workspace = true, features = ["fs"] } url = { workspace = true } vt_path = { workspace = true } vt_str = { workspace = true } diff --git a/crates/vt_remote_cache/README.md b/crates/vt_remote_cache/README.md index b24cd7a27..08650ece6 100644 --- a/crates/vt_remote_cache/README.md +++ b/crates/vt_remote_cache/README.md @@ -6,7 +6,7 @@ Client for the [remote cache server API](https://github.com/voidzero-dev/vite-ta `Client::fetch` sends a key and a secondary key as a CBOR map of byte strings. It decodes the response into `Fetched`: an exact match with its value and blob ID, a fallback match with the key it's stored under, or no match. An exact match must include `blob_id`, which is null when there's no blob. The fallback's value and blob ID aren't decoded. -`Client::download` gets a blob by its ID and streams it into a file. The file is created only after a 200 response. A download that fails after that can leave the file incomplete, so the caller removes it. +`Client::download` gets a blob by its ID. After a 200 response, it returns a `Download`, from which the caller reads the blob's chunks as they arrive. `Client::store` sends a multipart request: a CBOR `metadata` part with the key, secondary key, and value as byte strings, and an optional `blob` part streamed from a file. The response body isn't decoded. @@ -14,4 +14,4 @@ Only HTTP 200 counts as success for every operation. reqwest configures TLS. It uses the process's default rustls crypto provider, which the client installs as ring unless one is already installed, and verifies certificates with the operating system's verifier. Connections time out after 10 seconds. Reads time out after 60 seconds, and until the response headers arrive, that limit also covers sending the request. -`Error` names the kind of failure: an invalid endpoint, a client that couldn't be created, a blob file that couldn't be read or written, a network error (including timeouts and responses that end early), a status other than 200, or a malformed fetch response. Its messages contain no OS-specific details, so they can be shown to users as is. The details are in the source: the underlying error, the parse error for an endpoint that isn't a URL, or the message in an error response's body. +`Error` names the kind of failure: an invalid endpoint, a client that couldn't be created, a blob file that couldn't be read, a network error (including timeouts and responses that end early), a status other than 200, or a malformed fetch response. Its messages contain no OS-specific details, so they can be shown to users as is. The details are in the source: the underlying error, the parse error for an endpoint that isn't a URL, or the message in an error response's body. diff --git a/crates/vt_remote_cache/src/lib.rs b/crates/vt_remote_cache/src/lib.rs index 62c4b3538..d47a37d98 100644 --- a/crates/vt_remote_cache/src/lib.rs +++ b/crates/vt_remote_cache/src/lib.rs @@ -3,13 +3,13 @@ use std::time::Duration; +use bytes::Bytes; use reqwest::{ Response, StatusCode, header::CONTENT_TYPE, multipart::{Form, Part}, }; use serde::{Deserialize, Serialize}; -use tokio::io::AsyncWriteExt as _; use url::{ParseError, Url}; use vt_path::AbsolutePath; use vt_str::Str; @@ -37,9 +37,6 @@ pub enum Error { /// The blob file couldn't be opened. #[error("failed to read the blob")] ReadBlob(#[source] std::io::Error), - /// The downloaded blob couldn't be written to its file. - #[error("failed to write the blob")] - WriteBlob(#[source] std::io::Error), /// No complete response arrived, for example because the connection /// failed, timed out, or closed before the whole body arrived. #[error("network error")] @@ -58,6 +55,24 @@ pub enum Error { #[error("{0}")] pub struct ServerMessage(Str); +/// A blob being downloaded. +#[derive(Debug)] +pub struct Download { + response: Response, +} + +impl Download { + /// The next chunk of the blob, or `None` after the last one. + /// + /// # Errors + /// + /// Returns an error if the rest of the blob can't be read, for example + /// because the connection closes before it all arrives. + pub async fn chunk(&mut self) -> Result, Error> { + self.response.chunk().await.map_err(Error::Network) + } +} + /// The `metadata` part of a store request. #[derive(Serialize)] struct StoreMetadata<'a> { @@ -160,26 +175,19 @@ impl Client { decode_fetched(&body) } - /// Download the blob `blob_id` with `GET {endpoint}/blob/{blob_id}`, - /// writing it to the file at `path`. The file is created after a 200 - /// response. If the download fails after that, it may be left incomplete. + /// Start downloading the blob `blob_id` with + /// `GET {endpoint}/blob/{blob_id}`. Its chunks are read from the returned + /// [`Download`] as they arrive. /// /// # Errors /// - /// Returns an error if the request fails, the server responds with a - /// status other than 200, the body ends early, or the file can't be - /// written. - pub async fn download(&self, blob_id: &str, path: &AbsolutePath) -> Result<(), Error> { + /// Returns an error if the request fails or the server responds with a + /// status other than 200. + pub async fn download(&self, blob_id: &str) -> Result { let mut url = self.blob_url.clone(); url.path_segments_mut().map_err(|()| Error::InvalidEndpoint(None))?.push(blob_id); let response = self.http.get(url).send().await.map_err(Error::Network)?; - let mut response = check_status(response).await?; - let mut file = tokio::fs::File::create(path).await.map_err(Error::WriteBlob)?; - while let Some(chunk) = response.chunk().await.map_err(Error::Network)? { - file.write_all(&chunk).await.map_err(Error::WriteBlob)?; - } - file.flush().await.map_err(Error::WriteBlob)?; - Ok(()) + Ok(Download { response: check_status(response).await? }) } /// Store `value` under `key` with `POST {endpoint}/store`, uploading the @@ -467,42 +475,43 @@ mod tests { server.join().unwrap(); } + async fn read_to_end(mut download: Download) -> Result, Error> { + let mut blob = Vec::new(); + while let Some(chunk) = download.chunk().await? { + blob.extend_from_slice(&chunk); + } + Ok(blob) + } + #[tokio::test] - async fn download_writes_the_blob_to_a_file() { - let dir = tempfile::tempdir().unwrap(); - let path = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); + async fn download_streams_the_blob() { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); let client = client_for(&listener); let server = std::thread::spawn(move || serve_once(&listener, "HTTP/1.1 200 OK", b"archive bytes")); - client.download("7", &path).await.unwrap(); + let download = client.download("7").await.unwrap(); + assert_eq!(read_to_end(download).await.unwrap(), b"archive bytes"); let request = server.join().unwrap(); assert!(request.starts_with(b"GET /projects/test/blob/7 HTTP/1.1\r\n")); - assert_eq!(std::fs::read(path.as_path()).unwrap(), b"archive bytes"); } #[tokio::test] - async fn download_of_a_missing_blob_creates_no_file() { - let dir = tempfile::tempdir().unwrap(); - let path = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); + async fn download_of_a_missing_blob_fails() { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); let client = client_for(&listener); let server = std::thread::spawn(move || { serve_once(&listener, "HTTP/1.1 404 Not Found", b"Blob not found") }); - let error = client.download("7", &path).await.unwrap_err(); + let error = client.download("7").await.unwrap_err(); assert_eq!(error.to_string(), "HTTP status 404"); - assert!(!path.as_path().exists()); server.join().unwrap(); } #[tokio::test] async fn incomplete_download_is_a_network_error() { - let dir = tempfile::tempdir().unwrap(); - let path = AbsolutePathBuf::new(dir.path().join("archive.tar.zst")).unwrap(); let listener = TcpListener::bind("127.0.0.1:0").unwrap(); let client = client_for(&listener); // The connection closes before the announced length arrives. @@ -510,7 +519,8 @@ mod tests { serve_raw_once(&listener, b"HTTP/1.1 200 OK\r\ncontent-length: 100\r\n\r\npartial") }); - let error = client.download("7", &path).await.unwrap_err(); + let download = client.download("7").await.unwrap(); + let error = read_to_end(download).await.unwrap_err(); assert!(matches!(error, Error::Network(_)), "{error:?}"); server.join().unwrap(); } From 99659cab43bc9a95e2302334d0194b58781ad218 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Sun, 27 Sep 2026 21:28:08 +0800 Subject: [PATCH 6/8] feat(cache): save remote cache read failures with their causes Follow #761: a read failure's miss carries the error, and the run summary saves it as a `SavedError`. The miss reason is the error's message, such as `remote cache fetch failed`. `--verbose` and `--last-details` show each cause, such as the network error, on its own line below it. Co-authored-by: Claude Opus 5.5 --- crates/vt/src/session/cache/display.rs | 4 +- crates/vt/src/session/cache/mod.rs | 12 +++-- crates/vt/src/session/cache/remote.rs | 12 ++--- crates/vt/src/session/reporter/summary.rs | 51 ++++++++++++------- .../remote_cache/snapshots/corrupt_archive.md | 1 + .../snapshots/invalid_endpoint.md | 6 ++- .../snapshots/read_invalid_endpoint.md | 6 ++- .../snapshots/read_unreachable_endpoint.md | 9 +++- .../snapshots/unreachable_endpoint.md | 9 +++- 9 files changed, 71 insertions(+), 39 deletions(-) diff --git a/crates/vt/src/session/cache/display.rs b/crates/vt/src/session/cache/display.rs index 31e00752e..b7ba4eb18 100644 --- a/crates/vt/src/session/cache/display.rs +++ b/crates/vt/src/session/cache/display.rs @@ -193,8 +193,8 @@ pub fn format_cache_status_inline(cache_status: &CacheStatus) -> Option { }; Some(vt_str::format!("○ cache miss: {reason}, executing")) } - CacheStatus::Miss(CacheMiss::RemoteReadFailed(reason)) => { - Some(vt_str::format!("○ cache miss: {reason}, executing")) + CacheStatus::Miss(CacheMiss::RemoteReadFailed(error)) => { + Some(vt_str::format!("○ cache miss: {error}, executing")) } CacheStatus::Disabled(_) => Some(Str::from("⊘ cache disabled")), } diff --git a/crates/vt/src/session/cache/mod.rs b/crates/vt/src/session/cache/mod.rs index 1fba45d88..f0a05507c 100644 --- a/crates/vt/src/session/cache/mod.rs +++ b/crates/vt/src/session/cache/mod.rs @@ -30,7 +30,7 @@ use wincode::{ io::{Reader, Writer}, }; -use self::remote::{RemoteClients, Restore, UploadError}; +use self::remote::{ReadError, RemoteClients, Restore, UploadError}; use super::execute::{ fingerprint::{PostRunFingerprint, TrackedEnvQuery}, pipe::StdOutput, @@ -148,13 +148,17 @@ pub struct ExecutionCache { remote_clients: RemoteClients, } -#[derive(Debug, Clone, Serialize)] +#[derive(Debug, Clone)] +#[expect( + clippy::large_enum_variant, + reason = "FingerprintMismatch contains SpawnFingerprint which is intentionally large; boxing would add unnecessary indirection for a short-lived enum" +)] pub enum CacheMiss { NotFound, FingerprintMismatch(FingerprintMismatch), /// Reading the remote cache failed, and the local cache has no entry for - /// the task. The message names the cause. - RemoteReadFailed(Str), + /// the task. + RemoteReadFailed(Arc), } #[derive(Debug, Clone, Copy, Serialize, Deserialize)] diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index 16dbc60ae..db85cc5e9 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -49,9 +49,9 @@ pub enum UploadError { /// it's the same on every platform. #[derive(Debug, thiserror::Error)] pub enum ReadError { - #[error("remote cache fetch failed ({0})")] + #[error("remote cache fetch failed")] Fetch(#[source] vt_remote_cache::Error), - #[error("remote cache download failed ({0})")] + #[error("remote cache download failed")] Download(#[source] vt_remote_cache::Error), #[error("remote cache value is corrupt")] CorruptValue(#[source] wincode::error::ReadError), @@ -68,10 +68,8 @@ pub enum ReadError { } impl ReadError { - /// The miss for this failure. The full error is only logged. pub(super) fn into_miss(self) -> CacheMiss { - tracing::debug!(err = ?self, "remote cache read failed"); - CacheMiss::RemoteReadFailed(vt_str::format!("{self}")) + CacheMiss::RemoteReadFailed(Arc::new(self)) } } @@ -357,7 +355,7 @@ mod tests { fn read_failure(miss: CacheMiss) -> Str { match miss { - CacheMiss::RemoteReadFailed(reason) => reason, + CacheMiss::RemoteReadFailed(error) => vt_str::format!("{error}"), miss => panic!("expected a read failure, got {miss:?}"), } } @@ -439,7 +437,7 @@ mod tests { let key = cache_key(ResolvedGlobConfig::default_auto()); let fetched = Err(ReadError::Fetch(vt_remote_cache::Error::InvalidEndpoint(None))); let miss = resolve(fetched, &key, not_validated).unwrap_err(); - assert_eq!(read_failure(miss), "remote cache fetch failed (invalid endpoint)"); + assert_eq!(read_failure(miss), "remote cache fetch failed"); } #[test] diff --git a/crates/vt/src/session/reporter/summary.rs b/crates/vt/src/session/reporter/summary.rs index 2edf2e80f..12459cc20 100644 --- a/crates/vt/src/session/reporter/summary.rs +++ b/crates/vt/src/session/reporter/summary.rs @@ -151,8 +151,7 @@ pub enum SavedCacheMissReason { /// between runs. Carries the first differing entry. TrackedEnvQueryChanged { query: TrackedEnvQuery, mismatch: EnvMismatch }, /// Reading the remote cache failed, and the local cache had no entry. - /// Carries the failure's message. - RemoteReadFailed(Str), + RemoteReadFailed(SavedError), } /// An error's message and the messages of its causes, outermost first. @@ -277,7 +276,9 @@ impl SavedCacheMissReason { } } }, - CacheMiss::RemoteReadFailed(reason) => Self::RemoteReadFailed(reason.clone()), + CacheMiss::RemoteReadFailed(error) => { + Self::RemoteReadFailed(SavedError::new(error.as_ref())) + } } } } @@ -514,21 +515,22 @@ impl TaskResult { } } - /// Format the cache status detail line for the full summary. The caller - /// shows [`Self::ipc_server_error`] instead, if there is one. + /// Format the cache status detail line for the full summary, with the + /// causes to show below it. Only a remote read failure has causes. The + /// caller shows [`Self::ipc_server_error`] instead, if there is one. /// /// Examples: /// - "→ Cache hit - output replayed - 102.96ms saved" /// - "→ Cache miss: no previous cache entry found" /// - "→ Cache disabled in task configuration" - fn format_cache_detail(&self) -> Str { + fn format_cache_detail(&self) -> (Str, &[Str]) { // Tool-reported cache disable — the tool said it shouldn't be cached. if let Self::Spawned { outcome: SpawnOutcome::Success { tool_disabled_cache: true, .. }, .. } = self { - return Str::from("→ Not cached: the task opted out of caching"); + return (Str::from("→ Not cached: the task opted out of caching"), &[]); } // Check for input modification next — it overrides the cache miss reason @@ -537,7 +539,7 @@ impl TaskResult { .. } = self { - return vt_str::format!("→ Not cached: read and wrote '{path}'"); + return (vt_str::format!("→ Not cached: read and wrote '{path}'"), &[]); } // Tracking came up short, so the inferred inputs and outputs would // have been a subset of what the task touched. @@ -546,8 +548,11 @@ impl TaskResult { .. } = self { - return Str::from( - "→ Not cached: this task used more files than automatic tracking can record. Configure `input` and `output` manually to enable caching.", + return ( + Str::from( + "→ Not cached: this task used more files than automatic tracking can record. Configure `input` and `output` manually to enable caching.", + ), + &[], ); } // fspy-unsupported-on-this-OS message — same overrides precedence as above @@ -555,12 +560,15 @@ impl TaskResult { outcome: SpawnOutcome::Success { fspy_unsupported: true, .. }, .. } = self { - return Str::from( - "→ Not cached: `input` auto-inference isn't supported on this OS. Configure `input` manually to enable caching.", + return ( + Str::from( + "→ Not cached: `input` auto-inference isn't supported on this OS. Configure `input` manually to enable caching.", + ), + &[], ); } - match self { + let detail = match self { Self::CacheHit { saved_duration_ms } => { let d = Duration::from_millis(*saved_duration_ms); let formatted_duration = format_summary_duration(d); @@ -594,12 +602,13 @@ impl TaskResult { | SavedCacheMissReason::TrackedEnvQueryChanged { mismatch, .. } => { vt_str::format!("→ Cache miss: {mismatch}") } - SavedCacheMissReason::RemoteReadFailed(reason) => { - vt_str::format!("→ Cache miss: {reason}") + SavedCacheMissReason::RemoteReadFailed(error) => { + return (vt_str::format!("→ Cache miss: {}", error.message), &error.causes); } }, }, - } + }; + (detail, &[]) } /// The [`Style`] for the cache detail line. @@ -801,8 +810,9 @@ pub fn format_full_summary(summary: &LastRunSummary) -> Vec { detail_style, ); } else { - let cache_detail = task.result.format_cache_detail(); + let (cache_detail, causes) = task.result.format_cache_detail(); let _ = writeln!(buf, " {}", cache_detail.style(detail_style)); + write_causes(&mut buf, causes, detail_style); } if let Some(error) = task.result.upload_error() { @@ -848,7 +858,12 @@ pub fn format_full_summary(summary: &LastRunSummary) -> Vec { /// its causes on its own line below. fn write_error_lines(buf: &mut Vec, label: impl Display, error: &SavedError, style: Style) { let _ = writeln!(buf, " {label} {}", error.message.style(style)); - for cause in &error.causes { + write_causes(buf, &error.causes, style); +} + +/// Write each cause on its own line, below a task detail line. +fn write_causes(buf: &mut Vec, causes: &[Str], style: Style) { + for cause in causes { let _ = writeln!(buf, " {}", vt_str::format!("↳ {cause}").style(style)); } } diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md index d091f409a..94076a05c 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md @@ -49,6 +49,7 @@ Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ → Cache miss: downloaded archive is corrupt + ↳ Unknown frame descriptor ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ ``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md index 01d3f1414..067a8e20d 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/invalid_endpoint.md @@ -5,7 +5,7 @@ The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds. ``` -$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (invalid endpoint), executing +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed, executing --- vt run: remote-cache#build not uploaded to the remote cache: invalid endpoint. (Run `vt run --last-details` for full details) @@ -27,7 +27,9 @@ Performance: 0% cache hit rate Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ - → Cache miss: remote cache fetch failed (invalid endpoint) + → Cache miss: remote cache fetch failed + ↳ invalid endpoint + ↳ relative URL without a base ⚠ Not uploaded to the remote cache: invalid endpoint ↳ relative URL without a base ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md index a4bc8e07b..15fd38748 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_invalid_endpoint.md @@ -5,7 +5,7 @@ The failed fetch is the miss reason. Read failures aren't warnings. ``` -$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (invalid endpoint), executing +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed, executing ``` ## `vt run --last-details` @@ -22,6 +22,8 @@ Performance: 0% cache hit rate Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ - → Cache miss: remote cache fetch failed (invalid endpoint) + → Cache miss: remote cache fetch failed + ↳ invalid endpoint + ↳ relative URL without a base ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ ``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md index a941b5728..f61035964 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/read_unreachable_endpoint.md @@ -5,7 +5,7 @@ Nothing can listen on port 0. The failed fetch is the miss reason. Read failures aren't warnings. ``` -$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (network error), executing +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed, executing ``` ## `vt run --last-details` @@ -24,6 +24,11 @@ Performance: 0% cache hit rate Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ - → Cache miss: remote cache fetch failed (network error) + → Cache miss: remote cache fetch failed + ↳ network error + ↳ error sending request for url (http://127.0.0.1:0/projects/test/fetch) + ↳ client error (Connect) + ↳ tcp connect error + ↳ ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ ``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md index 047d22762..daa2d7e5d 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/unreachable_endpoint.md @@ -5,7 +5,7 @@ Nothing can listen on port 0. The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds. ``` -$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed (network error), executing +$ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed, executing --- vt run: remote-cache#build not uploaded to the remote cache: network error. (Run `vt run --last-details` for full details) @@ -27,7 +27,12 @@ Performance: 0% cache hit rate Task Details: ──────────────────────────────────────────────── [1] remote-cache#build: $ vtt write-file dist/output.txt built ✓ - → Cache miss: remote cache fetch failed (network error) + → Cache miss: remote cache fetch failed + ↳ network error + ↳ error sending request for url (http://127.0.0.1:0/projects/test/fetch) + ↳ client error (Connect) + ↳ tcp connect error + ↳ ⚠ Not uploaded to the remote cache: network error ↳ error sending request for url (http://127.0.0.1:0/projects/test/store) ↳ client error (Connect) From f5042c26fcd3a52153e4c48d683df5cb026e5f75 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Sun, 27 Sep 2026 21:38:21 +0800 Subject: [PATCH 7/8] refactor(cache): match every remote cache access mode Fetching and uploading each name every access mode, so adding one fails to compile until both decide what it does. Co-authored-by: Claude Opus 5.5 --- crates/vt/src/session/cache/mod.rs | 21 +++++++++++++++------ 1 file changed, 15 insertions(+), 6 deletions(-) diff --git a/crates/vt/src/session/cache/mod.rs b/crates/vt/src/session/cache/mod.rs index f0a05507c..7e1c887c7 100644 --- a/crates/vt/src/session/cache/mod.rs +++ b/crates/vt/src/session/cache/mod.rs @@ -353,8 +353,16 @@ impl ExecutionCache { Ok(cache_value) => return Ok(Ok(cache_value)), Err(miss) => miss, }; - let Some(ResolvedRemoteCacheConfig { url, .. }) = &cache_metadata.remote_cache else { - return Ok(Err(local_miss)); + #[expect( + clippy::manual_let_else, + reason = "naming every access mode makes adding one a compile error here" + )] + let url = match &cache_metadata.remote_cache { + Some(ResolvedRemoteCacheConfig { + access: RemoteCacheAccess::Read | RemoteCacheAccess::ReadWrite, + url, + }) => url, + None => return Ok(Err(local_miss)), }; let remote_miss = match self .try_hit_remote( @@ -504,10 +512,11 @@ impl ExecutionCache { self.record(&cache_key, execution_cache_key, &cache_value, cache_dir).await?; - let Some(ResolvedRemoteCacheConfig { access: RemoteCacheAccess::ReadWrite, url }) = - &cache_metadata.remote_cache - else { - return Ok(Ok(())); + let url = match &cache_metadata.remote_cache { + Some(ResolvedRemoteCacheConfig { access: RemoteCacheAccess::ReadWrite, url }) => url, + Some(ResolvedRemoteCacheConfig { access: RemoteCacheAccess::Read, .. }) | None => { + return Ok(Ok(())); + } }; let upload = self .remote_clients From aefe44a3a6ded332d532a52776270835c4ddd8d4 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Sun, 27 Sep 2026 21:49:51 +0800 Subject: [PATCH 8/8] refactor(cache): download remote cache archives through the check The archive is downloaded to a `.tmp` file and renamed once the check passes. The chunks flow one way: from the network to the check on a blocking thread, which writes each one to the file as it takes it. The file is closed before it's renamed or removed, even after a network error. Co-authored-by: Claude Opus 5.5 --- crates/vt/Cargo.toml | 1 - crates/vt/src/session/cache/remote.rs | 85 ++++++++++++------- .../fixtures/remote_cache/snapshots.toml | 4 +- .../remote_cache/snapshots/corrupt_archive.md | 5 +- 4 files changed, 60 insertions(+), 35 deletions(-) diff --git a/crates/vt/Cargo.toml b/crates/vt/Cargo.toml index 57bcabc7a..3759db96f 100644 --- a/crates/vt/Cargo.toml +++ b/crates/vt/Cargo.toml @@ -35,7 +35,6 @@ thiserror = { workspace = true } tar = { workspace = true } tokio = { workspace = true, features = [ "rt-multi-thread", - "fs", "io-std", "io-util", "macros", diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index db85cc5e9..b8d545a34 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -14,16 +14,17 @@ //! their keys differ. use std::{ - io, + fs::File, + io::{self, Write as _}, sync::{Arc, Mutex, PoisonError}, }; use bytes::Bytes; use rustc_hash::FxHashMap; -use tokio::{io::AsyncWriteExt as _, sync::mpsc}; +use tokio::sync::mpsc; use vt_path::AbsolutePath; use vt_plan::cache_metadata::ExecutionCacheKey; -use vt_remote_cache::{Client, Fetched}; +use vt_remote_cache::{Client, Download, Fetched}; use vt_str::Str; use wincode::{ SchemaWrite, @@ -146,8 +147,9 @@ impl RemoteClients { } /// Download the blob `blob_id` into `cache_dir`, checking that it decodes - /// as an output archive as it arrives. Returns the archive's file name. If - /// the download or the check fails, the file is removed. + /// as an output archive as it arrives. It's downloaded to a `.tmp` file, + /// which is renamed once the check passes and removed otherwise. Returns + /// the archive's file name. pub(super) async fn download_archive( &self, endpoint: &Arc, @@ -157,10 +159,14 @@ impl RemoteClients { let client = self.client(endpoint).map_err(ReadError::Download)?; let archive_name = vt_str::format!("{}.tar.zst", uuid::Uuid::new_v4()); let archive_path = cache_dir.join(archive_name.as_str()); - let result = download_checked(&client, blob_id, &archive_path).await; + let temp_path = cache_dir.join(vt_str::format!("{archive_name}.tmp").as_str()); + let result = download_checked(&client, blob_id, &temp_path).await.and_then(|()| { + std::fs::rename(temp_path.as_path(), archive_path.as_path()) + .map_err(ReadError::WriteArchive) + }); if result.is_err() { // Best-effort cleanup: the file may not have been created. - let _ = std::fs::remove_file(archive_path.as_path()); + let _ = std::fs::remove_file(temp_path.as_path()); } result.map(|()| archive_name) } @@ -188,50 +194,71 @@ impl RemoteClients { /// Chunks buffered between the download and the archive check. const CHECK_BUFFER_CHUNKS: usize = 16; -/// Download the blob `blob_id` to the file at `path`. Each chunk is written to -/// the file and passed to the archive check, which runs on a blocking thread. +/// Download the blob `blob_id` to the file at `path`. The chunks flow one way: +/// from the network to the archive check on a blocking thread, which writes +/// each one to the file as it takes it. async fn download_checked( client: &Client, blob_id: &str, path: &AbsolutePath, ) -> Result<(), ReadError> { - let mut download = client.download(blob_id).await.map_err(ReadError::Download)?; - let mut file = - tokio::fs::File::create(path.as_path()).await.map_err(ReadError::WriteArchive)?; + let download = client.download(blob_id).await.map_err(ReadError::Download)?; + let file = File::create(path.as_path()).map_err(ReadError::WriteArchive)?; let (sender, receiver) = mpsc::channel(CHECK_BUFFER_CHUNKS); let check = tokio::task::spawn_blocking(move || { - archive::check_output_archive(ChunkReader { receiver, chunk: Bytes::new() }) + let mut reader = DownloadReader { receiver, chunk: Bytes::new(), file, write_error: None }; + let checked = archive::check_output_archive(&mut reader); + if let Some(err) = reader.write_error { + return Err(ReadError::WriteArchive(err)); + } + checked.map_err(ReadError::CorruptArchive) }); - while let Some(chunk) = download.chunk().await.map_err(ReadError::Download)? { - file.write_all(&chunk).await.map_err(ReadError::WriteArchive)?; - // The check stops reading when it fails. + let sent = send_chunks(download, sender).await; + // Wait for the check even after a network error, so the file is closed + // before the caller removes it. + let checked = check.await.unwrap_or_else(|err| Err(ReadError::CorruptArchive(err.into()))); + sent.map_err(ReadError::Download)?; + checked +} + +/// Send the blob's chunks through `sender` until the blob ends or the receiver +/// is dropped. +async fn send_chunks( + mut download: Download, + sender: mpsc::Sender, +) -> Result<(), vt_remote_cache::Error> { + while let Some(chunk) = download.chunk().await? { + // The check drops the receiver when it fails. if sender.send(chunk).await.is_err() { break; } } - drop(sender); - check - .await - .map_err(io::Error::from) - .and_then(|checked| checked) - .map_err(ReadError::CorruptArchive)?; - file.flush().await.map_err(ReadError::WriteArchive) + Ok(()) } /// Reads the chunks sent through a channel, in order, until the sender is -/// dropped. Reading blocks, so it's only for blocking threads. -struct ChunkReader { +/// dropped, and writes each one to `file` as it takes it. Reading blocks, so +/// it's only for blocking threads. +struct DownloadReader { receiver: mpsc::Receiver, chunk: Bytes, + file: File, + /// Why writing to `file` failed. Reading fails after it, and the check's + /// error is then just a consequence. + write_error: Option, } -impl io::Read for ChunkReader { +impl io::Read for DownloadReader { fn read(&mut self, buf: &mut [u8]) -> io::Result { while self.chunk.is_empty() { - match self.receiver.blocking_recv() { - Some(chunk) => self.chunk = chunk, - None => return Ok(0), + let Some(chunk) = self.receiver.blocking_recv() else { + return Ok(0); + }; + if let Err(err) = self.file.write_all(&chunk) { + self.write_error = Some(err); + return Err(io::ErrorKind::Other.into()); } + self.chunk = chunk; } let len = buf.len().min(self.chunk.len()); buf[..len].copy_from_slice(&self.chunk.split_to(len)); diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml index 44b1c1182..93e714b8c 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots.toml @@ -207,9 +207,9 @@ steps = [ "list-dir", "node_modules/.vite/task-cache", "--ext", - ".tar.zst", + ".tmp", "--recursive", - ], comment = "Only the rerun's archive is on disk. The corrupt download was removed." }, + ], comment = "The corrupt download was removed." }, ] [[e2e]] diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md index 94076a05c..98ad169ed 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/corrupt_archive.md @@ -53,10 +53,9 @@ Task Details: ━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ ``` -## `vtt list-dir node_modules/.vite/task-cache --ext .tar.zst --recursive` +## `vtt list-dir node_modules/.vite/task-cache --ext .tmp --recursive` -Only the rerun's archive is on disk. The corrupt download was removed. +The corrupt download was removed. ``` -.tar.zst ```