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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
166 changes: 147 additions & 19 deletions crates/vt_bin/src/vtt/stalled_remote_cache.rs
Original file line number Diff line number Diff line change
@@ -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 `<command>` \[`<args>`...\]
/// stalled-remote-cache \[`--stall` `<route>`\]... `<command>` \[`<args>`...\]
///
/// Runs `<command>` 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 `<command>` with `VP_REMOTE_CACHE_URL` set to a loopback proxy for the
/// endpoint in `VP_REMOTE_CACHE_URL`, which must be
/// `http://<host>:<port>/<path>`, 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<dyn std::error::Error>> {
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 <command> [args...]".into());
return Err(
"Usage: vtt stalled-remote-cache [--stall <route>]... <command> [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://<host>:<port>/<path>")?;
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));
}
});

Expand All @@ -34,3 +59,106 @@ pub fn run(args: &[String]) -> Result<(), Box<dyn std::error::Error>> {
.status()?;
std::process::exit(status.code().unwrap_or(1));
}

struct Proxy {
/// `<host>:<port>` 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<String>,
}

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<u8>)> {
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
}
Original file line number Diff line number Diff line change
Expand Up @@ -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]]
Expand All @@ -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." },
]
Original file line number Diff line number Diff line change
@@ -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`

```
```
Expand Down
Original file line number Diff line number Diff line change
@@ -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.
```
Original file line number Diff line number Diff line change
@@ -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

Expand Down
2 changes: 1 addition & 1 deletion packages/tools/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:<port>/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:<port>/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:

Expand Down
2 changes: 2 additions & 0 deletions packages/tools/src/remote-cache/cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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}` },
Expand Down
Loading