diff --git a/CHANGELOG.md b/CHANGELOG.md index 5a035665e..61217dff1 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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)). diff --git a/Cargo.lock b/Cargo.lock index 67a03f1e3..d0d3f76bd 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1750,6 +1750,25 @@ dependencies = [ "regex-syntax", ] +[[package]] +name = "h2" +version = "0.4.19" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ef8e5e5a340588f4452631496976cf8636d4a7ecf600239fdc27615d2530bc16" +dependencies = [ + "atomic-waker", + "bytes", + "fnv", + "futures-core", + "futures-sink", + "http", + "indexmap", + "slab", + "tokio", + "tokio-util", + "tracing", +] + [[package]] name = "half" version = "2.7.1" @@ -1872,6 +1891,7 @@ dependencies = [ "bytes", "futures-channel", "futures-core", + "h2", "http", "http-body", "httparse", @@ -3629,6 +3649,7 @@ dependencies = [ "bytes", "futures-core", "futures-util", + "h2", "http", "http-body", "http-body-util", @@ -4632,6 +4653,7 @@ dependencies = [ "bytes", "futures-core", "futures-sink", + "futures-util", "pin-project-lite", "tokio", ] diff --git a/crates/vt/Cargo.toml b/crates/vt/Cargo.toml index 3759db96f..7b6aed2b8 100644 --- a/crates/vt/Cargo.toml +++ b/crates/vt/Cargo.toml @@ -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 } diff --git a/crates/vt/src/session/cache/mod.rs b/crates/vt/src/session/cache/mod.rs index c93d22235..4ff29bbba 100644 --- a/crates/vt/src/session/cache/mod.rs +++ b/crates/vt/src/session/cache/mod.rs @@ -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; @@ -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, @@ -112,6 +119,7 @@ pub struct CacheEntryValue { pub struct ExecutionCache { conn: Mutex, remote_clients: RemoteClients, + uploads: RemoteUploads, } /// A cache hit: the entry to replay, and the cache it came from. @@ -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] @@ -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> { + ) -> anyhow::Result>> { 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`. diff --git a/crates/vt/src/session/cache/remote.rs b/crates/vt/src/session/cache/remote.rs index 18008f28e..fb7a87d32 100644 --- a/crates/vt/src/session/cache/remote.rs +++ b/crates/vt/src/session/cache/remote.rs @@ -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}; @@ -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 @@ -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, cache_key: &CacheEntryKey, execution_cache_key: &ExecutionCacheKey, cache_value: &CacheEntryValue, cache_dir: &AbsolutePath, - cancel_token: &CancellationToken, - ) -> Result<(), UploadError> { + ) -> Result> + 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> + Send + 'static, + error: Arc>, + ) { + 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(); } } @@ -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, + ) -> Result> + 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) -> Option { + 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:?}"); } } diff --git a/crates/vt/src/session/event.rs b/crates/vt/src/session/event.rs index c52300569..caa54c000 100644 --- a/crates/vt/src/session/event.rs +++ b/crates/vt/src/session/event.rs @@ -1,4 +1,8 @@ -use std::{process::ExitStatus, time::Duration}; +use std::{ + process::ExitStatus, + sync::{Arc, OnceLock}, + time::Duration, +}; use vt_path::RelativePathBuf; use vt_server::Error as IpcServerError; @@ -126,9 +130,13 @@ pub enum CacheNotUpdatedReason { pub enum CacheUpdateStatus { /// Cache was successfully updated with new fingerprint and outputs Updated { - /// Why uploading the entry to the remote cache failed. `None` if the - /// upload succeeded or wasn't attempted. - upload_error: Option, + /// Why uploading the entry to the remote cache failed. It's set only + /// if the upload fails, which can happen in the background after the + /// task finishes. Empty if the upload succeeded or wasn't attempted. + /// Read it only after + /// [`ExecutionCache::wait_for_uploads`](super::cache::ExecutionCache::wait_for_uploads) + /// returns. + upload_error: Arc>, }, /// Cache was not updated (with reason). /// The reason is part of the `LeafExecutionReporter` trait contract — reporters diff --git a/crates/vt/src/session/execute/cache_update.rs b/crates/vt/src/session/execute/cache_update.rs index ee6a04f6e..9650c5af3 100644 --- a/crates/vt/src/session/execute/cache_update.rs +++ b/crates/vt/src/session/execute/cache_update.rs @@ -46,7 +46,7 @@ type TrackedEnvQueryValues = BTreeMap (CacheUpdateStatus::Updated { upload_error: upload.err() }, None), + match cache.update(metadata, new_cache_value, cache_dir).await { + Ok(upload_error) => (CacheUpdateStatus::Updated { upload_error }, None), Err(err) => ( CacheUpdateStatus::NotUpdated(CacheNotUpdatedReason::CacheDisabled), Some(ExecutionError::Cache { kind: CacheErrorKind::Update, source: err }), diff --git a/crates/vt/src/session/execute/mod.rs b/crates/vt/src/session/execute/mod.rs index 0426f060f..46dbacac1 100644 --- a/crates/vt/src/session/execute/mod.rs +++ b/crates/vt/src/session/execute/mod.rs @@ -339,7 +339,9 @@ impl Report { /// /// 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. +/// remote cache lookups and prevents caching. A cached run's upload to the +/// remote cache runs in the background, and the caller waits for it with +/// [`ExecutionCache::wait_for_uploads`]. /// /// Every path reports through the single `finish()` below — errors (cache /// lookup failure, spawn failure, cache update failure) do not abort the diff --git a/crates/vt/src/session/execute/scheduler.rs b/crates/vt/src/session/execute/scheduler.rs index 34180d62b..e01d4eb57 100644 --- a/crates/vt/src/session/execute/scheduler.rs +++ b/crates/vt/src/session/execute/scheduler.rs @@ -2,13 +2,14 @@ //! (dependency order, concurrency limits, fast-fail), and hands each leaf to //! [`execute_spawn`] which owns *how* a single spawn runs. -use std::{cell::RefCell, io::Write as _, sync::Arc}; +use std::{cell::RefCell, ffi::OsStr, io::Write as _, num::NonZeroUsize, sync::Arc}; use futures_util::{FutureExt, StreamExt, future::LocalBoxFuture, stream::FuturesUnordered}; use petgraph::Direction; use rustc_hash::FxHashMap; use tokio::sync::Semaphore; use tokio_util::sync::CancellationToken; +use vt_casefold::EnvName; use vt_path::AbsolutePath; use vt_plan::{ ExecutionGraph, ExecutionItemDisplay, ExecutionItemKind, LeafExecutionKind, @@ -25,6 +26,12 @@ use crate::{ }, }; +/// If set, the reporter isn't told about the uploads still running when the +/// graph is done, so there is no message about them. Whether an upload is +/// still running then depends on how fast the remote cache responds, so tests +/// set this to keep their output stable. +const HIDE_PENDING_UPLOADS_ENV: &str = "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS"; + /// Holds shared references needed during graph execution. /// /// The `reporter` field is wrapped in `RefCell` because concurrent futures @@ -52,7 +59,7 @@ struct ExecutionContext<'a> { fast_fail_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 + /// and stops remote cache lookups. On Ctrl-C, running processes are /// left to handle SIGINT naturally. cancel_token: CancellationToken, } @@ -233,6 +240,12 @@ impl Session<'_> { /// /// `fast_fail_token` and `cancel_token` are described on [`ExecutionContext`]. /// + /// Uploads to the remote cache run in the background and can outlive + /// their tasks. Once the graph is done, this waits for them before the + /// summary, telling the reporter how many are left. `interrupt_token`, + /// which only Ctrl-C cancels, cancels them instead. Fast-fail doesn't, so + /// tasks that succeeded before another one failed are still uploaded. + /// /// 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)] @@ -242,6 +255,7 @@ impl Session<'_> { builder: Box, fast_fail_token: CancellationToken, cancel_token: CancellationToken, + interrupt_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 @@ -271,8 +285,21 @@ impl Session<'_> { // are skipped. Leaf-level errors are reported through the reporter. execution_context.execute_expanded_graph(&execution_graph).await; - // Leaf-level errors and non-zero exit statuses are tracked internally - // by the reporter. - reporter.into_inner().finish() + // Nested graphs share the cache, so this waits for their uploads too. + // After Ctrl-C, the uploads are cancelled without a message. + let mut reporter = reporter.into_inner(); + let hide_pending_uploads = + self.envs.contains_key(EnvName::from_ref(OsStr::new(HIDE_PENDING_UPLOADS_ENV))); + if !interrupt_token.is_cancelled() + && !hide_pending_uploads + && let Some(pending) = NonZeroUsize::new(cache.pending_uploads()) + { + reporter.uploads_pending(pending); + } + cache.wait_for_uploads(&interrupt_token).await; + + // Leaf-level errors, non-zero exit statuses, and failed uploads are + // tracked internally by the reporter. + reporter.finish() } } diff --git a/crates/vt/src/session/mod.rs b/crates/vt/src/session/mod.rs index 6495080c7..a0f326747 100644 --- a/crates/vt/src/session/mod.rs +++ b/crates/vt/src/session/mod.rs @@ -383,8 +383,9 @@ impl<'a> Session<'a> { // the signal directly from the terminal driver and handle it // 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. + // remote cache lookups. It's a child of the fast-fail token, + // so fast-fail cancels it too. Only Ctrl-C cancels the + // interrupt token, which cancels remote cache uploads. // // On Windows, an ancestor process (e.g. cargo) may have been // created with CREATE_NEW_PROCESS_GROUP, which sets a per-process @@ -405,12 +406,15 @@ impl<'a> Session<'a> { } let fast_fail_token = tokio_util::sync::CancellationToken::new(); let cancel_token = fast_fail_token.child_token(); + let interrupt_token = tokio_util::sync::CancellationToken::new(); let ct = cancel_token.clone(); + let it = interrupt_token.clone(); ctrlc::set_handler(move || { ct.cancel(); + it.cancel(); })?; - self.execute_graph(graph, builder, fast_fail_token, cancel_token) + self.execute_graph(graph, builder, fast_fail_token, cancel_token, interrupt_token) .await .map_err(SessionError::EarlyExit) } @@ -735,6 +739,9 @@ impl<'a> Session<'a> { cancel_token, ) .await; + // Nothing interrupts the upload, if any. The plain reporter has no + // summary, so it isn't told about the upload or its result. + cache.wait_for_uploads(&tokio_util::sync::CancellationToken::new()).await; match outcome { // Cache hit — no process was spawned, success execute::SpawnOutcome::CacheHit => Ok(ExitStatus::SUCCESS), diff --git a/crates/vt/src/session/reporter/mod.rs b/crates/vt/src/session/reporter/mod.rs index 4a8522664..b4c03fc99 100644 --- a/crates/vt/src/session/reporter/mod.rs +++ b/crates/vt/src/session/reporter/mod.rs @@ -30,7 +30,7 @@ mod plain; pub mod summary; mod summary_reporter; -use std::{io::Write, process::ExitStatus as StdExitStatus}; +use std::{io::Write, num::NonZeroUsize, process::ExitStatus as StdExitStatus}; pub use grouped::GroupedReporterBuilder; pub use interleaved::InterleavedReporterBuilder; @@ -162,6 +162,11 @@ pub trait GraphExecutionReporter { leaf_kind: &LeafExecutionKind, ) -> Box; + /// Report that all tasks are done, but `count` uploads to the remote cache + /// are still running. The caller waits for them before calling + /// [`Self::finish`]. Not called after Ctrl-C, which cancels them instead. + fn uploads_pending(&mut self, _count: NonZeroUsize) {} + /// Finalize the graph execution session. /// /// Leaf-level errors are already tracked internally by the reporter via the diff --git a/crates/vt/src/session/reporter/summary.rs b/crates/vt/src/session/reporter/summary.rs index 14936fc4f..40ba11492 100644 --- a/crates/vt/src/session/reporter/summary.rs +++ b/crates/vt/src/session/reporter/summary.rs @@ -7,7 +7,12 @@ //! Both the live reporter and the `--last-details` display use the same rendering //! functions, ensuring consistent output. -use std::{fmt::Display, io::Write, num::NonZeroI32, time::Duration}; +use std::{ + fmt::Display, + io::Write, + num::{NonZeroI32, NonZeroUsize}, + time::Duration, +}; use owo_colors::Style; use serde::{Deserialize, Serialize}; @@ -349,10 +354,6 @@ impl TaskResult { cache_update_status, CacheUpdateStatus::NotUpdated(CacheNotUpdatedReason::TrackingIncomplete) ); - let upload_error = match cache_update_status { - CacheUpdateStatus::Updated { upload_error: Some(err) } => Some(SavedError::new(err)), - _ => None, - }; match cache_status { // The only error a cache hit can have is a failed restore. @@ -374,7 +375,6 @@ impl TaskResult { ipc_server_error, tool_disabled_cache, tracking_incomplete, - upload_error, ), }, CacheStatus::Miss(cache_miss) => Self::Spawned { @@ -389,18 +389,24 @@ impl TaskResult { ipc_server_error, tool_disabled_cache, tracking_incomplete, - upload_error, ), }, } } + + /// Record why uploading the entry to the remote cache failed. The upload + /// can fail after the task finishes, so this is set after + /// [`Self::from_execution`]. Only a successful spawned task uploads an + /// entry, so other results are left as they are. + pub fn set_upload_error(&mut self, error: SavedError) { + if let Self::Spawned { outcome: SpawnOutcome::Success { upload_error, .. }, .. } = self { + *upload_error = Some(error); + } + } } /// Build a [`SpawnOutcome`] from process exit status and optional pre-converted error. -#[expect( - clippy::too_many_arguments, - reason = "each cache update detail is extracted by the caller and passed through" -)] +/// A failed upload is set later, with [`TaskResult::set_upload_error`]. fn spawn_outcome_from_execution( exit_status: Option, saved_error: Option<&SavedError>, @@ -409,7 +415,6 @@ fn spawn_outcome_from_execution( ipc_server_error: Option, tool_disabled_cache: bool, tracking_incomplete: bool, - upload_error: Option, ) -> SpawnOutcome { match (exit_status, saved_error) { // Spawn error — process never ran @@ -422,7 +427,7 @@ fn spawn_outcome_from_execution( ipc_server_error, tool_disabled_cache, tracking_incomplete, - upload_error, + upload_error: None, }, // Process exited with non-zero code (Some(status), _) => { @@ -1143,6 +1148,22 @@ fn format_upload_failed_notice(buf: &mut Vec, failures: &[UploadFailure]) { let _ = write!(buf, "."); } +/// Render the line shown when all tasks are done, but `count` uploads to the +/// remote cache are still running. +pub fn format_uploads_pending(count: NonZeroUsize) -> Vec { + let uploads = if count.get() == 1 { "upload" } else { "uploads" }; + let mut buf = Vec::new(); + let _ = writeln!( + buf, + "{}", + vt_str::format!( + "Waiting for {count} remote cache {uploads} to finish (Ctrl-C to cancel)..." + ) + .style(Style::new().bright_black()) + ); + buf +} + #[cfg(test)] mod tests { use vt_path::RelativePathBuf; @@ -1399,6 +1420,19 @@ mod tests { ); } + #[test] + fn uploads_pending_names_the_count() { + let pending = |count| strip(&format_uploads_pending(NonZeroUsize::new(count).unwrap())); + assert_eq!( + pending(1).as_str(), + "Waiting for 1 remote cache upload to finish (Ctrl-C to cancel)...\n" + ); + assert_eq!( + pending(2).as_str(), + "Waiting for 2 remote cache uploads to finish (Ctrl-C to cancel)...\n" + ); + } + #[test] fn full_summary_shows_each_cause_on_its_own_line() { let task = upload_failed_task( diff --git a/crates/vt/src/session/reporter/summary_reporter.rs b/crates/vt/src/session/reporter/summary_reporter.rs index 1ab95de50..a406894f7 100644 --- a/crates/vt/src/session/reporter/summary_reporter.rs +++ b/crates/vt/src/session/reporter/summary_reporter.rs @@ -4,7 +4,14 @@ //! results, then renders a summary when the graph execution completes. The inner //! reporter handles all output formatting (interleaved, labeled, grouped). -use std::{cell::RefCell, io::Write, process::ExitStatus as StdExitStatus, rc::Rc, sync::Arc}; +use std::{ + cell::RefCell, + io::Write, + num::NonZeroUsize, + process::ExitStatus as StdExitStatus, + rc::Rc, + sync::{Arc, OnceLock}, +}; use vt_path::AbsolutePath; use vt_plan::{ExecutionItemDisplay, LeafExecutionKind}; @@ -15,10 +22,11 @@ use super::{ LeafExecutionReporter, StdioConfig, }; use crate::session::{ + cache::remote::UploadError, event::{CacheStatus, CacheUpdateStatus, ExecutionError}, reporter::summary::{ LastRunSummary, SavedError, SpawnOutcome, TaskResult, TaskSummary, format_compact_summary, - format_full_summary, + format_full_summary, format_uploads_pending, }, }; @@ -67,9 +75,18 @@ impl GraphExecutionReporterBuilder for SummaryReporterBuilder { } } +/// A finished task's summary, without the result of its upload to the remote +/// cache, which may still be running. +struct RecordedTask { + summary: TaskSummary, + /// Where the upload's error is set if it fails. `None` if the task didn't + /// update the cache. + upload_error: Option>>, +} + struct SummaryGraphReporter { inner: Box, - tasks: Rc>>, + tasks: Rc>>, workspace_path: Arc, writer: Box, show_details: bool, @@ -93,11 +110,29 @@ impl GraphExecutionReporter for SummaryGraphReporter { }) } + fn uploads_pending(&mut self, count: NonZeroUsize) { + let _ = self.writer.write_all(&format_uploads_pending(count)); + let _ = self.writer.flush(); + pty_terminal_test_client::mark_milestone("uploads-pending"); + } + + /// Called after the uploads to the remote cache finish, so their errors + /// are in the summary, and in the saved one. fn finish(self: Box) -> Result<(), ExitStatus> { // Let inner reporter finish first (flushes any pending output). let inner_result = self.inner.finish(); - let tasks = self.tasks.take(); + let tasks: Vec = self + .tasks + .take() + .into_iter() + .map(|RecordedTask { mut summary, upload_error }| { + if let Some(error) = upload_error.as_deref().and_then(OnceLock::get) { + summary.result.set_upload_error(SavedError::new(error)); + } + summary + }) + .collect(); let has_infra_errors = tasks.iter().any(|t| t.result.error().is_some()); @@ -159,7 +194,7 @@ impl GraphExecutionReporter for SummaryGraphReporter { /// Leaf reporter wrapper that records task results for the summary. struct SummaryLeafReporter { inner: Box, - tasks: Rc>>, + tasks: Rc>>, display: ExecutionItemDisplay, workspace_path: Arc, cache_status: Option, @@ -188,7 +223,7 @@ impl LeafExecutionReporter for SummaryLeafReporter { Str::default() }; - let task_summary = TaskSummary { + let summary = TaskSummary { package_name: self.display.task_display.package_name.clone(), task_name: self.display.task_display.task_name.clone(), command: self.display.command.clone(), @@ -202,10 +237,103 @@ impl LeafExecutionReporter for SummaryLeafReporter { &self.workspace_path, ), }; + let upload_error = match &cache_update_status { + CacheUpdateStatus::Updated { upload_error } => Some(Arc::clone(upload_error)), + CacheUpdateStatus::NotUpdated(_) => None, + }; - self.tasks.borrow_mut().push(task_summary); + self.tasks.borrow_mut().push(RecordedTask { summary, upload_error }); } self.inner.finish(status, cache_update_status, error); } } + +#[cfg(test)] +mod tests { + use vt_plan::ExecutionItemKind; + + use super::*; + use crate::session::{ + cache::CacheMiss, + reporter::{ + InterleavedReporterBuilder, + test_fixtures::{spawn_task, test_path}, + }, + }; + + /// A writer whose output stays readable after it's moved into a reporter. + #[derive(Clone, Default)] + struct SharedBuffer(Rc>>); + + impl Write for SharedBuffer { + fn write(&mut self, buf: &[u8]) -> std::io::Result { + self.0.borrow_mut().extend_from_slice(buf); + Ok(buf.len()) + } + + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } + } + + impl SharedBuffer { + fn text(&self) -> Str { + let bytes = self.0.borrow(); + let text = std::str::from_utf8(&bytes).unwrap(); + vt_str::format!("{}", anstream::adapter::strip_str(text)) + } + } + + #[test] + fn upload_error_set_after_the_task_finishes_is_in_the_summary() { + let task = spawn_task("build"); + let item = &task.items[0]; + let ExecutionItemKind::Leaf(leaf_kind) = &item.kind else { + panic!("test fixture item must be a Leaf"); + }; + let output = SharedBuffer::default(); + let saved = SharedBuffer::default(); + let write_summary: WriteSummaryFn = Box::new({ + let mut saved = saved.clone(); + move |summary| saved.write_all(&format_compact_summary(summary, "vp")).unwrap() + }); + let mut reporter = Box::new(SummaryReporterBuilder::new( + Box::new(InterleavedReporterBuilder::new( + test_path(), + Box::new(std::io::sink()), + ColorSupport::uniform(false), + )), + test_path(), + Box::new(output.clone()), + false, + Some(write_summary), + Str::from("vp"), + ColorSupport::uniform(false), + )) + .build(); + + let upload_error = Arc::new(OnceLock::new()); + let mut leaf = reporter.new_leaf_execution(&item.execution_item_display, leaf_kind); + leaf.start(CacheStatus::Miss(CacheMiss::NotFound)); + leaf.finish( + Some(StdExitStatus::default()), + CacheUpdateStatus::Updated { upload_error: Arc::clone(&upload_error) }, + None, + ); + reporter.uploads_pending(NonZeroUsize::MIN); + upload_error.set(UploadError::Interrupted).unwrap(); + reporter.finish().unwrap(); + + let summary = "---\nvp run: pkg#build not uploaded to the remote cache: interrupted. \ + (Run `vp run --last-details` for full details)\n"; + assert_eq!(saved.text().as_str(), summary); + assert_eq!( + output.text().as_str(), + vt_str::format!( + "Waiting for 1 remote cache upload to finish (Ctrl-C to cancel)...\n{summary}" + ) + .as_str() + ); + } +} 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 16045e01c..549ce0e5d 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 @@ -1,6 +1,8 @@ # Cases that use the remote cache backend are ignored because it runs on # Node.js. Windows is skipped because the PTY launcher cannot execute pnpm -# command shims. +# command shims. Whether an upload to the backend is still running when the +# tasks finish depends on timing, so steps that upload hide the message about +# pending uploads. The cases that test it stall the uploads instead. [[e2e]] name = "read_write" cfg = "not(windows)" @@ -16,7 +18,11 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], - ], comment = "The fetch finds no entry. The new execution is uploaded with one store request." }, + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], + ], comment = "The fetch finds no entry. The new execution is uploaded with one store request, and vt run waits for it after the task finishes." }, { argv = [ "remote-cache-server", "vt", @@ -45,6 +51,10 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], ] }, [ "vtt", @@ -75,6 +85,10 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], ] }, [ "vt", @@ -143,6 +157,10 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], ] }, [ "vt", @@ -182,6 +200,10 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], ] }, [ "vt", @@ -218,6 +240,10 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], ] }, { argv = [ "vtt", @@ -266,6 +292,10 @@ steps = [ "VP_REMOTE_CACHE", "read-write", ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], ] }, [ "vt", @@ -316,6 +346,34 @@ steps = [ ], comment = "Nothing was cached locally, so the remote entry is fetched and restored again." }, ] +[[e2e]] +name = "pending_uploads" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "--stall", + "/store", + "vt", + "run", + "all", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ], interactions = [ + { "expect-milestone" = "uploads-pending" }, + { "write-key" = "ctrl-c" }, + ], comment = "The backend never responds to the uploads. check doesn't wait for build's upload, so both are still running when check finishes, and vt run waits for them until Ctrl-C." }, + { argv = [ + "vt", + "run", + "--last-details", + ], comment = "Both uploads were cancelled." }, +] + [[e2e]] name = "invalid_endpoint" steps = [ @@ -370,7 +428,11 @@ steps = [ "VP_REMOTE_CACHE_URL", "http://127.0.0.1:0/projects/test", ], - ], comment = "Nothing can listen on port 0. The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds." }, + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], + ], comment = "Nothing can listen on port 0. The failed fetch is the miss reason, and the upload fails in the background with a warning. The task succeeds." }, { argv = [ "vt", "run", @@ -444,3 +506,86 @@ steps = [ "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." }, ] + +[[e2e]] +name = "ctrl_c_during_upload" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "--stall", + "/store", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ], interactions = [ + { "expect-milestone" = "uploads-pending" }, + { "write-key" = "ctrl-c" }, + ], comment = "The fetch is a miss, but the backend never responds to the upload. Ctrl-C cancels it while vt run waits." }, + { argv = [ + "vt", + "run", + "--last-details", + ], comment = "The details show why build wasn't uploaded." }, + { argv = [ + "vt", + "run", + "build", + ], comment = "The entry is still in the local cache." }, +] + +[[e2e]] +name = "fast_fail_during_upload" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "--stall", + "/store", + "vt", + "run", + "fail-after-build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ], interactions = [ + { "expect-milestone" = "uploads-pending" }, + { "write-key" = "ctrl-c" }, + ], comment = "The backend never responds to the upload. fail-after-build exits after build finishes, which doesn't cancel build's upload, so vt run waits for it until Ctrl-C." }, +] + +[[e2e]] +name = "hide_pending_uploads" +cfg = "not(windows)" +ignore = true +steps = [ + { argv = [ + "remote-cache-server", + "--stall", + "/store", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + [ + "VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS", + "1", + ], + ], interactions = [ + { "expect-milestone" = "stalled" }, + { "write-key" = "ctrl-c" }, + ], comment = "The backend never responds to the upload. vt run waits for it until Ctrl-C, without the message about pending uploads." }, +] 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 7e9a4e659..e6f510756 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 @@ -1,6 +1,6 @@ # corrupt_archive -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` ``` $ vtt write-file dist/output.txt built diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_upload.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_upload.md new file mode 100644 index 000000000..dedbb99be --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_upload.md @@ -0,0 +1,56 @@ +# ctrl_c_during_upload + +## `VP_REMOTE_CACHE=read-write remote-cache-server --stall /store vt run build` + +The fetch is a miss, but the backend never responds to the upload. Ctrl-C cancels it while vt run waits. + +**→ expect-milestone:** `uploads-pending` + +``` +$ vtt write-file dist/output.txt built + +Waiting for 1 remote cache upload to finish (Ctrl-C to cancel)... +``` + +**← write-key:** `ctrl-c` + +``` +$ vtt write-file dist/output.txt built + +Waiting for 1 remote cache upload to finish (Ctrl-C to cancel)... +--- +vt run: remote-cache#build not uploaded to the remote cache: interrupted. (Run `vt run --last-details` for full details) +[remote-cache] POST /fetch 404 +``` + +## `vt run --last-details` + +The details show why build wasn't uploaded. + +``` + +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + Vite+ Task Runner • Execution Summary +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + +Statistics: 1 task • 0 cache hits • 1 cache miss +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 + ⚠ Not uploaded to the remote cache: interrupted +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ +``` + +## `vt run build` + +The entry is still in the local cache. + +``` +$ vtt write-file dist/output.txt built ◉ cache hit, replaying + +--- +vt run: cache hit. +``` 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 index 25f5c1d5a..d2f8cffcd 100644 --- 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 @@ -1,6 +1,6 @@ # fallback -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` ``` $ vtt write-file dist/output.txt built diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_upload.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_upload.md new file mode 100644 index 000000000..8becb1e60 --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_upload.md @@ -0,0 +1,30 @@ +# fast_fail_during_upload + +## `VP_REMOTE_CACHE=read-write remote-cache-server --stall /store vt run fail-after-build` + +The backend never responds to the upload. fail-after-build exits after build finishes, which doesn't cancel build's upload, so vt run waits for it until Ctrl-C. + +**Exit code:** 1 + +**→ expect-milestone:** `uploads-pending` + +``` +$ vtt write-file dist/output.txt built + +$ vtt exit 1 ⊘ cache disabled + +Waiting for 1 remote cache upload to finish (Ctrl-C to cancel)... +``` + +**← write-key:** `ctrl-c` + +``` +$ vtt write-file dist/output.txt built + +$ vtt exit 1 ⊘ cache disabled + +Waiting for 1 remote cache upload to finish (Ctrl-C to cancel)... +--- +vt run: 0/2 cache hit (0%), 1 failed. remote-cache#build not uploaded to the remote cache: interrupted. (Run `vt run --last-details` for full details) +[remote-cache] POST /fetch 404 +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/hide_pending_uploads.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/hide_pending_uploads.md new file mode 100644 index 000000000..e5a1a0a14 --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/hide_pending_uploads.md @@ -0,0 +1,21 @@ +# hide_pending_uploads + +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server --stall /store vt run build` + +The backend never responds to the upload. vt run waits for it until Ctrl-C, without the message about pending uploads. + +**→ expect-milestone:** `stalled` + +``` +$ vtt write-file dist/output.txt built +``` + +**← write-key:** `ctrl-c` + +``` +$ vtt write-file dist/output.txt built + +--- +vt run: remote-cache#build not uploaded to the remote cache: interrupted. (Run `vt run --last-details` for full details) +[remote-cache] POST /fetch 404 +``` diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/local_and_remote_hits.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/local_and_remote_hits.md index c0a4b271f..ce193c69c 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/local_and_remote_hits.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/local_and_remote_hits.md @@ -1,6 +1,6 @@ # local_and_remote_hits -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` ``` $ vtt write-file dist/output.txt built diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/pending_uploads.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/pending_uploads.md new file mode 100644 index 000000000..d9c7b2848 --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/pending_uploads.md @@ -0,0 +1,56 @@ +# pending_uploads + +## `VP_REMOTE_CACHE=read-write remote-cache-server --stall /store vt run all` + +The backend never responds to the uploads. check doesn't wait for build's upload, so both are still running when check finishes, and vt run waits for them until Ctrl-C. + +**→ expect-milestone:** `uploads-pending` + +``` +$ vtt write-file dist/output.txt built + +$ vtt print checked +checked + +Waiting for 2 remote cache uploads to finish (Ctrl-C to cancel)... +``` + +**← write-key:** `ctrl-c` + +``` +$ vtt write-file dist/output.txt built + +$ vtt print checked +checked + +Waiting for 2 remote cache uploads to finish (Ctrl-C to cancel)... +--- +vt run: 0/2 cache hit (0%). remote-cache#build (and 1 more) not uploaded to the remote cache: interrupted. (Run `vt run --last-details` for full details) +[remote-cache] POST /fetch 404 +[remote-cache] POST /fetch 404 +``` + +## `vt run --last-details` + +Both uploads were cancelled. + +``` + +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + Vite+ Task Runner • Execution Summary +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ + +Statistics: 2 tasks • 0 cache hits • 2 cache misses +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 + ⚠ Not uploaded to the remote cache: interrupted + ······················································· + [2] remote-cache#check: $ vtt print checked ✓ + → Cache miss: no previous cache entry found + ⚠ Not uploaded to the remote cache: interrupted +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ +``` 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 aa4c5fe2c..f18b8d8c2 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 @@ -1,6 +1,6 @@ # read -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` ``` $ vtt write-file dist/output.txt built 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 699c9871f..677abcc78 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 @@ -1,8 +1,8 @@ # read_write -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` -The fetch finds no entry. The new execution is uploaded with one store request. +The fetch finds no entry. The new execution is uploaded with one store request, and vt run waits for it after the task finishes. ``` $ vtt write-file dist/output.txt built 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 index 16e50049a..763d33941 100644 --- 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 @@ -1,6 +1,6 @@ # restore -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` ``` $ vtt write-file dist/output.txt built diff --git a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore_failure.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore_failure.md index 3012d5dcf..63f5937cd 100644 --- a/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore_failure.md +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/restore_failure.md @@ -1,6 +1,6 @@ # restore_failure -## `VP_REMOTE_CACHE=read-write remote-cache-server vt run build` +## `VP_REMOTE_CACHE=read-write VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 remote-cache-server vt run build` ``` $ vtt write-file dist/output.txt built 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 91283286a..da7b5bdcb 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 @@ -1,8 +1,8 @@ # unreachable_endpoint -## `VP_REMOTE_CACHE=read-write VP_REMOTE_CACHE_URL=http://127.0.0.1:0/projects/test vt run build` +## `VP_REMOTE_CACHE=read-write VP_REMOTE_CACHE_URL=http://127.0.0.1:0/projects/test VP_RUN_INTERNAL_HIDE_PENDING_UPLOADS=1 vt run build` -Nothing can listen on port 0. The failed fetch is the miss reason, and the failed upload is a warning. The task succeeds. +Nothing can listen on port 0. The failed fetch is the miss reason, and the upload fails in the background with a warning. The task succeeds. ``` $ vtt write-file dist/output.txt built ○ cache miss: remote cache fetch failed, executing 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 14be7f2eb..893757ae0 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 @@ -25,6 +25,11 @@ "command": "vtt print done", "dependsOn": ["build", "fail"], "cache": false + }, + "fail-after-build": { + "command": "vtt exit 1", + "dependsOn": ["build"], + "cache": false } } } diff --git a/crates/vt_remote_cache/Cargo.toml b/crates/vt_remote_cache/Cargo.toml index 750b76f06..e8bcd52e3 100644 --- a/crates/vt_remote_cache/Cargo.toml +++ b/crates/vt_remote_cache/Cargo.toml @@ -11,6 +11,7 @@ rust-version.workspace = true bytes = { workspace = true } ciborium = { workspace = true } reqwest = { workspace = true, features = [ + "http2", "multipart", "rustls-no-provider", "stream", diff --git a/crates/vt_remote_cache/README.md b/crates/vt_remote_cache/README.md index 150cc62fb..632729b08 100644 --- a/crates/vt_remote_cache/README.md +++ b/crates/vt_remote_cache/README.md @@ -12,6 +12,6 @@ Client for the [remote cache server API](https://github.com/voidzero-dev/vite-ta Only HTTP 200 counts as success for every operation, except that a fetch also accepts 404 as no match. Redirects aren't followed, so a redirect fails like any other status. -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. Requests go through the proxy set in `HTTPS_PROXY`, `HTTP_PROXY`, or `ALL_PROXY`, except for hosts in `NO_PROXY`. Without those variables, the system proxy settings are used on macOS and Windows. 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. +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. HTTPS requests use HTTP/2 if the server accepts it during the TLS handshake and HTTP/1.1 otherwise. HTTP endpoints always use HTTP/1.1. Requests go through the proxy set in `HTTPS_PROXY`, `HTTP_PROXY`, or `ALL_PROXY`, except for hosts in `NO_PROXY`. Without those variables, the system proxy settings are used on macOS and Windows. 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 and responses that end early), a status other than 200 (or 404, for a fetch), 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. Network errors leave out the request URL, since the endpoint may contain credentials, such as a token in its query. diff --git a/crates/vt_remote_cache/src/lib.rs b/crates/vt_remote_cache/src/lib.rs index c5f91b2f9..be6006fab 100644 --- a/crates/vt_remote_cache/src/lib.rs +++ b/crates/vt_remote_cache/src/lib.rs @@ -148,7 +148,8 @@ impl Client { // ring too. let _ = rustls::crypto::ring::default_provider().install_default(); // A redirect fails like any other status. Following one could turn a - // store into a GET of a login page that responds with 200. + // store into a GET of a login page that responds with 200. HTTPS + // endpoints use HTTP/2 if the server accepts it in the TLS handshake. let http = reqwest::Client::builder() .connect_timeout(CONNECT_TIMEOUT) .read_timeout(READ_TIMEOUT) @@ -629,6 +630,34 @@ mod tests { server.join().unwrap(); } + /// Accept one connection and return the TLS record it starts with, which + /// is the `ClientHello`, then close the connection without responding. + fn read_client_hello(listener: &TcpListener) -> Vec { + let (mut stream, _) = listener.accept().unwrap(); + let mut record = vec![0; 5]; + stream.read_exact(&mut record).unwrap(); + assert_eq!(record[0], 0x16, "not a TLS handshake record"); + let length = usize::from(u16::from_be_bytes([record[3], record[4]])); + record.resize(5 + length, 0); + stream.read_exact(&mut record[5..]).unwrap(); + record + } + + #[tokio::test] + async fn offers_http2_over_tls() { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let port = listener.local_addr().unwrap().port(); + let client = + Client::new(&vt_str::format!("https://127.0.0.1:{port}/projects/test")).unwrap(); + let server = std::thread::spawn(move || read_client_hello(&listener)); + + let error = client.fetch(b"k", b"s").await.unwrap_err(); + assert!(matches!(error, Error::Network(_)), "{error:?}"); + + // The ALPN extension lists h2, then http/1.1. + assert!(contains(&server.join().unwrap(), b"\x02h2\x08http/1.1")); + } + #[tokio::test] async fn network_errors_leave_out_the_url() { let client = diff --git a/docs/cancellation.md b/docs/cancellation.md index 63cb816f4..ca2ec3dc3 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, prevent caching of in-flight results, and stop [remote cache requests](#remote-cache-requests), 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 lookups](#remote-cache-lookups), but they differ in how they treat running processes and [remote cache uploads](#remote-cache-uploads). ## Ctrl-C @@ -18,12 +18,16 @@ 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 +## Remote cache lookups -Both kinds of cancellation stop remote cache lookups, downloads, and uploads right away instead of waiting for them to finish or time out. +Both kinds of cancellation stop remote cache lookups and downloads 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 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. +## Remote cache uploads + +In `read-write` mode, a task's upload to the remote cache starts once its result is cached locally and keeps running after the task finishes, so tasks that depend on it don't wait for it. Once all tasks are done, `vp run` waits for the uploads still running before it prints the summary, with a message saying how many there are. + +- Ctrl-C cancels them. Pressed while `vp run` waits, it cancels the uploads right away. Pressed while tasks are running, it cancels the uploads still running once the tasks stop, without the message. Either way, the tasks keep their local cache entries, and the summary warns that they weren't uploaded because they were interrupted. +- Fast-fail doesn't cancel them. A task that succeeded before another task failed is still uploaded. ## Why interrupted tasks are not cached diff --git a/packages/tools/README.md b/packages/tools/README.md index 6d815e236..c282c02e0 100644 --- a/packages/tools/README.md +++ b/packages/tools/README.md @@ -9,9 +9,11 @@ remote-cache-server cbor-http POST /store --form-cbor "metadata={\"key\": 'A', \ remote-cache-server cbor-http POST /fetch --cbor "{\"key\": 'A', \"secondary_key\": 'S'}" ``` -`remote-cache-server COMMAND [ARGS...]` starts the backend on a free loopback port and runs the command with `VP_REMOTE_CACHE_URL` set to the endpoint, `http://127.0.0.1:/projects/test`. The fixed base path gives every endpoint a namespace path. The wrapper takes no options and passes all arguments to the command unchanged. The command inherits stdio. When it exits, the server stops and the wrapper exits with the command's exit code. +`remote-cache-server [--stall ROUTE]... COMMAND [ARGS...]` starts the backend on a free loopback port and runs the command with `VP_REMOTE_CACHE_URL` set to the endpoint, `http://127.0.0.1:/projects/test`. The fixed base path gives every endpoint a namespace path. The wrapper passes all arguments after `COMMAND` to the command unchanged. The command inherits stdio and handles Ctrl-C, which the wrapper ignores. When it exits, the server stops and the wrapper exits with the command's exit code. -After the command exits, the wrapper prints one line to stderr for each request it served, in the order of the responses. Each line has the method, the path below the base path, and the status. Successful fetch responses add their kind: +`--stall ROUTE` makes the backend read requests to the route below the base path, such as `/store`, but never answer them. It can repeat. Each stalled request emits a `stalled` milestone for the E2E harness when it arrives. `vp run` uploads in the background, so `--stall /store` keeps uploads running after their tasks finish, until the client gives up. + +After the command exits, the wrapper prints one line to stderr for each request it answered, in the order of the responses. Each line has the method, the path below the base path, and the status. Successful fetch responses add their kind: ```text [remote-cache] POST /fetch 404 diff --git a/packages/tools/src/remote-cache/cli.ts b/packages/tools/src/remote-cache/cli.ts index b118e8b18..cea722844 100755 --- a/packages/tools/src/remote-cache/cli.ts +++ b/packages/tools/src/remote-cache/cli.ts @@ -1,22 +1,47 @@ #!/usr/bin/env node import { spawn } from 'node:child_process'; +import { randomBytes } from 'node:crypto'; import { once } from 'node:events'; import type { AddressInfo } from 'node:net'; import { createCacheServer } from './server.ts'; const basePath = '/projects/test'; -const [command, ...args] = process.argv.slice(2); +/** + * Emit a milestone for the E2E harness: a window title that + * `pty_terminal_test` recognizes, in the OSC 2 form it gets on Unix. + */ +function markMilestone(name: string): void { + const id = randomBytes(16).toString('hex'); + process.stdout.write( + `\x1b]2;pty-terminal-test:${id}:${Buffer.from(name).toString('base64url')}\x1b\\`, + ); +} + +const args = process.argv.slice(2); +const stalledRoutes = new Set(); +while (args[0] === '--stall') { + args.shift(); + const route = args.shift(); + if (route === undefined) + throw new Error('Usage: remote-cache-server [--stall ROUTE]... COMMAND [ARGS...]'); + stalledRoutes.add(route); +} +const [command, ...commandArgs] = args; const requests: string[] = []; const server = createCacheServer({ basePath, directory: 'remote-cache', logRequest: (line) => requests.push(line), + stalledRoutes, + onStall: () => markMilestone('stalled'), }); server.listen(0, '127.0.0.1'); await once(server, 'listening'); const { port } = server.address() as AddressInfo; -const child = spawn(command!, args, { +// Ctrl-C is left to the command. +process.on('SIGINT', () => {}); +const child = spawn(command!, commandArgs, { stdio: 'inherit', env: { ...process.env, VP_REMOTE_CACHE_URL: `http://127.0.0.1:${port}${basePath}` }, }); diff --git a/packages/tools/src/remote-cache/server.ts b/packages/tools/src/remote-cache/server.ts index 76a14e33a..6125bdc8c 100644 --- a/packages/tools/src/remote-cache/server.ts +++ b/packages/tools/src/remote-cache/server.ts @@ -100,16 +100,21 @@ function cbor(response: ServerResponse, value: unknown): void { * blobs remain opaque bytes. A fetch that matches neither key gets a plain-text * 404. After each response, `logRequest` receives a line with the method, the * route below `basePath`, the status, and for successful fetch responses, the - * kind. + * kind. A request to one of the `stalledRoutes` is read but never answered or + * logged, and `onStall` is called when it arrives. */ export function createCacheServer({ basePath, directory, logRequest, + stalledRoutes = new Set(), + onStall = () => {}, }: { basePath: string; directory: string; logRequest: (line: string) => void; + stalledRoutes?: ReadonlySet; + onStall?: () => void; }) { const stateFile = join(directory, 'state.json'); const blobDirectory = join(directory, 'blobs'); @@ -194,6 +199,14 @@ export function createCacheServer({ return createServer((request, response) => { const path = new URL(request.url ?? '/', 'http://localhost').pathname; + if (path.startsWith(`${basePath}/`) && stalledRoutes.has(path.slice(basePath.length))) { + // Read the body, so the client can finish sending it, and ignore the + // error when the client gives up and closes the connection. + request.on('error', () => {}); + request.resume(); + onStall(); + return; + } const log = (kind?: string) => { const parts = [request.method, path.slice(basePath.length), response.statusCode, kind]; logRequest(parts.filter((part) => part !== undefined).join(' '));