From 524bc21ce9fd0c2b79e999304b165b54feb169f0 Mon Sep 17 00:00:00 2001 From: wan9chi Date: Mon, 28 Sep 2026 00:52:23 +0800 Subject: [PATCH] fix(cache): stop remote cache requests on Ctrl-C and fast-fail Remote cache lookups, downloads, and uploads now stop as soon as the run is cancelled, instead of holding `vp run` open until they finish or time out. A task whose lookup was cut short doesn't start or restore outputs. Co-authored-by: Claude Opus 5.5 --- CHANGELOG.md | 2 +- crates/vt/src/session/cache/mod.rs | 22 ++- crates/vt/src/session/cache/remote.rs | 163 ++++++++++++++---- crates/vt/src/session/execute/cache_update.rs | 10 +- crates/vt/src/session/execute/mod.rs | 48 +++++- crates/vt/src/session/execute/scheduler.rs | 29 ++-- crates/vt/src/session/mod.rs | 26 ++- crates/vt_bin/src/vtt/main.rs | 4 +- crates/vt_bin/src/vtt/stalled_remote_cache.rs | 36 ++++ .../fixtures/remote_cache/snapshots.toml | 27 +++ .../snapshots/ctrl_c_during_fetch.md | 17 ++ .../snapshots/fast_fail_during_fetch.md | 11 ++ .../fixtures/remote_cache/vite-task.json | 9 + docs/cancellation.md | 9 +- 14 files changed, 334 insertions(+), 79 deletions(-) create mode 100644 crates/vt_bin/src/vtt/stalled_remote_cache.rs create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_fetch.md create mode 100644 crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_fetch.md diff --git a/CHANGELOG.md b/CHANGELOG.md index 2faa74ede..44a7007eb 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. 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. 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)). +- **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), [#771](https://github.com/voidzero-dev/vite-task/pull/771)). - **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/mod.rs b/crates/vt/src/session/cache/mod.rs index 3b155c7bf..1f2d725ef 100644 --- a/crates/vt/src/session/cache/mod.rs +++ b/crates/vt/src/session/cache/mod.rs @@ -16,6 +16,7 @@ pub use display::{ use rusqlite::{Connection, OptionalExtension as _}; use serde::{Deserialize, Serialize}; use tokio::sync::Mutex; +use tokio_util::sync::CancellationToken; use vt_graph::config::ResolvedGlobConfig; use vt_path::{AbsolutePath, RelativePathBuf}; use vt_plan::{ @@ -352,7 +353,7 @@ impl ExecutionCache { /// 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. + /// remote cache. Remote requests stop when `cancel_token` is cancelled. #[tracing::instrument(level = "debug", skip_all)] pub async fn try_hit( &self, @@ -360,6 +361,7 @@ impl ExecutionCache { globbed_inputs: &BTreeMap, workspace_root: &AbsolutePath, cache_dir: &AbsolutePath, + cancel_token: &CancellationToken, ) -> anyhow::Result> { let cache_key = CacheEntryKey::from_metadata(cache_metadata); @@ -389,6 +391,7 @@ impl ExecutionCache { globbed_inputs, workspace_root, cache_dir, + cancel_token, ) .await? { @@ -444,6 +447,7 @@ impl ExecutionCache { /// 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. + #[expect(clippy::too_many_arguments, reason = "forwarded from `try_hit`")] async fn try_hit_remote( &self, endpoint: &Arc, @@ -452,10 +456,11 @@ impl ExecutionCache { globbed_inputs: &BTreeMap, workspace_root: &AbsolutePath, cache_dir: &AbsolutePath, + cancel_token: &CancellationToken, ) -> anyhow::Result> { let fetched = self .remote_clients - .fetch(endpoint, cache_key, &cache_metadata.execution_cache_key) + .fetch(endpoint, cache_key, &cache_metadata.execution_cache_key, cancel_token) .await; let validate = |cache_value: &CacheEntryValue| { cache_value.validate(&cache_metadata.unfiltered_envs, globbed_inputs, workspace_root) @@ -468,7 +473,11 @@ impl ExecutionCache { let output_archive = match blob_id { Some(blob_id) => { - match self.remote_clients.download_archive(endpoint, &blob_id, cache_dir).await { + match self + .remote_clients + .download_archive(endpoint, &blob_id, cache_dir, cancel_token) + .await + { Ok(archive_name) => Some(archive_name), Err(err) => return Ok(Err(err.into_miss())), } @@ -514,14 +523,15 @@ impl ExecutionCache { /// 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 - /// upload failed. + /// cache, until `cancel_token` is cancelled. Returns `Ok(Err(_))` if the + /// local update succeeded but the upload failed. #[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> { let execution_cache_key = &cache_metadata.execution_cache_key; @@ -537,7 +547,7 @@ impl ExecutionCache { }; let upload = self .remote_clients - .upload(url, &cache_key, execution_cache_key, &cache_value, cache_dir) + .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"); diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index a048c2caa..159563046 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -22,6 +22,7 @@ use std::{ use bytes::Bytes; use rustc_hash::FxHashMap; use tokio::sync::mpsc; +use tokio_util::sync::CancellationToken; use vt_path::AbsolutePath; use vt_plan::cache_metadata::ExecutionCacheKey; use vt_remote_cache::{Client, Download, Fetched}; @@ -43,6 +44,10 @@ 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, } /// Why no entry could be read from the remote cache. It's a cache miss, and @@ -70,6 +75,10 @@ pub enum ReadError { Validate(#[source] anyhow::Error), #[error("failed to encode the cache key")] Encode(#[from] WriteError), + /// The run was cancelled, by Ctrl-C or fast-fail, before the read + /// finished. The task doesn't start then, so this miss isn't reported. + #[error("cancelled")] + Cancelled, } impl ReadError { @@ -141,37 +150,45 @@ impl RemoteClients { } /// Fetch the entry stored under `cache_key`, falling back to the entry - /// last stored for `execution_cache_key`. + /// last stored for `execution_cache_key`. Stops when `cancel_token` is + /// cancelled. pub(super) async fn fetch( &self, endpoint: &Arc, cache_key: &CacheEntryKey, execution_cache_key: &ExecutionCacheKey, + cancel_token: &CancellationToken, ) -> Result { let client = self.client(endpoint).map_err(ReadError::Fetch)?; let key = encode_key(cache_key)?; let secondary_key = encode_key(execution_cache_key)?; - client.fetch(&key, &secondary_key).await.map_err(ReadError::Fetch) + cancel_token + .run_until_cancelled(client.fetch(&key, &secondary_key)) + .await + .ok_or(ReadError::Cancelled)? + .map_err(ReadError::Fetch) } /// Download the blob `blob_id` into `cache_dir`, checking that it decodes /// 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. + /// which is renamed once the check passes and removed otherwise, such as + /// when `cancel_token` is cancelled. Returns the archive's file name. pub(super) async fn download_archive( &self, endpoint: &Arc, blob_id: &str, cache_dir: &AbsolutePath, + cancel_token: &CancellationToken, ) -> 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 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) - }); + let result = + download_checked(&client, blob_id, &temp_path, cancel_token).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(temp_path.as_path()); @@ -180,7 +197,7 @@ impl RemoteClients { } /// Upload an entry that was just recorded locally, along with its output - /// archive in `cache_dir`. + /// archive in `cache_dir`. Stops when `cancel_token` is cancelled. pub(super) async fn upload( &self, endpoint: &Arc, @@ -188,13 +205,15 @@ impl RemoteClients { execution_cache_key: &ExecutionCacheKey, cache_value: &CacheEntryValue, cache_dir: &AbsolutePath, + cancel_token: &CancellationToken, ) -> Result<(), 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())); - client.store(&key, &secondary_key, &value, archive.as_deref()).await?; + let store = client.store(&key, &secondary_key, &value, archive.as_deref()); + cancel_token.run_until_cancelled(store).await.ok_or(UploadError::Cancelled)??; Ok(()) } } @@ -202,15 +221,20 @@ 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`. 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. +/// Download the blob `blob_id` to the file at `path`, until `cancel_token` is +/// cancelled. 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, + cancel_token: &CancellationToken, ) -> Result<(), ReadError> { - let download = client.download(blob_id).await.map_err(ReadError::Download)?; + let download = cancel_token + .run_until_cancelled(client.download(blob_id)) + .await + .ok_or(ReadError::Cancelled)? + .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 || { @@ -221,11 +245,12 @@ async fn download_checked( } checked.map_err(ReadError::CorruptArchive) }); - 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. + // Cancelling drops the sender, which ends the check. + let sent = cancel_token.run_until_cancelled(send_chunks(download, sender)).await; + // Wait for the check even after a network error or cancellation, 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)?; + sent.ok_or(ReadError::Cancelled)?.map_err(ReadError::Download)?; checked } @@ -313,6 +338,7 @@ fn decode_key(bytes: &[u8]) -> Result { mod tests { use std::{collections::BTreeMap, io::Read as _, net::TcpListener, time::Duration}; + use tokio::sync::oneshot; use vt_graph::config::ResolvedGlobConfig; use vt_path::{AbsolutePathBuf, RelativePathBuf}; use vt_plan::cache_metadata::{EnvValueHash, SpawnFingerprint}; @@ -522,38 +548,105 @@ mod tests { } } - #[tokio::test] - async fn failed_archive_check_stops_the_download() { + /// Serve one request on a loopback endpoint: once the request head + /// arrives, write `response`, then send nothing more and keep the + /// connection open until the client closes it. Returns the endpoint and a + /// receiver that resolves once `response` is written. + fn serve_stalled(response: &'static [u8]) -> (Arc, oneshot::Receiver<()>) { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); - let endpoint: Arc = + let endpoint = Arc::from(vt_str::format!("http://{}/projects/test", listener.local_addr().unwrap())); - let (done_sender, done_receiver) = std::sync::mpsc::channel::<()>(); - // The response announces more than it sends and stays open until the - // test is done, so only the failed check can end the download. - let server = std::thread::spawn(move || { + let (responded_sender, responded) = oneshot::channel(); + std::thread::spawn(move || { let (mut stream, _) = listener.accept().unwrap(); let mut request = Vec::new(); - while !request.ends_with(b"\r\n\r\n") { - let mut buf = [0; 1024]; + let mut buf = [0; 1024]; + while !request.windows(4).any(|window| window == b"\r\n\r\n") { let n = stream.read(&mut buf).unwrap(); - assert_ne!(n, 0, "connection closed before the request ended"); + assert_ne!(n, 0, "connection closed before the request head ended"); request.extend_from_slice(&buf[..n]); } - stream - .write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 1000\r\n\r\nnot an archive") - .unwrap(); - let _ = done_receiver.recv(); + stream.write_all(response).unwrap(); + let _ = responded_sender.send(()); + while stream.read(&mut buf).is_ok_and(|n| n > 0) {} }); + (endpoint, responded) + } + + #[tokio::test] + async fn failed_archive_check_stops_the_download() { + // The response announces more than it sends, so only the failed check + // can end the download. + let (endpoint, _) = + serve_stalled(b"HTTP/1.1 200 OK\r\ncontent-length: 1000\r\n\r\nnot an archive"); let dir = tempfile::tempdir().unwrap(); let cache_dir = AbsolutePathBuf::new(dir.path().to_path_buf()).unwrap(); let error = RemoteClients::default() - .download_archive(&endpoint, "1", &cache_dir) + .download_archive(&endpoint, "1", &cache_dir, &CancellationToken::new()) .await .unwrap_err(); assert!(matches!(error, ReadError::CorruptArchive(_)), "{error:?}"); assert_eq!(std::fs::read_dir(dir.path()).unwrap().count(), 0); - drop(done_sender); - server.join().unwrap(); + } + + #[tokio::test] + async fn cancelling_stops_a_fetch() { + let (endpoint, requested) = serve_stalled(b""); + let cancel_token = CancellationToken::new(); + let key = cache_key(ResolvedGlobConfig::default_auto()); + let execution_key = ExecutionCacheKey::ExecAPI(Arc::from([])); + + let clients = RemoteClients::default(); + let fetch = clients.fetch(&endpoint, &key, &execution_key, &cancel_token); + let (fetched, ()) = tokio::join!(fetch, async { + requested.await.unwrap(); + cancel_token.cancel(); + }); + assert!(matches!(fetched, Err(ReadError::Cancelled)), "{fetched:?}"); + } + + #[tokio::test] + async fn cancelling_stops_a_download_and_removes_it() { + // The response announces a body that never arrives, so only + // cancelling can end the download. + let (endpoint, responded) = + serve_stalled(b"HTTP/1.1 200 OK\r\ncontent-length: 1000\r\n\r\n"); + let cancel_token = CancellationToken::new(); + let dir = tempfile::tempdir().unwrap(); + let cache_dir = AbsolutePathBuf::new(dir.path().to_path_buf()).unwrap(); + + let clients = RemoteClients::default(); + let download = clients.download_archive(&endpoint, "1", &cache_dir, &cancel_token); + let (downloaded, ()) = tokio::join!(download, async { + responded.await.unwrap(); + // Cancel once the `.tmp` file exists, so the download has started + // writing it. + while std::fs::read_dir(dir.path()).unwrap().next().is_none() { + tokio::task::yield_now().await; + } + cancel_token.cancel(); + }); + assert!(matches!(downloaded, Err(ReadError::Cancelled)), "{downloaded:?}"); + 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(); + 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(); + + 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:?}"); } } diff --git a/crates/vt/src/session/execute/cache_update.rs b/crates/vt/src/session/execute/cache_update.rs index 029c48f52..ee6a04f6e 100644 --- a/crates/vt/src/session/execute/cache_update.rs +++ b/crates/vt/src/session/execute/cache_update.rs @@ -4,6 +4,7 @@ use std::{collections::BTreeMap, sync::Arc, time::Duration}; use rustc_hash::FxHashSet; +use tokio_util::sync::CancellationToken; use vt_path::{AbsolutePath, RelativePathBuf}; use vt_plan::cache_metadata::{CacheMetadata, EnvValueHash}; use vt_server::Reports; @@ -44,7 +45,8 @@ type TrackedEnvQueryValues = BTreeMap, duration: Duration, - cancelled: bool, + cancel_token: &CancellationToken, ) -> (CacheUpdateStatus, Option) { let CacheState { metadata, globbed_inputs, std_outputs, tracking } = state; let fspy = tracking.fspy.as_ref(); @@ -79,7 +81,7 @@ pub(super) async fn update_cache( .map(|r| normalize_ignored_paths(&r.ignored_outputs, workspace_root)) .unwrap_or_default(); - if cancelled { + if cancel_token.is_cancelled() { // Cancelled (Ctrl-C or sibling failure) — result is untrustworthy. return (CacheUpdateStatus::NotUpdated(CacheNotUpdatedReason::Cancelled), None); } @@ -185,7 +187,7 @@ pub(super) async fn update_cache( globbed_inputs, output_archive, }; - match cache.update(metadata, new_cache_value, cache_dir).await { + match cache.update(metadata, new_cache_value, cache_dir, cancel_token).await { Ok(upload) => (CacheUpdateStatus::Updated { upload_error: upload.err() }, None), Err(err) => ( CacheUpdateStatus::NotUpdated(CacheNotUpdatedReason::CacheDisabled), diff --git a/crates/vt/src/session/execute/mod.rs b/crates/vt/src/session/execute/mod.rs index 793bc69da..bb8bad7f8 100644 --- a/crates/vt/src/session/execute/mod.rs +++ b/crates/vt/src/session/execute/mod.rs @@ -53,6 +53,10 @@ pub enum SpawnOutcome { /// (cache lookup failure or spawn failure). /// Already reported through the leaf reporter. Failed, + /// The run was cancelled during the cache lookup, so the process didn't + /// start and no cached outputs were restored. Like a task that was never + /// scheduled, it isn't reported. + Cancelled, } /// All valid runtime configurations for a leaf execution, after the cache-hit @@ -273,6 +277,8 @@ enum Report { cache_update: CacheUpdateStatus, error: Option, }, + /// The run was cancelled before the process started. + Cancelled, } impl Report { @@ -304,6 +310,15 @@ impl Report { reporter.finish(Some(exit_status), cache_update, error); SpawnOutcome::Spawned(exit_status) } + // `start()` wasn't called, so the reporter shows and saves nothing. + Self::Cancelled => { + reporter.finish( + None, + CacheUpdateStatus::NotUpdated(CacheNotUpdatedReason::Cancelled), + None, + ); + SpawnOutcome::Cancelled + } } } } @@ -314,13 +329,18 @@ impl Report { /// reused from both graph-based execution and standalone synthetic execution. /// /// The full lifecycle is: -/// 1. Cache lookup (determines cache status) +/// 1. Cache lookup (determines cache status). If the run was cancelled by +/// the time it ends, nothing else happens. /// 2. `leaf_reporter.start(cache_status)` → `StdioConfig` /// 3. If cache hit: replay cached outputs via `StdioConfig` writers /// 4. Otherwise: `spawn()` with the chosen stdio mode, drain pipes and wait /// via [`run_child`], then decide the cache update /// ([`cache_update::update_cache`]) /// +/// Cancelling `fast_fail_token` kills the process. `cancel_token` must be a +/// child of `fast_fail_token` that Ctrl-C also cancels: cancelling it stops +/// remote cache requests and prevents caching. +/// /// Every path reports through the single `finish()` below — errors (cache /// lookup failure, spawn failure, cache update failure) do not abort the /// caller. @@ -337,7 +357,7 @@ pub async fn execute_spawn( cache_dir: &AbsolutePath, program_name: &str, fast_fail_token: CancellationToken, - interrupt_token: CancellationToken, + cancel_token: CancellationToken, ) -> SpawnOutcome { let pipeline = run( leaf_reporter.as_mut(), @@ -347,7 +367,7 @@ pub async fn execute_spawn( cache_dir, program_name, fast_fail_token, - interrupt_token, + cancel_token, ); let report = match pipeline.await { Ok(report) | Err(report) => report, @@ -371,14 +391,20 @@ async fn run( cache_dir: &AbsolutePath, program_name: &str, fast_fail_token: CancellationToken, - interrupt_token: CancellationToken, + cancel_token: CancellationToken, ) -> Result { let cache_metadata = spawn_execution.cache_metadata.as_ref(); // 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, cache_dir).await?; + // A lookup can take a while, especially a remote one. If the run was + // cancelled meanwhile, neither start the process nor restore outputs. + let lookup = + lookup_cache(cache_metadata, cache, workspace_root, cache_dir, &cancel_token).await?; + if cancel_token.is_cancelled() { + return Err(Report::Cancelled); + } // 2. Report execution start with the looked-up cache status (`start()` // runs exactly once on every arm) and either replay the hit — no need @@ -472,7 +498,6 @@ async fn run( // 7. Decide the cache update (only when we were in `Cached` mode). Cache // update errors are reported but do not affect the exit status we // return — the process ran, so we return its actual status. - let cancelled = fast_fail_token.is_cancelled() || interrupt_token.is_cancelled(); let (cache_update, error) = match mode { ExecutionMode::Cached { state, .. } => { cache_update::update_cache( @@ -483,7 +508,7 @@ async fn run( &outcome, reports.as_ref(), duration, - cancelled, + &cancel_token, ) .await } @@ -510,12 +535,14 @@ enum CacheLookup { } /// Phase 1: compute the globbed inputs and try to hit the cache. A remote hit -/// downloads its output archive into `cache_dir`. +/// downloads its output archive into `cache_dir`. Remote requests stop when +/// `cancel_token` is cancelled. async fn lookup_cache( cache_metadata: Option<&CacheMetadata>, cache: &ExecutionCache, workspace_root: &Arc, cache_dir: &AbsolutePath, + cancel_token: &CancellationToken, ) -> Result { let Some(cache_metadata) = cache_metadata else { return Ok(CacheLookup::Disabled); @@ -532,7 +559,10 @@ async fn lookup_cache( Report::failed(ExecutionError::Cache { kind: CacheErrorKind::Lookup, source: err }) })?; - match cache.try_hit(cache_metadata, &globbed_inputs, workspace_root, cache_dir).await { + match cache + .try_hit(cache_metadata, &globbed_inputs, workspace_root, cache_dir, cancel_token) + .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/execute/scheduler.rs b/crates/vt/src/session/execute/scheduler.rs index 44a3c6e37..34180d62b 100644 --- a/crates/vt/src/session/execute/scheduler.rs +++ b/crates/vt/src/session/execute/scheduler.rs @@ -48,20 +48,20 @@ struct ExecutionContext<'a> { /// messages that suggest a CLI command (e.g. `cache clean`). program_name: &'a str, /// Token cancelled when a task fails. Kills in-flight child processes - /// (via `start_kill` in spawn.rs), prevents scheduling new tasks, and - /// prevents caching results of concurrently-running tasks. + /// (via `start_kill` in spawn.rs). fast_fail_token: CancellationToken, - /// Token cancelled by Ctrl-C. Unlike `fast_fail_token` (which kills - /// children), this only prevents scheduling new tasks and caching - /// results — running processes are left to handle SIGINT naturally. - interrupt_token: CancellationToken, + /// Token cancelled by Ctrl-C, and by fast-fail as a child of + /// `fast_fail_token`. Prevents scheduling new tasks and caching results, + /// and stops remote cache requests. On Ctrl-C, running processes are + /// left to handle SIGINT naturally. + cancel_token: CancellationToken, } impl ExecutionContext<'_> { /// Returns true if execution has been cancelled, either by a task /// failure (fast-fail) or by Ctrl-C (interrupt). fn cancelled(&self) -> bool { - self.fast_fail_token.is_cancelled() || self.interrupt_token.is_cancelled() + self.cancel_token.is_cancelled() } /// Execute all tasks in an execution graph concurrently, respecting dependencies. @@ -72,7 +72,7 @@ impl ExecutionContext<'_> { /// semaphore, so nested graphs have independent concurrency limits. /// /// Fast-fail: if any task fails, `execute_leaf` cancels the `fast_fail_token` - /// (killing in-flight child processes). Ctrl-C cancels the `interrupt_token`. + /// (killing in-flight child processes). Ctrl-C cancels the `cancel_token`. /// Either cancellation causes this method to close the semaphore, drain /// remaining futures, and return. #[tracing::instrument(level = "debug", skip_all)] @@ -208,11 +208,11 @@ impl ExecutionContext<'_> { self.cache_dir, self.program_name, self.fast_fail_token.clone(), - self.interrupt_token.clone(), + self.cancel_token.clone(), ) .await; match outcome { - SpawnOutcome::CacheHit => false, + SpawnOutcome::CacheHit | SpawnOutcome::Cancelled => false, SpawnOutcome::Spawned(status) => !status.success(), SpawnOutcome::Failed => true, } @@ -231,6 +231,8 @@ impl Session<'_> { /// after cache initialization, so cache errors are reported directly to stderr /// without involving the reporter at all. /// + /// `fast_fail_token` and `cancel_token` are described on [`ExecutionContext`]. + /// /// Returns `Err(ExitStatus)` to indicate the caller should exit with the given status code. /// Returns `Ok(())` when all tasks succeeded. #[tracing::instrument(level = "debug", skip_all)] @@ -238,7 +240,8 @@ impl Session<'_> { &self, execution_graph: ExecutionGraph, builder: Box, - interrupt_token: CancellationToken, + fast_fail_token: CancellationToken, + cancel_token: CancellationToken, ) -> Result<(), ExitStatus> { // Initialize cache before building the reporter. Cache errors are reported // directly to stderr and cause an early exit, keeping the reporter flow clean @@ -260,8 +263,8 @@ impl Session<'_> { workspace_root: &self.workspace_path, cache_dir: &self.cache_path, program_name: self.program_name.as_str(), - fast_fail_token: CancellationToken::new(), - interrupt_token, + fast_fail_token, + cancel_token, }; // Execute the graph with fast-fail: if any task fails, remaining tasks diff --git a/crates/vt/src/session/mod.rs b/crates/vt/src/session/mod.rs index 6ec70e3e3..6495080c7 100644 --- a/crates/vt/src/session/mod.rs +++ b/crates/vt/src/session/mod.rs @@ -381,8 +381,10 @@ impl<'a> Session<'a> { )); // Don't let SIGINT/CTRL_C kill the runner. Child tasks receive // the signal directly from the terminal driver and handle it - // themselves. Cancelling the interrupt token prevents scheduling - // new tasks and caching results of in-flight tasks. + // themselves. Cancelling the cancel token prevents scheduling + // new tasks and caching results of in-flight tasks, and stops + // remote cache requests. It's a child of the fast-fail token, + // so fast-fail cancels it too. // // On Windows, an ancestor process (e.g. cargo) may have been // created with CREATE_NEW_PROCESS_GROUP, which sets a per-process @@ -401,13 +403,14 @@ impl<'a> Session<'a> { } SetConsoleCtrlHandler(None, 0); } - let interrupt_token = tokio_util::sync::CancellationToken::new(); - let ct = interrupt_token.clone(); + let fast_fail_token = tokio_util::sync::CancellationToken::new(); + let cancel_token = fast_fail_token.child_token(); + let ct = cancel_token.clone(); ctrlc::set_handler(move || { ct.cancel(); })?; - self.execute_graph(graph, builder, interrupt_token) + self.execute_graph(graph, builder, fast_fail_token, cancel_token) .await .map_err(SessionError::EarlyExit) } @@ -719,6 +722,8 @@ impl<'a> Session<'a> { ); // Execute the spawn directly using the free function, bypassing the graph pipeline + let fast_fail_token = tokio_util::sync::CancellationToken::new(); + let cancel_token = fast_fail_token.child_token(); let outcome = execute::execute_spawn( Box::new(plain_reporter), &spawn_execution, @@ -726,8 +731,8 @@ impl<'a> Session<'a> { &self.workspace_path, &self.cache_path, self.program_name.as_str(), - tokio_util::sync::CancellationToken::new(), - tokio_util::sync::CancellationToken::new(), + fast_fail_token, + cancel_token, ) .await; match outcome { @@ -744,8 +749,11 @@ impl<'a> Session<'a> { )] Ok(ExitStatus(code.clamp(1, 255) as u8)) } - // Infrastructure error — already reported through the reporter's finish() - execute::SpawnOutcome::Failed => Ok(ExitStatus::FAILURE), + // Infrastructure error — already reported through the reporter's finish(). + // Nothing cancels the tokens above, and a cancelled command wouldn't have run. + execute::SpawnOutcome::Failed | execute::SpawnOutcome::Cancelled => { + Ok(ExitStatus::FAILURE) + } } } diff --git a/crates/vt_bin/src/vtt/main.rs b/crates/vt_bin/src/vtt/main.rs index dd36cadba..e77c719df 100644 --- a/crates/vt_bin/src/vtt/main.rs +++ b/crates/vt_bin/src/vtt/main.rs @@ -25,6 +25,7 @@ mod replace_file_content; mod rm; #[cfg(target_os = "linux")] mod small_dev_shm; +mod stalled_remote_cache; mod stat_file; mod stat_many; mod touch_file; @@ -35,7 +36,7 @@ fn main() { if args.len() < 2 { eprintln!("Usage: vtt [args...]"); eprintln!( - "Subcommands: barrier, check-tty, cp, exit, exit-on-ctrlc, grep-file, list-dir, mkdir, pipe-stdin, print, print-color, print-cwd, print-env, print-file, read-stdin, replace-file-content, rm, small_dev_shm, stat-file, stat-many, touch-file, write-file" + "Subcommands: barrier, check-tty, cp, exit, exit-on-ctrlc, grep-file, list-dir, mkdir, pipe-stdin, print, print-color, print-cwd, print-env, print-file, read-stdin, replace-file-content, rm, small_dev_shm, stalled-remote-cache, stat-file, stat-many, touch-file, write-file" ); std::process::exit(1); } @@ -71,6 +72,7 @@ fn main() { "small_dev_shm" => small_dev_shm::run(&args[2..]).map_err(Into::into), #[cfg(not(target_os = "linux"))] "small_dev_shm" => Err("vtt small_dev_shm is only supported on Linux".into()), + "stalled-remote-cache" => stalled_remote_cache::run(&args[2..]), "stat-file" => { stat_file::run(&args[2..]); Ok(()) diff --git a/crates/vt_bin/src/vtt/stalled_remote_cache.rs b/crates/vt_bin/src/vtt/stalled_remote_cache.rs new file mode 100644 index 000000000..ab3d391c9 --- /dev/null +++ b/crates/vt_bin/src/vtt/stalled_remote_cache.rs @@ -0,0 +1,36 @@ +use std::io::Read as _; + +/// stalled-remote-cache `` \[``...\] +/// +/// Runs `` with `VP_REMOTE_CACHE_URL` set to a loopback endpoint +/// that accepts requests but never responds, then exits with the command's +/// exit code. Emits a "request" milestone when a request arrives. Ctrl-C is +/// left to the command. +pub fn run(args: &[String]) -> Result<(), Box> { + let [program, args @ ..] = args else { + return Err("Usage: vtt stalled-remote-cache [args...]".into()); + }; + ctrlc::set_handler(|| {})?; + + let listener = std::net::TcpListener::bind("127.0.0.1:0")?; + let endpoint = std::format!("http://{}/projects/test", listener.local_addr()?); + std::thread::spawn(move || { + for mut stream in listener.incoming().filter_map(Result::ok) { + std::thread::spawn(move || { + let mut buf = [0; 4096]; + if stream.read(&mut buf).is_ok_and(|n| n > 0) { + pty_terminal_test_client::mark_milestone("request"); + } + // Hold the connection without responding until the client + // closes it. + while stream.read(&mut buf).is_ok_and(|n| n > 0) {} + }); + } + }); + + let status = std::process::Command::new(program) + .args(args) + .env("VP_REMOTE_CACHE_URL", endpoint) + .status()?; + std::process::exit(status.code().unwrap_or(1)); +} 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 d9bea2423..ba61d0b67 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 @@ -352,3 +352,30 @@ steps = [ "--last-details", ], comment = "The details include the underlying error." }, ] + +[[e2e]] +name = "ctrl_c_during_fetch" +steps = [ + { argv = [ + "vtt", + "stalled-remote-cache", + "vt", + "run", + "build", + ], interactions = [ + { "expect-milestone" = "request" }, + { "write-key" = "ctrl-c" }, + ], comment = "The endpoint never responds. Ctrl-C stops the fetch, and the task doesn't start." }, +] + +[[e2e]] +name = "fast_fail_during_fetch" +steps = [ + { argv = [ + "vtt", + "stalled-remote-cache", + "vt", + "run", + "fail-during-build", + ], comment = "The endpoint never responds. fail exits while build's fetch is in flight, which stops the fetch, and build doesn't start." }, +] diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_fetch.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_fetch.md new file mode 100644 index 000000000..bd77d77bf --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_fetch.md @@ -0,0 +1,17 @@ +# ctrl_c_during_fetch + +## `vtt stalled-remote-cache vt run build` + +The endpoint never responds. Ctrl-C stops the fetch, and the task doesn't start. + +**→ expect-milestone:** `request` + +``` +``` + +**← write-key:** `ctrl-c` + +``` +--- +vt run: 0/0 cache hit (0%). (Run `vt run --last-details` for full details) +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_fetch.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_fetch.md new file mode 100644 index 000000000..d7b81dd79 --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_fetch.md @@ -0,0 +1,11 @@ +# fast_fail_during_fetch + +## `vtt stalled-remote-cache vt run fail-during-build` + +The endpoint never responds. fail exits while build's fetch is in flight, which stops the fetch, and build doesn't start. + +**Exit code:** 1 + +``` +$ vtt exit 1 ⊘ cache disabled +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/vite-task.json b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/vite-task.json index 3d2c4e967..14be7f2eb 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/vite-task.json +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/vite-task.json @@ -16,6 +16,15 @@ }, "all": { "command": "vt run build && vt run check" + }, + "fail": { + "command": "vtt exit 1", + "cache": false + }, + "fail-during-build": { + "command": "vtt print done", + "dependsOn": ["build", "fail"], + "cache": false } } } diff --git a/docs/cancellation.md b/docs/cancellation.md index 30d4b104f..63cb816f4 100644 --- a/docs/cancellation.md +++ b/docs/cancellation.md @@ -1,6 +1,6 @@ # Cancellation -`vp run` handles two kinds of cancellation: **Ctrl-C** (user interrupt) and **fast-fail** (a task exits with non-zero status). Both prevent new tasks from being scheduled and prevent caching of in-flight results, but they differ in how they treat running processes. +`vp run` handles two kinds of cancellation: **Ctrl-C** (user interrupt) and **fast-fail** (a task exits with non-zero status). Both prevent new tasks from being scheduled, prevent caching of in-flight results, and stop [remote cache requests](#remote-cache-requests), but they differ in how they treat running processes. ## Ctrl-C @@ -18,6 +18,13 @@ When any task exits with non-zero status: 2. No new tasks are scheduled. 3. Results of other in-flight tasks are **not cached** (they were killed mid-execution). +## Remote cache requests + +Both kinds of cancellation stop remote cache lookups, downloads, and uploads right away instead of waiting for them to finish or time out. + +- A task whose cache lookup was still in progress doesn't start and doesn't restore cached outputs, even if the lookup found them. Like a task that was never scheduled, it isn't shown in the summary. +- A task whose upload is stopped keeps its local cache entry, and the summary warns that it wasn't uploaded. + ## Why interrupted tasks are not cached A task that receives Ctrl-C might exit 0 after partial work (e.g., a build tool that flushes what it has so far). Caching this result would mean the next `vp run` replays incomplete output and skips the real execution. By never caching interrupted results, `vp run` guarantees that the next run starts fresh.