diff --git a/crates/vt_bin/src/vtt/stalled_remote_cache.rs b/crates/vt_bin/src/vtt/stalled_remote_cache.rs index ab3d391c9..eb468fa30 100644 --- a/crates/vt_bin/src/vtt/stalled_remote_cache.rs +++ b/crates/vt_bin/src/vtt/stalled_remote_cache.rs @@ -1,30 +1,55 @@ -use std::io::Read as _; +use std::{ + io::{Read as _, Write as _}, + net::{Shutdown, TcpListener, TcpStream}, + sync::Arc, +}; -/// stalled-remote-cache `` \[``...\] +/// stalled-remote-cache \[`--stall` ``\]... `` \[``...\] /// -/// 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. +/// Runs `` with `VP_REMOTE_CACHE_URL` set to a loopback proxy for the +/// endpoint in `VP_REMOTE_CACHE_URL`, which must be +/// `http://:/`, then exits with the command's exit code. +/// Requests to a stalled route below the endpoint, such as `/store`, are never +/// forwarded or answered: each emits a "stalled" milestone when it arrives, +/// and its connection is held until the client closes it. Other requests are +/// forwarded, each on its own connection. Ctrl-C is left to the command. +/// +/// On Windows a milestone is the console title, which `ConPTY` sends when it +/// next renders, so if the command sets one at about the same time, the +/// earlier title can be lost. pub fn run(args: &[String]) -> Result<(), Box> { + let mut args = args; + let mut stalled_routes = Vec::new(); + while let [flag, route, rest @ ..] = args + && flag == "--stall" + { + stalled_routes.push(route.clone()); + args = rest; + } let [program, args @ ..] = args else { - return Err("Usage: vtt stalled-remote-cache [args...]".into()); + return Err( + "Usage: vtt stalled-remote-cache [--stall ]... [args...]".into() + ); }; + let upstream = std::env::var("VP_REMOTE_CACHE_URL") + .map_err(|_| "VP_REMOTE_CACHE_URL must be set to the endpoint to proxy")?; + let (authority, path) = upstream + .strip_prefix("http://") + .and_then(|rest| rest.split_once('/')) + .ok_or("VP_REMOTE_CACHE_URL must be http://:/")?; + let proxy = Arc::new(Proxy { + upstream: authority.to_owned(), + base_path: std::format!("/{}", path.trim_end_matches('/')), + stalled_routes, + }); ctrlc::set_handler(|| {})?; - let listener = std::net::TcpListener::bind("127.0.0.1:0")?; - let endpoint = std::format!("http://{}/projects/test", listener.local_addr()?); + let listener = TcpListener::bind("127.0.0.1:0")?; + let endpoint = std::format!("http://{}{}", listener.local_addr()?, proxy.base_path); 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) {} - }); + for stream in listener.incoming().filter_map(Result::ok) { + let proxy = Arc::clone(&proxy); + std::thread::spawn(move || proxy.serve(stream)); } }); @@ -34,3 +59,106 @@ pub fn run(args: &[String]) -> Result<(), Box> { .status()?; std::process::exit(status.code().unwrap_or(1)); } + +struct Proxy { + /// `:` of the endpoint. + upstream: String, + /// The endpoint's path, without a trailing slash. + base_path: String, + /// Routes below `base_path` whose requests stall. + stalled_routes: Vec, +} + +impl Proxy { + /// Stall or forward the request on `client`. A forwarded request and its + /// response get `connection: close`, so the client sends its next request + /// on a new connection, which is served separately. + fn serve(&self, mut client: TcpStream) { + let Some((head, body_start)) = read_head(&mut client) else { + return; + }; + let target = head.split(' ').nth(1).unwrap_or_default(); + let path = target.split_once('?').map_or(target, |(path, _)| path); + if path + .strip_prefix(self.base_path.as_str()) + .is_some_and(|route| self.stalled_routes.iter().any(|stalled| stalled == route)) + { + pty_terminal_test_client::mark_milestone("stalled"); + let mut buf = [0; 4096]; + while client.read(&mut buf).is_ok_and(|n| n > 0) {} + return; + } + + let Ok(mut upstream) = TcpStream::connect(self.upstream.as_str()) else { + let _ = client.write_all( + b"HTTP/1.1 502 Bad Gateway\r\ncontent-length: 0\r\nconnection: close\r\n\r\n", + ); + return; + }; + let request_head = close_after_response(&head, Some(&self.upstream)); + if upstream.write_all(request_head.as_bytes()).is_err() + || upstream.write_all(&body_start).is_err() + { + return; + } + let (Ok(mut client_reader), Ok(mut upstream_writer)) = + (client.try_clone(), upstream.try_clone()) + else { + return; + }; + // The rest of the request body. + std::thread::spawn(move || { + let _ = std::io::copy(&mut client_reader, &mut upstream_writer); + let _ = upstream_writer.shutdown(Shutdown::Write); + }); + + if let Some((head, body_start)) = read_head(&mut upstream) + && client.write_all(close_after_response(&head, None).as_bytes()).is_ok() + && client.write_all(&body_start).is_ok() + { + let _ = std::io::copy(&mut upstream, &mut client); + } + let _ = client.shutdown(Shutdown::Both); + } +} + +/// Read an HTTP message's head, up to the blank line that ends it, and return +/// it with the bytes read after it. +fn read_head(stream: &mut TcpStream) -> Option<(String, Vec)> { + let mut bytes = Vec::new(); + let mut buf = [0; 4096]; + loop { + if let Some(end) = bytes.windows(4).position(|window| window == b"\r\n\r\n") { + let rest = bytes.split_off(end + 4); + return Some((String::from_utf8(bytes).ok()?, rest)); + } + match stream.read(&mut buf) { + Ok(n) if n > 0 => bytes.extend_from_slice(&buf[..n]), + _ => return None, + } + } +} + +/// `head` with `connection: close` in place of its `Connection` header, and +/// with `host` in place of its `Host` header if given. +fn close_after_response(head: &str, host: Option<&str>) -> String { + let mut lines = head.split("\r\n").filter(|line| !line.is_empty()); + let mut rewritten = std::format!("{}\r\n", lines.next().unwrap_or_default()); + for line in lines { + let name = line.split_once(':').map_or(line, |(name, _)| name).trim(); + if name.eq_ignore_ascii_case("connection") + || (host.is_some() && name.eq_ignore_ascii_case("host")) + { + continue; + } + rewritten.push_str(line); + rewritten.push_str("\r\n"); + } + if let Some(host) = host { + rewritten.push_str("host: "); + rewritten.push_str(host); + rewritten.push_str("\r\n"); + } + rewritten.push_str("connection: close\r\n\r\n"); + rewritten +} 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..10f5fd91c 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 @@ -424,13 +424,20 @@ steps = [ { argv = [ "vtt", "stalled-remote-cache", + "--stall", + "/fetch", "vt", "run", "build", + ], envs = [ + [ + "VP_REMOTE_CACHE_URL", + "http://127.0.0.1:0/projects/test", + ], ], interactions = [ - { "expect-milestone" = "request" }, + { "expect-milestone" = "stalled" }, { "write-key" = "ctrl-c" }, - ], comment = "The endpoint never responds. Ctrl-C stops the fetch, and the task doesn't start." }, + ], comment = "The proxy never answers the fetch, so nothing needs to listen behind it. Ctrl-C stops the fetch, and the task doesn't start." }, ] [[e2e]] @@ -439,8 +446,50 @@ steps = [ { argv = [ "vtt", "stalled-remote-cache", + "--stall", + "/fetch", "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." }, + ], envs = [ + [ + "VP_REMOTE_CACHE_URL", + "http://127.0.0.1:0/projects/test", + ], + ], comment = "The proxy never answers the fetch. 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", + "vtt", + "stalled-remote-cache", + "--stall", + "/store", + "vt", + "run", + "build", + ], envs = [ + [ + "VP_REMOTE_CACHE", + "read-write", + ], + ], interactions = [ + { "expect-milestone" = "stalled" }, + { "write-key" = "ctrl-c" }, + ], comment = "The proxy forwards the fetch to the backend, which has no entry, but never forwards the upload. Ctrl-C cancels it." }, + { 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." }, ] 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 index bd77d77bf..fe9a30759 100644 --- 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 @@ -1,10 +1,10 @@ # ctrl_c_during_fetch -## `vtt stalled-remote-cache vt run build` +## `VP_REMOTE_CACHE_URL=http://127.0.0.1:0/projects/test vtt stalled-remote-cache --stall /fetch vt run build` -The endpoint never responds. Ctrl-C stops the fetch, and the task doesn't start. +The proxy never answers the fetch, so nothing needs to listen behind it. Ctrl-C stops the fetch, and the task doesn't start. -**→ expect-milestone:** `request` +**→ expect-milestone:** `stalled` ``` ``` 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..251ca9a53 --- /dev/null +++ b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/ctrl_c_during_upload.md @@ -0,0 +1,53 @@ +# ctrl_c_during_upload + +## `VP_REMOTE_CACHE=read-write remote-cache-server vtt stalled-remote-cache --stall /store vt run build` + +The proxy forwards the fetch to the backend, which has no entry, but never forwards the upload. Ctrl-C cancels it. + +**→ 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: cancelled. (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: cancelled +━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━ +``` + +## `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/fast_fail_during_fetch.md b/crates/vt_bin/tests/e2e_snapshots/fixtures/remote_cache/snapshots/fast_fail_during_fetch.md index d7b81dd79..2946e7f65 100644 --- 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 @@ -1,8 +1,8 @@ # fast_fail_during_fetch -## `vtt stalled-remote-cache vt run fail-during-build` +## `VP_REMOTE_CACHE_URL=http://127.0.0.1:0/projects/test vtt stalled-remote-cache --stall /fetch 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. +The proxy never answers the fetch. fail exits while build's fetch is in flight, which stops the fetch, and build doesn't start. **Exit code:** 1 diff --git a/packages/tools/README.md b/packages/tools/README.md index 6d815e236..0e466a0b5 100644 --- a/packages/tools/README.md +++ b/packages/tools/README.md @@ -9,7 +9,7 @@ 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 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 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: diff --git a/packages/tools/src/remote-cache/cli.ts b/packages/tools/src/remote-cache/cli.ts index b118e8b18..82b00339a 100755 --- a/packages/tools/src/remote-cache/cli.ts +++ b/packages/tools/src/remote-cache/cli.ts @@ -16,6 +16,8 @@ const server = createCacheServer({ server.listen(0, '127.0.0.1'); await once(server, 'listening'); const { port } = server.address() as AddressInfo; +// Ctrl-C is left to the command. +process.on('SIGINT', () => {}); const child = spawn(command!, args, { stdio: 'inherit', env: { ...process.env, VP_REMOTE_CACHE_URL: `http://127.0.0.1:${port}${basePath}` },