mirror of
https://github.com/mountain-loop/yaak.git
synced 2026-08-25 21:04:04 +02:00
Add the browser send proxy and web sender (#572)
This commit is contained in:
@@ -0,0 +1,36 @@
|
||||
[package]
|
||||
name = "yaak-send-proxy"
|
||||
version = "0.1.0"
|
||||
edition = "2024"
|
||||
publish = false
|
||||
description = "Stateless HTTP send executor for Yaak in the browser"
|
||||
|
||||
# The send engine (yaak-http) and the model types it speaks (yaak-models, for
|
||||
# HttpRequest / Cookie / HttpResponseEventData). Deliberately NOT yaak (the
|
||||
# render + storage orchestration), yaak-plugins, or the RPC router: this binary
|
||||
# opens no database, runs no plugins, and renders nothing. yaak-models comes
|
||||
# along only because yaak-http's types are its types; nothing here calls into
|
||||
# its query layer.
|
||||
|
||||
[[bin]]
|
||||
name = "yaak-send-proxy"
|
||||
path = "src/main.rs"
|
||||
|
||||
[dependencies]
|
||||
async-trait = "0.1"
|
||||
axum = "0.7"
|
||||
base64 = "0.22.1"
|
||||
bytes = "1.11.1"
|
||||
clap = { version = "4.5", features = ["derive", "env"] }
|
||||
env_logger = "0.11"
|
||||
futures-util = "0.3"
|
||||
log = { workspace = true }
|
||||
serde = { workspace = true, features = ["derive"] }
|
||||
serde_json = { workspace = true }
|
||||
tokio = { workspace = true, features = ["rt-multi-thread", "macros", "signal", "sync", "io-util", "time", "net"] }
|
||||
tower-http = { version = "0.6", features = ["cors"] }
|
||||
ts-rs = { workspace = true }
|
||||
url = "2"
|
||||
uuid = { version = "1", features = ["v4"] }
|
||||
yaak-http = { workspace = true }
|
||||
yaak-models = { workspace = true }
|
||||
@@ -0,0 +1,139 @@
|
||||
# yaak-send-proxy
|
||||
|
||||
The network half of Yaak in a browser.
|
||||
|
||||
A tab can't see an HTTP response the way a desktop app can: CORS hides most
|
||||
headers (2 of 8 in a typical response), redirects are followed silently, and
|
||||
there is no timeline. So the tab renders the request and posts it here, and this
|
||||
process puts it on the network with the desktop's own engine (`yaak-http`) and
|
||||
streams back everything that happened — every header, every redirect hop, DNS
|
||||
timing, the body — for the tab to store.
|
||||
|
||||
It is a **stateless executor**. It keeps nothing: no database, no files, no
|
||||
sessions, no cookies between calls. Every byte it sees comes from the tab in the
|
||||
request, and every byte it returns is stored by the tab. Restart it any time.
|
||||
|
||||
## Running it
|
||||
|
||||
```shell
|
||||
cargo run -p yaak-send-proxy
|
||||
```
|
||||
|
||||
Listens on `127.0.0.1:9227`. Then run the web build against it:
|
||||
|
||||
```shell
|
||||
YAAK_TARGET=web npm run dev --workspace @yaakapp/yaak-client
|
||||
```
|
||||
|
||||
The tab looks for the proxy at `http://127.0.0.1:9227` unless
|
||||
`VITE_YAAK_SEND_PROXY_URL` says otherwise at build time.
|
||||
|
||||
Every flag has a `YAAK_PROXY_*` environment variable, so a container needs no
|
||||
arguments; `--help` lists them all.
|
||||
|
||||
| Flag | Default | What |
|
||||
| --- | --- | --- |
|
||||
| `--bind` | `127.0.0.1:9227` | Listen address. `0.0.0.0:9227` inside a container. |
|
||||
| `--allowed-origins` | `*` | CORS origins, comma-separated. A hosted instance should name its web origin. |
|
||||
| `--max-request-bytes` | 16 MiB | Largest rendered request accepted from the tab. |
|
||||
| `--max-response-bytes` | 64 MiB | Largest upstream body relayed before the send is cut off. |
|
||||
| `--max-timeout-secs` | 60 | Ceiling on a send's timeout; a request asking for more (or none) gets this. |
|
||||
| `--rate-limit-per-minute` | 120 | Sends per client IP per minute; 0 disables. |
|
||||
| `--max-concurrent` | 256 | Sends in flight at once. |
|
||||
| `--trust-forwarded-for` | off | Take the client IP from `X-Forwarded-For`. Only behind a load balancer that sets it. |
|
||||
|
||||
## What it refuses, and why
|
||||
|
||||
A hosted proxy is, by construction, a machine that makes HTTP requests on
|
||||
behalf of strangers. Left alone that is an open relay into whatever network it
|
||||
sits on. So it refuses, always, to connect to:
|
||||
|
||||
- loopback (`127/8`, `::1`), private (`10/8`, `172.16/12`, `192.168/16`,
|
||||
`fc00::/7`), link-local (`169.254/16` — where cloud metadata lives — and
|
||||
`fe80::/10`), carrier-grade NAT, multicast, reserved and unspecified ranges,
|
||||
IPv4 addresses carried inside IPv6 forms (`::ffff:a.b.c.d`, the well-known
|
||||
NAT64 prefix, 6to4), and the whole NAT64 local-use range;
|
||||
- anything not `http://` or `https://`.
|
||||
|
||||
The check runs **on the resolved addresses, after DNS**, for every hop of a
|
||||
redirect chain, so a public hostname that points at an internal address is
|
||||
caught, and so is a `Location:` header that points at one. It also refuses body
|
||||
types that would read files on the proxy's disk (`binary`, multipart file
|
||||
fields), since no browser tab could legitimately mean those.
|
||||
|
||||
Refusals are logged with the reason. There is no switch to turn this off: the
|
||||
proxy's private network is the cloud's, not the user's, so a `localhost` or LAN
|
||||
API can never be reached through it — that is what the desktop app is for.
|
||||
|
||||
## Deploying
|
||||
|
||||
One binary, no dependencies:
|
||||
|
||||
```shell
|
||||
cargo build --release -p yaak-send-proxy
|
||||
YAAK_PROXY_BIND=0.0.0.0:9227 \
|
||||
YAAK_PROXY_ALLOWED_ORIGINS=https://yaak.example.com \
|
||||
./target/release/yaak-send-proxy
|
||||
```
|
||||
|
||||
There is no authentication: an instance is anonymous and protected by the
|
||||
per-client rate limit and the destination policy, which is what the hosted
|
||||
funnel wants. Anything more (a shared token, per-user quotas) is a later slice
|
||||
and would sit in front of `send_http` in `main.rs`. Put TLS in front of it (a
|
||||
reverse proxy). If the reverse proxy buffers responses, tell it not to: the reply
|
||||
is a stream and the `X-Accel-Buffering: no` header it sets is honoured by
|
||||
nginx-shaped ones.
|
||||
|
||||
## The wire
|
||||
|
||||
`POST /v1/http/send` with a JSON body:
|
||||
|
||||
```json
|
||||
{
|
||||
"request": { "url": "https://…", "method": "GET", "headers": […], "body": {…}, "bodyType": null, "urlParameters": […] },
|
||||
"settings": { "validateCertificates": true, "followRedirects": true, "timeoutMs": 0, "sendCookies": true, "storeCookies": true },
|
||||
"cookies": [ … ]
|
||||
}
|
||||
```
|
||||
|
||||
`request` is a Yaak `HttpRequest` in the desktop's own model shape with every
|
||||
template already rendered by the tab; the proxy builds the URL, headers and
|
||||
body from it exactly the way the desktop does after rendering. `cookies` is the
|
||||
jar's contents (or `null` for no jar).
|
||||
|
||||
The reply is `application/x-ndjson`, one JSON frame per line, in the order things
|
||||
happened:
|
||||
|
||||
| `type` | When | Carries |
|
||||
| --- | --- | --- |
|
||||
| `event` | as the engine produces them | one timeline event, in the desktop's `http_response_event.event` shape |
|
||||
| `response` | once, when the final hop's headers arrive | status, all headers, request headers as sent, remote address, HTTP version, timing |
|
||||
| `body` | as the body is read | a decompressed chunk, base64 |
|
||||
| `done` | last, on success | elapsed, byte counts, and the cookie jar as the send left it |
|
||||
| `error` | last, on failure | the reason, and any cookies collected before the failure |
|
||||
|
||||
Refusals that happen before anything is sent (a blocked destination, a bad body,
|
||||
rate limit, capacity) are plain HTTP errors (`403`, `400`, `429`, `503`) with
|
||||
`{"error": "…"}`, not streams.
|
||||
|
||||
Why a streamed HTTP response and not a WebSocket: one `POST` is stateless by
|
||||
construction, cancellable by closing the connection, readable with `curl`, and
|
||||
needs no upgrade handling on either side. A WebSocket only earns its keep when
|
||||
traffic is bidirectional, which a single send is not.
|
||||
|
||||
The TypeScript side of this contract is generated from `src/wire.rs` by ts-rs
|
||||
into `bindings/` (run `cargo test -p yaak-send-proxy` after changing a frame)
|
||||
and published to the tab as `@yaakapp-internal/send-proxy`, so a change to the
|
||||
wire on one side is a type error on the other.
|
||||
|
||||
`GET /v1/health` reports the version and the effective limits.
|
||||
|
||||
## What comes later
|
||||
|
||||
Not built, by design, but the router is shaped for it: a WebSocket relay
|
||||
(`/v1/ws/relay`) and a gRPC relay (`/v1/grpc/relay`) would be long-lived,
|
||||
bidirectional endpoints on the same binary, behind the same destination policy
|
||||
and limits. They differ from this endpoint in holding per-connection
|
||||
in-memory state while a connection is open (never persisted), which brings
|
||||
connection limits and a larger abuse surface — the reason they are separate
|
||||
work.
|
||||
@@ -0,0 +1,48 @@
|
||||
// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually.
|
||||
|
||||
export type Cookie = { name: string, value: string, domain: CookieDomain, expires: CookieExpires, path: string, secure: boolean, httpOnly: boolean, sameSite: CookieSameSite | null, };
|
||||
|
||||
export type CookieDomain = { "HostOnly": string } | { "Suffix": string } | "NotPresent" | "Empty";
|
||||
|
||||
export type CookieExpires = { "AtUtc": string } | "SessionEnd";
|
||||
|
||||
export type CookieSameSite = "Strict" | "Lax" | "None";
|
||||
|
||||
export type HttpRequest = { model: "http_request", id: string, createdAt: string, updatedAt: string, workspaceId: string, folderId: string | null, authentication: Record<string, any>, authenticationType: string | null, body: Record<string, any>, bodyType: string | null, description: string, headers: Array<HttpRequestHeader>, method: string, name: string, sortPriority: number, url: string,
|
||||
/**
|
||||
* URL parameters used for both path placeholders (`:id`) and query string entries.
|
||||
*/
|
||||
urlParameters: Array<HttpUrlParameter>, settingSendCookies: InheritedBoolSetting, settingStoreCookies: InheritedBoolSetting, settingValidateCertificates: InheritedBoolSetting, settingFollowRedirects: InheritedBoolSetting, settingRequestTimeout: InheritedIntSetting, };
|
||||
|
||||
export type HttpRequestHeader = { enabled?: boolean, name: string, value: string, id?: string, };
|
||||
|
||||
/**
|
||||
* Serializable representation of HTTP response events for DB storage.
|
||||
* This mirrors `yaak_http::sender::HttpResponseEvent` but with serde support.
|
||||
* The `From` impl is in yaak-http to avoid circular dependencies.
|
||||
*/
|
||||
export type HttpResponseEventData = { "type": "setting", name: string, value: string, source_model?: string, source_id?: string, source_name?: string, } | { "type": "info", message: string, } | { "type": "redirect", url: string, status: number, behavior: string, dropped_body: boolean, dropped_headers: Array<string>, } | { "type": "send_url", method: string, scheme: string, username: string, password: string, host: string, port: number, path: string, query: string, fragment: string, } | { "type": "receive_url", version: string, status: string, } | { "type": "header_up", name: string, value: string, } | { "type": "header_down", name: string, value: string, } | { "type": "chunk_sent", bytes: number, } | { "type": "chunk_received", bytes: number, } | { "type": "dns_resolved", hostname: string, addresses: Array<string>, duration: bigint, overridden: boolean, };
|
||||
|
||||
export type HttpResponseHeader = { name: string, value: string, };
|
||||
|
||||
/**
|
||||
* The resolved send settings, values only: what an executor has to obey, with the sources
|
||||
* (which model each came from) left behind in [`ResolvedHttpRequestSettings`]. This is what
|
||||
* crosses from a tab to the send proxy, and what the proxy reads.
|
||||
*/
|
||||
export type HttpSendSettings = { validateCertificates: boolean, followRedirects: boolean,
|
||||
/**
|
||||
* Milliseconds. Zero or negative means no timeout.
|
||||
*/
|
||||
timeoutMs: number, sendCookies: boolean, storeCookies: boolean, };
|
||||
|
||||
export type HttpUrlParameter = { enabled?: boolean,
|
||||
/**
|
||||
* Colon-prefixed parameters are treated as path parameters if they match, like `/users/:id`
|
||||
* Other entries are appended as query parameters
|
||||
*/
|
||||
name: string, value: string, id?: string, };
|
||||
|
||||
export type InheritedBoolSetting = { enabled?: boolean, value: boolean, };
|
||||
|
||||
export type InheritedIntSetting = { enabled?: boolean, value: number, };
|
||||
@@ -0,0 +1,64 @@
|
||||
// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually.
|
||||
import type { Cookie, HttpRequest, HttpResponseEventData, HttpResponseHeader, HttpSendSettings } from "./gen_models";
|
||||
|
||||
/**
|
||||
* One line of the reply stream. Tags are snake_case like the timeline event tags; fields are
|
||||
* camelCase like every model the tab stores.
|
||||
*/
|
||||
export type Frame = { "type": "event", event: HttpResponseEventData, } | { "type": "response", status: number, statusReason: string | null,
|
||||
/**
|
||||
* The URL that answered, after redirects.
|
||||
*/
|
||||
url: string, remoteAddr: string | null, version: string | null, headers: Array<HttpResponseHeader>,
|
||||
/**
|
||||
* The headers that were actually sent on the final hop, cookies and all.
|
||||
*/
|
||||
requestHeaders: Array<HttpResponseHeader>,
|
||||
/**
|
||||
* `Content-Length` as declared by the server, if it declared one.
|
||||
*/
|
||||
contentLength: number | null,
|
||||
/**
|
||||
* Milliseconds from the start of the send to the response head.
|
||||
*/
|
||||
elapsedHeaders: number,
|
||||
/**
|
||||
* Milliseconds spent in DNS on the last lookup, or zero.
|
||||
*/
|
||||
elapsedDns: number, } | { "type": "body", data: string, } | { "type": "done",
|
||||
/**
|
||||
* Milliseconds from the start of the send to the end of the body.
|
||||
*/
|
||||
elapsed: number,
|
||||
/**
|
||||
* Bytes of body relayed, after decompression.
|
||||
*/
|
||||
contentLength: number,
|
||||
/**
|
||||
* Bytes on the wire as declared by the server, or the relayed size when unknown.
|
||||
*/
|
||||
contentLengthCompressed: number,
|
||||
/**
|
||||
* The jar as the send left it, for the tab to persist. `None` when the tab sent none.
|
||||
*/
|
||||
cookies: Array<Cookie> | null, } | { "type": "error", message: string, cookies: Array<Cookie> | null, };
|
||||
|
||||
/**
|
||||
* The body of `POST /v1/http/send`.
|
||||
*/
|
||||
export type SendRequest = {
|
||||
/**
|
||||
* The request to send, in the desktop's own model shape but with every template already
|
||||
* rendered by the tab. The proxy builds the URL, headers and body from it exactly the way
|
||||
* the desktop does after rendering.
|
||||
*/
|
||||
request: HttpRequest,
|
||||
/**
|
||||
* The resolved settings, values only. Where they came from is the tab's to record in
|
||||
* its timeline; the proxy only needs to obey them.
|
||||
*/
|
||||
settings: HttpSendSettings,
|
||||
/**
|
||||
* The cookies to start with. `None` means no jar at all: nothing sent, nothing kept.
|
||||
*/
|
||||
cookies: Array<Cookie> | null, };
|
||||
@@ -0,0 +1,4 @@
|
||||
// The send proxy's wire contract, generated by ts-rs from src/wire.rs
|
||||
// (`cargo test -p yaak-send-proxy`). The tab imports these so a change to a
|
||||
// frame on the Rust side is a type error in packages/platform/src/web.
|
||||
export type { Frame, SendRequest } from "./bindings/gen_send_proxy";
|
||||
@@ -0,0 +1,6 @@
|
||||
{
|
||||
"name": "@yaakapp-internal/send-proxy",
|
||||
"version": "1.0.0",
|
||||
"private": true,
|
||||
"main": "index.ts"
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
use clap::Parser;
|
||||
use std::net::SocketAddr;
|
||||
|
||||
/// A stateless HTTP send executor for Yaak running in a browser.
|
||||
///
|
||||
/// The tab renders the request and owns the data; this binary only puts bytes on the
|
||||
/// network and streams back what came back. Nothing is written to disk or a database.
|
||||
#[derive(Parser, Debug, Clone)]
|
||||
#[command(name = "yaak-send-proxy", version, about, long_about = None)]
|
||||
pub struct Config {
|
||||
/// Address to listen on. 127.0.0.1 for a local instance; 0.0.0.0 inside a container.
|
||||
#[arg(long, env = "YAAK_PROXY_BIND", default_value = "127.0.0.1:9227")]
|
||||
pub bind: SocketAddr,
|
||||
|
||||
/// Browser origins allowed to call this proxy (CORS), comma-separated. `*` allows any.
|
||||
/// A local dev instance wants the Vite origin; a hosted instance wants its own web origin.
|
||||
#[arg(
|
||||
long,
|
||||
env = "YAAK_PROXY_ALLOWED_ORIGINS",
|
||||
default_value = "*",
|
||||
value_delimiter = ','
|
||||
)]
|
||||
pub allowed_origins: Vec<String>,
|
||||
|
||||
/// Largest request the proxy accepts from the tab (the rendered request JSON, body included).
|
||||
#[arg(long, env = "YAAK_PROXY_MAX_REQUEST_BYTES", default_value_t = 16 * 1024 * 1024)]
|
||||
pub max_request_bytes: usize,
|
||||
|
||||
/// Largest upstream response body the proxy will relay before cutting the send off.
|
||||
#[arg(long, env = "YAAK_PROXY_MAX_RESPONSE_BYTES", default_value_t = 64 * 1024 * 1024)]
|
||||
pub max_response_bytes: usize,
|
||||
|
||||
/// Ceiling on a send's timeout, in seconds. A request asking for longer (or for no timeout)
|
||||
/// gets this instead.
|
||||
#[arg(long, env = "YAAK_PROXY_MAX_TIMEOUT_SECS", default_value_t = 60)]
|
||||
pub max_timeout_secs: u64,
|
||||
|
||||
/// Sends allowed per client IP per minute. 0 disables the limit. This and the concurrency
|
||||
/// cap are the whole of what protects an instance: there is no authentication.
|
||||
#[arg(long, env = "YAAK_PROXY_RATE_LIMIT_PER_MINUTE", default_value_t = 120)]
|
||||
pub rate_limit_per_minute: u32,
|
||||
|
||||
/// Sends in flight at once across all clients.
|
||||
#[arg(long, env = "YAAK_PROXY_MAX_CONCURRENT", default_value_t = 256)]
|
||||
pub max_concurrent: usize,
|
||||
|
||||
/// Take the client IP from `X-Forwarded-For` (first hop) instead of the socket. Only turn
|
||||
/// this on behind a load balancer that sets the header; otherwise anyone can spoof their way
|
||||
/// past the rate limit.
|
||||
#[arg(long, env = "YAAK_PROXY_TRUST_FORWARDED_FOR", default_value_t = false)]
|
||||
pub trust_forwarded_for: bool,
|
||||
}
|
||||
@@ -0,0 +1,252 @@
|
||||
//! Where a send may go.
|
||||
//!
|
||||
//! A hosted sender is, by construction, a machine that makes HTTP requests on
|
||||
//! behalf of strangers. Left alone that is an open relay into whatever network
|
||||
//! it sits on: cloud metadata endpoints, internal admin panels, the database
|
||||
//! next door. So every destination is checked twice — once on the URL before a
|
||||
//! hop is attempted (literal IPs, host allow/deny lists) and once on the
|
||||
//! addresses a hostname actually resolves to, right before the connection is
|
||||
//! made. The second check is the one that matters for a hostname pointing at
|
||||
//! an internal address, and it runs on every redirect hop because the engine
|
||||
//! resolves every hop.
|
||||
|
||||
use async_trait::async_trait;
|
||||
use log::warn;
|
||||
use std::net::{IpAddr, Ipv4Addr, Ipv6Addr};
|
||||
use std::sync::Arc;
|
||||
use tokio::sync::mpsc;
|
||||
use url::Url;
|
||||
use yaak_http::dns::AddressFilter;
|
||||
use yaak_http::sender::{HttpResponse, HttpResponseEvent, HttpSender};
|
||||
use yaak_http::types::SendableHttpRequest;
|
||||
|
||||
/// The destination policy, shared by every send: public addresses only, always. A hosted
|
||||
/// proxy's "private network" is the cloud's, not the user's, so there is no configuration
|
||||
/// that makes reaching it right.
|
||||
#[derive(Clone, Default)]
|
||||
pub struct DestinationPolicy;
|
||||
|
||||
impl DestinationPolicy {
|
||||
/// Check a URL before a hop is attempted: scheme and literal IPs. A hostname that passes
|
||||
/// here still has its resolved addresses checked by [`Self::address_filter`].
|
||||
pub fn check_url(&self, raw: &str) -> Result<(), String> {
|
||||
let url = Url::parse(raw).map_err(|e| format!("Invalid URL {raw:?}: {e}"))?;
|
||||
match url.scheme() {
|
||||
"http" | "https" => {}
|
||||
other => return Err(format!("Refusing to send over {other:?}; only http and https")),
|
||||
}
|
||||
let host = url.host_str().ok_or_else(|| format!("URL {raw:?} has no host"))?;
|
||||
let host = host.trim_matches(|c| c == '[' || c == ']');
|
||||
|
||||
// A literal IP never reaches the resolver, so it is checked here. Hostnames are checked
|
||||
// where their addresses become known.
|
||||
if let Ok(ip) = host.parse::<IpAddr>() {
|
||||
self.check_ip(ip)?;
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// The veto the engine's resolver applies to every address a hostname resolves to.
|
||||
pub fn address_filter(&self) -> AddressFilter {
|
||||
let policy = self.clone();
|
||||
Arc::new(move |ip| policy.check_ip(ip))
|
||||
}
|
||||
|
||||
pub fn check_ip(&self, ip: IpAddr) -> Result<(), String> {
|
||||
match non_public_reason(ip) {
|
||||
Some(reason) => Err(format!(
|
||||
"Refusing to connect to {ip}: {reason}. This proxy only sends to public addresses"
|
||||
)),
|
||||
None => Ok(()),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Why an address is not a public internet address, or `None` if it is one.
|
||||
///
|
||||
/// Every range here is one a hosted relay must never be talked into reaching: the machine
|
||||
/// itself, the network it sits on, and the link-local range where cloud metadata services
|
||||
/// (169.254.169.254) live. IPv4 addresses carried inside fixed-layout IPv6 forms — IPv4-mapped,
|
||||
/// the well-known NAT64 prefix, 6to4 — are unwrapped and judged as IPv4, since that is where
|
||||
/// the packets end up; the NAT64 local-use range is refused outright. This is the stable-Rust
|
||||
/// stand-in for `IpAddr::is_global`, which is still behind `#![feature(ip)]`; a network-specific
|
||||
/// NAT64 prefix is not knowable here.
|
||||
pub fn non_public_reason(ip: IpAddr) -> Option<&'static str> {
|
||||
match ip {
|
||||
IpAddr::V4(v4) => non_public_v4(v4),
|
||||
IpAddr::V6(v6) => {
|
||||
if let Some(v4) = v6.to_ipv4_mapped() {
|
||||
return non_public_v4(v4);
|
||||
}
|
||||
if let Some(v4) = embedded_v4(&v6) {
|
||||
return non_public_v4(v4);
|
||||
}
|
||||
if v6.is_loopback() {
|
||||
Some("loopback")
|
||||
} else if v6.is_unspecified() {
|
||||
Some("unspecified")
|
||||
} else if v6.is_unique_local() {
|
||||
Some("unique local (fc00::/7)")
|
||||
} else if v6.is_unicast_link_local() {
|
||||
Some("link-local (fe80::/10)")
|
||||
} else if v6.is_multicast() {
|
||||
Some("multicast")
|
||||
} else if v6.segments()[..3] == [0x64, 0xff9b, 1] {
|
||||
Some("NAT64 local-use (64:ff9b:1::/48)")
|
||||
} else if v6.segments()[..4] == [0x100, 0, 0, 0] {
|
||||
Some("discard-only (100::/64)")
|
||||
} else if (v6.segments()[0] & 0xffc0) == 0xfec0 {
|
||||
Some("site-local (fec0::/10)")
|
||||
} else if v6.segments()[0] == 0x2001 && v6.segments()[1] == 0x0db8 {
|
||||
Some("documentation (2001:db8::/32)")
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn non_public_v4(v4: Ipv4Addr) -> Option<&'static str> {
|
||||
let o = v4.octets();
|
||||
if v4.is_loopback() {
|
||||
Some("loopback (127.0.0.0/8)")
|
||||
} else if v4.is_private() {
|
||||
Some("private (10/8, 172.16/12, 192.168/16)")
|
||||
} else if v4.is_link_local() {
|
||||
Some("link-local (169.254.0.0/16, where cloud metadata lives)")
|
||||
} else if v4.is_unspecified() || o[0] == 0 {
|
||||
Some("this network (0.0.0.0/8)")
|
||||
} else if o[0] == 100 && (o[1] & 0xc0) == 64 {
|
||||
Some("carrier-grade NAT (100.64.0.0/10)")
|
||||
} else if v4.is_broadcast() {
|
||||
Some("broadcast")
|
||||
} else if v4.is_multicast() {
|
||||
Some("multicast (224.0.0.0/4)")
|
||||
} else if o[0] >= 240 {
|
||||
Some("reserved (240.0.0.0/4)")
|
||||
} else if v4.is_documentation() {
|
||||
Some("documentation")
|
||||
} else if o[0] == 192 && o[1] == 0 && o[2] == 0 {
|
||||
Some("IETF protocol assignments (192.0.0.0/24)")
|
||||
} else if o[0] == 198 && (o[1] & 0xfe) == 18 {
|
||||
Some("benchmarking (198.18.0.0/15)")
|
||||
} else {
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
/// The IPv4 address an IPv6 address stands for, when it is one of the fixed-layout translation
|
||||
/// forms: the NAT64 well-known prefix (64:ff9b::/96) or 6to4 (2002::/16, IPv4 in the next 32
|
||||
/// bits). The NAT64 local-use range (64:ff9b:1::/48) is a pool operators carve their own
|
||||
/// prefix from, at a length only they know, so it is refused wholesale in [`non_public_reason`]
|
||||
/// rather than decoded — the same call `std`'s (still unstable) `Ipv6Addr::is_global` makes.
|
||||
fn embedded_v4(v6: &Ipv6Addr) -> Option<Ipv4Addr> {
|
||||
let s = v6.segments();
|
||||
let o = v6.octets();
|
||||
if s[0] == 0x64 && s[1] == 0xff9b && s[2..6].iter().all(|x| *x == 0) {
|
||||
return Some(Ipv4Addr::new(o[12], o[13], o[14], o[15]));
|
||||
}
|
||||
if s[0] == 0x2002 {
|
||||
return Some(Ipv4Addr::new(o[2], o[3], o[4], o[5]));
|
||||
}
|
||||
None
|
||||
}
|
||||
|
||||
/// An [`HttpSender`] that checks each hop's URL against the policy before delegating.
|
||||
///
|
||||
/// The engine's redirect loop calls the sender once per hop with the hop's URL, so wrapping
|
||||
/// the sender is what makes `Location:` headers subject to the same rules as the first URL —
|
||||
/// including a redirect to a literal internal IP, which the resolver would never see.
|
||||
pub struct GuardedSender<S> {
|
||||
inner: S,
|
||||
policy: DestinationPolicy,
|
||||
}
|
||||
|
||||
impl<S: HttpSender> GuardedSender<S> {
|
||||
pub fn new(inner: S, policy: DestinationPolicy) -> Self {
|
||||
Self { inner, policy }
|
||||
}
|
||||
}
|
||||
|
||||
#[async_trait]
|
||||
impl<S: HttpSender> HttpSender for GuardedSender<S> {
|
||||
async fn send(
|
||||
&self,
|
||||
request: SendableHttpRequest,
|
||||
event_tx: mpsc::Sender<HttpResponseEvent>,
|
||||
) -> yaak_http::error::Result<HttpResponse> {
|
||||
if let Err(reason) = self.policy.check_url(&request.url) {
|
||||
warn!("Refused {} {}: {reason}", request.method, request.url);
|
||||
return Err(yaak_http::error::Error::RequestError(reason));
|
||||
}
|
||||
self.inner.send(request, event_tx).await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn ip(s: &str) -> IpAddr {
|
||||
s.parse().unwrap()
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn refuses_the_ranges_a_relay_must_never_reach() {
|
||||
for addr in [
|
||||
"127.0.0.1",
|
||||
"127.9.9.9",
|
||||
"10.0.0.1",
|
||||
"172.16.0.1",
|
||||
"172.31.255.255",
|
||||
"192.168.1.1",
|
||||
"169.254.169.254",
|
||||
"169.254.0.1",
|
||||
"0.0.0.0",
|
||||
"100.64.0.1",
|
||||
"255.255.255.255",
|
||||
"224.0.0.1",
|
||||
"240.0.0.1",
|
||||
"::1",
|
||||
"::",
|
||||
"fc00::1",
|
||||
"fd12::1",
|
||||
"fe80::1",
|
||||
"::ffff:127.0.0.1",
|
||||
"::ffff:169.254.169.254",
|
||||
"64:ff9b::7f00:1",
|
||||
"ff02::1",
|
||||
] {
|
||||
assert!(non_public_reason(ip(addr)).is_some(), "{addr} should be refused");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn allows_public_addresses() {
|
||||
for addr in [
|
||||
"1.1.1.1",
|
||||
"8.8.8.8",
|
||||
"93.184.216.34",
|
||||
"172.32.0.1",
|
||||
"2606:4700:4700::1111",
|
||||
] {
|
||||
assert!(non_public_reason(ip(addr)).is_none(), "{addr} should be allowed");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn literal_private_addresses_in_urls_are_refused() {
|
||||
let policy = DestinationPolicy;
|
||||
assert!(policy.check_url("http://127.0.0.1/").is_err());
|
||||
assert!(policy.check_url("http://[::1]/").is_err());
|
||||
assert!(policy.check_url("http://169.254.169.254/latest/meta-data").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_http_schemes() {
|
||||
let policy = DestinationPolicy;
|
||||
assert!(policy.check_url("ftp://example.com/").is_err());
|
||||
assert!(policy.check_url("file:///etc/passwd").is_err());
|
||||
assert!(policy.check_url("https://example.com/").is_ok());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
//! Per-client rate limiting, kept deliberately small.
|
||||
//!
|
||||
//! One token bucket per client IP, refilled continuously, in a mutex-guarded
|
||||
//! map that is swept of idle entries as it goes. Good enough to keep one
|
||||
//! caller from monopolising a hosted instance; not a substitute for whatever
|
||||
//! sits in front of it in production.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::net::IpAddr;
|
||||
use std::sync::Mutex;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
pub struct RateLimiter {
|
||||
per_minute: u32,
|
||||
buckets: Mutex<HashMap<IpAddr, Bucket>>,
|
||||
}
|
||||
|
||||
struct Bucket {
|
||||
tokens: f64,
|
||||
last: Instant,
|
||||
}
|
||||
|
||||
impl RateLimiter {
|
||||
/// `per_minute == 0` disables limiting.
|
||||
pub fn new(per_minute: u32) -> Self {
|
||||
Self { per_minute, buckets: Mutex::new(HashMap::new()) }
|
||||
}
|
||||
|
||||
/// Take one token for `client`, or say how long until one is available.
|
||||
pub fn check(&self, client: IpAddr) -> Result<(), Duration> {
|
||||
if self.per_minute == 0 {
|
||||
return Ok(());
|
||||
}
|
||||
let capacity = self.per_minute as f64;
|
||||
let per_second = capacity / 60.0;
|
||||
let now = Instant::now();
|
||||
|
||||
let mut buckets = self.buckets.lock().unwrap_or_else(|e| e.into_inner());
|
||||
|
||||
// Sweep buckets that have been idle long enough to be full again; there is nothing
|
||||
// to remember about them.
|
||||
if buckets.len() > 1024 {
|
||||
buckets.retain(|_, b| now.duration_since(b.last).as_secs_f64() * per_second < capacity);
|
||||
}
|
||||
|
||||
let bucket = buckets.entry(client).or_insert(Bucket { tokens: capacity, last: now });
|
||||
let elapsed = now.duration_since(bucket.last).as_secs_f64();
|
||||
bucket.tokens = (bucket.tokens + elapsed * per_second).min(capacity);
|
||||
bucket.last = now;
|
||||
|
||||
if bucket.tokens >= 1.0 {
|
||||
bucket.tokens -= 1.0;
|
||||
Ok(())
|
||||
} else {
|
||||
let wait = (1.0 - bucket.tokens) / per_second;
|
||||
Err(Duration::from_secs_f64(wait.max(0.001)))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn a_full_bucket_then_a_wait() {
|
||||
let limiter = RateLimiter::new(3);
|
||||
let ip: IpAddr = "203.0.113.5".parse().unwrap();
|
||||
assert!(limiter.check(ip).is_ok());
|
||||
assert!(limiter.check(ip).is_ok());
|
||||
assert!(limiter.check(ip).is_ok());
|
||||
let wait = limiter.check(ip).expect_err("fourth call in a burst should wait");
|
||||
assert!(wait > Duration::ZERO && wait <= Duration::from_secs(20));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn clients_are_independent_and_zero_disables() {
|
||||
let limiter = RateLimiter::new(1);
|
||||
let a: IpAddr = "203.0.113.5".parse().unwrap();
|
||||
let b: IpAddr = "203.0.113.6".parse().unwrap();
|
||||
assert!(limiter.check(a).is_ok());
|
||||
assert!(limiter.check(a).is_err());
|
||||
assert!(limiter.check(b).is_ok());
|
||||
|
||||
let unlimited = RateLimiter::new(0);
|
||||
for _ in 0..1000 {
|
||||
assert!(unlimited.check(a).is_ok());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,191 @@
|
||||
//! yaak-send-proxy: the network half of Yaak in a browser.
|
||||
//!
|
||||
//! A tab can't see a response the way a desktop app can — CORS hides most
|
||||
//! headers, redirects are followed silently, there is no timeline. So the tab
|
||||
//! renders the request and hands it here; this process puts it on the network
|
||||
//! with the desktop's own engine and streams back everything that happened,
|
||||
//! for the tab to store. It keeps nothing: no database, no files, no session.
|
||||
//!
|
||||
//! One binary, configured by flags or `YAAK_PROXY_*` environment variables.
|
||||
//! See README.md for running and deploying it, and `guard.rs` for what it
|
||||
//! refuses to talk to.
|
||||
|
||||
mod config;
|
||||
mod guard;
|
||||
mod limits;
|
||||
mod send;
|
||||
mod wire;
|
||||
|
||||
use axum::Router;
|
||||
use axum::body::Body;
|
||||
use axum::extract::{ConnectInfo, DefaultBodyLimit, State};
|
||||
use axum::http::{HeaderMap, HeaderValue, Method, StatusCode, header};
|
||||
use axum::response::{IntoResponse, Json, Response};
|
||||
use axum::routing::{get, post};
|
||||
use clap::Parser;
|
||||
use config::Config;
|
||||
use guard::DestinationPolicy;
|
||||
use limits::RateLimiter;
|
||||
use log::{info, warn};
|
||||
use send::{Refusal, SendLimits};
|
||||
use serde_json::json;
|
||||
use std::net::{IpAddr, SocketAddr};
|
||||
use std::sync::Arc;
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::sync::Semaphore;
|
||||
use tower_http::cors::{AllowOrigin, CorsLayer};
|
||||
use wire::SendRequest;
|
||||
|
||||
#[derive(Clone)]
|
||||
struct AppState {
|
||||
config: Arc<Config>,
|
||||
limits: Arc<SendLimits>,
|
||||
rate_limiter: Arc<RateLimiter>,
|
||||
in_flight: Arc<Semaphore>,
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() {
|
||||
env_logger::Builder::from_env(env_logger::Env::default().default_filter_or("info")).init();
|
||||
let config = Config::parse();
|
||||
|
||||
let policy = DestinationPolicy;
|
||||
let state = AppState {
|
||||
limits: Arc::new(SendLimits {
|
||||
policy,
|
||||
max_response_bytes: config.max_response_bytes,
|
||||
max_timeout: Duration::from_secs(config.max_timeout_secs),
|
||||
}),
|
||||
rate_limiter: Arc::new(RateLimiter::new(config.rate_limit_per_minute)),
|
||||
in_flight: Arc::new(Semaphore::new(config.max_concurrent)),
|
||||
config: Arc::new(config),
|
||||
};
|
||||
|
||||
let cors = CorsLayer::new()
|
||||
.allow_methods([Method::GET, Method::POST, Method::OPTIONS])
|
||||
.allow_headers([header::CONTENT_TYPE])
|
||||
.allow_origin(allowed_origins(&state.config.allowed_origins));
|
||||
|
||||
let app = Router::new()
|
||||
.route("/v1/health", get(health))
|
||||
// A WebSocket or gRPC relay would sit beside this as `/v1/ws/relay` and `/v1/grpc/relay`
|
||||
// on the same router, behind the same policy, limits and auth. Not built; see README.
|
||||
.route("/v1/http/send", post(send_http))
|
||||
.layer(DefaultBodyLimit::max(state.config.max_request_bytes))
|
||||
.layer(cors)
|
||||
.with_state(state.clone());
|
||||
|
||||
let bind = state.config.bind;
|
||||
let listener = tokio::net::TcpListener::bind(bind).await.unwrap_or_else(|e| {
|
||||
eprintln!("Failed to bind {bind}: {e}");
|
||||
std::process::exit(1);
|
||||
});
|
||||
info!(
|
||||
"yaak-send-proxy listening on http://{bind} (rate limit: {}/min)",
|
||||
state.config.rate_limit_per_minute,
|
||||
);
|
||||
|
||||
axum::serve(listener, app.into_make_service_with_connect_info::<SocketAddr>())
|
||||
.with_graceful_shutdown(async {
|
||||
let _ = tokio::signal::ctrl_c().await;
|
||||
info!("Shutting down");
|
||||
})
|
||||
.await
|
||||
.expect("server error");
|
||||
}
|
||||
|
||||
fn allowed_origins(origins: &[String]) -> AllowOrigin {
|
||||
if origins.iter().any(|o| o.trim() == "*") {
|
||||
return AllowOrigin::any();
|
||||
}
|
||||
let parsed: Vec<HeaderValue> =
|
||||
origins.iter().filter_map(|o| HeaderValue::from_str(o.trim()).ok()).collect();
|
||||
AllowOrigin::list(parsed)
|
||||
}
|
||||
|
||||
async fn health(State(state): State<AppState>) -> impl IntoResponse {
|
||||
Json(json!({
|
||||
"ok": true,
|
||||
"version": env!("CARGO_PKG_VERSION"),
|
||||
"maxResponseBytes": state.config.max_response_bytes,
|
||||
"maxTimeoutSecs": state.config.max_timeout_secs,
|
||||
}))
|
||||
}
|
||||
|
||||
fn error_response(status: StatusCode, message: impl Into<String>) -> Response {
|
||||
let message = message.into();
|
||||
(status, Json(json!({ "error": message }))).into_response()
|
||||
}
|
||||
|
||||
/// The client's address for rate limiting: the socket peer, or the first `X-Forwarded-For`
|
||||
/// hop when the operator has said the header can be trusted.
|
||||
fn client_ip(config: &Config, headers: &HeaderMap, peer: SocketAddr) -> IpAddr {
|
||||
if config.trust_forwarded_for
|
||||
&& let Some(forwarded) = headers.get("x-forwarded-for").and_then(|v| v.to_str().ok())
|
||||
&& let Some(first) = forwarded.split(',').next()
|
||||
&& let Ok(ip) = first.trim().parse::<IpAddr>()
|
||||
{
|
||||
return ip;
|
||||
}
|
||||
peer.ip()
|
||||
}
|
||||
|
||||
async fn send_http(
|
||||
State(state): State<AppState>,
|
||||
ConnectInfo(peer): ConnectInfo<SocketAddr>,
|
||||
headers: HeaderMap,
|
||||
Json(body): Json<SendRequest>,
|
||||
) -> Response {
|
||||
let ip = client_ip(&state.config, &headers, peer);
|
||||
if let Err(wait) = state.rate_limiter.check(ip) {
|
||||
warn!("Rate limited {ip}");
|
||||
let mut res = error_response(
|
||||
StatusCode::TOO_MANY_REQUESTS,
|
||||
format!("Rate limit reached; try again in {}s", wait.as_secs().max(1)),
|
||||
);
|
||||
res.headers_mut().insert(header::RETRY_AFTER, HeaderValue::from(wait.as_secs().max(1)));
|
||||
return res;
|
||||
}
|
||||
|
||||
let Ok(permit) = state.in_flight.clone().try_acquire_owned() else {
|
||||
warn!("At capacity; refusing {ip}");
|
||||
return error_response(StatusCode::SERVICE_UNAVAILABLE, "This proxy is at capacity");
|
||||
};
|
||||
|
||||
let prepared = match send::prepare(state.limits.clone(), body).await {
|
||||
Ok(p) => p,
|
||||
Err(Refusal::Unsupported(m)) => return error_response(StatusCode::BAD_REQUEST, m),
|
||||
Err(Refusal::Invalid(m)) => return error_response(StatusCode::BAD_REQUEST, m),
|
||||
Err(Refusal::Destination(m)) => {
|
||||
warn!("Refused send from {ip}: {m}");
|
||||
return error_response(StatusCode::FORBIDDEN, m);
|
||||
}
|
||||
};
|
||||
|
||||
let description = prepared.describe();
|
||||
info!("{ip} -> {description}");
|
||||
let started = Instant::now();
|
||||
|
||||
let (tx, rx) = tokio::sync::mpsc::channel(send::FRAME_CHANNEL_CAPACITY);
|
||||
tokio::spawn(async move {
|
||||
prepared.run(tx).await;
|
||||
send::log_outcome(&description, started, "finished");
|
||||
drop(permit);
|
||||
});
|
||||
|
||||
let stream = tokio_stream_from(rx);
|
||||
Response::builder()
|
||||
.status(StatusCode::OK)
|
||||
.header(header::CONTENT_TYPE, "application/x-ndjson")
|
||||
.header(header::CACHE_CONTROL, "no-store")
|
||||
// Some reverse proxies buffer streamed responses unless told not to
|
||||
.header("x-accel-buffering", "no")
|
||||
.body(Body::from_stream(stream))
|
||||
.expect("valid response")
|
||||
}
|
||||
|
||||
fn tokio_stream_from<T: Send + 'static>(
|
||||
mut rx: tokio::sync::mpsc::Receiver<T>,
|
||||
) -> impl futures_util::Stream<Item = T> + Send + 'static {
|
||||
futures_util::stream::poll_fn(move |cx| rx.poll_recv(cx))
|
||||
}
|
||||
@@ -0,0 +1,357 @@
|
||||
//! The one thing this binary does: execute a rendered request and stream back what happened.
|
||||
//!
|
||||
//! This is the "execute" half of the desktop's `send_http_request` — the part after rendering
|
||||
//! and before storage — driven through the same `HttpTransaction` the desktop drives, with the
|
||||
//! same redirect loop, cookie jar, decompression and timeline events. Everything the desktop
|
||||
//! would write to its database is written to the reply stream instead, and the tab stores it.
|
||||
|
||||
use crate::guard::{DestinationPolicy, GuardedSender};
|
||||
use crate::wire::{Frame, SendRequest};
|
||||
use base64::Engine;
|
||||
use bytes::Bytes;
|
||||
use log::{info, warn};
|
||||
use std::convert::Infallible;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::{Duration, Instant};
|
||||
use tokio::io::AsyncReadExt;
|
||||
use tokio::sync::{mpsc, watch};
|
||||
use yaak_http::client::{HttpConnectionOptions, HttpConnectionProxySetting};
|
||||
use yaak_http::cookies::CookieStore;
|
||||
use yaak_http::sender::{HttpResponseEvent, ReqwestSender};
|
||||
use yaak_http::transaction::HttpTransaction;
|
||||
use yaak_http::types::{SendableHttpRequest, SendableHttpRequestOptions};
|
||||
use yaak_models::models::HttpResponseHeader;
|
||||
|
||||
/// How many frames may sit unread by the client before body reading pauses. Backpressure, so a
|
||||
/// slow tab slows the upstream read rather than filling memory.
|
||||
pub const FRAME_CHANNEL_CAPACITY: usize = 64;
|
||||
const EVENT_CHANNEL_CAPACITY: usize = 256;
|
||||
const BODY_READ_CHUNK: usize = 64 * 1024;
|
||||
|
||||
/// What a send needs from the process, beyond the request itself.
|
||||
pub struct SendLimits {
|
||||
pub policy: DestinationPolicy,
|
||||
pub max_response_bytes: usize,
|
||||
pub max_timeout: Duration,
|
||||
}
|
||||
|
||||
/// Why a send was refused before anything was put on the network. Distinct from a failure
|
||||
/// mid-stream: these become a plain HTTP error, not a stream with an error frame.
|
||||
#[derive(Debug)]
|
||||
pub enum Refusal {
|
||||
/// The request asks for something a browser-originated send cannot mean.
|
||||
Unsupported(String),
|
||||
/// The destination is not one this proxy will talk to.
|
||||
Destination(String),
|
||||
/// The request could not be turned into something sendable.
|
||||
Invalid(String),
|
||||
}
|
||||
|
||||
pub type FrameSender = mpsc::Sender<Result<Bytes, Infallible>>;
|
||||
|
||||
/// Check and prepare a send, then hand back the task that runs it. Refusals happen here, before
|
||||
/// the caller has committed to a streaming response.
|
||||
pub async fn prepare(limits: Arc<SendLimits>, send: SendRequest) -> Result<PreparedSend, Refusal> {
|
||||
let request = send.request;
|
||||
|
||||
// The engine reads files for these body types. There are no files here that a browser tab
|
||||
// could legitimately mean, and letting a request name a path on this machine would be a
|
||||
// local file read for anyone who can reach the proxy.
|
||||
if request.body_type.as_deref() == Some("binary") {
|
||||
return Err(Refusal::Unsupported(
|
||||
"Binary file bodies can't be sent from the browser: the proxy has no access to your files"
|
||||
.to_string(),
|
||||
));
|
||||
}
|
||||
if request.body_type.as_deref() == Some("multipart/form-data") {
|
||||
let names_a_file =
|
||||
request.body.get("form").and_then(|f| f.as_array()).is_some_and(|entries| {
|
||||
entries.iter().any(|e| {
|
||||
e.get("enabled").and_then(|v| v.as_bool()).unwrap_or(true)
|
||||
&& e.get("file").and_then(|v| v.as_str()).is_some_and(|f| !f.is_empty())
|
||||
})
|
||||
});
|
||||
if names_a_file {
|
||||
return Err(Refusal::Unsupported(
|
||||
"Multipart file fields can't be sent from the browser: the proxy has no access to your files"
|
||||
.to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
|
||||
// The tab's requested timeout, capped. Zero means "none", which here means the cap.
|
||||
let requested = if send.settings.timeout_ms > 0 {
|
||||
Some(Duration::from_millis(send.settings.timeout_ms as u64))
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let timeout = requested.map_or(limits.max_timeout, |t| t.min(limits.max_timeout));
|
||||
let timeout_capped = requested.is_none_or(|t| t > limits.max_timeout);
|
||||
|
||||
let sendable = SendableHttpRequest::from_http_request(
|
||||
&request,
|
||||
SendableHttpRequestOptions {
|
||||
timeout: Some(timeout),
|
||||
follow_redirects: send.settings.follow_redirects,
|
||||
},
|
||||
)
|
||||
.await
|
||||
.map_err(|e| Refusal::Invalid(e.to_string()))?;
|
||||
|
||||
// The first hop, checked up front so a bad destination is a clean refusal rather than a
|
||||
// stream that opens and immediately errors. Every later hop is checked by GuardedSender.
|
||||
limits.policy.check_url(&sendable.url).map_err(Refusal::Destination)?;
|
||||
|
||||
Ok(PreparedSend {
|
||||
limits,
|
||||
sendable,
|
||||
settings: send.settings,
|
||||
cookies: send.cookies,
|
||||
timeout,
|
||||
timeout_capped,
|
||||
})
|
||||
}
|
||||
|
||||
pub struct PreparedSend {
|
||||
limits: Arc<SendLimits>,
|
||||
sendable: SendableHttpRequest,
|
||||
settings: yaak_models::models::HttpSendSettings,
|
||||
cookies: Option<Vec<yaak_models::models::Cookie>>,
|
||||
timeout: Duration,
|
||||
timeout_capped: bool,
|
||||
}
|
||||
|
||||
impl PreparedSend {
|
||||
pub fn describe(&self) -> String {
|
||||
format!("{} {}", self.sendable.method, self.sendable.url)
|
||||
}
|
||||
|
||||
/// Run the send, writing frames to `frames` until the terminal frame. Returns when the
|
||||
/// stream is complete or the client has gone away.
|
||||
pub async fn run(mut self, frames: FrameSender) {
|
||||
let cookie_store = self.cookies.take().map(CookieStore::from_cookies);
|
||||
let store_for_result = cookie_store.clone();
|
||||
let outcome = self.execute(frames.clone(), cookie_store).await;
|
||||
|
||||
let cookies = store_for_result.as_ref().map(|s| s.get_all_cookies());
|
||||
let terminal = match outcome {
|
||||
Ok(done) => Frame::Done {
|
||||
elapsed: done.elapsed,
|
||||
content_length: done.content_length,
|
||||
content_length_compressed: done.content_length_compressed,
|
||||
cookies,
|
||||
},
|
||||
Err(message) => Frame::Error { message, cookies },
|
||||
};
|
||||
let _ = write_frame(&frames, &terminal).await;
|
||||
}
|
||||
|
||||
async fn execute(
|
||||
self,
|
||||
frames: FrameSender,
|
||||
cookie_store: Option<CookieStore>,
|
||||
) -> Result<DoneStats, String> {
|
||||
let limits = self.limits;
|
||||
|
||||
let (client, resolver) = HttpConnectionOptions {
|
||||
id: uuid::Uuid::new_v4().to_string(),
|
||||
validate_certificates: self.settings.validate_certificates,
|
||||
// The proxy connects directly. Going through a system proxy would move DNS, and
|
||||
// therefore the address check, somewhere this process can't see.
|
||||
proxy: HttpConnectionProxySetting::Disabled,
|
||||
client_certificate: None,
|
||||
dns_overrides: Vec::new(),
|
||||
address_filter: Some(limits.policy.address_filter()),
|
||||
}
|
||||
.build_client()
|
||||
.map_err(|e| format!("Failed to build HTTP client: {e}"))?;
|
||||
|
||||
// Timeline events go into the same frame stream as everything else, as they happen.
|
||||
// The desktop persists them from a task like this one; here the task serialises them.
|
||||
let (event_tx, mut event_rx) = mpsc::channel::<HttpResponseEvent>(EVENT_CHANNEL_CAPACITY);
|
||||
resolver.set_event_sender(Some(event_tx.clone())).await;
|
||||
let dns_elapsed = Arc::new(AtomicU64::new(0));
|
||||
let event_frames = frames.clone();
|
||||
let event_dns = dns_elapsed.clone();
|
||||
let event_task = tokio::spawn(async move {
|
||||
while let Some(event) = event_rx.recv().await {
|
||||
if let HttpResponseEvent::DnsResolved { duration, .. } = &event {
|
||||
event_dns.store(*duration, Ordering::Relaxed);
|
||||
}
|
||||
let frame = Frame::Event { event: event.into() };
|
||||
if write_frame(&event_frames, &frame).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Cancellation: the client hanging up, or the overall deadline. The deadline exists
|
||||
// because a per-hop timeout times each hop separately; ten slow redirects must not add
|
||||
// up to ten timeouts.
|
||||
let (cancel_tx, cancel_rx) = watch::channel(false);
|
||||
let deadline = self.timeout * 2 + Duration::from_secs(5);
|
||||
let deadline_cancel = cancel_tx.clone();
|
||||
let deadline_task = tokio::spawn(async move {
|
||||
tokio::time::sleep(deadline).await;
|
||||
let _ = deadline_cancel.send(true);
|
||||
});
|
||||
let hangup_frames = frames.clone();
|
||||
let hangup_task = tokio::spawn(async move {
|
||||
hangup_frames.closed().await;
|
||||
let _ = cancel_tx.send(true);
|
||||
});
|
||||
|
||||
if self.timeout_capped {
|
||||
let _ = event_tx.try_send(HttpResponseEvent::Info(format!(
|
||||
"Timeout set to {:?} (this proxy's ceiling)",
|
||||
self.timeout
|
||||
)));
|
||||
}
|
||||
|
||||
let sender = GuardedSender::new(ReqwestSender::with_client(client), limits.policy.clone());
|
||||
let transaction = match cookie_store {
|
||||
Some(store) => HttpTransaction::with_cookie_behavior(
|
||||
sender,
|
||||
store,
|
||||
self.settings.send_cookies,
|
||||
self.settings.store_cookies,
|
||||
),
|
||||
None => HttpTransaction::new(sender),
|
||||
};
|
||||
|
||||
let started_at = Instant::now();
|
||||
let result = transaction
|
||||
.execute_with_cancellation(self.sendable, cancel_rx.clone(), event_tx.clone())
|
||||
.await;
|
||||
resolver.set_event_sender(None).await;
|
||||
|
||||
let mut response = match result {
|
||||
Ok(response) => response,
|
||||
Err(err) => {
|
||||
drop(event_tx);
|
||||
let _ = event_task.await;
|
||||
deadline_task.abort();
|
||||
hangup_task.abort();
|
||||
return Err(describe_error(&err));
|
||||
}
|
||||
};
|
||||
|
||||
let elapsed_headers = started_at.elapsed().as_millis() as u64;
|
||||
let head = Frame::Response {
|
||||
status: response.status,
|
||||
status_reason: response.status_reason.clone(),
|
||||
url: response.url.clone(),
|
||||
remote_addr: response.remote_addr.clone(),
|
||||
version: response.version.clone(),
|
||||
headers: to_wire_headers(&response.headers),
|
||||
request_headers: to_wire_headers(&response.request_headers),
|
||||
content_length: response.content_length,
|
||||
elapsed_headers,
|
||||
elapsed_dns: dns_elapsed.load(Ordering::Relaxed),
|
||||
};
|
||||
write_frame(&frames, &head).await.map_err(|_| "Client went away".to_string())?;
|
||||
|
||||
let declared_length = response.content_length;
|
||||
let mut body = response
|
||||
.into_body_stream()
|
||||
.map_err(|e| format!("Failed to read response body: {e}"))?;
|
||||
let mut buf = vec![0u8; BODY_READ_CHUNK];
|
||||
let mut total: usize = 0;
|
||||
let mut cancel_rx = cancel_rx;
|
||||
let base64 = base64::engine::general_purpose::STANDARD;
|
||||
|
||||
let read_result: Result<(), String> = loop {
|
||||
if *cancel_rx.borrow() {
|
||||
break Err("Request canceled".to_string());
|
||||
}
|
||||
let read = tokio::select! {
|
||||
biased;
|
||||
_ = cancel_rx.changed() => break Err("Request canceled".to_string()),
|
||||
r = body.read(&mut buf) => r,
|
||||
};
|
||||
match read {
|
||||
Ok(0) => break Ok(()),
|
||||
Ok(n) => {
|
||||
total += n;
|
||||
if total > limits.max_response_bytes {
|
||||
break Err(format!(
|
||||
"Response body exceeds this proxy's limit of {} bytes",
|
||||
limits.max_response_bytes
|
||||
));
|
||||
}
|
||||
let frame = Frame::Body { data: base64.encode(&buf[..n]) };
|
||||
if write_frame(&frames, &frame).await.is_err() {
|
||||
break Err("Client went away".to_string());
|
||||
}
|
||||
}
|
||||
Err(e) => break Err(format!("Failed to read response body: {e}")),
|
||||
}
|
||||
};
|
||||
drop(body);
|
||||
|
||||
// Let the timeline drain before the terminal frame, so nothing arrives after "done".
|
||||
drop(event_tx);
|
||||
let _ = event_task.await;
|
||||
deadline_task.abort();
|
||||
hangup_task.abort();
|
||||
|
||||
read_result?;
|
||||
Ok(DoneStats {
|
||||
elapsed: started_at.elapsed().as_millis() as u64,
|
||||
content_length: total as u64,
|
||||
content_length_compressed: declared_length.unwrap_or(total as u64),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
/// A send error as a sentence, not a debug dump.
|
||||
///
|
||||
/// A connection error from reqwest arrives wrapped several layers deep, and the layer that
|
||||
/// says something useful — "Refusing to connect to ::1: loopback" — is the innermost. The
|
||||
/// desktop shows the outer `Debug`; a stranger reading a proxy's reply deserves the reason.
|
||||
fn describe_error(err: &yaak_http::error::Error) -> String {
|
||||
match err {
|
||||
yaak_http::error::Error::Client(e) => {
|
||||
let mut leaf: &dyn std::error::Error = e;
|
||||
while let Some(next) = leaf.source() {
|
||||
leaf = next;
|
||||
}
|
||||
let outer = e.to_string();
|
||||
let inner = leaf.to_string();
|
||||
if inner == outer { outer } else { format!("{outer}: {inner}") }
|
||||
}
|
||||
yaak_http::error::Error::RequestError(message) => format!("Request failed: {message}"),
|
||||
other => other.to_string(),
|
||||
}
|
||||
}
|
||||
|
||||
struct DoneStats {
|
||||
elapsed: u64,
|
||||
content_length: u64,
|
||||
content_length_compressed: u64,
|
||||
}
|
||||
|
||||
fn to_wire_headers(headers: &[(String, String)]) -> Vec<HttpResponseHeader> {
|
||||
headers
|
||||
.iter()
|
||||
.map(|(name, value)| HttpResponseHeader { name: name.clone(), value: value.clone() })
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn write_frame(frames: &FrameSender, frame: &Frame) -> Result<(), ()> {
|
||||
let mut line = match serde_json::to_vec(frame) {
|
||||
Ok(v) => v,
|
||||
Err(e) => {
|
||||
warn!("Failed to serialize frame: {e}");
|
||||
return Err(());
|
||||
}
|
||||
};
|
||||
line.push(b'\n');
|
||||
frames.send(Ok(Bytes::from(line))).await.map_err(|_| ())
|
||||
}
|
||||
|
||||
/// Log a finished send at info: destination, outcome, and how long, never the content.
|
||||
pub fn log_outcome(description: &str, started: Instant, outcome: &str) {
|
||||
info!("{description} -> {outcome} in {:?}", started.elapsed());
|
||||
}
|
||||
@@ -0,0 +1,97 @@
|
||||
//! What crosses the wire between a tab and this proxy.
|
||||
//!
|
||||
//! One `POST /v1/http/send` carries a request the tab has already rendered —
|
||||
//! templates resolved, inheritance applied — plus the send settings and the
|
||||
//! cookies the send starts with. The reply is a stream of newline-delimited
|
||||
//! JSON frames: timeline events as they happen, the response head as soon as
|
||||
//! headers arrive, body chunks as they are read, and one terminal frame.
|
||||
//!
|
||||
//! Nothing here names a workspace, a request id, or a response id. The proxy
|
||||
//! does not know what the tab will call this response; it only knows what came
|
||||
//! back.
|
||||
//!
|
||||
//! The TypeScript side of this contract is generated from these types into
|
||||
//! `bindings/` (`cargo test -p yaak-send-proxy`) and published to the tab as
|
||||
//! `@yaakapp-internal/send-proxy`, so a change here is a type error there.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use ts_rs::TS;
|
||||
use yaak_models::models::{
|
||||
Cookie, HttpRequest, HttpResponseEventData, HttpResponseHeader, HttpSendSettings,
|
||||
};
|
||||
|
||||
/// The body of `POST /v1/http/send`.
|
||||
#[derive(Deserialize, Debug, TS)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
#[ts(export, export_to = "gen_send_proxy.ts")]
|
||||
pub struct SendRequest {
|
||||
/// The request to send, in the desktop's own model shape but with every template already
|
||||
/// rendered by the tab. The proxy builds the URL, headers and body from it exactly the way
|
||||
/// the desktop does after rendering.
|
||||
pub request: HttpRequest,
|
||||
/// The resolved settings, values only. Where they came from is the tab's to record in
|
||||
/// its timeline; the proxy only needs to obey them.
|
||||
pub settings: HttpSendSettings,
|
||||
/// The cookies to start with. `None` means no jar at all: nothing sent, nothing kept.
|
||||
#[serde(default)]
|
||||
pub cookies: Option<Vec<Cookie>>,
|
||||
}
|
||||
|
||||
/// One line of the reply stream. Tags are snake_case like the timeline event tags; fields are
|
||||
/// camelCase like every model the tab stores.
|
||||
#[derive(Serialize, Debug, TS)]
|
||||
#[serde(
|
||||
tag = "type",
|
||||
rename_all = "snake_case",
|
||||
rename_all_fields = "camelCase"
|
||||
)]
|
||||
#[ts(export, export_to = "gen_send_proxy.ts")]
|
||||
pub enum Frame {
|
||||
/// A timeline event, in the same shape the desktop stores. Interleaved with everything
|
||||
/// else in the order the engine produced it.
|
||||
Event { event: HttpResponseEventData },
|
||||
/// The response head. Sent once, as soon as the final hop's headers are in — before any of
|
||||
/// the body — so the tab can show status and headers while the body streams.
|
||||
Response {
|
||||
status: u16,
|
||||
status_reason: Option<String>,
|
||||
/// The URL that answered, after redirects.
|
||||
url: String,
|
||||
remote_addr: Option<String>,
|
||||
version: Option<String>,
|
||||
headers: Vec<HttpResponseHeader>,
|
||||
/// The headers that were actually sent on the final hop, cookies and all.
|
||||
request_headers: Vec<HttpResponseHeader>,
|
||||
/// `Content-Length` as declared by the server, if it declared one.
|
||||
#[ts(type = "number | null")]
|
||||
content_length: Option<u64>,
|
||||
/// Milliseconds from the start of the send to the response head.
|
||||
#[ts(type = "number")]
|
||||
elapsed_headers: u64,
|
||||
/// Milliseconds spent in DNS on the last lookup, or zero.
|
||||
#[ts(type = "number")]
|
||||
elapsed_dns: u64,
|
||||
},
|
||||
/// A piece of the response body, decompressed, base64-encoded.
|
||||
Body { data: String },
|
||||
/// The send finished. The last frame on a successful stream.
|
||||
Done {
|
||||
/// Milliseconds from the start of the send to the end of the body.
|
||||
#[ts(type = "number")]
|
||||
elapsed: u64,
|
||||
/// Bytes of body relayed, after decompression.
|
||||
#[ts(type = "number")]
|
||||
content_length: u64,
|
||||
/// Bytes on the wire as declared by the server, or the relayed size when unknown.
|
||||
#[ts(type = "number")]
|
||||
content_length_compressed: u64,
|
||||
/// The jar as the send left it, for the tab to persist. `None` when the tab sent none.
|
||||
cookies: Option<Vec<Cookie>>,
|
||||
},
|
||||
/// The send failed. The last frame on a failed stream. Cookies collected before the failure
|
||||
/// still come back — the transaction may have set some before the hop that failed.
|
||||
Error {
|
||||
message: String,
|
||||
cookies: Option<Vec<Cookie>>,
|
||||
},
|
||||
}
|
||||
Reference in New Issue
Block a user