Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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)).
Expand Down
2 changes: 2 additions & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
1 change: 1 addition & 0 deletions crates/vt/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down
45 changes: 44 additions & 1 deletion crates/vt/src/session/cache/archive.rs
Original file line number Diff line number Diff line change
@@ -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};

Expand Down Expand Up @@ -61,3 +64,43 @@ pub fn extract_output_archive(
archive.unpack(workspace_root.as_path())?;
Ok(())
}

/// Read a tar.zst archive from `reader` to the end without writing any files,
/// to check that it decodes.
///
/// # Errors
///
/// 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())?;
}
// 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();
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"] {
assert!(check_output_archive(corrupt).is_err());
}
}
}
3 changes: 3 additions & 0 deletions crates/vt/src/session/cache/display.rs
Original file line number Diff line number Diff line change
Expand Up @@ -193,6 +193,9 @@ pub fn format_cache_status_inline(cache_status: &CacheStatus) -> Option<Str> {
};
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")),
}
}
Expand Down
171 changes: 142 additions & 29 deletions crates/vt/src/session/cache/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ use wincode::{
io::{Reader, Writer},
};

use self::remote::{RemoteClients, UploadError};
use self::remote::{ReadError, RemoteClients, Restore, UploadError};
use super::execute::{
fingerprint::{PostRunFingerprint, TrackedEnvQuery},
pipe::StdOutput,
Expand Down Expand Up @@ -148,14 +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.
RemoteReadFailed(Arc<ReadError>),
}

#[derive(Debug, Clone, Copy, Serialize, Deserialize)]
Expand Down Expand Up @@ -327,27 +330,81 @@ 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<RelativePathBuf, u64>,
workspace_root: &AbsolutePath,
cache_dir: &AbsolutePath,
) -> anyhow::Result<Result<CacheEntryValue, CacheMiss>> {
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,
};
#[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(
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<RelativePathBuf, u64>,
workspace_root: &AbsolutePath,
) -> anyhow::Result<Result<CacheEntryValue, CacheMiss>> {
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(mismatch) =
cache_value.validate(cache_metadata, globbed_inputs, workspace_root)?
{
if let Some(cache_value) = self.get_by_cache_key(cache_key).await? {
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,
// 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));
}

Expand All @@ -358,18 +415,86 @@ 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. 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<str>,
cache_metadata: &CacheMetadata,
cache_key: &CacheEntryKey,
globbed_inputs: &BTreeMap<RelativePathBuf, u64>,
workspace_root: &AbsolutePath,
cache_dir: &AbsolutePath,
) -> anyhow::Result<Result<CacheEntryValue, CacheMiss>> {
let fetched = self
.remote_clients
.fetch(endpoint, cache_key, &cache_metadata.execution_cache_key)
.await;
let validate = |cache_value: &CacheEntryValue| {
cache_value.validate(&cache_metadata.unfiltered_envs, globbed_inputs, workspace_root)
};
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(err.into_miss())),
}
}
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
Expand All @@ -385,25 +510,13 @@ 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.record(&cache_key, execution_cache_key, &cache_value, cache_dir).await?;

self.upsert_cache_entry(&cache_key, &cache_value).await?;
self.upsert_task_fingerprint(execution_cache_key, &cache_key).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
Expand Down
Loading
Loading