diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 9842e19..3dfab35 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -105,7 +105,7 @@ The daemon's shared state is wrapped in `Arc` for concurrent access across conne |-----------|------|---------| | Config | `Arc` | Parsed `airlock.toml` (immutable after startup) | | Secrets | `SecretStore` = `Arc>>` | Per-label slot holding `Arc>` plus refresh health; the map is fixed at startup, slot contents swap on refresh | -| Redactor | `Arc` | Aho-Corasick automaton for output redaction | +| Redactor | `Arc>>` | Aho-Corasick automaton for output redaction; refresh tasks swap the inner `Arc`. A connection snapshots it for the child's stdout/stderr; a proxy session carries the handle itself and snapshots per response | | Ring buffer | `Arc>` | Last 1000 log entries (`VecDeque`) | | Child registry | `Arc>>` | PIDs of currently running children | @@ -189,16 +189,30 @@ child (curl) daemon upstream │ │ inject header + hop-by-hop │ │ │ secret store lookup (Stale→502)│ │ │ attach prefix+secret+suffix │ + │ │ force Accept-Encoding:identity,│ + │ │ strip Range / If-Range │ │ │ resolve host, refuse non- │ │ │ routable addrs, dial that │ │ │ exact SocketAddr │ │ ├───── TLS ≥1.2, public roots ──►│ - │◄─────────────────────────────┤◄────── response streamed ──────┤ + │ │◄──────── response head ────────┤ + │ │ content-encoded / odd framing /│ + │ │ 206? → 502, body unread │ + │ │ redact every header value │ + │ │ drop Content-Length unless the │ + │ │ response is bodiless │ + │◄──── redacted, chunked ──────┤◄──── body frames streamed ─────┤ │ │ audit: method, host, path, │ - │ │ decision, status (no query, │ - │ │ no header values) │ + │ │ decision, status, redaction │ + │ │ counts (no query, no header │ + │ │ values, no matched bytes) │ ``` +The response never reaches the tool unexamined. Header values and body both go +through the redactor, so what `curl -o` writes into the sandbox was already +redacted — see [SECURITY.md](SECURITY.md#response-redaction) for what that +covers and what fails closed. + Meanwhile the sandbox holds the other end: the profile permits a TCP connect to that port and nothing else — no DNS, no other destination — so a tool that ignores `HTTPS_PROXY` gets nowhere. @@ -228,6 +242,21 @@ select! loop → NDJSON → Unix socket → client This design keeps the automaton's streaming state machine on a dedicated blocking thread (via `spawn_blocking`) while the daemon's main loop remains fully async. +The proxy's response path does not use this bridge. It already holds the bytes +as owned frames handed to it by hyper, so it drives a `StreamRedactor` — an +incremental redactor that keeps between chunks only the bytes a pattern could +still be starting in, and whose output for any chunking is what `redact_bytes` +makes of the whole input — directly from `poll_frame`. No thread and no channel +per response: hyper's polling is the backpressure, and dropping the response +stops the upstream read. + +Both paths take their redactor from the same `Arc>>`. A +connection snapshots it once for the child's stdout and stderr, which are +framed against the secrets the child was spawned with. A proxy response +snapshots it per response, because the proxy injects whatever the store holds +at that moment and a token refreshed mid-exec must be redacted on the way +back. + ## Wire protocol Communication uses **NDJSON** (newline-delimited JSON) over the Unix domain socket. Each message is a single JSON line with a `"type"` discriminator field. diff --git a/CLAUDE.md b/CLAUDE.md index b098d6b..c3cd8e7 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -23,7 +23,7 @@ Read [README.md](README.md), [ARCHITECTURE.md](ARCHITECTURE.md), and [SECURITY.m ## Tests -- Most logic lives in `cargo test --lib` (432 tests). All hermetic. +- Most logic lives in `cargo test --lib` (458 tests). All hermetic. - `tests/cli_integration.rs` spawns the real `airlock` binary and runs `daemon start/stop/status`. **These might fail in the Claude Code sandbox** - Some `client.rs` tests read the process's real stdin. Run `cargo test` with `< /dev/null` or they can hang on an inherited pipe that never closes. diff --git a/README.md b/README.md index 6f47d2c..d412d3e 100644 --- a/README.md +++ b/README.md @@ -303,6 +303,8 @@ airlock exec -- curl -s https://run.googleapis.com/v2/projects/my-project/locati What the daemon does for that invocation: binds a proxy on an ephemeral loopback port, points the tool at it with `HTTPS_PROXY` and `CURL_CA_BUNDLE`, pins the tool's egress to that one port with Seatbelt (macOS) or Landlock (Linux), and tears the whole thing down when the child exits. Per request it checks the host against the route table, the method and path against the rules, attaches the credential, and forwards over a verified TLS connection to a host it has confirmed is publicly routable. Unlisted hosts are simply unreachable — deny by default. +The response is redacted on the way back — header values and body both — so an API that echoes the credential cannot hand it to the agent even via `curl -o file`. Compressed and partial responses are refused rather than forwarded unread: see [SECURITY.md](SECURITY.md#response-redaction). + Caveats worth knowing up front: HTTP/1.1 only (no gRPC or HTTP/2-only endpoints), certificate-pinned clients break under interception, and the agent gets the credential's full API authority on the routed hosts — scope the service account narrowly. Full threat model and residual risks: [SECURITY.md](SECURITY.md#proxy-tools); design rationale: [docs/proxy-tools-design.md](docs/proxy-tools-design.md). ### `[agent]` — for `airlock run` diff --git a/SECURITY.md b/SECURITY.md index 42fc2d6..c940aa9 100644 --- a/SECURITY.md +++ b/SECURITY.md @@ -340,9 +340,35 @@ Egress restriction is the second layer, not the first. It is what makes the tool | Path contains `.`/`..` segments, `//`, a backslash, or an encoded `/`, `.`, `\` or NUL | `403` — refused rather than normalized, because the upstream's normalization may differ from the matcher's | | The injected secret's slot is `Stale` | `502` | | Host resolves to any private, loopback, link-local (incl. `169.254.169.254`), CGNAT, ULA, multicast, documentation or otherwise non-routable address | `502` | +| The upstream answers with a `Content-Encoding` other than `identity`, a transfer coding other than `chunked`, or a partial representation (`206` / `Content-Range`) | `502`, body dropped unread — see [Response redaction](#response-redaction) | On the way through, every client-supplied copy of the injected header is removed before the credential is attached, as are `Proxy-Authorization`, `Proxy-Connection` and the other hop-by-hop headers. The secret is read from the secret store **per request**, so a background refresh applies to the next one. The header value is assembled into a buffer that is zeroized, never through `format!`, and is marked sensitive. +The request is also rewritten so that the response is something the redactor can read: `Accept-Encoding` is forced to `identity` whatever the tool asked for, and `Range` / `If-Range` are removed. + +### Response redaction + +Everything the upstream sends back is redacted before it reaches the tool, with the same automaton and the same secret set (raw, base64, URL-encoded and hex variants of *every* declared secret, not just this route's) that the tool's stdout goes through: + +- **All response header values** — every one of them, including `Location`, `Set-Cookie` and `WWW-Authenticate`. A value that will not rebuild after replacement is dropped rather than forwarded. +- **The body**, streamed. Nothing is buffered beyond the partial match at the end of a frame, so a multi-gigabyte download costs what a small one costs, and the tool's own read rate is what drives the upstream read. A secret split across two upstream writes is still caught. +- **Trailers are dropped**, not forwarded. +- **The upstream's reason phrase is dropped.** `HTTP/1.1 200 ` is a legal status line and sits outside the header map, so the tool sees the status code with the standard phrase, never the upstream's text. + +The redactor is taken per response from the daemon's live handle, not snapshotted when the exec started: a tool runs for minutes, the proxy injects whatever the store holds *now*, and the two generations a refresh leaves behind cover a swap that lands mid-response. + +Because a `[REDACTED:name]` placeholder is not the length of the secret it replaced, an upstream `Content-Length` is wrong whenever anything matches — and which it is cannot be known before the body has been read. The proxy therefore drops it for any response that has a body and lets hyper frame the response as chunked (HTTP/1.1 always supports it). A bodiless response — HEAD, `1xx`, `204`, `304` — keeps its length, which describes the representation rather than bytes on the wire, so `curl -I` still reports one. + +Three things fail closed rather than being handled, all for the same reason — the redactor reads bytes, not formats, and Airlock adds no decoder to the response path: + +- **Compressed responses.** `gzip`, `br`, `zstd` and `deflate` are opaque to a byte-pattern scanner. The request demands `identity`; an upstream that compresses anyway gets a `502` and its body is dropped unread. +- **Unknown transfer codings**, for the same reason. +- **Byte ranges.** A range may begin in the middle of a secret, which would split the pattern across two responses the proxy never sees together while the tool reassembles the plaintext in a file. `Range` and `If-Range` are stripped from the request so the upstream sends the whole representation, and a `206` or `Content-Range` that arrives anyway is refused. Resumed and parallel-chunked downloads therefore do not work through a proxy tool. + +There is no configuration that turns any of this off. Redaction on the output path is mandatory in Airlock, and the proxy is an output path. + +When a response had anything replaced, the audit line says so: a count of header values on the line written when the headers arrive, and a second line when the body ends carrying the body's count. Counts only — never the matched bytes. + The CONNECT authority is the single source of truth: it selects the route, names the leaf certificate the tool is shown, is the name resolved and dialled, and is the name the upstream certificate is verified against (TLS ≥ 1.2, public roots). The client's SNI is ignored entirely, so `curl --resolve`, `--connect-to`, a forged `Host` and a forged SNI cannot make any two of those disagree. DNS is resolved once and the concrete `SocketAddr` that passed the address check is the one dialled, so rebinding cannot slip between check and use. Each request is logged to the ring buffer: tool, method, host, path, decision and upstream status. Never a header value, and never the query string — it may carry data. @@ -360,7 +386,7 @@ Each request is logged to the ring buffer: tool, method, host, path, decision an - **Misuse, not leakage.** The agent gets the credential's full API authority on routed hosts — broader than a purpose-built CLI. Mitigate with a narrowly scoped service account first and method/path rules second. - **Data exfiltration to co-tenants.** Anything the tool can read can be uploaded to an attacker's project on an allowed multi-tenant host (`storage.googleapis.com` serves every GCP customer). The credential cannot. - **Path rules are a convenience layer, not an authorization system.** They see the path, not the body; a `POST` allowed for one purpose may do another (`:batchUpdate`, GraphQL). IAM is the authority boundary. -- **`-o` and friends bypass the stdout redactor.** Response bodies reach the agent through the tool's stdout, which passes through the Aho-Corasick redactor — so an API that echoes the bearer token back is covered there. A body written to a file (`curl -o`, `--dump-header`, `--trace`) is not. This gap exists for every Airlock tool, but it is more reachable here. +- **An upstream that *transforms* the secret is not caught.** [Response redaction](#response-redaction) closes the `curl -o` / `--dump-header` / `--trace` path for response bytes: what the tool writes to a file was already redacted, so the plaintext credential never exists inside the sandbox. What it does not catch is an upstream that reflects the secret reversed, re-chunked, or encoded in a scheme the redactor does not know — the same limitation the stdout path has always had. It also says nothing about the agent's own request data: a query string or request body the agent chose is forwarded as sent. - **Linux egress pinning is port-scoped and TCP-only** — see [Linux — Landlock LSM](#linux--landlock-lsm). - **HTTP/1.1 only.** ALPN offers `http/1.1` and nothing else; gRPC and HTTP/2-only endpoints will not work. - **Certificate-pinned clients break** under interception. By design. diff --git a/SKILL.md b/SKILL.md index 2277f99..aea7d17 100644 --- a/SKILL.md +++ b/SKILL.md @@ -81,7 +81,7 @@ airlock exec -- curl -s https://run.googleapis.com/v2/projects/my-project/locati airlock exec -- curl -s 'https://storage.googleapis.com/storage/v1/b?project=my-project' ``` -Four things to know: +Things to know: - **Do not pass authentication headers.** Airlock attaches the credential itself. An `-H 'Authorization: ...'` you supply is removed before the request @@ -93,9 +93,18 @@ Four things to know: **not** retry with `--noproxy`, `--insecure`/`-k`, a different port, or a rewritten URL. Those either fail the same way or fail harder — the sandbox blocks direct connections and DNS outright. Report the refusal to the user. -- **`-o file` skips redaction.** Output written to a file does not pass through - the redactor, so a response that echoes a credential lands on disk in the - clear. Prefer stdout. +- **Responses are redacted, including to a file.** Header values and body are + redacted inside the daemon, so `-o file`, `--dump-header` and `-D` write + bytes that already have `[REDACTED:NAME]` in place of any credential. Seeing + that in a downloaded file is expected — the API echoed a secret, and Airlock + replaced it. +- **Compressed transfer is not available.** The request always asks the API for + an uncompressed response, so `--compressed` gets you plain bytes, and an API + that compresses anyway produces `502 ... content-encoded`. Nothing to work + around; just read the plain response. +- **Range requests and resumed downloads do not work.** `--range` / `-r` and + `-C -` are stripped, so the API sends the whole resource. A resume attempt + will fail or re-download from the start. Download in one go. ### Check daemon status diff --git a/docs/proxy-tools-design.md b/docs/proxy-tools-design.md index 21d1b3b..b47a70a 100644 --- a/docs/proxy-tools-design.md +++ b/docs/proxy-tools-design.md @@ -1,6 +1,6 @@ # Proxy tools — design proposal -**Status:** implemented through phase 2. The config schema and route matcher +**Status:** implemented through phase 3. The config schema and route matcher live in [src/proxy.rs](../src/proxy.rs) and [src/config.rs](../src/config.rs); the runtime is [src/proxy/server.rs](../src/proxy/server.rs) and [src/proxy/ca.rs](../src/proxy/ca.rs), wired into `handle_exec_request` in @@ -68,7 +68,7 @@ what it does and what this proposal takes from it: | Proxy auth: random token, Basic auth, constant-time compare. | Same, but per-exec (see below for why it is mandatory here). | | SSRF: private-range check in the dialer's `Control` callback — i.e. *after* DNS resolution, which defeats rebinding. | Same: check the resolved `SocketAddr` immediately before `connect()`. | | Host header vs. CONNECT authority consistency check. | Same, and the client's SNI is ignored entirely. | -| No response redaction, no audit trail for proxied requests. | Tool stdout already flows through the redactor; proxied requests are logged to the ring buffer. | +| No response redaction, no audit trail for proxied requests. | **Responses are redacted in the proxy** — header values and body — so nothing the tool writes to a file holds a secret; proxied requests are logged to the ring buffer. | ## Decisions taken @@ -187,8 +187,12 @@ CONNECT run.googleapis.com:443 ├─ resolve host; any private / loopback / link-local / │ CGNAT / ULA / metadata (169.254.0.0/16) address → 403 │ (checked on the SocketAddr passed to connect(), post-DNS) + ├─ force Accept-Encoding: identity; strip Range / If-Range ├─ upstream TLS ≥ 1.2, verified against public roots for `host` - └─ stream response back unmodified + ├─ response Content-Encoding ≠ identity, odd transfer coding, + │ or 206 / Content-Range → 502 (fail closed) + └─ redact every header value; drop Content-Length unless the + response is bodiless; stream the body through the redactor plain `GET http://…` (non-CONNECT) → 403 (never send credentials in cleartext) ``` @@ -258,13 +262,55 @@ backend closes both gaps and is the intended follow-up. ### Output -Response bodies reach the agent via curl's stdout, which already passes -through the Aho-Corasick redactor — so an API that echoes the bearer token is -covered on that path. It is **not** covered when curl writes to a file -(`-o`). That gap exists for every Airlock tool today, but it is more -reachable here. Options, deferred: redact inside the proxy (requires forcing -`Accept-Encoding: identity` and re-framing), or deny the tool write access -outside a scratch directory. +Response bodies reach the agent via curl's stdout, which passes through the +Aho-Corasick redactor — but not when curl writes to a file (`-o`, +`--dump-header`, `--trace`), and that is exactly what an HTTP client is for. +Rather than fence the tool out of the filesystem, the redaction moved into the +proxy: **every response header value and every body byte is redacted before it +reaches the tool**, so the plaintext secret never exists inside the sandbox at +all. The two options the first draft weighed against each other turned out not +to be alternatives — the second only narrows where the plaintext can land, +while the first stops it being produced. + +Consequences, all accepted deliberately: + +- **Streaming, not buffering.** The body is redacted frame by frame by an + incremental redactor ([src/redact.rs](../src/redact.rs)) that holds back only + the bytes a pattern could still be starting in — never more than the longest + pattern — and whose output for any chunking equals what the single-shot + redactor makes of the whole input. It runs inside `poll_frame` with no thread + and no channel behind it, so hyper's own polling is the backpressure and + dropping the response stops the upstream read. The `spawn_blocking` bridge the + stdout path uses would have cost a thread per response and would have had to + be cancelled by hand. +- **The redactor is taken per response from the live handle**, not snapshotted + at session start. A tool runs for minutes, the proxy injects whatever the + store holds *now*, and a session snapshot would not know a token minted after + the exec began. The two generations a refresh leaves behind cover a swap that + lands mid-response. +- **`Content-Length` is dropped whenever there is a body.** A placeholder is not + the length of the secret it replaced, and which it is cannot be known before + the body has been read; hyper frames the response as chunked instead, which + HTTP/1.1 always supports. A bodiless response (HEAD, `1xx`, `204`, `304`) + keeps its length — there it is metadata about the representation, and `curl + -I` must still report one. +- **Compression fails closed.** The request forces `Accept-Encoding: identity`; + an upstream that answers with a content coding (or a transfer coding other + than chunked) gets a `502` and its body is dropped unread. No decompressor is + added: it would be a second parser of attacker-supplied bytes in the response + path for no security gain. +- **Ranges are stripped, not supported.** A range may begin in the middle of a + secret, splitting the pattern across two responses the proxy never sees + together while the tool reassembles the plaintext in a file. `Range` and + `If-Range` are removed so the upstream sends the whole representation, and a + `206` arriving anyway is refused. Resumed downloads therefore do not work. +- **Trailers are dropped.** +- **No opt-out.** Redaction on the output path is mandatory in Airlock; the + proxy is an output path. + +What it does not catch is an upstream that *transforms* the secret — reversed, +re-encoded in a scheme the redactor does not know — which is the same +limitation the stdout path has always had. Each proxied request is logged to the ring buffer: method, host, path (no query string — it may carry data), route decision, upstream status. @@ -301,7 +347,8 @@ request path is security-critical and small enough to audit. | **0** | Design; `proxy` / `routes` schema, validation, matcher, tests; daemon fails closed; `airlock list` shows routes. | landed | | **1** | Runtime on macOS: per-exec listener, CA, interception, injection, SSRF dial check, Seatbelt `ProxyOnly`. Curl guidance flipped in SECURITY.md / SKILL.md / README. | landed | | **2** | Linux: Landlock ABI v4 network rules, fail-closed kernel check. | landed; exercised by `tests/proxy_e2e_integration.rs` on the Linux CI runner (proxy port reachable, direct TCP connect refused). The fail-closed path for kernels older than 6.7 has not been run on such a kernel. | -| 3 | In-proxy response redaction, HTTP/2, per-route upstream port, network-namespace backend. | open | +| **3** | In-proxy response redaction: header values and body, streaming; identity encoding forced; compressed, oddly-framed and partial responses refused. | landed | +| 4 | HTTP/2, per-route upstream port, network-namespace backend. | open | Phase 1 landed with request auditing included rather than deferred to phase 3 — the ring-buffer line (tool, method, host, path, decision, upstream status) falls diff --git a/src/daemon.rs b/src/daemon.rs index aad9270..2708280 100644 --- a/src/daemon.rs +++ b/src/daemon.rs @@ -791,10 +791,15 @@ pub(crate) async fn run_embedded( // Snapshot the redactor at accept time so in-flight // connections are not affected by concurrent refreshes. let red = redactor.read().unwrap_or_else(|e| e.into_inner()).clone(); + // The live handle travels alongside the snapshot: a + // proxy response is redacted with whatever the daemon + // holds *now*, because the proxy injects whatever it + // holds now. + let live = Arc::clone(&redactor); let ca = proxy_ca.clone(); tokio::spawn(async move { - handle_connection(stream, cfg, sec, red, rb.clone(), cr, ca).await; + handle_connection(stream, cfg, sec, red, live, rb.clone(), cr, ca).await; rb.log(format!("connection closed ({peer_info})")); }); } @@ -976,10 +981,15 @@ async fn async_main( // may swap the inner Arc later; this connection keeps // its snapshot for its full lifetime. let red = redactor.read().unwrap_or_else(|e| e.into_inner()).clone(); + // The live handle travels alongside the snapshot: a + // proxy response is redacted with whatever the daemon + // holds *now*, because the proxy injects whatever it + // holds now. + let live = Arc::clone(&redactor); let ca = proxy_ca.clone(); tokio::spawn(async move { - handle_connection(stream, cfg, sec, red, rb.clone(), cr, ca).await; + handle_connection(stream, cfg, sec, red, live, rb.clone(), cr, ca).await; rb.log(format!("connection closed ({peer_info})")); }); } @@ -1032,6 +1042,7 @@ async fn handle_connection( config: Arc, secrets: SecretStore, redactor: Arc, + live_redactor: Arc>>, ring_buffer: RingBuffer, child_registry: ChildRegistry, proxy_ca: Option>, @@ -1117,6 +1128,7 @@ async fn handle_connection( config, secrets, redactor, + live_redactor, ring_buffer, child_registry, proxy_ca, @@ -1199,6 +1211,7 @@ async fn handle_exec_request( config: Arc, secrets: SecretStore, redactor: Arc, + live_redactor: Arc>>, ring_buffer: RingBuffer, child_registry: ChildRegistry, proxy_ca: Option>, @@ -1317,6 +1330,7 @@ async fn handle_exec_request( ca, config.ca_path.clone(), Arc::clone(&secrets), + live_redactor, ring_buffer.clone(), ) { Ok(session) => { diff --git a/src/proxy/server.rs b/src/proxy/server.rs index 0f2e90f..2cf0f96 100644 --- a/src/proxy/server.rs +++ b/src/proxy/server.rs @@ -21,12 +21,14 @@ use std::collections::HashMap; use std::convert::Infallible; use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr}; use std::path::PathBuf; -use std::sync::Arc; +use std::pin::Pin; +use std::sync::{Arc, RwLock}; +use std::task::{Context, Poll}; use std::time::Duration; use bytes::Bytes; use http_body_util::{BodyExt, Full, combinators::BoxBody}; -use hyper::body::Incoming; +use hyper::body::{Body, Frame, Incoming}; use hyper::header::{self, HeaderMap, HeaderName, HeaderValue}; use hyper::service::service_fn; use hyper::{Method, Request, Response, StatusCode}; @@ -41,6 +43,7 @@ use tokio_util::sync::{CancellationToken, DropGuard}; use zeroize::Zeroize; use crate::daemon::RingBuffer; +use crate::redact::{Redactor, StreamRedactor}; use crate::secrets::{Health, SecretStore}; use super::ca::ProxyCa; @@ -138,12 +141,14 @@ impl Drop for ProxySession { impl ProxySession { /// Bind a listener and start serving `policy` on it. + #[allow(clippy::too_many_arguments)] pub fn start( tool: String, policy: ProxyPolicy, ca: Arc, ca_path: PathBuf, secrets: SecretStore, + redactor: Arc>>, ring_buffer: RingBuffer, ) -> std::io::Result { Self::start_with_upstream( @@ -152,17 +157,20 @@ impl ProxySession { ca, ca_path, secrets, + redactor, ring_buffer, Upstream::public(), ) } + #[allow(clippy::too_many_arguments)] fn start_with_upstream( tool: String, policy: ProxyPolicy, ca: Arc, ca_path: PathBuf, secrets: SecretStore, + redactor: Arc>>, ring_buffer: RingBuffer, upstream: Upstream, ) -> std::io::Result { @@ -182,6 +190,7 @@ impl ProxySession { policy, ca, secrets, + redactor, ring_buffer, expected_auth: basic_auth_header(&token), upstream, @@ -242,6 +251,11 @@ struct ProxyContext { policy: ProxyPolicy, ca: Arc, secrets: SecretStore, + /// The daemon's live redactor, not a snapshot of it. A tool runs for + /// minutes and a refreshed token is injected from the *next* request on, + /// so a redactor snapshotted when the session started would not know the + /// value the proxy is now attaching. + redactor: Arc>>, ring_buffer: RingBuffer, /// The full `Proxy-Authorization` value this exec accepts. expected_auth: String, @@ -260,6 +274,13 @@ impl ProxyContext { self.tool )); } + + /// The redactor to apply to one response, taken when that response's + /// headers arrive. The two generations a refresh leaves behind cover a + /// swap that lands between this snapshot and the end of the body. + fn redactor(&self) -> Arc { + Arc::clone(&self.redactor.read().unwrap_or_else(|e| e.into_inner())) + } } // ─── Accept loop ────────────────────────────────────────────────────────────── @@ -539,6 +560,7 @@ async fn handle_tunneled_request( } strip_forbidden_headers(&mut parts.headers, route.inject.as_ref()); + demand_a_plain_full_response(&mut parts.headers); if let Some(inject) = &route.inject { match self::inject_credential(&ctx.secrets, inject, &mut parts.headers) { @@ -556,15 +578,7 @@ async fn handle_tunneled_request( .send(&host, Request::from_parts(parts, body)) .await { - Ok(response) => { - ctx.audit( - &method, - &host, - &path, - &format!("allowed ({})", response.status().as_u16()), - ); - Ok(response.map(|body| body.boxed())) - } + Ok(response) => Ok(forward_response(response, method, host, path, &ctx)), Err(e) => { ctx.audit(&method, &host, &path, &format!("upstream error: {e}")); Ok(refuse( @@ -575,8 +589,263 @@ async fn handle_tunneled_request( } } +// ─── Response handling ──────────────────────────────────────────────────────── + +/// Hand one upstream response back to the tool, redacted. +/// +/// Every header value and every body byte goes through the same automaton the +/// tool's stdout goes through, so the plaintext secret never exists inside the +/// sandbox at all — not in a `-o` file, not in `--dump-header` output, not in +/// a trace. Redaction on the way out is not optional here any more than it is +/// on stdout, so there is no configuration that turns it off. +fn forward_response( + response: Response, + method: Method, + host: String, + path: String, + ctx: &Arc, +) -> Response { + let (mut parts, body) = response.into_parts(); + + if let Err(reason) = vet_response(&parts) { + // `body` is dropped unread. An opaque body is exactly the case where + // forwarding would put bytes the redactor cannot see into the tool's + // hands, so the response is refused rather than passed through. + ctx.audit(&method, &host, &path, reason); + return refuse(StatusCode::BAD_GATEWAY, reason); + } + + // Extensions are how hyper carries the upstream's reason phrase (and its + // original header casing) from the client half to the server half, which + // would write them back out verbatim. A reason phrase is free text — + // `HTTP/1.1 200 ` is a legal status line — and it is outside the + // header map the redactor is about to walk, so none of it is forwarded. + parts.extensions.clear(); + + for name in HOP_BY_HOP_HEADERS { + parts.headers.remove(name); + } + let redactor = ctx.redactor(); + let header_redactions = redact_header_values(&mut parts.headers, &redactor); + + let bodiless = carries_no_body(&method, parts.status); + if !bodiless { + // A placeholder is not the length of the secret it replaces, and the + // body has not been read yet, so an upstream `Content-Length` is + // unknowable here and wrong the moment anything matches. Dropping it + // leaves hyper to frame the response as chunked, which is always + // available on HTTP/1.1. On a bodiless response the length describes + // the representation rather than bytes on the wire, so it is kept. + parts.headers.remove(header::CONTENT_LENGTH); + } + + let mut decision = format!("allowed ({})", parts.status.as_u16()); + if header_redactions > 0 { + decision.push_str(&format!( + " with {header_redactions} header value(s) redacted" + )); + } + ctx.audit(&method, &host, &path, &decision); + + let body: ProxyBody = if bodiless { + empty_body() + } else { + RedactedBody { + inner: body, + stream: StreamRedactor::new(redactor), + ended: false, + audit: ResponseAudit { + ctx: Arc::clone(ctx), + method, + host, + path, + }, + } + .boxed() + }; + Response::from_parts(parts, body) +} + +/// Whether an upstream response may be forwarded at all. +/// +/// The redactor reads bytes, not formats. Anything that leaves the body as +/// something other than its plain, whole representation — a compressed +/// `Content-Encoding`, a transfer coding hyper has not already undone, or a +/// byte range that could begin in the middle of a secret — fails closed +/// instead of reaching the tool unexamined. Airlock does not decompress: a +/// decoder in the response path would be a second parser of attacker-supplied +/// bytes for no security gain, since the request already demands `identity`. +fn vet_response(parts: &hyper::http::response::Parts) -> Result<(), &'static str> { + // Every copy, not the first: two `Content-Encoding` lines are one list, + // and the coding that matters could be in either. + for encoding in parts.headers.get_all(header::CONTENT_ENCODING) { + if !encoding.as_bytes().eq_ignore_ascii_case(b"identity") { + return Err("denied: upstream response is content-encoded and cannot be redacted"); + } + } + for coding in parts.headers.get_all(header::TRANSFER_ENCODING) { + if !coding.as_bytes().eq_ignore_ascii_case(b"chunked") { + return Err("denied: upstream response uses a transfer coding that cannot be redacted"); + } + } + if parts.status == StatusCode::PARTIAL_CONTENT + || parts.headers.contains_key(header::CONTENT_RANGE) + { + return Err("denied: upstream response is a byte range, which may split a secret"); + } + Ok(()) +} + +/// Replace every secret occurrence in every response header value. +/// +/// All values, not a chosen subset: a `Location` carrying the token in a +/// query string, a `Set-Cookie` minted from it, and a debug header some API +/// adds are the same problem, and the set of header names an upstream may use +/// is not knowable in advance. +fn redact_header_values(headers: &mut HeaderMap, redactor: &Redactor) -> usize { + let mut redactions = 0; + let mut clean = HeaderMap::with_capacity(headers.len()); + for (name, value) in headers.iter() { + let redacted = redactor.redact_bytes(value.as_bytes()); + if redacted == value.as_bytes() { + clean.append(name.clone(), value.clone()); + continue; + } + redactions += 1; + // A redacted value that will not rebuild is dropped. The one outcome + // that must not happen is the original going out instead. + if let Ok(value) = HeaderValue::from_bytes(&redacted) { + clean.append(name.clone(), value); + } + } + *headers = clean; + redactions +} + +/// Whether the response has no body on the wire, in which case there is +/// nothing to redact and a `Content-Length` describes the representation the +/// request asked about rather than bytes being sent. +fn carries_no_body(method: &Method, status: StatusCode) -> bool { + *method == Method::HEAD + || status.is_informational() + || status == StatusCode::NO_CONTENT + || status == StatusCode::NOT_MODIFIED +} + +/// The upstream body, redacted frame by frame on its way to the tool. +/// +/// Nothing is buffered beyond the partial match at the end of a frame, so a +/// multi-gigabyte download costs what a small one costs, and no thread or task +/// sits behind it: the tool's own read rate drives the polls, the polls drive +/// the upstream reads, and a slow tool slows the upstream instead of filling +/// the daemon's memory. Dropping it — a disconnected tool, a cancelled +/// session — drops the upstream body with it. +struct RedactedBody { + inner: Incoming, + stream: StreamRedactor, + ended: bool, + audit: ResponseAudit, +} + +/// What the body needs to report its own redaction count once it ends. The +/// count only — never a matched byte. +struct ResponseAudit { + ctx: Arc, + method: Method, + host: String, + path: String, +} + +impl Body for RedactedBody { + type Data = Bytes; + type Error = hyper::Error; + + fn poll_frame( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + ) -> Poll, hyper::Error>>> { + let this = self.get_mut(); + loop { + if this.ended { + return Poll::Ready(None); + } + match Pin::new(&mut this.inner).poll_frame(cx) { + Poll::Pending => return Poll::Pending, + Poll::Ready(Some(Err(e))) => { + this.ended = true; + return Poll::Ready(Some(Err(e))); + } + Poll::Ready(Some(Ok(frame))) => { + // Trailers are dropped rather than forwarded: they arrive + // after the tool has already been handed the body, and no + // request reaches the upstream asking for them. + let Ok(data) = frame.into_data() else { + continue; + }; + let redacted = this.stream.push(&data); + // A frame that was entirely held back as a possible + // partial match yields nothing yet; poll again rather + // than emit an empty frame. + if redacted.is_empty() { + continue; + } + return Poll::Ready(Some(Ok(Frame::data(Bytes::from(redacted))))); + } + Poll::Ready(None) => { + this.ended = true; + let tail = this.stream.finish(); + if tail.is_empty() { + return Poll::Ready(None); + } + return Poll::Ready(Some(Ok(Frame::data(Bytes::from(tail))))); + } + } + } + } +} + +impl Drop for RedactedBody { + /// The second audit line for a request, written when the body ends. + /// + /// It belongs here rather than at end-of-stream so that a body the tool + /// abandoned half-way still reports what it had replaced. The first line + /// went out when the headers arrived and is not held back for this. + fn drop(&mut self) { + let redactions = self.stream.redactions(); + if redactions == 0 { + return; + } + self.audit.ctx.audit( + &self.audit.method, + &self.audit.host, + &self.audit.path, + &format!("response body: {redactions} secret occurrence(s) redacted"), + ); + } +} + // ─── Header handling ────────────────────────────────────────────────────────── +/// Constrain the upstream request so that its response is something the +/// redactor can read. +/// +/// `Accept-Encoding: identity` replaces whatever the tool asked for: a gzip, +/// br or zstd body is opaque to a byte-pattern scanner, and an upstream that +/// compresses anyway is refused rather than forwarded. `Range` and `If-Range` +/// go because a range may begin in the middle of a secret — the pattern would +/// be split across two responses the proxy never sees together, and the tool +/// would reassemble the plaintext in a file. Stripping them makes the +/// upstream send the whole representation, which is the form redaction is +/// sound on. +fn demand_a_plain_full_response(headers: &mut HeaderMap) { + headers.insert( + header::ACCEPT_ENCODING, + HeaderValue::from_static("identity"), + ); + headers.remove(header::RANGE); + headers.remove(header::IF_RANGE); +} + /// Remove hop-by-hop headers and every client-supplied copy of the header this /// route injects. /// diff --git a/src/proxy/server/tests.rs b/src/proxy/server/tests.rs index 6c4f251..78f40f3 100644 --- a/src/proxy/server/tests.rs +++ b/src/proxy/server/tests.rs @@ -22,6 +22,12 @@ use crate::proxy::{HostPattern, Inject, PathRule, ProxyRoute}; use crate::secrets::{Secret, SecretSlot}; const UPSTREAM_HOST: &str = "upstream.test"; +const SECRET_LABEL: &str = "api_key"; +const SECRET_VALUE: &str = "s3cret-value"; +const PLACEHOLDER: &str = "[REDACTED:api_key]"; +/// Literal text the injected header wraps the secret in, so a test can take +/// the secret back out of what the upstream received. +const INJECT_PREFIX: &str = "tok-"; // ─── Fixtures ───────────────────────────────────────────────────────────────── @@ -35,7 +41,7 @@ fn route(host: &str, allow: &[&str], inject: Option) -> ProxyRoute { } fn bearer() -> Inject { - Inject::parse("X-Test-Secret", "tok-{secret}", "api_key").unwrap() + Inject::parse("X-Test-Secret", "tok-{secret}", SECRET_LABEL).unwrap() } fn secret_store(label: &str, value: &str, healthy: bool) -> SecretStore { @@ -58,6 +64,40 @@ fn secret_store(label: &str, value: &str, healthy: bool) -> SecretStore { Arc::new(map) } +/// The daemon's live redactor handle, built over the same store the proxy +/// injects from — which is how the real daemon builds it. +fn live_redactor(secrets: &SecretStore) -> Arc>> { + let values: Vec<(String, Arc>)> = secrets + .iter() + .map(|(name, slot)| { + let slot = slot.read().unwrap(); + (name.clone(), Arc::clone(&slot.value)) + }) + .collect(); + let refs: Vec<(&str, &Secret)> = values + .iter() + .map(|(name, value)| (name.as_str(), value.as_ref())) + .collect(); + Arc::new(RwLock::new(Arc::new(Redactor::new(refs).unwrap()))) +} + +fn encode_base64(value: &str) -> String { + use base64::Engine; + base64::engine::general_purpose::STANDARD.encode(value.as_bytes()) +} + +fn encode_url(value: &str) -> String { + percent_encoding::utf8_percent_encode(value, percent_encoding::NON_ALPHANUMERIC).to_string() +} + +fn encode_hex(value: &str) -> String { + value + .as_bytes() + .iter() + .map(|b| format!("{b:02x}")) + .collect() +} + fn pem_to_der(pem: &str) -> Vec { use base64::Engine; let body: String = pem @@ -89,12 +129,172 @@ fn ca_for(hosts: &[&str]) -> Arc { // ─── The local "upstream" ───────────────────────────────────────────────────── +/// Every request the upstream received, recorded out of band. +/// +/// The proxy redacts what it forwards, so an echo that comes back through the +/// proxy can no longer tell a test what the upstream actually saw. This can. +type Seen = Arc>>; + /// A TLS server that reports back what it received: the request target and one -/// line per header. That is what lets a test see whether the credential -/// arrived and whether the query string survived. +/// line per header. Paths other than the default select a canned response +/// shape — a compressed body, a leaking header, a body delivered in pieces — +/// so that each thing the response path has to cope with has an upstream that +/// produces it. struct TestUpstream { addr: SocketAddr, ca: Arc, + seen: Seen, +} + +/// What the upstream received, as `METHOD target` followed by sorted headers. +fn request_summary(req: &Request) -> String { + let mut out = format!( + "{} {}\n", + req.method(), + req.uri().path_and_query().map(|p| p.as_str()).unwrap_or("") + ); + let mut headers: Vec<_> = req + .headers() + .iter() + .map(|(n, v)| format!("{}: {}\n", n, v.to_str().unwrap_or(""))) + .collect(); + headers.sort(); + out.extend(headers); + out +} + +/// The secret the proxy injected, recovered from the credential header. The +/// upstream learns it the same way a real one would: it was sent the thing. +fn injected_secret(req: &Request) -> String { + req.headers() + .get("x-test-secret") + .and_then(|v| v.to_str().ok()) + .and_then(|v| v.strip_prefix(INJECT_PREFIX)) + .unwrap_or("") + .to_string() +} + +/// A body whose frames arrive one at a time from a task, so that a secret can +/// be made to straddle a frame boundary and a long download can be stopped +/// half-way. The channel holds one frame, so a reader that stops reading stops +/// the sender. +struct FramedBody(tokio::sync::mpsc::Receiver>); + +impl Body for FramedBody { + type Data = Bytes; + type Error = hyper::Error; + + fn poll_frame( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + ) -> Poll, hyper::Error>>> { + match self.get_mut().0.poll_recv(cx) { + Poll::Ready(Some(chunk)) => Poll::Ready(Some(Ok(Frame::data(Bytes::from(chunk))))), + Poll::Ready(None) => Poll::Ready(None), + Poll::Pending => Poll::Pending, + } + } +} + +fn framed_body(chunks: Vec>) -> ProxyBody { + let (tx, rx) = tokio::sync::mpsc::channel::>(1); + tokio::spawn(async move { + for chunk in chunks { + if tx.send(chunk).await.is_err() { + return; + } + tokio::time::sleep(Duration::from_millis(5)).await; + } + }); + FramedBody(rx).boxed() +} + +fn status_only(status: StatusCode, headers: &[(&str, &str)]) -> Response { + let mut response = Response::new(empty_body()); + *response.status_mut() = status; + for (name, value) in headers { + response.headers_mut().insert( + HeaderName::from_bytes(name.as_bytes()).unwrap(), + HeaderValue::from_str(value).unwrap(), + ); + } + response +} + +async fn upstream_response(req: Request, seen: Seen) -> Response { + let summary = request_summary(&req); + seen.lock().unwrap().push(summary.clone()); + let secret = injected_secret(&req); + + match req.uri().path() { + // The credential in every encoding the redactor knows. + "/encoded" => Response::new(text_body(&format!( + "b64 {} url {} hex {}", + encode_base64(&secret), + encode_url(&secret), + encode_hex(&secret) + ))), + // The credential cut in half across two frames. + "/frames" => { + let (head, tail) = secret.split_at(secret.len() / 2); + Response::new(framed_body(vec![ + b"start ".to_vec(), + head.as_bytes().to_vec(), + format!("{tail} end").into_bytes(), + ])) + } + // The credential in header values rather than the body. + "/header-leak" => { + let mut response = Response::new(text_body("see the headers")); + response.headers_mut().insert( + "location", + HeaderValue::from_str(&format!("https://{UPSTREAM_HOST}/next?token={secret}")) + .unwrap(), + ); + response.headers_mut().insert( + "set-cookie", + HeaderValue::from_str(&format!("session={secret}; Path=/")).unwrap(), + ); + response + } + // The credential in the status line's reason phrase. + "/reason-leak" => { + let mut response = Response::new(text_body("see the status line")); + response.extensions_mut().insert( + hyper::ext::ReasonPhrase::try_from(format!("OK {secret}").into_bytes()).unwrap(), + ); + response + } + // An upstream that ignores `Accept-Encoding: identity`. + "/gzip" => { + let mut response = Response::new(text_body(&format!("not really gzip, but {secret}"))); + response + .headers_mut() + .insert(header::CONTENT_ENCODING, HeaderValue::from_static("gzip")); + response + } + // A partial representation, which a range request would have asked for. + "/partial" => { + let mut response = Response::new(text_body(&secret)); + *response.status_mut() = StatusCode::PARTIAL_CONTENT; + response.headers_mut().insert( + header::CONTENT_RANGE, + HeaderValue::from_static("bytes 0-11/24"), + ); + response + } + // Tens of megabytes with the credential buried in the middle. + "/large" => { + let filler = vec![b'x'; 4 * 1024 * 1024]; + let mut chunks: Vec> = vec![filler.clone(); 4]; + chunks.push(secret.into_bytes()); + chunks.extend(std::iter::repeat_n(filler, 4)); + Response::new(framed_body(chunks)) + } + "/no-content" => status_only(StatusCode::NO_CONTENT, &[]), + "/not-modified" => status_only(StatusCode::NOT_MODIFIED, &[("content-length", "42")]), + _ => Response::new(text_body(&summary)), + } } async fn start_upstream() -> TestUpstream { @@ -102,31 +302,23 @@ async fn start_upstream() -> TestUpstream { let listener = TcpListener::bind((Ipv4Addr::LOCALHOST, 0)).await.unwrap(); let addr = listener.local_addr().unwrap(); let acceptor = TlsAcceptor::from(ca.server_config(UPSTREAM_HOST).unwrap()); + let seen: Seen = Arc::new(std::sync::Mutex::new(Vec::new())); + let accepted = Arc::clone(&seen); tokio::spawn(async move { loop { let Ok((stream, _)) = listener.accept().await else { return; }; let acceptor = acceptor.clone(); + let seen = Arc::clone(&accepted); tokio::spawn(async move { let Ok(tls) = acceptor.accept(stream).await else { return; }; - let service = service_fn(|req: Request| async move { - let mut out = format!( - "{} {}\n", - req.method(), - req.uri().path_and_query().map(|p| p.as_str()).unwrap_or("") - ); - let mut names: Vec<_> = req - .headers() - .iter() - .map(|(n, v)| format!("{}: {}\n", n, v.to_str().unwrap_or(""))) - .collect(); - names.sort(); - out.extend(names); - Ok::<_, Infallible>(Response::new(text_body(&out))) + let service = service_fn(move |req: Request| { + let seen = Arc::clone(&seen); + async move { Ok::<_, Infallible>(upstream_response(req, seen).await) } }); let _ = hyper::server::conn::http1::Builder::new() .serve_connection(TokioIo::new(tls), service) @@ -135,7 +327,7 @@ async fn start_upstream() -> TestUpstream { } }); - TestUpstream { addr, ca } + TestUpstream { addr, ca, seen } } // ─── The harness ────────────────────────────────────────────────────────────── @@ -144,6 +336,9 @@ struct Harness { session: ProxySession, ca: Arc, ring_buffer: RingBuffer, + secrets: SecretStore, + redactor: Arc>>, + seen: Seen, } impl Harness { @@ -151,12 +346,14 @@ impl Harness { let upstream = start_upstream().await; let ca = ca_for(&[UPSTREAM_HOST, "other.test"]); let ring_buffer = RingBuffer::new(); + let redactor = live_redactor(&secrets); let session = ProxySession::start_with_upstream( "curl".to_string(), ProxyPolicy { routes }, Arc::clone(&ca), PathBuf::from("/nonexistent/airlock-ca.pem"), - secrets, + Arc::clone(&secrets), + Arc::clone(&redactor), ring_buffer.clone(), Upstream::fixed(upstream.addr, root_store(&upstream.ca)), ) @@ -165,6 +362,9 @@ impl Harness { session, ca, ring_buffer, + secrets, + redactor, + seen: upstream.seen, } } @@ -173,11 +373,26 @@ impl Harness { async fn default() -> Self { Self::new( vec![route(UPSTREAM_HOST, &["GET /**"], Some(bearer()))], - secret_store("api_key", "s3cret-value", true), + secret_store(SECRET_LABEL, SECRET_VALUE, true), ) .await } + /// What the upstream received, as it received it. + fn seen(&self) -> String { + self.seen.lock().unwrap().join("\n") + } + + /// Stand in for a background refresh: swap the stored value and rebuild + /// the redactor, exactly as [`crate::refresh::refresh_once`] does. + fn refresh_secret(&self, value: &str) { + { + let mut slot = self.secrets[SECRET_LABEL].write().unwrap(); + slot.value = Arc::new(Secret::new(value.to_string())); + } + *self.redactor.write().unwrap() = Arc::clone(&live_redactor(&self.secrets).read().unwrap()); + } + fn logs(&self) -> String { self.ring_buffer .entries() @@ -489,8 +704,13 @@ async fn a_permitted_request_reaches_the_upstream_with_the_credential() { assert_eq!(status, StatusCode::OK); assert!( - body.contains("x-test-secret: tok-s3cret-value"), - "the upstream should have seen the injected credential, got:\n{body}" + h.seen().contains("x-test-secret: tok-s3cret-value"), + "the upstream should have seen the injected credential, got:\n{}", + h.seen() + ); + assert!( + body.contains(&format!("x-test-secret: tok-{PLACEHOLDER}")), + "the echo of it must come back redacted, got:\n{body}" ); assert!(h.logs().contains("allowed (200)"), "{}", h.logs()); } @@ -511,19 +731,21 @@ async fn a_client_supplied_copy_of_the_injected_header_is_replaced() { ) .await; + let seen = h.seen(); assert!( - body.contains("x-test-secret: tok-s3cret-value"), - "got:\n{body}" + seen.contains("x-test-secret: tok-s3cret-value"), + "got:\n{seen}" ); assert!( - !body.contains("attacker-chosen"), - "the client's copy must not survive, got:\n{body}" + !seen.contains("attacker-chosen"), + "the client's copy must not survive, got:\n{seen}" ); assert_eq!( - body.matches("x-test-secret").count(), + seen.matches("x-test-secret").count(), 1, - "the upstream must see exactly one copy, got:\n{body}" + "the upstream must see exactly one copy, got:\n{seen}" ); + assert!(!body.contains(SECRET_VALUE), "got:\n{body}"); } #[tokio::test] @@ -652,9 +874,11 @@ async fn one_tunnel_serves_several_requests() { assert_eq!(second_status, StatusCode::OK); assert!(first.starts_with("GET /v1/one\n"), "{first}"); assert!(second.starts_with("GET /v1/two\n"), "{second}"); - assert!( - second.contains("x-test-secret: tok-s3cret-value"), - "the credential is looked up per request, got:\n{second}" + assert_eq!( + h.seen().matches("x-test-secret: tok-s3cret-value").count(), + 2, + "the credential is looked up per request, got:\n{}", + h.seen() ); assert_eq!( h.logs().matches("allowed (200)").count(), @@ -687,6 +911,452 @@ async fn the_secret_never_reaches_the_ring_buffer() { assert!(!logs.contains(&h.token()), "{logs}"); } +// ─── Response redaction ─────────────────────────────────────────────────────── + +// Everything the upstream sends back goes through the redactor, so a test +// asserting on a response is asserting on what the *tool* would have written +// to a file. `h.seen()` is the only window onto the unredacted truth. + +fn request(method: Method, target: &str, extra: &[(&str, &str)]) -> Request { + let mut builder = Request::builder() + .method(method) + .uri(target) + .header(header::HOST, UPSTREAM_HOST); + for (name, value) in extra { + builder = builder.header(*name, *value); + } + builder.body(empty_body()).unwrap() +} + +/// A harness whose route permits every method, for HEAD and the canned +/// response endpoints. +async fn any_method_harness() -> Harness { + Harness::new( + vec![route(UPSTREAM_HOST, &["* /**"], Some(bearer()))], + secret_store(SECRET_LABEL, SECRET_VALUE, true), + ) + .await +} + +/// Read a response body frame by frame, returning the byte count and whether +/// any window of the stream ever spelled the secret out. +async fn scan_body(response: Response) -> (usize, bool) { + let mut body = response.into_body(); + let mut total = 0; + let mut leaked = false; + // A window wide enough that a secret straddling two frames is still seen. + let mut window: Vec = Vec::new(); + while let Some(frame) = body.frame().await { + let frame = frame.unwrap(); + let Ok(data) = frame.into_data() else { + continue; + }; + total += data.len(); + window.extend_from_slice(&data); + if String::from_utf8_lossy(&window).contains(SECRET_VALUE) { + leaked = true; + } + let keep = window.len().saturating_sub(64); + window.drain(..keep); + } + (total, leaked) +} + +#[tokio::test] +async fn an_echoed_credential_is_redacted_in_the_response_body() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let (status, body) = read_body( + tunnel + .send_request(get("/v1/things", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::OK); + assert!(!body.contains(SECRET_VALUE), "got:\n{body}"); + assert!(body.contains(PLACEHOLDER), "got:\n{body}"); + // Byte-for-byte what the single-shot redactor makes of what the upstream + // sent: streaming must not change the result. + let expected = h + .redactor + .read() + .unwrap() + .redact_bytes(format!("{}\n", h.seen()).as_bytes()); + assert_eq!(body.as_bytes(), expected.as_slice()); +} + +#[tokio::test] +async fn encoded_spellings_of_the_secret_are_redacted() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let (_, body) = read_body( + tunnel + .send_request(get("/encoded", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + + for encoded in [ + encode_base64(SECRET_VALUE), + encode_url(SECRET_VALUE), + encode_hex(SECRET_VALUE), + ] { + assert!(!body.contains(&encoded), "{encoded} survived in:\n{body}"); + } + assert_eq!(body.matches(PLACEHOLDER).count(), 3, "got:\n{body}"); +} + +#[tokio::test] +async fn a_secret_split_across_upstream_frames_is_redacted() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let (_, body) = read_body( + tunnel + .send_request(get("/frames", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + + assert_eq!(body, format!("start {PLACEHOLDER} end"), "got:\n{body}"); +} + +#[tokio::test] +async fn secrets_in_response_header_values_are_redacted() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let response = tunnel + .send_request(get("/header-leak", UPSTREAM_HOST, &[])) + .await + .unwrap(); + + let location = response.headers()[header::LOCATION].to_str().unwrap(); + let cookie = response.headers()["set-cookie"].to_str().unwrap(); + assert_eq!( + location, + format!("https://{UPSTREAM_HOST}/next?token={PLACEHOLDER}") + ); + assert_eq!(cookie, format!("session={PLACEHOLDER}; Path=/")); + assert!( + h.logs().contains("with 2 header value(s) redacted"), + "{}", + h.logs() + ); +} + +#[tokio::test] +async fn the_upstream_content_length_is_not_forwarded() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let response = tunnel + .send_request(get("/v1/things", UPSTREAM_HOST, &[])) + .await + .unwrap(); + + assert!( + response.headers().get(header::CONTENT_LENGTH).is_none(), + "a length computed before redaction would be wrong; hyper must chunk instead" + ); + // Reading to the end must neither truncate nor hang. + let (_, body) = read_body(response).await; + assert!(body.ends_with('\n'), "got:\n{body}"); +} + +#[tokio::test] +async fn a_head_response_keeps_its_length_and_carries_no_body() { + let h = any_method_harness().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let response = tunnel + .send_request(request(Method::HEAD, "/v1/things", &[])) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + assert!( + response.headers().get(header::CONTENT_LENGTH).is_some(), + "on a bodiless response the length is metadata and must survive" + ); + let (_, body) = read_body(response).await; + assert!(body.is_empty(), "a HEAD response has no body, got:\n{body}"); +} + +#[tokio::test] +async fn bodiless_statuses_pass_through() { + let h = any_method_harness().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + + let (status, body) = read_body( + tunnel + .send_request(get("/no-content", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + assert_eq!(status, StatusCode::NO_CONTENT); + assert!(body.is_empty()); + + // The proxy keeps the upstream's `Content-Length` on a 304 — hyper then + // drops it on the wire, as it does for every 204 and 304 it writes. What + // matters here is that nothing is framed as a body and nothing hangs. + let response = tunnel + .send_request(get("/not-modified", UPSTREAM_HOST, &[])) + .await + .unwrap(); + assert_eq!(response.status(), StatusCode::NOT_MODIFIED); + let (_, body) = read_body(response).await; + assert!(body.is_empty()); +} + +#[tokio::test] +async fn a_compressed_response_fails_closed() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let (status, body) = read_body( + tunnel + .send_request(get("/gzip", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::BAD_GATEWAY); + assert!(body.contains("content-encoded"), "got:\n{body}"); + assert!( + !body.contains(SECRET_VALUE) && !body.contains("not really gzip"), + "no byte of an unreadable body may be forwarded, got:\n{body}" + ); + assert!(h.logs().contains("content-encoded"), "{}", h.logs()); +} + +#[tokio::test] +async fn a_partial_response_fails_closed() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let (status, body) = read_body( + tunnel + .send_request(get("/partial", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + + assert_eq!(status, StatusCode::BAD_GATEWAY); + assert!(body.contains("byte range"), "got:\n{body}"); + assert!(!body.contains(SECRET_VALUE), "got:\n{body}"); +} + +#[tokio::test] +async fn the_client_accept_encoding_is_replaced_and_ranges_are_stripped() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let _ = tunnel + .send_request(get( + "/v1/things", + UPSTREAM_HOST, + &[ + ("accept-encoding", "gzip, deflate, br"), + ("range", "bytes=0-99"), + ("if-range", "\"etag\""), + ], + )) + .await + .unwrap(); + + let seen = h.seen(); + assert!(seen.contains("accept-encoding: identity"), "got:\n{seen}"); + assert!(!seen.contains("gzip"), "got:\n{seen}"); + assert!(!seen.contains("range:"), "got:\n{seen}"); +} + +#[tokio::test] +async fn a_large_body_streams_through_redacted() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let response = tunnel + .send_request(get("/large", UPSTREAM_HOST, &[])) + .await + .unwrap(); + + let (total, leaked) = scan_body(response).await; + assert!(!leaked, "the secret must not survive a 33 MB stream"); + assert_eq!( + total, + 32 * 1024 * 1024 + PLACEHOLDER.len(), + "every filler byte must arrive, and the secret must arrive replaced" + ); + assert!( + h.logs() + .contains("response body: 1 secret occurrence(s) redacted"), + "{}", + h.logs() + ); +} + +#[tokio::test] +async fn a_secret_refreshed_mid_session_is_redacted_in_the_next_response() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + + let (_, first) = read_body( + tunnel + .send_request(get("/v1/one", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + assert!(first.contains(PLACEHOLDER), "got:\n{first}"); + + h.refresh_secret("rotated-value"); + + let (_, second) = read_body( + tunnel + .send_request(get("/v1/two", UPSTREAM_HOST, &[])) + .await + .unwrap(), + ) + .await; + assert!( + h.seen().contains("x-test-secret: tok-rotated-value"), + "the proxy must inject the fresh value, got:\n{}", + h.seen() + ); + assert!( + !second.contains("rotated-value"), + "and must redact the value it just injected, got:\n{second}" + ); + assert!(second.contains(PLACEHOLDER), "got:\n{second}"); +} + +#[tokio::test] +async fn dropping_the_session_ends_a_response_body_mid_stream() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let response = tunnel + .send_request(get("/large", UPSTREAM_HOST, &[])) + .await + .unwrap(); + let mut body = response.into_body(); + + let first = body.frame().await.expect("a first frame").unwrap(); + assert!(first.data_ref().is_some()); + + drop(h); + tokio::time::sleep(Duration::from_millis(50)).await; + + let mut delivered = 0; + while let Some(frame) = body.frame().await { + match frame { + Ok(frame) => delivered += frame.data_ref().map_or(0, |d| d.len()), + Err(_) => break, + } + } + assert!( + delivered < 32 * 1024 * 1024, + "a cancelled session must stop the upstream read, not drain it" + ); +} + +// ─── Response handling (pure) ───────────────────────────────────────────────── + +fn response_parts(status: StatusCode, headers: &[(&str, &str)]) -> hyper::http::response::Parts { + let mut builder = Response::builder().status(status); + for (name, value) in headers { + builder = builder.header(*name, *value); + } + builder.body(()).unwrap().into_parts().0 +} + +#[test] +fn only_a_plain_whole_representation_is_forwarded() { + assert!(vet_response(&response_parts(StatusCode::OK, &[])).is_ok()); + assert!( + vet_response(&response_parts( + StatusCode::OK, + &[("content-encoding", "identity")] + )) + .is_ok() + ); + assert!( + vet_response(&response_parts( + StatusCode::OK, + &[("transfer-encoding", "chunked")] + )) + .is_ok() + ); + + for bad in [ + vec![("content-encoding", "gzip")], + vec![("content-encoding", "br")], + vec![("content-encoding", "zstd")], + vec![("content-encoding", "deflate")], + vec![("content-encoding", "gzip, identity")], + vec![ + ("content-encoding", "identity"), + ("content-encoding", "gzip"), + ], + vec![("transfer-encoding", "gzip, chunked")], + vec![("content-range", "bytes 0-9/100")], + ] { + assert!( + vet_response(&response_parts(StatusCode::OK, &bad)).is_err(), + "{bad:?} should fail closed" + ); + } + assert!(vet_response(&response_parts(StatusCode::PARTIAL_CONTENT, &[])).is_err()); +} + +#[test] +fn only_bodiless_responses_keep_their_length() { + assert!(carries_no_body(&Method::HEAD, StatusCode::OK)); + assert!(carries_no_body(&Method::GET, StatusCode::NO_CONTENT)); + assert!(carries_no_body(&Method::GET, StatusCode::NOT_MODIFIED)); + assert!(carries_no_body(&Method::GET, StatusCode::CONTINUE)); + assert!(!carries_no_body(&Method::GET, StatusCode::OK)); + assert!(!carries_no_body(&Method::POST, StatusCode::CREATED)); +} + +#[test] +fn every_header_value_is_scanned_and_an_unrebuildable_one_is_dropped() { + let secret = Secret::new("abc".to_string()); + let redactor = Redactor::new([("KEY", &secret)]).unwrap(); + + let mut headers = HeaderMap::new(); + headers.insert(header::LOCATION, HeaderValue::from_static("/next?t=abc")); + headers.append("set-cookie", HeaderValue::from_static("a=abc")); + headers.append("set-cookie", HeaderValue::from_static("b=plain")); + headers.insert("x-clean", HeaderValue::from_static("nothing here")); + + assert_eq!(redact_header_values(&mut headers, &redactor), 2); + assert_eq!(headers[header::LOCATION], "/next?t=[REDACTED:KEY]"); + let cookies: Vec<_> = headers + .get_all("set-cookie") + .iter() + .map(|v| v.to_str().unwrap()) + .collect(); + assert_eq!(cookies, vec!["a=[REDACTED:KEY]", "b=plain"]); + assert_eq!(headers["x-clean"], "nothing here"); +} + +#[test] +fn the_upstream_request_asks_for_a_plain_whole_representation() { + let mut headers = HeaderMap::new(); + headers.insert( + header::ACCEPT_ENCODING, + HeaderValue::from_static("gzip, br"), + ); + headers.insert(header::RANGE, HeaderValue::from_static("bytes=0-99")); + headers.insert(header::IF_RANGE, HeaderValue::from_static("\"etag\"")); + + demand_a_plain_full_response(&mut headers); + + assert_eq!(headers[header::ACCEPT_ENCODING], "identity"); + assert!(headers.get(header::RANGE).is_none()); + assert!(headers.get(header::IF_RANGE).is_none()); +} + // ─── Header assembly ────────────────────────────────────────────────────────── #[test] @@ -900,8 +1570,13 @@ async fn system_curl_completes_a_request_through_the_proxy() { "got:\n{stdout}" ); assert!( - stdout.contains("x-test-secret: tok-s3cret-value"), - "the credential must reach the upstream, got:\n{stdout}" + h.seen().contains("x-test-secret: tok-s3cret-value"), + "the credential must reach the upstream, got:\n{}", + h.seen() + ); + assert!( + !stdout.contains(SECRET_VALUE), + "but must not come back, got:\n{stdout}" ); assert!( !stdout.contains("proxy-authorization"), @@ -909,6 +1584,71 @@ async fn system_curl_completes_a_request_through_the_proxy() { ); } +/// The file-write path the design record left open: `curl -o` puts the +/// response body straight into the sandbox root, where the agent reads it +/// without the daemon's stdout redaction ever touching it. What lands there +/// must already be redacted, because the proxy redacted it. +#[tokio::test] +async fn system_curl_writes_a_file_that_does_not_contain_the_secret() { + let curl = std::path::Path::new("/usr/bin/curl"); + if !curl.exists() { + return; + } + + let h = Harness::default().await; + let dir = tempfile::tempdir().unwrap(); + let ca_path = dir.path().join("airlock-ca.pem"); + h.ca.write_cert_pem(&ca_path).unwrap(); + let out_path = dir.path().join("out.json"); + let header_path = dir.path().join("headers.txt"); + + let mut env: HashMap = HashMap::new(); + h.session.apply_env(&mut env); + for name in CA_BUNDLE_VARS { + env.insert((*name).to_string(), ca_path.display().to_string()); + } + + let args = vec![ + "-sS".to_string(), + "--max-time".to_string(), + "20".to_string(), + "--compressed".to_string(), + "-o".to_string(), + out_path.display().to_string(), + "--dump-header".to_string(), + header_path.display().to_string(), + format!("https://{UPSTREAM_HOST}/header-leak"), + ]; + let output = tokio::task::spawn_blocking(move || { + std::process::Command::new("/usr/bin/curl") + .args(&args) + .env_clear() + .envs(&env) + .output() + .unwrap() + }) + .await + .unwrap(); + + let stderr = String::from_utf8_lossy(&output.stderr); + assert!(output.status.success(), "curl failed: {stderr}"); + + let written = std::fs::read_to_string(&out_path).unwrap(); + let headers = std::fs::read_to_string(&header_path).unwrap(); + assert!( + h.seen().contains("x-test-secret: tok-s3cret-value"), + "the upstream must still have been authenticated" + ); + assert!( + !written.contains(SECRET_VALUE) && !headers.contains(SECRET_VALUE), + "nothing curl wrote may hold the secret:\nbody: {written}\nheaders: {headers}" + ); + assert!( + headers.contains(PLACEHOLDER), + "the leaking headers must arrive redacted, got:\n{headers}" + ); +} + /// The same client, asking for a host no route covers. #[tokio::test] async fn system_curl_is_refused_for_an_unrouted_host() { @@ -1000,3 +1740,22 @@ fn connect_authority_is_canonicalized() { let no_port = "api.example.com".parse().unwrap(); assert_eq!(split_authority(&no_port), None); } + +#[tokio::test] +async fn a_secret_in_the_reason_phrase_is_not_forwarded() { + let h = Harness::default().await; + let mut tunnel = h.tunnel(UPSTREAM_HOST).await; + let response = tunnel + .send_request(get("/reason-leak", UPSTREAM_HOST, &[])) + .await + .unwrap(); + + assert_eq!(response.status(), StatusCode::OK); + assert!( + response + .extensions() + .get::() + .is_none(), + "the upstream's reason phrase is free text the redactor never sees" + ); +} diff --git a/src/redact.rs b/src/redact.rs index 880a6dd..9222ebb 100644 --- a/src/redact.rs +++ b/src/redact.rs @@ -228,6 +228,110 @@ impl Redactor { None => 0, } } + + /// The length of the longest prefix of `buf` whose redaction is already + /// decided, and the number of matches inside it. + /// + /// The automaton reports a match at the position its last byte lands on, + /// so a match found in `buf` can never be created, extended or displaced + /// by bytes that arrive later. What later bytes *can* do is complete a + /// pattern that starts near the end of `buf`, and such a pattern must + /// begin within the last `max_pattern_len - 1` bytes. Everything before + /// that — and everything up to the end of the last match, which the + /// non-overlapping scan resumes from — is settled. + fn settled_prefix(&self, buf: &[u8]) -> (usize, usize) { + let Some(automaton) = &self.automaton else { + return (buf.len(), 0); + }; + let mut matches = 0; + let mut end_of_last_match = 0; + for m in automaton.find_iter(buf) { + matches += 1; + end_of_last_match = m.end(); + } + let unfinishable = buf + .len() + .saturating_sub(automaton.max_pattern_len().saturating_sub(1)); + (end_of_last_match.max(unfinishable), matches) + } +} + +// ─── Incremental streaming redaction ────────────────────────────────────────── + +/// Chunk-wise redaction for callers that hold the pieces themselves. +/// +/// Feeding a stream through [`StreamRedactor::push`] and finally +/// [`StreamRedactor::finish`] yields, concatenated, exactly what +/// [`Redactor::redact_bytes`] produces for the whole stream — for *any* +/// chunking, including one that cuts a secret in half. Chunks that could still +/// turn out to be the start of a pattern are held back until the next one +/// arrives, so the memory held between chunks never exceeds the longest +/// pattern. +/// +/// This is the counterpart of [`Redactor::redact_stream`] for data that +/// arrives as owned buffers rather than through a blocking `Read`: the proxy's +/// response path drives it from a `poll` with no thread and no channel behind +/// it, so backpressure is whatever the caller's own polling imposes. +pub struct StreamRedactor { + redactor: Arc, + /// Trailing bytes a pattern could still be starting in. + carry: Vec, + redactions: usize, +} + +impl StreamRedactor { + /// Start a stream against a snapshot of the redactor. + pub fn new(redactor: Arc) -> Self { + StreamRedactor { + redactor, + carry: Vec::new(), + redactions: 0, + } + } + + /// Feed the next chunk and take back whatever is now settled. + pub fn push(&mut self, chunk: &[u8]) -> Vec { + self.carry.extend_from_slice(chunk); + let (settled, matches) = self.redactor.settled_prefix(&self.carry); + self.redactions += matches; + let out = self.redactor.redact_bytes(&self.carry[..settled]); + self.carry.drain(..settled); + out + } + + /// Flush the held-back tail. Nothing more may be pushed afterwards — at + /// end of stream there is no later byte to hold anything back for. + pub fn finish(&mut self) -> Vec { + let (_, matches) = self.redactor.settled_prefix(&self.carry); + self.redactions += matches; + let out = self.redactor.redact_bytes(&self.carry); + self.carry.clear(); + out + } + + /// How many secret occurrences have been replaced so far. Counted, never + /// the bytes themselves. + pub fn redactions(&self) -> usize { + self.redactions + } + + /// Bytes currently held back. Tests assert the bound; nothing else needs + /// to know. + #[cfg(test)] + fn carry_len(&self) -> usize { + self.carry.len() + } +} + +#[cfg(test)] +impl Redactor { + /// The longest pattern in the automaton, which bounds what a + /// [`StreamRedactor`] can be holding between chunks. + fn longest_pattern(&self) -> usize { + self.automaton + .as_ref() + .map_or(0, |automaton| automaton.max_pattern_len()) + } } // ─── Lossy UTF-8 conversion ────────────────────────────────────────────────── @@ -654,6 +758,127 @@ mod tests { ); } + // ── Incremental streaming redaction ─────────────────────────────────── + + /// Drive a [`StreamRedactor`] over `input` cut at `splits` and return the + /// concatenated output. + fn stream_in_chunks(redactor: &Arc, input: &[u8], chunks: &[&[u8]]) -> Vec { + let mut stream = StreamRedactor::new(Arc::clone(redactor)); + let mut out = Vec::new(); + for chunk in chunks { + out.extend_from_slice(&stream.push(chunk)); + assert!( + stream.carry_len() < redactor.longest_pattern().max(1), + "the held-back tail must stay under the longest pattern" + ); + } + out.extend_from_slice(&stream.finish()); + assert_eq!( + redactor.redact_bytes(input), + out, + "chunked output must equal whole-input redaction" + ); + out + } + + fn streaming_fixture() -> (Arc, Vec) { + let secret = Secret::new("s3cret-value".to_string()); + let other = Secret::new("s3cret".to_string()); + let redactor = Arc::new(Redactor::new([("API_KEY", &secret), ("SHORT", &other)]).unwrap()); + + let mut input = Vec::new(); + input.extend_from_slice(b"prefix s3cret-value middle "); + input.extend_from_slice(encode_base64("s3cret-value").as_bytes()); + input.extend_from_slice(b" then "); + input.extend_from_slice(encode_hex("s3cret-value").as_bytes()); + input.extend_from_slice(b" and s3cret alone, s3cret-valu (partial), tail"); + (redactor, input) + } + + #[test] + fn streaming_any_two_way_split_equals_whole_input() { + let (redactor, input) = streaming_fixture(); + for split in 0..=input.len() { + let (a, b) = input.split_at(split); + stream_in_chunks(&redactor, &input, &[a, b]); + } + } + + #[test] + fn streaming_any_fixed_chunk_size_equals_whole_input() { + let (redactor, input) = streaming_fixture(); + for size in 1..=input.len() { + let chunks: Vec<&[u8]> = input.chunks(size).collect(); + stream_in_chunks(&redactor, &input, &chunks); + } + } + + #[test] + fn streaming_splits_inside_a_match_and_around_a_replacement() { + let secret = Secret::new("abcdef".to_string()); + let redactor = Arc::new(Redactor::new([("KEY", &secret)]).unwrap()); + // Two adjacent occurrences: a split anywhere between them lands inside + // a match, and the emitted replacement is longer than what it + // replaced, so output and input offsets no longer line up. + let input = b"xxabcdefabcdefyy".to_vec(); + for split in 0..=input.len() { + let (a, b) = input.split_at(split); + let out = stream_in_chunks(&redactor, &input, &[a, b]); + assert_eq!(out, b"xx[REDACTED:KEY][REDACTED:KEY]yy"); + } + } + + #[test] + fn streaming_byte_at_a_time_redacts_every_variant() { + let (redactor, input) = streaming_fixture(); + let chunks: Vec<&[u8]> = input.chunks(1).collect(); + let out = stream_in_chunks(&redactor, &input, &chunks); + let text = String::from_utf8_lossy(&out); + assert!( + !text.contains("s3cret"), + "no spelling of the secret may survive a one-byte-at-a-time stream: {text}" + ); + } + + #[test] + fn streaming_counts_redactions_without_logging_bytes() { + let secret = Secret::new("abcdef".to_string()); + let redactor = Arc::new(Redactor::new([("KEY", &secret)]).unwrap()); + let mut stream = StreamRedactor::new(redactor); + stream.push(b"one abc"); + stream.push(b"def two abcdef"); + stream.finish(); + assert_eq!(stream.redactions(), 2); + } + + #[test] + fn streaming_with_no_secrets_passes_everything_through() { + let redactor = + Arc::new(Redactor::new(std::iter::empty::<(&str, &Secret)>()).unwrap()); + let mut stream = StreamRedactor::new(redactor); + let mut out = stream.push(b"hello "); + out.extend_from_slice(&stream.push(b"world")); + out.extend_from_slice(&stream.finish()); + assert_eq!(out, b"hello world"); + assert_eq!(stream.redactions(), 0); + } + + #[test] + fn streaming_holds_back_a_partial_match_at_the_end_of_a_chunk() { + let secret = Secret::new("tail-secret".to_string()); + let redactor = Arc::new(Redactor::new([("KEY", &secret)]).unwrap()); + let mut stream = StreamRedactor::new(redactor); + + let mut out = stream.push(b"ends with tail-se"); + assert!( + !out.ends_with(b"tail-se"), + "bytes that could still be the start of a secret must not be handed on" + ); + out.extend_from_slice(&stream.push(b"cret and more")); + out.extend_from_slice(&stream.finish()); + assert_eq!(out, b"ends with [REDACTED:KEY] and more"); + } + // ── Lossy UTF-8 conversion utility ──────────────────────────────────── #[test]