Skip to content
Open
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 @@ -4,7 +4,7 @@
- **Changed** The run summary now says a task that wrote a file it also read was `not cached because it modified its inputs`, and the statistics in `vp run --verbose` and `vp run --last-details` use the singular for a count of one, e.g. `1 task • 1 cache miss` ([#783](https://github.com/voidzero-dev/vite-task/pull/783)).
- **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. 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. The task output and the run summary show which hits came from the remote cache. 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. Ctrl-C, or a failing task, stops remote cache requests right away, and a task still being looked up doesn't start. Tasks can opt out with `cache: { remote: false }`. Requests use the proxy environment variables or, on macOS and Windows, the system proxy settings ([#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), [#757](https://github.com/voidzero-dev/vite-task/pull/757), [#764](https://github.com/voidzero-dev/vite-task/pull/764), [#770](https://github.com/voidzero-dev/vite-task/pull/770), [#771](https://github.com/voidzero-dev/vite-task/pull/771), [#772](https://github.com/voidzero-dev/vite-task/pull/772), [#786](https://github.com/voidzero-dev/vite-task/pull/786)).
- **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. The task output and the run summary show which hits came from the remote cache. 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. Uploads run in the background, so tasks that depend on an uploading task don't wait for it; once all tasks are done, `vp run` waits for the uploads still running and says so. A failed upload doesn't fail the task; the run summary shows a warning instead. Ctrl-C, or a failing task, stops remote cache lookups right away, and a task still being looked up doesn't start. Ctrl-C also cancels the uploads, but a failing task doesn't. Tasks can opt out with `cache: { remote: false }`. Requests to an HTTPS endpoint use HTTP/2 if the server supports it. Requests use the proxy environment variables or, on macOS and Windows, the system proxy settings ([#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), [#757](https://github.com/voidzero-dev/vite-task/pull/757), [#764](https://github.com/voidzero-dev/vite-task/pull/764), [#770](https://github.com/voidzero-dev/vite-task/pull/770), [#771](https://github.com/voidzero-dev/vite-task/pull/771), [#772](https://github.com/voidzero-dev/vite-task/pull/772), [#786](https://github.com/voidzero-dev/vite-task/pull/786), [#787](https://github.com/voidzero-dev/vite-task/pull/787)).
- **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
22 changes: 22 additions & 0 deletions Cargo.lock

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

2 changes: 1 addition & 1 deletion crates/vt/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,7 @@ tokio = { workspace = true, features = [
"process",
"sync",
] }
tokio-util = { workspace = true }
tokio-util = { workspace = true, features = ["rt"] }
tracing = { workspace = true }
twox-hash = { workspace = true }
materialized_artifact = { workspace = true }
Expand Down
63 changes: 48 additions & 15 deletions crates/vt/src/session/cache/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,14 @@ pub mod display;
pub mod remote;
mod validation;

use std::{collections::BTreeMap, fmt::Display, fs::File, io::Write, sync::Arc, time::Duration};
use std::{
collections::BTreeMap,
fmt::Display,
fs::File,
io::Write,
sync::{Arc, OnceLock},
time::Duration,
};

// Re-export display functions for convenience
pub use display::format_cache_status_inline;
Expand All @@ -30,7 +37,7 @@ use wincode::{
error::{ReadResult, WriteResult},
};

use self::remote::{ReadError, RemoteClients, Restore, UploadError};
use self::remote::{ReadError, RemoteClients, RemoteUploads, Restore, UploadError};
use super::execute::{
fingerprint::{PostRunFingerprint, TrackedEnvQuery},
pipe::StdOutput,
Expand Down Expand Up @@ -112,6 +119,7 @@ pub struct CacheEntryValue {
pub struct ExecutionCache {
conn: Mutex<Connection>,
remote_clients: RemoteClients,
uploads: RemoteUploads,
}

/// A cache hit: the entry to replay, and the cache it came from.
Expand Down Expand Up @@ -303,7 +311,11 @@ impl ExecutionCache {
CREATE TABLE IF NOT EXISTS task_fingerprints (key BLOB PRIMARY KEY, value BLOB);",
)?;
// Lock is released when lock_file is dropped
Ok(Self { conn: Mutex::new(conn), remote_clients: RemoteClients::default() })
Ok(Self {
conn: Mutex::new(conn),
remote_clients: RemoteClients::default(),
uploads: RemoteUploads::default(),
})
}

#[tracing::instrument]
Expand Down Expand Up @@ -487,36 +499,57 @@ impl ExecutionCache {
/// as [`Self::record`] does.
///
/// In `read-write` remote mode, the entry is then uploaded to the remote
/// cache, until `cancel_token` is cancelled. Returns `Ok(Err(_))` if the
/// local update succeeded but the upload failed.
/// cache in the background, and the task doesn't wait for it. Returns
/// where the upload's error is set if it fails: right away if the upload
/// can't start, or later by the background upload. Read it after
/// [`Self::wait_for_uploads`] returns.
#[tracing::instrument(level = "debug", skip_all)]
pub async fn update(
&self,
cache_metadata: &CacheMetadata,
cache_value: CacheEntryValue,
cache_dir: &AbsolutePath,
cancel_token: &CancellationToken,
) -> anyhow::Result<Result<(), UploadError>> {
) -> anyhow::Result<Arc<OnceLock<UploadError>>> {
let execution_cache_key = &cache_metadata.execution_cache_key;

let cache_key = CacheEntryKey::from_metadata(cache_metadata);

self.record(&cache_key, execution_cache_key, &cache_value, cache_dir).await?;

let upload_error = Arc::new(OnceLock::new());
let url = match &cache_metadata.remote_cache {
Some(ResolvedRemoteCacheConfig { access: RemoteCacheAccess::ReadWrite, url }) => url,
Some(ResolvedRemoteCacheConfig { access: RemoteCacheAccess::Read, .. }) | None => {
return Ok(Ok(()));
return Ok(upload_error);
}
};
let upload = self
.remote_clients
.upload(url, &cache_key, execution_cache_key, &cache_value, cache_dir, cancel_token)
.await;
if let Err(err) = &upload {
tracing::debug!(?err, "remote cache upload failed");
match self.remote_clients.prepare_upload(
url,
&cache_key,
execution_cache_key,
&cache_value,
cache_dir,
) {
Ok(upload) => self.uploads.spawn(upload, Arc::clone(&upload_error)),
Err(err) => {
tracing::debug!(?err, "remote cache upload failed");
let _ = upload_error.set(err);
}
}
Ok(upload)
Ok(upload_error)
}

/// The number of uploads to the remote cache still running.
pub fn pending_uploads(&self) -> usize {
self.uploads.pending()
}

/// Wait for the uploads to the remote cache to finish. If
/// `interrupt_token` is cancelled first, or already was, the uploads are
/// cancelled instead, each with [`UploadError::Interrupted`] as its error.
/// The entries stay in the local cache.
pub async fn wait_for_uploads(&self, interrupt_token: &CancellationToken) {
self.uploads.wait(interrupt_token).await;
}

/// Restore the output files of `hit` into `workspace_root`.
Expand Down
149 changes: 124 additions & 25 deletions crates/vt/src/session/cache/remote.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,13 +16,13 @@
use std::{
fs::File,
io::{self, Write as _},
sync::{Arc, Mutex, PoisonError},
sync::{Arc, Mutex, OnceLock, PoisonError},
};

use bytes::Bytes;
use rustc_hash::FxHashMap;
use tokio::sync::mpsc;
use tokio_util::sync::CancellationToken;
use tokio_util::{sync::CancellationToken, task::TaskTracker};
use vt_path::AbsolutePath;
use vt_plan::cache_metadata::ExecutionCacheKey;
use vt_remote_cache::{Client, Download, Fetched};
Expand All @@ -44,10 +44,9 @@ pub enum UploadError {
Remote(#[from] vt_remote_cache::Error),
#[error("failed to encode the cache entry")]
Encode(#[from] WriteError),
/// The run was cancelled, by Ctrl-C or fast-fail, before the upload
/// finished.
#[error("cancelled")]
Cancelled,
/// Ctrl-C cancelled the upload before it finished.
#[error("interrupted")]
Interrupted,
}

/// Why no entry could be read from the remote cache. It's a cache miss, and
Expand Down Expand Up @@ -196,25 +195,80 @@ impl RemoteClients {
result.map(|()| archive_name)
}

/// Upload an entry that was just recorded locally, along with its output
/// archive in `cache_dir`. Stops when `cancel_token` is cancelled.
pub(super) async fn upload(
/// Prepare to upload an entry that was just recorded locally, along with
/// its output archive in `cache_dir`: get the endpoint's client and encode
/// the entry. Nothing is sent until the returned future is polled, so an
/// invalid endpoint or an entry that doesn't encode fails here. The future
/// owns everything the upload needs, so it can run in a spawned task.
pub(super) fn prepare_upload(
&self,
endpoint: &Arc<str>,
cache_key: &CacheEntryKey,
execution_cache_key: &ExecutionCacheKey,
cache_value: &CacheEntryValue,
cache_dir: &AbsolutePath,
cancel_token: &CancellationToken,
) -> Result<(), UploadError> {
) -> Result<impl Future<Output = Result<(), UploadError>> + Send + use<>, UploadError> {
let client = self.client(endpoint)?;
let key = encode_key(cache_key)?;
let secondary_key = encode_key(execution_cache_key)?;
let value = serialize_cache(cache_value)?;
let archive = cache_value.output_archive.as_ref().map(|name| cache_dir.join(name.as_str()));
let store = client.store(&key, &secondary_key, &value, archive.as_deref());
cancel_token.run_until_cancelled(store).await.ok_or(UploadError::Cancelled)??;
Ok(())
Ok(async move {
client.store(&key, &secondary_key, &value, archive.as_deref()).await?;
Ok(())
})
}
}

/// Uploads running in the background. Each keeps running after its task
/// finishes, until [`Self::wait`] waits for all of them.
#[derive(Debug, Default)]
pub(super) struct RemoteUploads {
tracker: TaskTracker,
/// Cancelled when the wait is interrupted, which stops the uploads. It
/// stays cancelled, since the run ends after Ctrl-C.
cancel: CancellationToken,
}

impl RemoteUploads {
/// Send `upload` in the background. If it fails or is cancelled, the
/// error is set in `error`.
pub(super) fn spawn(
&self,
upload: impl Future<Output = Result<(), UploadError>> + Send + 'static,
error: Arc<OnceLock<UploadError>>,
) {
let cancel = self.cancel.clone();
self.tracker.spawn(async move {
// A cancelled upload sets its error itself, instead of being
// aborted, so every cancelled upload has one.
let result =
cancel.run_until_cancelled(upload).await.unwrap_or(Err(UploadError::Interrupted));
if let Err(err) = result {
tracing::debug!(?err, "remote cache upload failed");
let _ = error.set(err);
}
});
}

/// The number of uploads still running.
pub(super) fn pending(&self) -> usize {
self.tracker.len()
}

/// Wait for all uploads to finish. If `interrupt_token` is cancelled
/// first, or already was, cancel them and wait for them to stop.
pub(super) async fn wait(&self, interrupt_token: &CancellationToken) {
self.tracker.close();
tokio::select! {
biased;
() = self.tracker.wait() => {}
() = interrupt_token.cancelled() => {
self.cancel.cancel();
self.tracker.wait().await;
}
}
self.tracker.reopen();
}
}

Expand Down Expand Up @@ -633,22 +687,67 @@ mod tests {
assert_eq!(std::fs::read_dir(dir.path()).unwrap().count(), 0);
}

#[tokio::test]
async fn cancelling_stops_an_upload() {
let (endpoint, requested) = serve_stalled(b"");
let cancel_token = CancellationToken::new();
/// Prepare to upload an entry without an output archive to `endpoint`.
fn prepare_upload(
clients: &RemoteClients,
endpoint: &Arc<str>,
) -> Result<impl Future<Output = Result<(), UploadError>> + use<>, UploadError> {
let key = cache_key(ResolvedGlobConfig::default_auto());
let execution_key = ExecutionCacheKey::ExecAPI(Arc::from([]));
let value = CacheEntryValue { output_archive: None, ..cache_value() };
let cache_dir = vt_path::current_dir().unwrap();
clients.prepare_upload(endpoint, &key, &execution_key, &value, &cache_dir)
}

fn upload_error(error: &OnceLock<UploadError>) -> Option<Str> {
error.get().map(|error| vt_str::format!("{error}"))
}

#[test]
fn upload_to_an_invalid_endpoint_fails_before_it_starts() {
let endpoint = Arc::from("cache.example/projects/test");
let Err(error) = prepare_upload(&RemoteClients::default(), &endpoint) else {
panic!("an invalid endpoint should fail");
};
assert!(
matches!(error, UploadError::Remote(vt_remote_cache::Error::InvalidEndpoint(_))),
"{error:?}"
);
}

#[tokio::test]
async fn failed_upload_sets_its_error() {
let (error_status, _) =
serve_stalled(b"HTTP/1.1 500 Internal Server Error\r\ncontent-length: 0\r\n\r\n");
// Nothing can listen on port 0.
let unreachable = Arc::from("http://127.0.0.1:0/projects/test");
let clients = RemoteClients::default();
let upload =
clients.upload(&endpoint, &key, &execution_key, &value, &cache_dir, &cancel_token);
let (uploaded, ()) = tokio::join!(upload, async {
requested.await.unwrap();
cancel_token.cancel();
});
assert!(matches!(uploaded, Err(UploadError::Cancelled)), "{uploaded:?}");
let uploads = RemoteUploads::default();

for (endpoint, message) in
[(error_status, "HTTP status 500"), (unreachable, "network error")]
{
let error = Arc::new(OnceLock::new());
uploads.spawn(prepare_upload(&clients, &endpoint).unwrap(), Arc::clone(&error));
uploads.wait(&CancellationToken::new()).await;
assert_eq!(uploads.pending(), 0);
assert_eq!(upload_error(&error).as_deref(), Some(message));
}
}

#[tokio::test]
async fn interrupting_the_wait_cancels_the_uploads() {
let (endpoint, requested) = serve_stalled(b"");
let clients = RemoteClients::default();
let uploads = RemoteUploads::default();
let error = Arc::new(OnceLock::new());
uploads.spawn(prepare_upload(&clients, &endpoint).unwrap(), Arc::clone(&error));
requested.await.unwrap();
assert_eq!(uploads.pending(), 1);

let interrupt_token = CancellationToken::new();
tokio::join!(uploads.wait(&interrupt_token), async { interrupt_token.cancel() });
assert_eq!(uploads.pending(), 0);
assert!(matches!(error.get(), Some(UploadError::Interrupted)), "{error:?}");
}
}
Loading
Loading