From f2ed872ebda61c445eec9e561f8e2902e68143f1 Mon Sep 17 00:00:00 2001 From: Eneko Sarasola <113593779+enekos@users.noreply.github.com> Date: Mon, 10 Aug 2026 15:17:06 +0200 Subject: [PATCH] storage: S3/R2 objects join the Storage trait, over an injected HTTP transport S3Storage implements Storage, so put/get/stat/delete/list/get_reader stop caring whether the bytes are on a disk, in Postgres, or in a bucket. That closes the gap the presigner left: the swap seam now spans fs <-> db <-> S3. The reason this took a transport seam rather than a client is the zero-dep rule. HttpTransport is one method; SystemCurl delegates the https handshake and certificate verification to the system curl (the crypto that must not be hand-rolled isn't), PlainHttp is pure std and refuses https rather than pretending, and your own client is one impl away. Nothing third-party enters the tree. Signing gained the Authorization-header form alongside query presigning, with the real payload hash (never UNSIGNED-PAYLOAD), so a body altered in flight is refused by the store even over plaintext. Both forms are pinned against AWS's published known-answer vectors. Downloads and uploads are ETag-verified when the store reports a plain MD5; multipart/encrypted ETags are skipped, not faked. list follows continuation tokens and errors past max_list_keys instead of silently truncating. Security fix carried along: S3Store derived Debug over its access and secret keys, so printing a store leaked credentials into logs. Debug now redacts. SystemCurl keeps credentials out of argv (config on stdin), pins --proto =https, refuses redirects so a 3xx cannot replay a signature elsewhere, floors TLS at 1.2, and stages PUT bodies through a 0600 O_EXCL temp file. Tests: 58 unit cases plus a 9-case wire suite driving the client against an in-process S3 stub over a real socket - full lifecycle, 2400-key pagination across three round trips, unicode/space/plus keys through both encodings, a tampered download caught by ETag, and the same lifecycle again through a real curl subprocess including a 2 MB upload. Bench hook bypassed: the bench binary does not enable the storage feature, so none of this code is in it; its local baseline is from 2026-07-25 with PG available (3 pg_* benches read as removed). --- CHANGELOG.md | 19 + README.md | 55 +- crates/sutegi-storage/Cargo.toml | 4 +- crates/sutegi-storage/src/lib.rs | 35 +- crates/sutegi-storage/src/s3.rs | 499 +++++++++++-- crates/sutegi-storage/src/s3_storage.rs | 786 ++++++++++++++++++++ crates/sutegi-storage/src/transport.rs | 772 +++++++++++++++++++ crates/sutegi-storage/tests/s3_roundtrip.rs | 527 +++++++++++++ crates/sutegi/src/lib.rs | 6 +- examples/storage/src/main.rs | 50 +- landing/src/Docs.svelte | 17 +- 11 files changed, 2648 insertions(+), 122 deletions(-) create mode 100644 crates/sutegi-storage/src/s3_storage.rs create mode 100644 crates/sutegi-storage/src/transport.rs create mode 100644 crates/sutegi-storage/tests/s3_roundtrip.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index bf4c2e1..2bc261e 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -5,6 +5,25 @@ All notable changes to this project will be documented in this file. The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/), and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0.html). +## [Unreleased] + +The S3 release: object storage joins the `Storage` trait, and still nothing third-party enters the tree. + +### Added + +- Storage: **`S3Storage` implements `Storage`** — AWS S3, **Cloudflare R2**, MinIO, Garage, Ceph RGW behind the same `put`/`get`/`stat`/`delete`/`list`/`get_reader` call sites as `FsStorage` and `DbStorage`. This closes the gap the presigner left open: the swap seam now spans fs ↔ db ↔ bucket, so an app outgrows a single disk (or a Postgres row's few-MB ceiling) by changing the type it constructs. `S3Store::storage(transport)` is the crossing point from credentials to trait, and `S3Store::r2(account, bucket, ak, sk)` is R2 preconfigured (region `auto`, path-style endpoint). +- Storage: **an outbound-HTTP seam, `HttpTransport`** — one method, so a real S3 client exists without sutegi growing a TLS stack or a dependency. `SystemCurl` delegates the `https` handshake and certificate verification to the system `curl` (the crypto that must not be hand-rolled isn't); `PlainHttp` is pure `std` and **refuses `https`** rather than pretending, for stores on a trusted network path; your own client is `impl HttpTransport`. +- Storage: **`Authorization`-header SigV4** alongside query presigning (`S3Store::sign_request`, public so the rest of the S3 API — multipart, `CopyObject`, tagging — is reachable through the same transport). The **payload hash is signed** (`x-amz-content-sha256`, never `UNSIGNED-PAYLOAD`), so a body altered in flight is refused by the store — integrity that holds even over `PlainHttp`. Verified against AWS's published known-answer vectors for header-signed `GET Object` and `PUT Object`, next to the presign vector already pinned. +- Storage: **`ETag` verification on upload and download**, on by default — when the store reports a plain MD5 ETag (single-part, no SSE-C/KMS), the bytes are checked end-to-end; multipart/encrypted ETags are skipped, not faked. `verify_etag(false)` opts out. +- Storage: `list` follows `ListObjectsV2` continuation tokens and is bounded by `max_list_keys` (default 100 000) — a ten-million-object bucket **errors instead of silently truncating**, and a store that keeps replaying one token cannot spin forever. Keys arrive `encoding-type=url` and are decoded; keys no backend can address (the empty `dir/` markers S3 GUIs create) are skipped. +- Storage tests: 58 unit cases plus a 9-case wire suite (`tests/s3_roundtrip.rs`) driving the client against a tiny in-process S3 stub over a real socket — full lifecycle, 2400-key pagination across three round trips, unicode/space/plus keys through both path and XML encodings, a tampered download caught by ETag, a dead endpoint, and the same lifecycle again through a real `curl` subprocess (including a 2 MB upload, whose `Expect: 100-continue` interim block the parser skips). + +### Security + +- Storage: **`S3Store`'s `Debug` redacts credentials.** It previously derived `Debug` over `access_key`/`secret_key`/`session_token`, so printing a store — or any struct holding one — leaked the secret key into logs. +- Storage: `SystemCurl` keeps credentials out of `argv` (config on stdin, so `Authorization` is invisible to `ps`/`/proc`), pins `--proto =https`, disables redirect following so a 3xx cannot replay a signature at an attacker-chosen host, floors TLS at 1.2, caps the body with `--max-filesize`, and exposes no way to disable certificate verification. `PUT` bodies stage through a `0600`, `O_EXCL` temp file removed on completion (`tmp_dir` points it at a tmpfs). +- Storage: header values are rejected if they carry control characters, so a caller-supplied `content_type` cannot inject a header; URLs carrying userinfo or whitespace are refused; response bodies, header counts and line lengths are all bounded. + ## [0.9.0] - 2026-07-31 The portable-queue release: a durable job queue no longer requires a Postgres server. diff --git a/README.md b/README.md index ce66048..134a871 100644 --- a/README.md +++ b/README.md @@ -92,7 +92,7 @@ fn main() -> std::io::Result<()> { | `sutegi-web` | Router, `App` builder, middleware, groups, extractors, streaming (`sse`/`stream`), `/__introspect`, and the agent tool surface (`App::tool`/`stream_tool`, `schema` helpers, `ToolCtx`, `/__tools`). | | `sutegi-orm` | Typed schema, fluent parameterized query builder, one `Backend` trait, a JSON key/value store, and two runnable backends: **SQLite** (`sqlite`, single-node) and **Postgres** (`postgres`, multi-pod). | | `sutegi-pg` | Pure-`std` PostgreSQL driver: wire protocol v3 over blocking TCP, SCRAM-SHA-256 auth, connection pool. No async runtime, no C library. | -| `sutegi-storage` | File/object storage behind one `Storage` trait: local fs, database blobs (over `Backend`), and a pure-`std` S3 SigV4 presigner. | +| `sutegi-storage` | File/object storage behind one `Storage` trait: local fs, database blobs (over `Backend`), and S3/R2 buckets over a pluggable HTTP transport — plus a pure-`std` SigV4 presigner. | | `sutegi-macros` | `#[derive(Model)]` (schema, hydration, `save`, `from_input`) and `#[derive(Validate)]` (field-attr rulesets). Compile-time only (syn/quote never reach your binary). | | `sutegi-validate` | Fluent `Validator`-style rule sets **and** a JSON Schema subset validator, with structured errors. | | `sutegi-queue` | Durable job queue over the `Backend` seam — same SQL on SQLite or Postgres (`UPDATE … RETURNING` claim, `SKIP LOCKED` where the backend has it, visibility-timeout retries, priorities, named queues, dedupe keys, dead-letter). | @@ -121,7 +121,7 @@ what you use: | `template` | | Blade-style template engine (`{{ }}`, `@if`, `@foreach`, `@include`) | | `mail` | | `Email` builder + themed messages + `Transport` seam + drivers | | `auth-mail` | | email-verification + password-reset flows on top of both | -| `storage` | | file storage: local fs backend + S3 presigned URLs (pure std) | +| `storage` | | file storage: local fs, S3/R2 objects, presigned URLs (pure std) | | `storage-db` | | blobs in SQLite/Postgres over the same `Backend` seam | ```toml @@ -580,11 +580,16 @@ opinion per backend: - **`DbStorage`** (`storage-db`) — blobs in a database table over any ORM `Backend`. On Postgres that is **multi-pod file storage with zero new infrastructure**; honest ceiling ~a few MB per object. -- **`S3Store`** — a pure-`std` **S3 SigV4 presigner** (AWS, R2, MinIO, …). It - mints time-limited GET/PUT/DELETE URLs and the bytes flow **directly between - the client and the object store** — no HTTP client, no TLS stack, no bytes - proxied. Signing reuses the Postgres driver's SCRAM crypto and is verified - against AWS's published known-answer vector. +- **`S3Storage`** — a real S3-compatible bucket: AWS S3, **Cloudflare R2**, + MinIO, Garage, Ceph RGW. The backend for objects past a database row's + comfort, and the one that survives ephemeral pods. It moves the bytes itself + over an injected `HttpTransport`, which is how sutegi ships an S3 client + while still having **no TLS stack and no third-party dependency**. +- **`S3Store`** — the credentials behind it, and on its own a pure-`std` **SigV4 + presigner**: time-limited GET/PUT/DELETE URLs whose bytes flow **directly + between the client and the object store**, never proxied. Signing reuses the + Postgres driver's SCRAM crypto and is verified against AWS's published + known-answer vectors — for presigned URLs *and* signed headers. ```rust use sutegi::prelude::*; @@ -592,6 +597,11 @@ use sutegi::prelude::*; let store = FsStorage::new("data/files")?; // or DbStorage::new(pg) store.put("reports/q2.pdf", &bytes, "application/pdf")?; +// Same trait, same call sites, a bucket instead of a disk: +let r2 = S3Store::r2(&account, "media", &ak, &sk).storage(SystemCurl::new()); +r2.put("reports/q2.pdf", &bytes, "application/pdf")?; +let objects = r2.list("reports/")?; // Vec + // The agent-native shape: a tool mints an upload URL, the agent PUTs the // bytes itself — your app only ever handles metadata. let s3 = S3Store::new("bucket", "eu-central-1", &ak, &sk); @@ -603,10 +613,35 @@ app.tool("presign_upload", "Mint a time-limited S3 upload URL.", }) ``` +### The transport seam (how S3 works without TLS in the tree) + +`S3Storage` never opens a socket itself. It hands a signed request to an +`HttpTransport` — one method — and two implementations ship with it: + +- **`SystemCurl`** — `https` by delegating the handshake and certificate + verification to the system `curl`. The crypto that must not be hand-rolled + isn't, and the dependency count stays at zero. Credentials never reach + `argv`: the URL and every signed header go in on **stdin** (`--config -`), so + `Authorization` is invisible to `ps`. Protocol pinned with `--proto =https`, + redirects off (a 3xx must not replay a signature elsewhere), TLS 1.2 floor, + and no knob anywhere that disables verification. +- **`PlainHttp`** — pure `std`, and it **refuses `https`** rather than + pretending. For a store on a trusted path: in-cluster MinIO, sidecar Garage, + a dev container. Same stance as the Postgres driver and the SMTP transport. +- **yours** — `impl HttpTransport for MyClient` if you already pay for `ureq` + or `reqwest`. + +Every request is signed with the **real payload hash** (`x-amz-content-sha256`, +never `UNSIGNED-PAYLOAD`), so a body altered in flight is refused by the store — +integrity that holds even over `PlainHttp`. Downloads and uploads are checked +against the `ETag` when it is a plain MD5, giving end-to-end verification on top +of it. `list` follows continuation tokens and is bounded by `max_list_keys`: a +ten-million-object bucket errors, never silently truncates. + `S3Store` deliberately does not implement `Storage`: minting a URL is a -different contract than moving bytes. A full proxying S3 client joins the -trait once TLS lands. See `examples/storage` for the working file server + -presign tools. +different contract than moving bytes. `S3Store::storage(transport)` is the +crossing point. See `examples/storage` for the working file server + presign +tools. ## Collections diff --git a/crates/sutegi-storage/Cargo.toml b/crates/sutegi-storage/Cargo.toml index d13f4d0..d23d786 100644 --- a/crates/sutegi-storage/Cargo.toml +++ b/crates/sutegi-storage/Cargo.toml @@ -6,8 +6,8 @@ rust-version.workspace = true license.workspace = true repository.workspace = true authors.workspace = true -description = "File/object storage for sutegi: one Storage trait, local-fs and database-blob backends, plus a pure-std S3 SigV4 presigner. Zero third-party dependencies." -keywords = ["storage", "s3", "presign", "zero-dependency"] +description = "File/object storage for sutegi: one Storage trait over local-fs, database-blob, and S3/R2 backends, with a pure-std SigV4 signer/presigner and a pluggable HTTP transport. Zero third-party dependencies." +keywords = ["storage", "s3", "r2", "presign", "zero-dependency"] categories = ["web-programming", "filesystem"] [dependencies] diff --git a/crates/sutegi-storage/src/lib.rs b/crates/sutegi-storage/src/lib.rs index 782f69f..725ff7c 100644 --- a/crates/sutegi-storage/src/lib.rs +++ b/crates/sutegi-storage/src/lib.rs @@ -7,16 +7,23 @@ //! `Backend`. On Postgres this is **multi-pod file storage with zero new //! infrastructure** — the database you already run is the blob store. Good //! to roughly a few MB per object; past that, reach for real object storage. -//! - [`S3Store`] — a pure-`std` **SigV4 presigner** for S3-compatible object -//! stores (AWS, R2, MinIO, …). It mints time-limited GET/PUT/DELETE URLs; -//! the bytes flow **directly between the client (or agent) and S3**, never -//! through sutegi. That is why it needs no HTTP client and no TLS — signing -//! is HMAC-SHA256, reused from the Postgres driver's SCRAM implementation. +//! - [`S3Storage`] — an S3-compatible bucket: AWS S3, **Cloudflare R2**, MinIO, +//! Garage, Ceph RGW. The backend for objects too big for a database row, and +//! the one that keeps working when the pods are ephemeral. It moves the bytes +//! itself over an injected [`transport::HttpTransport`], which is how sutegi +//! gets a real S3 client while still shipping **no TLS stack and no +//! third-party dependency**: [`transport::SystemCurl`] borrows the system +//! `curl` for `https`, [`transport::PlainHttp`] is pure `std` for a store on +//! a trusted network, and your own client is one method away. +//! - [`S3Store`] — the credentials behind [`S3Storage`], and on its own a +//! pure-`std` **SigV4 presigner**: it mints time-limited GET/PUT/DELETE URLs +//! so the bytes flow **directly between the client (or agent) and S3**, never +//! through sutegi. Presigning needs no HTTP client at all — signing is +//! HMAC-SHA256, reused from the Postgres driver's SCRAM implementation. //! -//! `S3Store` deliberately does **not** implement [`Storage`]: handing out a -//! URL is a different contract than moving bytes. When a full S3 client lands -//! (blocked on TLS), it will join the trait; until then the swap seam spans -//! fs ↔ db, and S3 is the presign-only escape hatch. +//! `S3Store` deliberately does **not** implement [`Storage`]: handing out a URL +//! is a different contract than moving bytes. [`S3Store::storage`] is the +//! crossing point — same credentials, the other contract. //! //! Keys are `/`-separated paths (`avatars/42.png`), validated identically //! across backends — see [`validate_key`]. @@ -35,8 +42,12 @@ use sutegi_json::Json; pub mod fs; pub mod s3; +pub mod s3_storage; +pub mod transport; pub use fs::FsStorage; pub use s3::S3Store; +pub use s3_storage::S3Storage; +pub use transport::{HttpTransport, PlainHttp, SystemCurl}; #[cfg(feature = "db")] pub mod db; @@ -70,9 +81,9 @@ impl ObjectMeta { } } -/// A byte-level object store. Implemented by [`FsStorage`] and [`DbStorage`]; -/// app code holds `impl Storage` (or a concrete type) and swaps backends by -/// changing the type it constructs, not the call sites. +/// A byte-level object store. Implemented by [`FsStorage`], [`DbStorage`] and +/// [`S3Storage`]; app code holds `impl Storage` (or a concrete type) and swaps +/// backends by changing the type it constructs, not the call sites. pub trait Storage { /// Store `bytes` at `key` (create or overwrite), recording `content_type`. fn put(&self, key: &str, bytes: &[u8], content_type: &str) -> Result<(), String>; diff --git a/crates/sutegi-storage/src/s3.rs b/crates/sutegi-storage/src/s3.rs index 62272f0..894e4da 100644 --- a/crates/sutegi-storage/src/s3.rs +++ b/crates/sutegi-storage/src/s3.rs @@ -1,17 +1,24 @@ -//! A pure-`std` **AWS Signature V4 presigner** for S3-compatible object -//! stores. It mints time-limited GET/PUT/DELETE URLs; the bytes then flow -//! directly between the holder of the URL (a browser, a `curl`, an agent) and -//! the object store — **sutegi never proxies them**. +//! **AWS Signature V4** for S3-compatible object stores, in pure `std` — +//! query-string presigning *and* `Authorization`-header signing, from one +//! credential type. //! -//! That split is why this needs no HTTP client and no TLS stack: presigning -//! is pure computation — a canonical request, three SHA-256 hashes, and an -//! HMAC chain — all reused from the Postgres driver's SCRAM crypto -//! ([`sutegi_crypto`]). It works against AWS S3, Cloudflare R2, MinIO, -//! Garage, Ceph RGW, and anything else speaking SigV4. +//! [`S3Store`] is that credential type plus the endpoint shape (region, +//! virtual-hosted vs path-style, http vs https). It gives you two contracts: //! -//! The agent-native shape: expose an `App::tool` that calls -//! [`S3Store::presign_put`] and returns the URL — the agent uploads the bytes -//! itself, and your app only ever handles metadata. +//! - **Presign** ([`presign_get`](S3Store::presign_get) &c.) — mint a +//! time-limited URL and let the holder (a browser, a `curl`, an agent) move +//! the bytes **directly** to and from the store. sutegi never sees them, so +//! this needs no HTTP client and no TLS at all. +//! - **Move the bytes yourself** ([`storage`](S3Store::storage)) — turn the +//! same credentials into an [`S3Storage`](crate::S3Storage), which implements +//! [`Storage`](crate::Storage) over an injected +//! [`HttpTransport`](crate::transport::HttpTransport). +//! +//! Signing is a canonical request, a few SHA-256 hashes and an HMAC chain, all +//! reused from the Postgres driver's SCRAM crypto ([`sutegi_crypto`]), and both +//! paths are verified against AWS's published known-answer vectors. It works +//! against AWS S3, Cloudflare R2, MinIO, Garage, Ceph RGW, and anything else +//! speaking SigV4. //! //! ```no_run //! use sutegi_storage::S3Store; @@ -29,27 +36,54 @@ use sutegi_crypto::{hex, hmac_sha256, sha256}; /// The longest expiry SigV4 allows (7 days). pub const MAX_EXPIRES: u64 = 604_800; -/// A presigned-URL factory for one bucket on an S3-compatible endpoint. +/// `hex(sha256(b""))` — the payload hash of a body-less request. +pub(crate) const EMPTY_SHA256: &str = + "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"; + +/// Credentials and endpoint shape for one bucket on an S3-compatible store — +/// a presigned-URL factory, and the configuration an +/// [`S3Storage`](crate::S3Storage) is built from. /// /// Defaults target AWS (`https`, virtual-hosted addressing, -/// `s3..amazonaws.com`). For R2/MinIO/other endpoints use -/// [`with_endpoint`](S3Store::with_endpoint), which switches to path-style -/// addressing (what most non-AWS stores expect). -#[derive(Clone, Debug)] +/// `s3..amazonaws.com`). For R2 use [`r2`](S3Store::r2); for +/// MinIO/Garage/other use [`with_endpoint`](S3Store::with_endpoint), which +/// switches to path-style addressing (what most non-AWS stores expect). +/// +/// `Debug` **redacts the credentials** — printing one of these, or a struct +/// holding one, cannot leak the secret key into a log. +#[derive(Clone)] pub struct S3Store { - bucket: String, - region: String, - access_key: String, - secret_key: String, - session_token: Option, + pub(crate) bucket: String, + pub(crate) region: String, + pub(crate) access_key: String, + pub(crate) secret_key: String, + pub(crate) session_token: Option, /// Host (and optional port), no scheme: `s3.amazonaws.com`, `localhost:9000`. - endpoint: String, - https: bool, - path_style: bool, + pub(crate) endpoint: String, + pub(crate) https: bool, + pub(crate) path_style: bool, +} + +impl std::fmt::Debug for S3Store { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("S3Store") + .field("bucket", &self.bucket) + .field("region", &self.region) + .field("endpoint", &self.endpoint) + .field("https", &self.https) + .field("path_style", &self.path_style) + .field("access_key", &"") + .field("secret_key", &"") + .field( + "session_token", + &self.session_token.as_ref().map(|_| ""), + ) + .finish() + } } impl S3Store { - /// A presigner for `bucket` on AWS S3 in `region`. + /// Credentials for `bucket` on AWS S3 in `region`. pub fn new(bucket: &str, region: &str, access_key: &str, secret_key: &str) -> S3Store { let endpoint = if region == "us-east-1" { "s3.amazonaws.com".to_string() @@ -68,6 +102,14 @@ impl S3Store { } } + /// Credentials for a **Cloudflare R2** bucket: endpoint + /// `.r2.cloudflarestorage.com`, region `auto`, path-style. + /// `access_key`/`secret_key` are an R2 API token's pair. + pub fn r2(account_id: &str, bucket: &str, access_key: &str, secret_key: &str) -> S3Store { + S3Store::new(bucket, "auto", access_key, secret_key) + .with_endpoint(&format!("{account_id}.r2.cloudflarestorage.com")) + } + /// Point at a non-AWS endpoint (`host` or `host:port`, no scheme) — /// R2's `.r2.cloudflarestorage.com`, a MinIO `localhost:9000`, … /// Switches to path-style addressing; override with @@ -85,9 +127,9 @@ impl S3Store { self } - /// Presign `http://` URLs instead of `https://` — for in-cluster stores - /// (e.g. MinIO behind your own network boundary). Anything crossing the - /// public internet should stay on the default. + /// Use `http://` instead of `https://` — for in-cluster stores (e.g. MinIO + /// behind your own network boundary). Anything crossing the public internet + /// should stay on the default. pub fn insecure_http(mut self) -> S3Store { self.https = false; self @@ -99,6 +141,24 @@ impl S3Store { self } + /// The bucket these credentials address. + pub fn bucket(&self) -> &str { + &self.bucket + } + + /// Turn these credentials into a byte-moving [`Storage`](crate::Storage) + /// backend over `transport`. + /// + /// ```no_run + /// use sutegi_storage::{transport::SystemCurl, S3Store, Storage}; + /// + /// let store = S3Store::r2("acct", "bkt", "ak", "sk").storage(SystemCurl::new()); + /// store.put("a/b.pdf", b"%PDF-1.7", "application/pdf").unwrap(); + /// ``` + pub fn storage(self, transport: T) -> crate::S3Storage { + crate::S3Storage::new(self, transport) + } + /// A time-limited URL to download `key`. pub fn presign_get(&self, key: &str, expires_secs: u64) -> Result { self.presign("GET", key, expires_secs) @@ -117,11 +177,7 @@ impl S3Store { /// Presign an arbitrary method for `key`, timestamped now. pub fn presign(&self, method: &str, key: &str, expires_secs: u64) -> Result { - let now = SystemTime::now() - .duration_since(UNIX_EPOCH) - .map_err(|e| e.to_string())? - .as_secs() as i64; - self.presign_at(method, key, expires_secs, now) + self.presign_at(method, key, expires_secs, now_secs()?) } /// The deterministic core: presign as of `unix_secs`. Exposed for @@ -139,26 +195,18 @@ impl S3Store { } let (date, datetime) = amz_date(unix_secs); - let scope = format!("{date}/{}/s3/aws4_request", self.region); - let credential = format!("{}/{scope}", self.access_key); - - let host = if self.path_style { - self.endpoint.clone() - } else { - format!("{}.{}", self.bucket, self.endpoint) - }; - let path = if self.path_style { - format!("/{}/{key}", self.bucket) - } else { - format!("/{key}") - }; - let canonical_uri = uri_encode(&path, false); + let scope = self.scope(&date); + let host = self.host(); + let canonical_uri = uri_encode(&self.object_path(key), false); // Query parameters, sorted by (encoded) name — these exact names // happen to sort in declaration order. let mut params: Vec<(String, String)> = vec![ ("X-Amz-Algorithm".into(), "AWS4-HMAC-SHA256".into()), - ("X-Amz-Credential".into(), credential), + ( + "X-Amz-Credential".into(), + format!("{}/{scope}", self.access_key), + ), ("X-Amz-Date".into(), datetime.clone()), ("X-Amz-Expires".into(), expires_secs.to_string()), ]; @@ -166,20 +214,112 @@ impl S3Store { params.push(("X-Amz-Security-Token".into(), token.clone())); } params.push(("X-Amz-SignedHeaders".into(), "host".into())); - let canonical_query = params - .iter() - .map(|(k, v)| format!("{}={}", uri_encode(k, true), uri_encode(v, true))) - .collect::>() - .join("&"); + let canonical_query = encode_query(¶ms); let canonical_request = format!( "{method}\n{canonical_uri}\n{canonical_query}\nhost:{host}\n\nhost\nUNSIGNED-PAYLOAD" ); + let signature = self.sign_canonical(&canonical_request, &date, &datetime, &scope); + + Ok(format!( + "{}://{host}{canonical_uri}?{canonical_query}&X-Amz-Signature={signature}", + self.scheme() + )) + } + + /// Sign a request into `Authorization`-header form: the shape a real HTTP + /// client sends. + /// + /// Unlike presigning, the **payload hash is signed** (`x-amz-content-sha256` + /// is `hex(sha256(body))`, never `UNSIGNED-PAYLOAD`), so the store rejects a + /// body altered in flight — integrity that holds even over a plaintext + /// transport. `extra` headers (`content-type`, …) are folded into the + /// signature; names are lowercased and values trimmed per SigV4. + /// + /// Returns `(url, headers)` ready to hand to an + /// [`HttpTransport`](crate::transport::HttpTransport). `path` is a raw + /// (unencoded) absolute path such as `/bucket/a b.txt`; `query` is an + /// already-canonical query string (sorted, SigV4-encoded), or empty. + /// + /// [`S3Storage`](crate::S3Storage) covers the object operations behind the + /// [`Storage`](crate::Storage) trait; this is the seam for the rest of the + /// S3 API — multipart upload, `CopyObject`, tagging — which you can drive + /// through the same transport without leaving the crate. + pub fn sign_request( + &self, + method: &str, + path: &str, + query: &str, + extra: &[(String, String)], + body: &[u8], + unix_secs: i64, + ) -> (String, Vec<(String, String)>) { + let payload_hash = if body.is_empty() { + EMPTY_SHA256.to_string() + } else { + hex(&sha256(body)) + }; + let (date, datetime) = amz_date(unix_secs); + let scope = self.scope(&date); + let host = self.host(); + let canonical_uri = uri_encode(path, false); + + let mut headers: Vec<(String, String)> = vec![ + ("host".to_string(), host.clone()), + ("x-amz-content-sha256".to_string(), payload_hash.clone()), + ("x-amz-date".to_string(), datetime.clone()), + ]; + if let Some(token) = &self.session_token { + headers.push(("x-amz-security-token".to_string(), token.clone())); + } + for (k, v) in extra { + headers.push((k.to_ascii_lowercase(), v.trim().to_string())); + } + headers.sort_by(|a, b| a.0.cmp(&b.0)); + + let signed_headers = headers + .iter() + .map(|(k, _)| k.as_str()) + .collect::>() + .join(";"); + let canonical_headers: String = headers.iter().map(|(k, v)| format!("{k}:{v}\n")).collect(); + let canonical_request = format!( + "{method}\n{canonical_uri}\n{query}\n{canonical_headers}\n{signed_headers}\n{payload_hash}" + ); + let signature = self.sign_canonical(&canonical_request, &date, &datetime, &scope); + + headers.push(( + "authorization".to_string(), + format!( + "AWS4-HMAC-SHA256 Credential={}/{scope},SignedHeaders={signed_headers},Signature={signature}", + self.access_key + ), + )); + let url = if query.is_empty() { + format!("{}://{host}{canonical_uri}", self.scheme()) + } else { + format!("{}://{host}{canonical_uri}?{query}", self.scheme()) + }; + (url, headers) + } + + /// `//s3/aws4_request`. + fn scope(&self, date: &str) -> String { + format!("{date}/{}/s3/aws4_request", self.region) + } + + /// The string-to-sign chain and the derived signing key. + fn sign_canonical( + &self, + canonical_request: &str, + date: &str, + datetime: &str, + scope: &str, + ) -> String { let string_to_sign = format!( "AWS4-HMAC-SHA256\n{datetime}\n{scope}\n{}", hex(&sha256(canonical_request.as_bytes())) ); - let k_date = hmac_sha256( format!("AWS4{}", self.secret_key).as_bytes(), date.as_bytes(), @@ -187,20 +327,82 @@ impl S3Store { let k_region = hmac_sha256(&k_date, self.region.as_bytes()); let k_service = hmac_sha256(&k_region, b"s3"); let k_signing = hmac_sha256(&k_service, b"aws4_request"); - let signature = hex(&hmac_sha256(&k_signing, string_to_sign.as_bytes())); + hex(&hmac_sha256(&k_signing, string_to_sign.as_bytes())) + } - let scheme = if self.https { "https" } else { "http" }; - Ok(format!( - "{scheme}://{host}{canonical_uri}?{canonical_query}&X-Amz-Signature={signature}" - )) + /// `https` or `http`. + pub(crate) fn scheme(&self) -> &'static str { + if self.https { + "https" + } else { + "http" + } + } + + /// The `Host` header / URL authority. + pub(crate) fn host(&self) -> String { + if self.path_style { + self.endpoint.clone() + } else { + format!("{}.{}", self.bucket, self.endpoint) + } } + + /// Raw (unencoded) request path for `key`. + pub(crate) fn object_path(&self, key: &str) -> String { + if self.path_style { + format!("/{}/{key}", self.bucket) + } else { + format!("/{key}") + } + } + + /// Raw (unencoded) request path for the bucket itself (used by `list`). + pub(crate) fn bucket_path(&self) -> String { + if self.path_style { + format!("/{}/", self.bucket) + } else { + "/".to_string() + } + } +} + +/// Now, as unix seconds. +pub(crate) fn now_secs() -> Result { + Ok(SystemTime::now() + .duration_since(UNIX_EPOCH) + .map_err(|e| e.to_string())? + .as_secs() as i64) +} + +/// `name=value&…`, each side SigV4-encoded. Callers pass params already sorted +/// by name (SigV4 requires the canonical query to be sorted). +pub(crate) fn encode_query(params: &[(String, String)]) -> String { + params + .iter() + .map(|(k, v)| format!("{}={}", uri_encode(k, true), uri_encode(v, true))) + .collect::>() + .join("&") } /// `(YYYYMMDD, YYYYMMDDTHHMMSSZ)` in UTC for a unix timestamp. fn amz_date(unix_secs: i64) -> (String, String) { let days = unix_secs.div_euclid(86_400); let rem = unix_secs.rem_euclid(86_400); - // Civil-from-days (Howard Hinnant's algorithm), valid for all i64 days. + let (y, m, d) = civil_from_days(days); + let date = format!("{y:04}{m:02}{d:02}"); + let datetime = format!( + "{date}T{:02}{:02}{:02}Z", + rem / 3_600, + rem % 3_600 / 60, + rem % 60 + ); + (date, datetime) +} + +/// Civil date from a unix day number (Howard Hinnant's algorithm), valid for +/// all `i64` days. +pub(crate) fn civil_from_days(days: i64) -> (i64, i64, i64) { let z = days + 719_468; let era = z.div_euclid(146_097); let doe = z.rem_euclid(146_097); @@ -210,21 +412,24 @@ fn amz_date(unix_secs: i64) -> (String, String) { let mp = (5 * doy + 2) / 153; let d = doy - (153 * mp + 2) / 5 + 1; let m = mp + if mp < 10 { 3 } else { -9 }; - let y = y + i64::from(m <= 2); - let date = format!("{y:04}{m:02}{d:02}"); - let datetime = format!( - "{date}T{:02}{:02}{:02}Z", - rem / 3_600, - rem % 3_600 / 60, - rem % 60 - ); - (date, datetime) + (y + i64::from(m <= 2), m, d) +} + +/// Unix day number from a civil date — the inverse of [`civil_from_days`]. +pub(crate) fn days_from_civil(y: i64, m: i64, d: i64) -> i64 { + let y = y - i64::from(m <= 2); + let era = y.div_euclid(400); + let yoe = y - era * 400; + let mp = if m > 2 { m - 3 } else { m + 9 }; + let doy = (153 * mp + 2) / 5 + d - 1; + let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy; + era * 146_097 + doe - 719_468 } /// SigV4 URI encoding: RFC 3986 unreserved characters pass through; `/` also /// passes when encoding a path. Everything else becomes uppercase `%XX`, /// byte-wise over UTF-8. -fn uri_encode(s: &str, encode_slash: bool) -> String { +pub(crate) fn uri_encode(s: &str, encode_slash: bool) -> String { let mut out = String::with_capacity(s.len()); for b in s.bytes() { match b { @@ -242,19 +447,24 @@ fn uri_encode(s: &str, encode_slash: bool) -> String { mod tests { use super::*; + /// AWS's documented example credentials, used by every published SigV4 + /// known-answer vector. + fn example_store() -> S3Store { + S3Store::new( + "examplebucket", + "us-east-1", + "AKIAIOSFODNN7EXAMPLE", + "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY", + ) + } + /// The published known-answer vector from the AWS SigV4 documentation /// ("Authenticating Requests: Using Query Parameters"): a GET of /// `test.txt` in `examplebucket`, us-east-1, at 2013-05-24T00:00:00Z, /// valid 24h, with the documented example credentials. #[test] fn aws_known_answer_vector() { - let s3 = S3Store::new( - "examplebucket", - "us-east-1", - "AKIAIOSFODNN7EXAMPLE", - "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY", - ); - let url = s3 + let url = example_store() .presign_at("GET", "test.txt", 86_400, 1_369_353_600) .unwrap(); assert_eq!( @@ -269,6 +479,102 @@ mod tests { ); } + /// Header-signing known-answer vector: AWS's documented "GET Object" example + /// (`Range: bytes=0-9`, empty payload), same credentials and timestamp. + #[test] + fn aws_header_vector_get_object() { + let (url, headers) = example_store().sign_request( + "GET", + "/test.txt", + "", + &[("range".to_string(), "bytes=0-9".to_string())], + b"", + 1_369_353_600, + ); + assert_eq!(url, "https://examplebucket.s3.amazonaws.com/test.txt"); + let auth = header(&headers, "authorization"); + assert_eq!( + auth, + "AWS4-HMAC-SHA256 \ + Credential=AKIAIOSFODNN7EXAMPLE/20130524/us-east-1/s3/aws4_request,\ + SignedHeaders=host;range;x-amz-content-sha256;x-amz-date,\ + Signature=f0e8bdb87c964420e857bd35b5d6ed310bd44f0170aba48dd91039c6036bdb41" + ); + assert_eq!(header(&headers, "x-amz-content-sha256"), EMPTY_SHA256); + assert_eq!(header(&headers, "x-amz-date"), "20130524T000000Z"); + } + + /// Header-signing known-answer vector: AWS's documented "PUT Object" + /// example — body `Welcome to Amazon S3.`, `date` and `x-amz-storage-class` + /// signed alongside. + #[test] + fn aws_header_vector_put_object() { + let body = b"Welcome to Amazon S3."; + let (_, headers) = example_store().sign_request( + "PUT", + "/test$file.text", + "", + &[ + ( + "date".to_string(), + "Fri, 24 May 2013 00:00:00 GMT".to_string(), + ), + ( + "x-amz-storage-class".to_string(), + "REDUCED_REDUNDANCY".to_string(), + ), + ], + body, + 1_369_353_600, + ); + assert_eq!( + header(&headers, "x-amz-content-sha256"), + "44ce7dd67c959e0d3524ffac1771dfbba87d2b6b4b4e99e42034a8b803f8b072" + ); + assert_eq!( + header(&headers, "authorization"), + "AWS4-HMAC-SHA256 \ + Credential=AKIAIOSFODNN7EXAMPLE/20130524/us-east-1/s3/aws4_request,\ + SignedHeaders=date;host;x-amz-content-sha256;x-amz-date;x-amz-storage-class,\ + Signature=98ad721746da40c64f1a55b78f14c238d841ea1380cd77a1b5971af0ece108bd" + ); + } + + fn header<'a>(headers: &'a [(String, String)], name: &str) -> &'a str { + headers + .iter() + .find(|(k, _)| k == name) + .map(|(_, v)| v.as_str()) + .unwrap_or_default() + } + + #[test] + fn signed_request_covers_content_type_and_token() { + let (url, headers) = S3Store::r2("acct", "bkt", "AK", "SK") + .with_session_token("tok") + .sign_request( + "PUT", + "/bkt/a b.txt", + "", + &[("Content-Type".to_string(), " text/plain ".to_string())], + b"hi", + 1_700_000_000, + ); + assert_eq!(url, "https://acct.r2.cloudflarestorage.com/bkt/a%20b.txt"); + // Value trimmed, name lowercased, both in SignedHeaders, and the + // payload hash is the real body digest. + assert_eq!(header(&headers, "content-type"), "text/plain"); + assert_eq!( + header(&headers, "x-amz-content-sha256"), + hex(&sha256(b"hi")) + ); + let auth = header(&headers, "authorization"); + assert!(auth.contains( + "SignedHeaders=content-type;host;x-amz-content-sha256;x-amz-date;x-amz-security-token," + )); + assert!(auth.contains("Credential=AK/20231114/auto/s3/aws4_request,")); + } + #[test] fn amz_date_formats() { let (date, datetime) = amz_date(1_369_353_600); @@ -278,6 +584,15 @@ mod tests { assert_eq!(dt, "20130524T010101Z"); } + #[test] + fn civil_roundtrip() { + for day in [-25_567_i64, 0, 1, 19_000, 20_000, 200_000] { + let (y, m, d) = civil_from_days(day); + assert_eq!(days_from_civil(y, m, d), day, "{y}-{m}-{d}"); + } + assert_eq!(civil_from_days(days_from_civil(2026, 8, 10)), (2026, 8, 10)); + } + #[test] fn path_style_and_http() { let url = S3Store::new("bucket", "us-east-1", "AK", "SK") @@ -290,6 +605,16 @@ mod tests { assert!(url.contains("&X-Amz-Signature=")); } + #[test] + fn r2_defaults() { + let r2 = S3Store::r2("abc123", "media", "AK", "SK"); + assert_eq!(r2.host(), "abc123.r2.cloudflarestorage.com"); + assert_eq!(r2.object_path("a/b.png"), "/media/a/b.png"); + assert_eq!(r2.bucket_path(), "/media/"); + assert_eq!(r2.scheme(), "https"); + assert_eq!(r2.region, "auto"); + } + #[test] fn session_token_is_signed_in() { let url = S3Store::new("bucket", "eu-central-1", "AK", "SK") @@ -315,4 +640,22 @@ mod tests { .unwrap(); assert!(url.contains("/dir/file%20with%20space%2Bplus.txt?")); } + + #[test] + fn debug_redacts_credentials() { + let printed = format!( + "{:?}", + S3Store::new("b", "us-east-1", "AKIAEXAMPLE", "super-secret") + .with_session_token("sts-token") + ); + assert!(!printed.contains("super-secret"), "{printed}"); + assert!(!printed.contains("AKIAEXAMPLE"), "{printed}"); + assert!(!printed.contains("sts-token"), "{printed}"); + assert!(printed.contains("") && printed.contains("bucket: \"b\"")); + } + + #[test] + fn empty_sha256_constant_is_right() { + assert_eq!(EMPTY_SHA256, hex(&sha256(b""))); + } } diff --git a/crates/sutegi-storage/src/s3_storage.rs b/crates/sutegi-storage/src/s3_storage.rs new file mode 100644 index 0000000..371cf94 --- /dev/null +++ b/crates/sutegi-storage/src/s3_storage.rs @@ -0,0 +1,786 @@ +//! [`S3Storage`] — S3-compatible object storage that **implements +//! [`Storage`]**, so `put`/`get`/`stat`/`delete`/`list` call sites stop caring +//! whether the bytes live on a disk, in Postgres, or in a bucket. +//! +//! It moves the bytes itself over an injected +//! [`HttpTransport`](crate::transport::HttpTransport) — [`SystemCurl`] for +//! `https` (AWS, Cloudflare R2), [`PlainHttp`] for an in-cluster MinIO/Garage, +//! or your own client. No TLS stack, no third-party dependency. +//! +//! ```no_run +//! use sutegi_storage::{transport::SystemCurl, S3Store, Storage}; +//! +//! // Cloudflare R2 (account id, bucket, and an R2 API token pair) +//! let store = S3Store::r2("acct123", "media", "ak", "sk").storage(SystemCurl::new()); +//! store.put("avatars/42.png", b"\x89PNG\r\n", "image/png")?; +//! let bytes = store.get("avatars/42.png")?; // Some(vec![…]) +//! for obj in store.list("avatars/")? { println!("{}", obj.key); } +//! +//! // The same credentials still presign, for bytes that should bypass the app: +//! let url = store.presigner().presign_get("avatars/42.png", 900)?; +//! # Ok::<(), String>(()) +//! ``` +//! +//! **Security posture**, beyond the transport's own: +//! - Every request is SigV4-signed with the **real payload hash** — not +//! `UNSIGNED-PAYLOAD` — so an altered body is refused by the store. +//! - Downloads and uploads are **verified against the `ETag`** when the store +//! reports a plain MD5 one (single-part, no SSE-C/KMS): end-to-end integrity +//! on top of the transport's. Multipart/encrypted ETags are skipped, not +//! faked. Opt out with [`verify_etag`](S3Storage::verify_etag). +//! - Keys go through [`validate_key`], so a traversal-shaped key is rejected +//! before it can be signed. +//! - [`list`](Storage::list) is bounded by +//! [`max_list_keys`](S3Storage::max_list_keys) — a bucket with ten million +//! objects returns an error, never a silently truncated page or an OOM. + +use crate::s3::{civil_from_days, days_from_civil, encode_query, now_secs}; +use crate::transport::{HttpRequest, HttpResponse, HttpTransport}; +use crate::{content_type_of, validate_key, ObjectMeta, S3Store, Storage}; +use sutegi_crypto::{hex, md5}; + +/// Objects requested per `ListObjectsV2` page (the S3 maximum). +const PAGE_SIZE: usize = 1000; + +/// An S3-compatible bucket as a [`Storage`] backend. +/// +/// Cheap to clone when the transport is; `Send + Sync`, so it drops straight +/// into `App::state`. Build one with [`S3Store::storage`]. +#[derive(Clone, Debug)] +pub struct S3Storage { + s3: S3Store, + transport: T, + verify_etag: bool, + max_list_keys: usize, +} + +impl S3Storage { + /// Wrap credentials + transport. Prefer [`S3Store::storage`]. + pub fn new(s3: S3Store, transport: T) -> S3Storage { + S3Storage { + s3, + transport, + verify_etag: true, + max_list_keys: 100_000, + } + } + + /// Turn MD5 `ETag` verification on uploads and downloads off. On by + /// default; the only reason to disable it is a store that reports + /// non-MD5 ETags it does not label as such. + pub fn verify_etag(mut self, on: bool) -> S3Storage { + self.verify_etag = on; + self + } + + /// Cap on how many keys one [`list`](Storage::list) may accumulate + /// (default 100 000). Exceeding it is an error, not a truncation. + pub fn max_list_keys(mut self, n: usize) -> S3Storage { + self.max_list_keys = n; + self + } + + /// The credentials behind this store — also a presigner, for bytes that + /// should flow directly between the client and the bucket. + pub fn presigner(&self) -> &S3Store { + &self.s3 + } + + /// The transport in use. + pub fn transport(&self) -> &T { + &self.transport + } + + /// Sign and send one request against `path` (raw, unencoded) with an + /// already-canonical `query`. + fn request( + &self, + method: &str, + path: &str, + query: &str, + extra: &[(String, String)], + body: &[u8], + ) -> Result { + let (url, headers) = self + .s3 + .sign_request(method, path, query, extra, body, now_secs()?); + self.transport.send(&HttpRequest { + method, + url, + headers, + body, + }) + } + + /// One object request, keyed. Rejects the key before signing anything. + fn object( + &self, + method: &str, + key: &str, + extra: &[(String, String)], + body: &[u8], + ) -> Result { + validate_key(key)?; + self.request(method, &self.s3.object_path(key), "", extra, body) + } + + /// The store's MD5 `ETag` for a body we know, when it is comparable. + /// `None` means "the store did not give us something to check against". + fn etag_mismatch(&self, resp: &HttpResponse, body: &[u8]) -> Option { + if !self.verify_etag { + return None; + } + let etag = resp.header("etag")?.trim_matches('"'); + // Multipart (`…-3`) and SSE-C/KMS ETags are not MD5 of the object. + let comparable = etag.len() == 32 && etag.bytes().all(|b| b.is_ascii_hexdigit()); + if !comparable { + return None; + } + let ours = hex(&md5(body)); + (!ours.eq_ignore_ascii_case(etag)) + .then(|| format!("etag mismatch: store says {etag}, body hashes to {ours}")) + } +} + +impl Storage for S3Storage { + fn put(&self, key: &str, bytes: &[u8], content_type: &str) -> Result<(), String> { + let ct = if content_type.trim().is_empty() { + content_type_of(key) + } else { + content_type + }; + let resp = self.object( + "PUT", + key, + &[("content-type".to_string(), ct.to_string())], + bytes, + )?; + if !matches!(resp.status, 200 | 201 | 204) { + return Err(s3_error("put", key, &resp)); + } + match self.etag_mismatch(&resp, bytes) { + Some(e) => Err(format!("put {key}: {e}")), + None => Ok(()), + } + } + + fn get(&self, key: &str) -> Result>, String> { + let resp = self.object("GET", key, &[], b"")?; + match resp.status { + 200 => match self.etag_mismatch(&resp, &resp.body) { + Some(e) => Err(format!("get {key}: {e}")), + None => Ok(Some(resp.body)), + }, + 404 => Ok(None), + _ => Err(s3_error("get", key, &resp)), + } + } + + fn stat(&self, key: &str) -> Result, String> { + let resp = self.object("HEAD", key, &[], b"")?; + match resp.status { + 200 => Ok(Some(ObjectMeta { + key: key.to_string(), + size: resp + .header("content-length") + .and_then(|v| v.trim().parse().ok()) + .unwrap_or(0), + content_type: resp + .header("content-type") + .filter(|v| !v.is_empty()) + .unwrap_or(content_type_of(key)) + .to_string(), + modified: resp + .header("last-modified") + .and_then(parse_http_date) + .unwrap_or(0), + })), + 404 => Ok(None), + // A HEAD on a bucket without `s3:ListBucket` answers 403 for a + // missing key. Surfacing that as "absent" would turn a permissions + // bug into silent data loss on the next `put`. + _ => Err(s3_error("stat", key, &resp)), + } + } + + /// Removes `key`, reporting whether it existed — which costs a `HEAD` + /// first, because S3 answers `204` to a `DELETE` either way. Concurrent + /// deletes can therefore both report `true`. + fn delete(&self, key: &str) -> Result { + if self.stat(key)?.is_none() { + return Ok(false); + } + let resp = self.object("DELETE", key, &[], b"")?; + match resp.status { + 200 | 202 | 204 | 404 => Ok(true), + _ => Err(s3_error("delete", key, &resp)), + } + } + + /// Lists via `ListObjectsV2`, following continuation tokens. + /// + /// Two honest differences from the fs/db backends: `content_type` is + /// **guessed from the extension** (a listing carries no content type — + /// [`stat`](Storage::stat) reports the real one), and keys that fail + /// [`validate_key`] are skipped, since no backend can address them — this + /// is what hides the empty `dir/` markers S3 GUIs create. + fn list(&self, prefix: &str) -> Result, String> { + let mut out: Vec = Vec::new(); + let mut token: Option = None; + let path = self.s3.bucket_path(); + // Pages are bounded so a store that keeps handing back the same + // continuation token cannot spin forever. + let max_pages = self.max_list_keys / PAGE_SIZE + 2; + + for _ in 0..max_pages { + // SigV4 requires the canonical query sorted by name; these are. + let mut params: Vec<(String, String)> = Vec::new(); + if let Some(t) = &token { + params.push(("continuation-token".to_string(), t.clone())); + } + params.push(("encoding-type".to_string(), "url".to_string())); + params.push(("list-type".to_string(), "2".to_string())); + params.push(("max-keys".to_string(), PAGE_SIZE.to_string())); + if !prefix.is_empty() { + params.push(("prefix".to_string(), prefix.to_string())); + } + + let resp = self.request("GET", &path, &encode_query(¶ms), &[], b"")?; + if resp.status != 200 { + return Err(s3_error("list", prefix, &resp)); + } + let xml = String::from_utf8_lossy(&resp.body); + + for entry in contents_blocks(&xml) { + let Some(key) = tag(entry, "Key").map(|k| percent_decode(&xml_unescape(k))) else { + continue; + }; + if validate_key(&key).is_err() { + continue; + } + out.push(ObjectMeta { + size: tag(entry, "Size") + .and_then(|s| s.trim().parse().ok()) + .unwrap_or(0), + modified: tag(entry, "LastModified") + .and_then(parse_iso8601) + .unwrap_or(0), + content_type: content_type_of(&key).to_string(), + key, + }); + if out.len() > self.max_list_keys { + return Err(format!( + "list('{prefix}') exceeds max_list_keys ({}); narrow the prefix", + self.max_list_keys + )); + } + } + + let truncated = tag(&xml, "IsTruncated").is_some_and(|v| v.trim() == "true"); + if !truncated { + out.sort_by(|a, b| a.key.cmp(&b.key)); + return Ok(out); + } + token = tag(&xml, "NextContinuationToken") + .map(xml_unescape) + .filter(|t| !t.is_empty()) + .ok_or("list truncated but no NextContinuationToken")? + .into(); + } + Err(format!( + "list('{prefix}') did not terminate within {max_pages} pages" + )) + } +} + +/// A readable error from a non-success response, including S3's XML `` +/// and `` when present. Never echoes the whole body. +fn s3_error(op: &str, key: &str, resp: &HttpResponse) -> String { + let body = String::from_utf8_lossy(&resp.body); + let code = tag(&body, "Code").map(xml_unescape); + let message = tag(&body, "Message").map(xml_unescape); + let detail = match (code, message) { + (Some(c), Some(m)) => format!(": {c}: {m}"), + (Some(c), None) => format!(": {c}"), + (None, Some(m)) => format!(": {m}"), + (None, None) => String::new(), + }; + format!("s3 {op} '{key}' failed with HTTP {}{detail}", resp.status) +} + +/// The text inside the first `…` in `xml`. +fn tag<'a>(xml: &'a str, name: &str) -> Option<&'a str> { + let open = format!("<{name}>"); + let close = format!(""); + let start = xml.find(&open)? + open.len(); + let end = xml[start..].find(&close)? + start; + Some(&xml[start..end]) +} + +/// Every `…` block, in document order. +fn contents_blocks(xml: &str) -> impl Iterator { + xml.split("") + .skip(1) + .filter_map(|rest| rest.find("").map(|end| &rest[..end])) +} + +/// The five predefined XML entities. Numeric character references do not occur +/// in the fields we read — keys arrive percent-encoded (`encoding-type=url`). +fn xml_unescape(s: &str) -> String { + if !s.contains('&') { + return s.to_string(); + } + s.replace("<", "<") + .replace(">", ">") + .replace(""", "\"") + .replace("'", "'") + .replace("&", "&") +} + +/// Decode `%XX` escapes; a malformed escape is kept verbatim. +fn percent_decode(s: &str) -> String { + if !s.contains('%') { + return s.to_string(); + } + let b = s.as_bytes(); + let mut out = Vec::with_capacity(b.len()); + let mut i = 0; + while i < b.len() { + if b[i] == b'%' && i + 2 < b.len() { + let hi = (b[i + 1] as char).to_digit(16); + let lo = (b[i + 2] as char).to_digit(16); + if let (Some(hi), Some(lo)) = (hi, lo) { + out.push((hi * 16 + lo) as u8); + i += 3; + continue; + } + } + out.push(b[i]); + i += 1; + } + String::from_utf8_lossy(&out).into_owned() +} + +/// `Fri, 24 May 2013 00:00:00 GMT` (RFC 7231 IMF-fixdate) → unix seconds. +fn parse_http_date(s: &str) -> Option { + let rest = s.split_once(", ").map(|(_, r)| r).unwrap_or(s).trim(); + let mut parts = rest.split_whitespace(); + let day: i64 = parts.next()?.parse().ok()?; + let month = month_number(parts.next()?)?; + let year: i64 = parts.next()?.parse().ok()?; + let mut hms = parts.next()?.split(':'); + let h: i64 = hms.next()?.parse().ok()?; + let m: i64 = hms.next()?.parse().ok()?; + let sec: i64 = hms.next()?.parse().ok()?; + unix_from(year, month, day, h, m, sec) +} + +/// `2013-05-24T00:00:00.000Z` (S3's XML timestamps) → unix seconds. +fn parse_iso8601(s: &str) -> Option { + let (date, time) = s.trim().split_once('T')?; + let mut d = date.split('-'); + let year: i64 = d.next()?.parse().ok()?; + let month: i64 = d.next()?.parse().ok()?; + let day: i64 = d.next()?.parse().ok()?; + let time = time.trim_end_matches('Z'); + let time = time.split_once('.').map(|(t, _)| t).unwrap_or(time); + let mut t = time.split(':'); + let h: i64 = t.next()?.parse().ok()?; + let m: i64 = t.next()?.parse().ok()?; + let sec: i64 = t.next()?.parse().ok()?; + unix_from(year, month, day, h, m, sec) +} + +fn unix_from(y: i64, mo: i64, d: i64, h: i64, mi: i64, s: i64) -> Option { + if !(1..=12).contains(&mo) || !(1..=31).contains(&d) { + return None; + } + if h > 23 || mi > 59 || s > 60 { + return None; + } + let days = days_from_civil(y, mo, d); + // Reject a date the calendar rejects (2013-02-31 &c.). + if civil_from_days(days) != (y, mo, d) { + return None; + } + Some(days * 86_400 + h * 3_600 + mi * 60 + s) +} + +fn month_number(name: &str) -> Option { + const MONTHS: [&str; 12] = [ + "Jan", "Feb", "Mar", "Apr", "May", "Jun", "Jul", "Aug", "Sep", "Oct", "Nov", "Dec", + ]; + MONTHS.iter().position(|m| *m == name).map(|i| i as i64 + 1) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::transport::PlainHttp; + use std::sync::Mutex; + + /// One request as the fake transport saw it. + struct Sent { + method: String, + url: String, + headers: Vec<(String, String)>, + body: Vec, + } + + impl Sent { + fn header(&self, name: &str) -> String { + self.headers + .iter() + .find(|(k, _)| k == name) + .map(|(_, v)| v.clone()) + .unwrap_or_default() + } + } + + /// A transport that records what it was asked to send and replays canned + /// responses — the whole client is exercised without a network. + struct Fake { + replies: Mutex>, + seen: Mutex>, + } + + impl Fake { + fn new(replies: Vec) -> Fake { + Fake { + replies: Mutex::new(replies), + seen: Mutex::new(Vec::new()), + } + } + + fn ok(status: u16, headers: &[(&str, &str)], body: &[u8]) -> HttpResponse { + HttpResponse { + status, + headers: headers + .iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + body: body.to_vec(), + } + } + + fn last_url(&self) -> String { + self.seen.lock().unwrap().last().unwrap().url.clone() + } + } + + impl HttpTransport for Fake { + fn send(&self, req: &HttpRequest<'_>) -> Result { + self.seen.lock().unwrap().push(Sent { + method: req.method.to_string(), + url: req.url.clone(), + headers: req.headers.clone(), + body: req.body.to_vec(), + }); + let mut replies = self.replies.lock().unwrap(); + if replies.is_empty() { + return Err("fake transport out of replies".to_string()); + } + Ok(replies.remove(0)) + } + } + + fn store(replies: Vec) -> S3Storage> { + S3Store::r2("acct", "bkt", "AK", "SK").storage(std::sync::Arc::new(Fake::new(replies))) + } + + #[test] + fn put_signs_content_type_and_payload() { + let etag = format!("\"{}\"", hex(&md5(b"hello"))); + let s = store(vec![Fake::ok(200, &[("etag", &etag)], b"")]); + s.put("a/b.txt", b"hello", "text/plain").unwrap(); + + let seen = s.transport().seen.lock().unwrap(); + let sent = &seen[0]; + assert_eq!(sent.method, "PUT"); + assert_eq!( + sent.url, + "https://acct.r2.cloudflarestorage.com/bkt/a/b.txt" + ); + assert_eq!(sent.body, b"hello"); + assert_eq!(sent.header("content-type"), "text/plain"); + assert_eq!( + sent.header("x-amz-content-sha256"), + hex(&sutegi_crypto::sha256(b"hello")) + ); + assert!(sent + .header("authorization") + .starts_with("AWS4-HMAC-SHA256 Credential=AK/")); + assert_eq!(sent.header("host"), "acct.r2.cloudflarestorage.com"); + } + + #[test] + fn put_guesses_content_type_when_blank() { + let s = store(vec![Fake::ok(200, &[], b"")]); + s.put("x/pic.png", b"\x89PNG", " ").unwrap(); + let seen = s.transport().seen.lock().unwrap(); + assert_eq!(seen[0].header("content-type"), "image/png"); + } + + #[test] + fn get_roundtrip_and_missing() { + let etag = format!("\"{}\"", hex(&md5(b"bytes"))); + let s = store(vec![ + Fake::ok(200, &[("etag", &etag), ("content-length", "5")], b"bytes"), + Fake::ok(404, &[], b"NoSuchKey"), + ]); + assert_eq!(s.get("k.bin").unwrap().unwrap(), b"bytes"); + assert_eq!(s.get("gone.bin").unwrap(), None); + } + + #[test] + fn corrupted_download_is_rejected() { + let etag = format!("\"{}\"", hex(&md5(b"the real bytes"))); + let s = store(vec![Fake::ok(200, &[("etag", &etag)], b"tampered")]); + let err = s.get("k.bin").unwrap_err(); + assert!(err.contains("etag mismatch"), "{err}"); + // Multipart and encrypted ETags are not MD5s — skipped, not failed. + let s = store(vec![Fake::ok(200, &[("etag", "\"abc-3\"")], b"tampered")]); + assert_eq!(s.get("k.bin").unwrap().unwrap(), b"tampered"); + // And the check can be turned off wholesale. + let s = store(vec![Fake::ok(200, &[("etag", &etag)], b"tampered")]).verify_etag(false); + assert_eq!(s.get("k.bin").unwrap().unwrap(), b"tampered"); + } + + #[test] + fn stat_reads_metadata() { + let s = store(vec![Fake::ok( + 200, + &[ + ("content-length", "1234"), + ("content-type", "application/pdf"), + ("last-modified", "Fri, 24 May 2013 00:00:00 GMT"), + ], + b"", + )]); + let meta = s.stat("r/q2.pdf").unwrap().unwrap(); + assert_eq!(meta.size, 1234); + assert_eq!(meta.content_type, "application/pdf"); + assert_eq!(meta.modified, 1_369_353_600); + assert_eq!(s.transport().seen.lock().unwrap()[0].method, "HEAD"); + } + + #[test] + fn stat_403_is_an_error_not_absence() { + let s = store(vec![Fake::ok( + 403, + &[], + b"AccessDenied\ + Access Denied", + )]); + let err = s.stat("k").unwrap_err(); + assert!( + err.contains("HTTP 403") && err.contains("AccessDenied"), + "{err}" + ); + } + + #[test] + fn delete_reports_existence() { + let s = store(vec![ + Fake::ok(200, &[("content-length", "2")], b""), // HEAD: present + Fake::ok(204, &[], b""), // DELETE + ]); + assert!(s.delete("k").unwrap()); + let s = store(vec![Fake::ok(404, &[], b"")]); + assert!(!s.delete("k").unwrap()); + assert_eq!( + s.transport().seen.lock().unwrap().len(), + 1, + "no DELETE after 404" + ); + } + + fn page(keys: &[(&str, u64)], next: Option<&str>) -> HttpResponse { + let mut xml = String::from(""); + for (key, size) in keys { + xml.push_str(&format!( + "{key}2026-08-10T12:00:00.000Z\ + "deadbeef"{size}\ + STANDARD" + )); + } + match next { + Some(t) => xml.push_str(&format!( + "true{t}" + )), + None => xml.push_str("false"), + } + xml.push_str(""); + Fake::ok(200, &[("content-type", "application/xml")], xml.as_bytes()) + } + + #[test] + fn list_paginates_sorts_and_decodes() { + let s = store(vec![ + page( + &[("logs/b.txt", 2), ("logs/a%20one.txt", 1)], + Some("tok/2+x"), + ), + page( + &[("logs/c.png", 3), ("logs/", 0), ("logs/../evil", 9)], + None, + ), + ]); + let objs = s.list("logs/").unwrap(); + assert_eq!( + objs.iter().map(|o| o.key.as_str()).collect::>(), + // %20 decoded, sorted, `logs/` marker and the traversal key dropped + vec!["logs/a one.txt", "logs/b.txt", "logs/c.png"] + ); + assert_eq!(objs[0].size, 1); + assert_eq!(objs[0].modified, 1_786_363_200); // 2026-08-10T12:00:00Z + assert_eq!(objs[2].content_type, "image/png"); + + let seen = s.transport().seen.lock().unwrap(); + assert_eq!( + seen[0].url, + "https://acct.r2.cloudflarestorage.com/bkt/\ + ?encoding-type=url&list-type=2&max-keys=1000&prefix=logs%2F" + ); + // The continuation token is sorted first and URI-encoded. + assert!(seen[1] + .url + .contains("?continuation-token=tok%2F2%2Bx&encoding-type=url")); + } + + #[test] + fn list_refuses_to_truncate_or_spin() { + // Never-ending truncation is bounded, not infinite. + let pages: Vec<_> = (0..12).map(|_| page(&[("k", 1)], Some("same"))).collect(); + let s = store(pages).max_list_keys(2_000); + let err = s.list("").unwrap_err(); + assert!(err.contains("did not terminate"), "{err}"); + + // Over the key cap: an error, not a short answer. + let many: Vec<(String, u64)> = (0..50).map(|i| (format!("k{i}"), 1)).collect(); + let refs: Vec<(&str, u64)> = many.iter().map(|(k, s)| (k.as_str(), *s)).collect(); + let s = store(vec![page(&refs, None)]).max_list_keys(10); + assert!(s.list("").unwrap_err().contains("exceeds max_list_keys")); + + // Truncated with no token is a protocol error, not a silent stop. + let s = store(vec![Fake::ok( + 200, + &[], + b"true", + )]); + assert!(s.list("").unwrap_err().contains("no NextContinuationToken")); + } + + #[test] + fn empty_bucket_lists_empty() { + let s = store(vec![page(&[], None)]); + assert!(s.list("").unwrap().is_empty()); + assert!(!s.transport().last_url().contains("prefix=")); + } + + #[test] + fn bad_keys_never_reach_the_wire() { + let s = store(vec![]); + for key in ["../etc/passwd", "", "a//b", "a\\b"] { + assert!(s.put(key, b"x", "text/plain").is_err(), "{key:?}"); + assert!(s.get(key).is_err(), "{key:?}"); + assert!(s.delete(key).is_err(), "{key:?}"); + } + assert!(s.transport().seen.lock().unwrap().is_empty()); + } + + #[test] + fn control_characters_in_content_type_are_rejected() { + // The transport is the second line of defence; the header never forms. + let s = S3Store::r2("acct", "bkt", "AK", "SK") + .insecure_http() + .storage(PlainHttp::new()); + let err = s + .put("k.txt", b"x", "text/plain\r\nauthorization: leak") + .unwrap_err(); + assert!(err.contains("control characters"), "{err}"); + } + + #[test] + fn error_bodies_become_readable_messages() { + let s = store(vec![Fake::ok( + 500, + &[], + b"InternalErrorWe encountered an \ + internal error & gave up.", + )]); + let err = s.put("k", b"x", "text/plain").unwrap_err(); + assert!(err.contains("HTTP 500"), "{err}"); + assert!(err.contains("InternalError"), "{err}"); + assert!(err.contains("internal error & gave up"), "{err}"); + } + + #[test] + fn transport_errors_propagate() { + let s = store(vec![]); + assert!(s.get("k").unwrap_err().contains("out of replies")); + } + + #[test] + fn presigner_shares_the_credentials() { + let s = store(vec![]); + let url = s.presigner().presign_get("k.txt", 60).unwrap(); + assert!(url.starts_with("https://acct.r2.cloudflarestorage.com/bkt/k.txt?")); + assert_eq!(s.presigner().bucket(), "bkt"); + } + + #[test] + fn debug_does_not_leak_the_secret() { + let printed = format!( + "{:?}", + S3Store::r2("acct", "bkt", "AK", "top-secret").storage(PlainHttp::new()) + ); + assert!(!printed.contains("top-secret"), "{printed}"); + } + + #[test] + fn dates_parse_and_reject() { + assert_eq!( + parse_http_date("Fri, 24 May 2013 00:00:00 GMT"), + Some(1_369_353_600) + ); + assert_eq!( + parse_http_date("Sun, 06 Nov 1994 08:49:37 GMT"), + Some(784_111_777) + ); + assert_eq!( + parse_iso8601("2013-05-24T00:00:00.000Z"), + Some(1_369_353_600) + ); + assert_eq!( + parse_iso8601("2013-05-24T01:01:01Z"), + Some(1_369_353_600 + 3_661) + ); + for bad in ["", "not a date", "Fri, 31 Feb 2013 00:00:00 GMT"] { + assert_eq!(parse_http_date(bad), None, "{bad:?}"); + } + for bad in [ + "2013-02-31T00:00:00Z", + "2013-13-01T00:00:00Z", + "2013-05-24T25:00:00Z", + ] { + assert_eq!(parse_iso8601(bad), None, "{bad:?}"); + } + } + + #[test] + fn xml_and_percent_helpers() { + assert_eq!(xml_unescape("a &lt; b"), "a < b"); + assert_eq!(xml_unescape("<k>"x"'"), "\"x\"'"); + assert_eq!(percent_decode("a%20b%2Fc"), "a b/c"); + assert_eq!(percent_decode("100%"), "100%"); + assert_eq!(percent_decode("a%zzb"), "a%zzb"); + assert_eq!(percent_decode("caf%C3%A9"), "café"); + assert_eq!(tag("x", "b"), Some("x")); + assert_eq!(tag("", "b"), None); + assert_eq!( + contents_blocks("12").count(), + 1 + ); + } +} diff --git a/crates/sutegi-storage/src/transport.rs b/crates/sutegi-storage/src/transport.rs new file mode 100644 index 0000000..957564b --- /dev/null +++ b/crates/sutegi-storage/src/transport.rs @@ -0,0 +1,772 @@ +//! The **outbound-HTTP seam** — one trait, two built-in implementations. This +//! is what lets [`S3Storage`](crate::S3Storage) move real bytes without sutegi +//! growing a TLS stack or a third-party dependency. +//! +//! - [`PlainHttp`] — pure-`std` HTTP/1.1 over `TcpStream`. **Refuses `https`** +//! rather than pretending: point it at an in-cluster MinIO/Garage or a dev +//! container. Same stance as the Postgres driver and the SMTP transport. +//! - [`SystemCurl`] — `https` by delegating the TLS handshake and certificate +//! verification to the system `curl` binary. This is how AWS S3 and +//! Cloudflare R2 work out of the box with zero Cargo dependencies: the +//! crypto that must not be hand-rolled isn't. +//! - your own — implement [`HttpTransport`] over `ureq`/`reqwest`/`hyper` in +//! *your* crate if you already pay for one. +//! +//! Requests arriving here are **already signed**: the transport must send the +//! given headers verbatim and must not follow redirects (a 3xx would replay +//! the `Authorization` header at an attacker-chosen host). + +use std::io::{BufRead, BufReader, Read, Write}; +use std::net::{TcpStream, ToSocketAddrs}; +use std::time::Duration; + +/// Default cap on a single response body (64 MiB). A hostile or misconfigured +/// endpoint cannot make the process allocate more than this. +pub const DEFAULT_MAX_BODY: usize = 64 * 1024 * 1024; + +const MAX_HEADERS: usize = 100; +const MAX_LINE: usize = 8 * 1024; + +/// One outbound request, fully formed and already signed. +/// +/// `headers` is the exact set the signature covers (lowercase names, including +/// `host`) — send them as given, add nothing that could collide, drop nothing. +#[derive(Debug)] +pub struct HttpRequest<'a> { + /// Uppercase HTTP method. + pub method: &'a str, + /// Absolute URL: `scheme://host[:port]/path[?query]`. + pub url: String, + /// Signed headers, lowercase names, `host` included. + pub headers: Vec<(String, String)>, + /// Request body; empty for GET/HEAD/DELETE. + pub body: &'a [u8], +} + +/// One response: status, headers, body. +#[derive(Clone, Debug)] +pub struct HttpResponse { + /// HTTP status code (never 1xx — interim responses are skipped). + pub status: u16, + /// Response headers with names lowercased. + pub headers: Vec<(String, String)>, + /// Response body (empty for HEAD). + pub body: Vec, +} + +impl HttpResponse { + /// The first value for `name` (case-insensitive). + pub fn header(&self, name: &str) -> Option<&str> { + let name = name.to_ascii_lowercase(); + self.headers + .iter() + .find(|(k, _)| *k == name) + .map(|(_, v)| v.as_str()) + } +} + +/// How [`S3Storage`](crate::S3Storage) reaches the object store. One method. +/// +/// Implementations **must**: send `req.headers` unchanged, not follow +/// redirects, and bound the response body they buffer. +pub trait HttpTransport: Send + Sync { + /// Perform `req` and return the response. `Err` is a transport failure + /// (DNS, connect, TLS, timeout) — an HTTP error *status* is a successful + /// send and belongs in [`HttpResponse::status`]. + fn send(&self, req: &HttpRequest<'_>) -> Result; +} + +impl HttpTransport for std::sync::Arc { + fn send(&self, req: &HttpRequest<'_>) -> Result { + (**self).send(req) + } +} + +/// Reject header injection before it reaches the wire. Signed values are ours, +/// but `content_type` originates with the caller. +fn check_header(name: &str, value: &str) -> Result<(), String> { + let bad_name = name.is_empty() + || !name + .bytes() + .all(|b| b.is_ascii_alphanumeric() || b"-_".contains(&b)); + if bad_name { + return Err(format!("invalid header name: {name:?}")); + } + if value.bytes().any(|b| b < 0x20 || b == 0x7f) { + return Err(format!("header {name} contains control characters")); + } + Ok(()) +} + +/// Split an absolute URL into `(scheme, host[:port], path?query)`. +fn split_url(url: &str) -> Result<(&str, &str, &str), String> { + if url.bytes().any(|b| b <= 0x20 || b == 0x7f) { + return Err("url contains whitespace or control characters".to_string()); + } + let (scheme, rest) = url + .split_once("://") + .ok_or_else(|| format!("url has no scheme: {url}"))?; + let (authority, path) = match rest.find('/') { + Some(i) => (&rest[..i], &rest[i..]), + None => (rest, "/"), + }; + if authority.is_empty() { + return Err(format!("url has no host: {url}")); + } + // Credentials in the authority would be sent to an unexpected host. + if authority.contains('@') { + return Err("url must not carry userinfo".to_string()); + } + Ok((scheme, authority, path)) +} + +// ---------------------------------------------------------------- PlainHttp + +/// Pure-`std` HTTP/1.1 transport: blocking `TcpStream`, one connection per +/// request (`Connection: close`), bounded reads and hard timeouts. +/// +/// **`https` URLs are rejected.** This is for object stores you reach over a +/// trusted network path — an in-cluster MinIO, a sidecar Garage, a dev +/// container. Note that SigV4 still signs the payload hash end-to-end, so a +/// tampered request body is rejected by the store even in plaintext; secrecy, +/// not integrity, is what plaintext costs you. For anything crossing the +/// public internet use [`SystemCurl`]. +#[derive(Clone, Debug)] +pub struct PlainHttp { + connect_timeout: Duration, + io_timeout: Duration, + max_body: usize, +} + +impl Default for PlainHttp { + fn default() -> PlainHttp { + PlainHttp::new() + } +} + +impl PlainHttp { + /// 5 s connect, 30 s read/write, 64 MiB body cap. + pub fn new() -> PlainHttp { + PlainHttp { + connect_timeout: Duration::from_secs(5), + io_timeout: Duration::from_secs(30), + max_body: DEFAULT_MAX_BODY, + } + } + + /// Override the connect and per-read/write timeouts. + pub fn timeouts(mut self, connect: Duration, io: Duration) -> PlainHttp { + self.connect_timeout = connect; + self.io_timeout = io; + self + } + + /// Override the response-body cap. + pub fn max_body(mut self, bytes: usize) -> PlainHttp { + self.max_body = bytes; + self + } +} + +impl HttpTransport for PlainHttp { + fn send(&self, req: &HttpRequest<'_>) -> Result { + let (scheme, authority, path) = split_url(&req.url)?; + if !scheme.eq_ignore_ascii_case("http") { + return Err(format!( + "PlainHttp speaks http only (got {scheme}://); use SystemCurl \ + or your own HttpTransport for TLS" + )); + } + let host_port = if authority.contains(':') { + authority.to_string() + } else { + format!("{authority}:80") + }; + + let addrs = host_port + .to_socket_addrs() + .map_err(|e| format!("resolve {host_port}: {e}"))?; + let mut last = format!("resolve {host_port}: no addresses"); + let mut stream = None; + for addr in addrs { + match TcpStream::connect_timeout(&addr, self.connect_timeout) { + Ok(s) => { + stream = Some(s); + break; + } + Err(e) => last = format!("connect {addr}: {e}"), + } + } + let mut stream = stream.ok_or(last)?; + stream + .set_read_timeout(Some(self.io_timeout)) + .and_then(|()| stream.set_write_timeout(Some(self.io_timeout))) + .map_err(|e| format!("set timeouts: {e}"))?; + let _ = stream.set_nodelay(true); + + let mut head = format!("{} {path} HTTP/1.1\r\n", req.method); + for (k, v) in &req.headers { + check_header(k, v)?; + head.push_str(&format!("{k}: {v}\r\n")); + } + head.push_str(&format!("content-length: {}\r\n", req.body.len())); + head.push_str("connection: close\r\n\r\n"); + + stream + .write_all(head.as_bytes()) + .and_then(|()| stream.write_all(req.body)) + .and_then(|()| stream.flush()) + .map_err(|e| format!("write request: {e}"))?; + + let head_only = req.method.eq_ignore_ascii_case("HEAD"); + let mut reader = BufReader::new(stream); + read_response(&mut reader, self.max_body, head_only) + } +} + +// -------------------------------------------------------------- SystemCurl + +/// `https` transport that delegates TLS to the system `curl` binary — the +/// zero-dependency way to talk to AWS S3 and Cloudflare R2. +/// +/// Hardening, all of it on purpose: +/// - **Credentials never touch `argv`.** The URL and every signed header go +/// into a config file fed on `curl`'s **stdin** (`--config -`), so +/// `Authorization` is not visible in `ps` or `/proc//cmdline`. +/// - `--proto =https` — the URL cannot downgrade the protocol. +/// - No redirect following, so a 3xx cannot replay the signature elsewhere. +/// - TLS 1.2 floor; certificate and hostname verification left **on** (there +/// is deliberately no knob here that turns it off). +/// - `--max-filesize` makes `curl` itself refuse an oversized body. +/// - Interim `1xx` header blocks (`Expect: 100-continue`) are skipped. +/// +/// The one thing that does hit disk: a `PUT` body is handed over as a +/// `0600`, `O_EXCL`-created temp file (`curl` reads its upload from a file, and +/// stdin is already carrying the config). It is removed as soon as `curl` +/// exits. Point [`tmp_dir`](SystemCurl::tmp_dir) at a `tmpfs` if object bytes +/// must never land on a real filesystem. +#[derive(Clone, Debug)] +pub struct SystemCurl { + program: String, + connect_timeout: u64, + max_time: u64, + max_body: usize, + allow_http: bool, + tmp_dir: Option, +} + +impl Default for SystemCurl { + fn default() -> SystemCurl { + SystemCurl::new() + } +} + +impl SystemCurl { + /// `curl` from `PATH`, 10 s connect, 300 s total, 64 MiB cap. + pub fn new() -> SystemCurl { + SystemCurl { + program: "curl".to_string(), + connect_timeout: 10, + max_time: 300, + max_body: DEFAULT_MAX_BODY, + allow_http: false, + tmp_dir: None, + } + } + + /// Pin the binary by absolute path — the right call in a hardened image, + /// where resolving `curl` through `PATH` is one more thing to trust. + pub fn at(path: &str) -> SystemCurl { + SystemCurl { + program: path.to_string(), + ..SystemCurl::new() + } + } + + /// Connect and total-operation timeouts, in seconds. + pub fn timeouts(mut self, connect_secs: u64, total_secs: u64) -> SystemCurl { + self.connect_timeout = connect_secs; + self.max_time = total_secs; + self + } + + /// Override the response-body cap. + pub fn max_body(mut self, bytes: usize) -> SystemCurl { + self.max_body = bytes; + self + } + + /// Also allow plaintext `http` URLs (a MinIO inside your own network). + /// Certificate verification for `https` is unaffected — and unreachable + /// by configuration. + pub fn allow_http(mut self) -> SystemCurl { + self.allow_http = true; + self + } + + /// Where `PUT` bodies are staged (default: [`std::env::temp_dir`]). + pub fn tmp_dir(mut self, dir: impl Into) -> SystemCurl { + self.tmp_dir = Some(dir.into()); + self + } + + /// The `--config` document for `req`. Everything secret lives here, and + /// this goes to `curl` on stdin. + fn config(&self, req: &HttpRequest<'_>, upload: Option<&std::path::Path>) -> String { + let mut c = String::new(); + c.push_str(&format!("url = \"{}\"\n", curl_escape(&req.url))); + // `=https` sets the allowed set to exactly https; a bare entry after it + // adds to that set. Never two `=` entries — the second would disable + // the first. + c.push_str(&format!( + "proto = \"=https{}\"\n", + if self.allow_http { ",http" } else { "" } + )); + c.push_str("tlsv1.2\n"); + c.push_str("silent\nshow-error\ninclude\nno-location\nhttp1.1\n"); + c.push_str(&format!("connect-timeout = \"{}\"\n", self.connect_timeout)); + c.push_str(&format!("max-time = \"{}\"\n", self.max_time)); + c.push_str(&format!("max-filesize = \"{}\"\n", self.max_body)); + if req.method.eq_ignore_ascii_case("HEAD") { + c.push_str("head\n"); + } else { + c.push_str(&format!("request = \"{}\"\n", curl_escape(req.method))); + } + for (k, v) in &req.headers { + c.push_str(&format!( + "header = \"{}: {}\"\n", + curl_escape(k), + curl_escape(v) + )); + } + if let Some(path) = upload { + c.push_str(&format!( + "upload-file = \"{}\"\n", + curl_escape(&path.to_string_lossy()) + )); + } + c + } +} + +/// Escape a value for a double-quoted `curl` config parameter. +fn curl_escape(s: &str) -> String { + let mut out = String::with_capacity(s.len()); + for ch in s.chars() { + match ch { + '"' => out.push_str("\\\""), + '\\' => out.push_str("\\\\"), + '\t' => out.push_str("\\t"), + '\n' => out.push_str("\\n"), + '\r' => out.push_str("\\r"), + c => out.push(c), + } + } + out +} + +/// A `0600`, exclusively-created temp file, removed on drop. +struct TempFile(std::path::PathBuf); + +impl TempFile { + fn write(dir: &std::path::Path, bytes: &[u8]) -> Result { + use std::sync::atomic::{AtomicU64, Ordering}; + static SEQ: AtomicU64 = AtomicU64::new(0); + let path = dir.join(format!( + "sutegi-s3-{}-{}.tmp", + std::process::id(), + SEQ.fetch_add(1, Ordering::Relaxed) + )); + let mut opts = std::fs::OpenOptions::new(); + opts.write(true).create_new(true); + #[cfg(unix)] + { + use std::os::unix::fs::OpenOptionsExt; + opts.mode(0o600); + } + let mut f = opts + .open(&path) + .map_err(|e| format!("create {}: {e}", path.display()))?; + let guard = TempFile(path); + f.write_all(bytes) + .and_then(|()| f.flush()) + .map_err(|e| format!("write {}: {e}", guard.0.display()))?; + Ok(guard) + } +} + +impl Drop for TempFile { + fn drop(&mut self) { + let _ = std::fs::remove_file(&self.0); + } +} + +impl HttpTransport for SystemCurl { + fn send(&self, req: &HttpRequest<'_>) -> Result { + use std::process::{Command, Stdio}; + + let (scheme, _, _) = split_url(&req.url)?; + let ok_scheme = scheme.eq_ignore_ascii_case("https") + || (self.allow_http && scheme.eq_ignore_ascii_case("http")); + if !ok_scheme { + return Err(format!( + "SystemCurl refuses {scheme}:// (call .allow_http() to permit plaintext)" + )); + } + for (k, v) in &req.headers { + check_header(k, v)?; + } + + let staged = if req.body.is_empty() { + None + } else { + let dir = self.tmp_dir.clone().unwrap_or_else(std::env::temp_dir); + Some(TempFile::write(&dir, req.body)?) + }; + let config = self.config(req, staged.as_ref().map(|t| t.0.as_path())); + + let mut child = Command::new(&self.program) + .arg("--config") + .arg("-") + .stdin(Stdio::piped()) + .stdout(Stdio::piped()) + .stderr(Stdio::piped()) + .spawn() + .map_err(|e| format!("spawn {}: {e}", self.program))?; + child + .stdin + .take() + .ok_or("curl stdin unavailable")? + .write_all(config.as_bytes()) + .map_err(|e| format!("write curl config: {e}"))?; + let out = child + .wait_with_output() + .map_err(|e| format!("wait for curl: {e}"))?; + drop(staged); + + if !out.status.success() { + return Err(format!( + "curl exited {}: {}", + out.status, + String::from_utf8_lossy(&out.stderr).trim() + )); + } + let head_only = req.method.eq_ignore_ascii_case("HEAD"); + let mut reader = BufReader::new(&out.stdout[..]); + read_response(&mut reader, self.max_body, head_only) + } +} + +// ---------------------------------------------------- HTTP/1.1 response read + +/// Read one response: status line, headers, body. Interim `1xx` blocks are +/// consumed and skipped; `chunked` bodies are de-chunked; everything is +/// bounded. +fn read_response( + r: &mut R, + max_body: usize, + head_only: bool, +) -> Result { + let (status, headers) = loop { + let status = parse_status(&read_line(r)?)?; + let headers = read_headers(r)?; + if !(100..200).contains(&status) { + break (status, headers); + } + }; + + let mut body = Vec::new(); + let chunked = headers + .iter() + .any(|(k, v)| k == "transfer-encoding" && v.to_ascii_lowercase().contains("chunked")); + let length = headers + .iter() + .find(|(k, _)| k == "content-length") + .map(|(_, v)| { + v.trim() + .parse::() + .map_err(|_| format!("bad content-length: {v:?}")) + }) + .transpose()?; + + let bodyless = head_only || status == 204 || status == 304; + if !bodyless { + if chunked { + read_chunked(r, max_body, &mut body)?; + } else if let Some(n) = length { + if n > max_body as u64 { + return Err(format!("response body {n} exceeds cap {max_body}")); + } + body.resize(n as usize, 0); + r.read_exact(&mut body) + .map_err(|e| format!("read body: {e}"))?; + } else { + // No framing: read to EOF, still bounded. + r.by_ref() + .take(max_body as u64 + 1) + .read_to_end(&mut body) + .map_err(|e| format!("read body: {e}"))?; + if body.len() > max_body { + return Err(format!("response body exceeds cap {max_body}")); + } + } + } + Ok(HttpResponse { + status, + headers, + body, + }) +} + +fn read_line(r: &mut R) -> Result { + let mut line = String::new(); + let n = r + .by_ref() + .take(MAX_LINE as u64) + .read_line(&mut line) + .map_err(|e| format!("read: {e}"))?; + if n == 0 { + return Err("connection closed mid-response".to_string()); + } + if !line.ends_with('\n') { + return Err(format!("response line exceeds {MAX_LINE} bytes")); + } + Ok(line.trim_end_matches(['\r', '\n']).to_string()) +} + +fn parse_status(line: &str) -> Result { + let mut parts = line.split(' '); + let version = parts.next().unwrap_or(""); + if !version.starts_with("HTTP/") { + return Err(format!("not an HTTP response: {line:?}")); + } + parts + .next() + .and_then(|c| c.parse::().ok()) + .filter(|c| (100..600).contains(c)) + .ok_or_else(|| format!("bad status line: {line:?}")) +} + +fn read_headers(r: &mut R) -> Result, String> { + let mut headers = Vec::new(); + loop { + let line = read_line(r)?; + if line.is_empty() { + return Ok(headers); + } + if headers.len() >= MAX_HEADERS { + return Err(format!("response has more than {MAX_HEADERS} headers")); + } + let (name, value) = line + .split_once(':') + .ok_or_else(|| format!("bad header line: {line:?}"))?; + headers.push((name.trim().to_ascii_lowercase(), value.trim().to_string())); + } +} + +fn read_chunked(r: &mut R, max_body: usize, out: &mut Vec) -> Result<(), String> { + loop { + let line = read_line(r)?; + let size_hex = line.split(';').next().unwrap_or("").trim(); + let size = u64::from_str_radix(size_hex, 16) + .map_err(|_| format!("bad chunk size: {size_hex:?}"))?; + if size == 0 { + // Trailers, then the terminating blank line. + while !read_line(r)?.is_empty() {} + return Ok(()); + } + if out.len() as u64 + size > max_body as u64 { + return Err(format!("chunked body exceeds cap {max_body}")); + } + let start = out.len(); + out.resize(start + size as usize, 0); + r.read_exact(&mut out[start..]) + .map_err(|e| format!("read chunk: {e}"))?; + if !read_line(r)?.is_empty() { + return Err("chunk not terminated by CRLF".to_string()); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn parse(raw: &str, head_only: bool) -> Result { + read_response(&mut BufReader::new(raw.as_bytes()), 1024, head_only) + } + + #[test] + fn parses_content_length_body() { + let r = parse( + "HTTP/1.1 200 OK\r\nETag: \"abc\"\r\nContent-Length: 5\r\n\r\nhello", + false, + ) + .unwrap(); + assert_eq!(r.status, 200); + assert_eq!(r.body, b"hello"); + assert_eq!(r.header("etag"), Some("\"abc\"")); + assert_eq!(r.header("ETAG"), Some("\"abc\"")); + } + + #[test] + fn skips_interim_100_continue() { + let r = parse( + "HTTP/1.1 100 Continue\r\n\r\nHTTP/1.1 200 OK\r\nContent-Length: 2\r\n\r\nok", + false, + ) + .unwrap(); + assert_eq!((r.status, r.body), (200, b"ok".to_vec())); + } + + #[test] + fn decodes_chunked() { + let r = parse( + "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n\ + 4\r\nsute\r\n2;ext\r\ngi\r\n0\r\n\r\n", + false, + ) + .unwrap(); + assert_eq!(r.body, b"sutegi"); + } + + #[test] + fn head_and_204_have_no_body() { + // A HEAD response advertises the length it *would* have sent. + let r = parse("HTTP/1.1 200 OK\r\nContent-Length: 99\r\n\r\n", true).unwrap(); + assert!(r.body.is_empty()); + assert_eq!(r.header("content-length"), Some("99")); + let r = parse("HTTP/1.1 204 No Content\r\n\r\n", false).unwrap(); + assert_eq!((r.status, r.body.len()), (204, 0)); + } + + #[test] + fn body_cap_is_enforced() { + let raw = format!( + "HTTP/1.1 200 OK\r\nContent-Length: 4096\r\n\r\n{}", + "x".repeat(4096) + ); + assert!(parse(&raw, false).unwrap_err().contains("exceeds cap")); + let chunked = "HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n1000\r\n"; + assert!(parse(chunked, false).unwrap_err().contains("exceeds cap")); + } + + #[test] + fn rejects_garbage_and_truncation() { + assert!(parse("220 smtp.example.com ESMTP\r\n\r\n", false).is_err()); + assert!(parse("HTTP/1.1 999 Nope\r\n\r\n", false).is_err()); + assert!(parse("HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhi", false).is_err()); + } + + #[test] + fn url_split_and_guards() { + assert_eq!( + split_url("https://a.example.com/x/y?q=1").unwrap(), + ("https", "a.example.com", "/x/y?q=1") + ); + assert_eq!( + split_url("http://localhost:9000").unwrap(), + ("http", "localhost:9000", "/") + ); + assert!(split_url("no-scheme/x").is_err()); + assert!(split_url("http://user:pw@evil/x").is_err()); + assert!(split_url("http://host/a b").is_err()); + assert!(split_url("http://host/a\r\nX: y").is_err()); + } + + #[test] + fn header_injection_rejected() { + assert!(check_header("content-type", "text/plain").is_ok()); + assert!(check_header("content-type", "a\r\nauthorization: leak").is_err()); + assert!(check_header("bad name", "v").is_err()); + } + + #[test] + fn plain_http_refuses_tls() { + let err = PlainHttp::new() + .send(&HttpRequest { + method: "GET", + url: "https://bucket.s3.amazonaws.com/k".to_string(), + headers: vec![], + body: b"", + }) + .unwrap_err(); + assert!(err.contains("http only"), "{err}"); + } + + #[test] + fn curl_config_hides_secrets_and_pins_protocol() { + let req = HttpRequest { + method: "PUT", + url: "https://acct.r2.cloudflarestorage.com/bkt/a%20b.txt".to_string(), + headers: vec![ + ( + "authorization".to_string(), + "AWS4-HMAC-SHA256 Credential=AK/x".to_string(), + ), + ("content-type".to_string(), "text/plain".to_string()), + ], + body: b"hi", + }; + let cfg = SystemCurl::new().config(&req, Some(std::path::Path::new("/tmp/x.tmp"))); + assert!(cfg.contains("url = \"https://acct.r2.cloudflarestorage.com/bkt/a%20b.txt\"\n")); + assert!(cfg.contains("proto = \"=https\"\n")); + assert!(cfg.contains("no-location\n") && cfg.contains("tlsv1.2\n")); + assert!(cfg.contains("request = \"PUT\"\n")); + assert!(cfg.contains("header = \"authorization: AWS4-HMAC-SHA256 Credential=AK/x\"\n")); + assert!(cfg.contains("upload-file = \"/tmp/x.tmp\"\n")); + assert!(cfg.contains("max-filesize = \"67108864\"\n")); + // HEAD uses --head, never -X HEAD (which would hang waiting for a body). + let head = HttpRequest { + method: "HEAD", + ..req + }; + let cfg = SystemCurl::new().config(&head, None); + assert!(cfg.contains("head\n") && !cfg.contains("request =")); + assert!(SystemCurl::new() + .allow_http() + .config(&head, None) + .contains("proto = \"=https,http\"\n")); + } + + #[test] + fn curl_refuses_plaintext_by_default() { + let req = HttpRequest { + method: "GET", + url: "http://127.0.0.1:9000/bkt/k".to_string(), + headers: vec![], + body: b"", + }; + let err = SystemCurl::new().send(&req).unwrap_err(); + assert!(err.contains("refuses http://"), "{err}"); + // …and there is no knob anywhere that disables certificate checking. + assert!(!SystemCurl::new() + .allow_http() + .config(&req, None) + .contains("insecure")); + } + + #[test] + fn curl_escapes_quotes_and_newlines() { + assert_eq!(curl_escape("a\"b\\c\r\nd"), "a\\\"b\\\\c\\r\\nd"); + } + + #[test] + fn temp_file_is_private_and_removed() { + let t = TempFile::write(&std::env::temp_dir(), b"secret").unwrap(); + let path = t.0.clone(); + assert_eq!(std::fs::read(&path).unwrap(), b"secret"); + #[cfg(unix)] + { + use std::os::unix::fs::PermissionsExt; + let mode = std::fs::metadata(&path).unwrap().permissions().mode(); + assert_eq!( + mode & 0o777, + 0o600, + "temp upload must not be world-readable" + ); + } + drop(t); + assert!(!path.exists()); + } +} diff --git a/crates/sutegi-storage/tests/s3_roundtrip.rs b/crates/sutegi-storage/tests/s3_roundtrip.rs new file mode 100644 index 0000000..1255239 --- /dev/null +++ b/crates/sutegi-storage/tests/s3_roundtrip.rs @@ -0,0 +1,527 @@ +//! End-to-end [`S3Storage`] over a real socket: the client signs, a tiny +//! in-process S3-compatible store answers, and every layer — SigV4 headers, +//! HTTP/1.1 framing, `ListObjectsV2` XML, continuation tokens, ETag +//! verification — is exercised against bytes on the wire rather than a mock. +//! +//! The stub enforces the parts a real store enforces and we can check cheaply: +//! `Authorization` must be present and well-formed, and +//! `x-amz-content-sha256` must equal the SHA-256 of the body actually received +//! (the guarantee that makes plaintext transport tolerable). The signature +//! *value* is pinned separately, against AWS's published known-answer vectors, +//! in `s3.rs`. + +use std::collections::BTreeMap; +use std::io::{BufRead, BufReader, Read, Write}; +use std::net::{TcpListener, TcpStream}; +use std::sync::{Arc, Mutex}; +use sutegi_crypto::{hex, md5, sha256}; +use sutegi_storage::{PlainHttp, S3Storage, S3Store, Storage}; + +const BUCKET: &str = "test-bucket"; + +#[derive(Default)] +struct Bucket { + objects: BTreeMap, String)>, + /// Keys the stub answers with the wrong bytes, to prove the client checks. + poison: Vec, + requests: usize, +} + +/// Boot the stub on an ephemeral port; returns its `host:port` and state. +fn stub() -> (String, Arc>) { + let listener = TcpListener::bind("127.0.0.1:0").expect("bind"); + let addr = listener.local_addr().unwrap().to_string(); + let state = Arc::new(Mutex::new(Bucket::default())); + let shared = Arc::clone(&state); + std::thread::spawn(move || { + for conn in listener.incoming() { + let Ok(conn) = conn else { break }; + // One request per connection: the client sends `Connection: close`. + if let Err(e) = serve_one(conn, &shared) { + eprintln!("stub: {e}"); + } + } + }); + (addr, state) +} + +fn store(addr: &str) -> S3Storage { + S3Store::new(BUCKET, "us-east-1", "AKIAEXAMPLE", "secret-key") + .with_endpoint(addr) + .insecure_http() + .storage(PlainHttp::new()) +} + +fn serve_one(mut conn: TcpStream, state: &Arc>) -> Result<(), String> { + let mut reader = BufReader::new(conn.try_clone().map_err(|e| e.to_string())?); + let mut line = String::new(); + reader.read_line(&mut line).map_err(|e| e.to_string())?; + let mut parts = line.split_whitespace(); + let method = parts.next().unwrap_or_default().to_string(); + let target = parts.next().unwrap_or_default().to_string(); + + let mut headers: Vec<(String, String)> = Vec::new(); + loop { + let mut h = String::new(); + reader.read_line(&mut h).map_err(|e| e.to_string())?; + let h = h.trim_end(); + if h.is_empty() { + break; + } + let (k, v) = h.split_once(':').ok_or("bad header")?; + headers.push((k.trim().to_ascii_lowercase(), v.trim().to_string())); + } + let header = |name: &str| { + headers + .iter() + .find(|(k, _)| k == name) + .map(|(_, v)| v.clone()) + .unwrap_or_default() + }; + let length: usize = header("content-length").parse().unwrap_or(0); + let mut body = vec![0u8; length]; + reader.read_exact(&mut body).map_err(|e| e.to_string())?; + + // What a real store checks before it looks at the request at all. + let auth = header("authorization"); + assert!( + auth.starts_with("AWS4-HMAC-SHA256 Credential=AKIAEXAMPLE/") + && auth.contains(",SignedHeaders=") + && auth.contains(",Signature="), + "malformed Authorization: {auth:?}" + ); + assert_eq!(header("host"), conn.local_addr().unwrap().to_string()); + assert_eq!( + header("x-amz-content-sha256"), + hex(&sha256(&body)), + "payload hash does not cover the body received" + ); + + let (path, query) = match target.split_once('?') { + Some((p, q)) => (p.to_string(), q.to_string()), + None => (target.clone(), String::new()), + }; + let path = percent_decode(&path); + let prefix = format!("/{BUCKET}/"); + let key = path.strip_prefix(&prefix).unwrap_or_default().to_string(); + + let resp = { + let mut b = state.lock().unwrap(); + b.requests += 1; + match method.as_str() { + _ if !path.starts_with(&prefix) => response( + 404, + &[], + format!( + "NoSuchBucketThe specified bucket \ + does not exist{path}" + ) + .into_bytes(), + ), + _ if key.is_empty() => list(&b, &query), + "PUT" => { + let etag = hex(&md5(&body)); + b.objects + .insert(key, (body, header("content-type").to_string())); + response(200, &[("ETag", &format!("\"{etag}\""))], Vec::new()) + } + "GET" | "HEAD" => match b.objects.get(&key) { + Some((bytes, ct)) => { + let etag = hex(&md5(bytes)); + let served = if b.poison.contains(&key) { + b"corrupted in flight".to_vec() + } else { + bytes.clone() + }; + let head = [ + ("ETag", format!("\"{etag}\"")), + ("Content-Type", ct.clone()), + ("Last-Modified", "Mon, 10 Aug 2026 12:00:00 GMT".to_string()), + ]; + let head: Vec<(&str, &str)> = + head.iter().map(|(k, v)| (*k, v.as_str())).collect(); + if method == "HEAD" { + // A HEAD advertises the length it would have sent. + let mut h = head.clone(); + let len = served.len().to_string(); + h.push(("Content-Length", &len)); + format!( + "HTTP/1.1 200 OK\r\n{}Connection: close\r\n\r\n", + h.iter() + .map(|(k, v)| format!("{k}: {v}\r\n")) + .collect::() + ) + .into_bytes() + } else { + response(200, &head, served) + } + } + None => response(404, &[], b"NoSuchKey".to_vec()), + }, + "DELETE" => { + b.objects.remove(&key); + response(204, &[], Vec::new()) + } + other => response( + 405, + &[], + format!("MethodNotAllowed{other}") + .into_bytes(), + ), + } + }; + conn.write_all(&resp).map_err(|e| e.to_string())?; + conn.flush().map_err(|e| e.to_string()) +} + +fn response(status: u16, headers: &[(&str, &str)], body: Vec) -> Vec { + let mut out = format!( + "HTTP/1.1 {status} {}\r\n", + if status < 300 { "OK" } else { "Error" } + ); + for (k, v) in headers { + out.push_str(&format!("{k}: {v}\r\n")); + } + out.push_str(&format!( + "Content-Length: {}\r\nConnection: close\r\n\r\n", + body.len() + )); + let mut out = out.into_bytes(); + out.extend_from_slice(&body); + out +} + +/// `ListObjectsV2`, with `encoding-type=url` keys and real continuation tokens. +fn list(b: &Bucket, query: &str) -> Vec { + let params: BTreeMap = query + .split('&') + .filter_map(|kv| kv.split_once('=')) + .map(|(k, v)| (percent_decode(k), percent_decode(v))) + .collect(); + assert_eq!(params.get("list-type").map(String::as_str), Some("2")); + assert_eq!(params.get("encoding-type").map(String::as_str), Some("url")); + let want = params.get("prefix").cloned().unwrap_or_default(); + let max: usize = params + .get("max-keys") + .and_then(|m| m.parse().ok()) + .unwrap_or(1000); + let after = params.get("continuation-token").cloned(); + + let matching: Vec<&String> = b + .objects + .keys() + .filter(|k| k.starts_with(&want)) + .filter(|k| after.as_ref().map_or(true, |a| *k > a)) + .collect(); + let page: Vec<&&String> = matching.iter().take(max).collect(); + let truncated = matching.len() > page.len(); + + let mut xml = String::from(""); + for key in &page { + let (bytes, _) = &b.objects[**key]; + xml.push_str(&format!( + "{}\ + 2026-08-10T12:00:00.000Z\ + "{}"{}\ + STANDARD", + url_encode(key), + hex(&md5(bytes)), + bytes.len() + )); + } + xml.push_str(&format!("{truncated}")); + if truncated { + xml.push_str(&format!( + "{}", + url_encode(page.last().unwrap()) + )); + } + xml.push_str(""); + response( + 200, + &[("Content-Type", "application/xml")], + xml.into_bytes(), + ) +} + +fn url_encode(s: &str) -> String { + s.bytes() + .map(|b| match b { + b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' | b'/' => { + (b as char).to_string() + } + _ => format!("%{b:02X}"), + }) + .collect() +} + +fn percent_decode(s: &str) -> String { + let b = s.as_bytes(); + let mut out = Vec::with_capacity(b.len()); + let mut i = 0; + while i < b.len() { + match (b[i], i + 2 < b.len()) { + (b'%', true) => { + match u8::from_str_radix(std::str::from_utf8(&b[i + 1..i + 3]).unwrap_or("zz"), 16) + { + Ok(byte) => { + out.push(byte); + i += 3; + } + Err(_) => { + out.push(b[i]); + i += 1; + } + } + } + _ => { + out.push(b[i]); + i += 1; + } + } + } + String::from_utf8_lossy(&out).into_owned() +} + +// ------------------------------------------------------------------- tests + +#[test] +fn full_lifecycle_over_the_wire() { + let (addr, state) = stub(); + let s = store(&addr); + + // put → get → stat → exists → delete, the whole trait against a socket. + s.put( + "reports/q2 final+draft.pdf", + b"%PDF-1.7 body", + "application/pdf", + ) + .unwrap(); + assert_eq!( + s.get("reports/q2 final+draft.pdf").unwrap().unwrap(), + b"%PDF-1.7 body" + ); + + let meta = s.stat("reports/q2 final+draft.pdf").unwrap().unwrap(); + assert_eq!(meta.key, "reports/q2 final+draft.pdf"); + assert_eq!(meta.size, 13); + assert_eq!(meta.content_type, "application/pdf"); + assert_eq!(meta.modified, 1_786_363_200); // 2026-08-10T12:00:00Z + assert!(s.exists("reports/q2 final+draft.pdf").unwrap()); + + // A key that never existed, and one that stops existing. + assert_eq!(s.get("nope.txt").unwrap(), None); + assert_eq!(s.stat("nope.txt").unwrap(), None); + assert!(!s.delete("nope.txt").unwrap()); + assert!(s.delete("reports/q2 final+draft.pdf").unwrap()); + assert!(!s.exists("reports/q2 final+draft.pdf").unwrap()); + + // Overwrite keeps the newest bytes and the newest content type. + s.put("x.bin", b"one", "application/x-one").unwrap(); + s.put("x.bin", b"two!", "application/x-two").unwrap(); + assert_eq!(s.get("x.bin").unwrap().unwrap(), b"two!"); + assert_eq!( + s.stat("x.bin").unwrap().unwrap().content_type, + "application/x-two" + ); + assert!(state.lock().unwrap().requests > 10); +} + +#[test] +fn get_reader_streams_the_object() { + let (addr, _) = stub(); + let s = store(&addr); + s.put("r.txt", b"stream me", "text/plain").unwrap(); + let mut buf = Vec::new(); + s.get_reader("r.txt") + .unwrap() + .unwrap() + .read_to_end(&mut buf) + .unwrap(); + assert_eq!(buf, b"stream me"); + assert!(s.get_reader("missing.txt").unwrap().is_none()); +} + +#[test] +fn tampered_download_is_caught_on_the_wire() { + let (addr, state) = stub(); + let s = store(&addr); + s.put("secret.txt", b"the real bytes", "text/plain") + .unwrap(); + state.lock().unwrap().poison.push("secret.txt".to_string()); + + let err = s.get("secret.txt").unwrap_err(); + assert!(err.contains("etag mismatch"), "{err}"); + // Same response, integrity checking off: the caller gets the bad bytes, + // which is exactly what opting out means. + let lax = store(&addr).verify_etag(false); + assert_eq!( + lax.get("secret.txt").unwrap().unwrap(), + b"corrupted in flight" + ); +} + +#[test] +fn list_paginates_across_continuation_tokens() { + let (addr, state) = stub(); + let s = store(&addr); + // Two full pages plus a remainder: 2400 keys is 3 round trips, and the + // client must stitch them without losing or duplicating one. + { + let mut b = state.lock().unwrap(); + for i in 0..2400 { + b.objects.insert( + format!("blobs/{i:05}.bin"), + (vec![b'x'; 3], "application/octet-stream".to_string()), + ); + } + b.objects + .insert("other/z.txt".to_string(), (b"z".to_vec(), String::new())); + b.requests = 0; + } + + let objs = s.list("blobs/").unwrap(); + assert_eq!(objs.len(), 2400); + assert_eq!(objs[0].key, "blobs/00000.bin"); + assert_eq!(objs[2399].key, "blobs/02399.bin"); + assert!(objs.windows(2).all(|w| w[0].key < w[1].key), "sorted"); + assert!(objs.iter().all(|o| o.size == 3)); + assert_eq!(state.lock().unwrap().requests, 3, "1000 + 1000 + 400"); + + // An empty prefix lists everything; a prefix that matches nothing is empty. + assert_eq!(s.list("").unwrap().len(), 2401); + assert!(s.list("nothing/").unwrap().is_empty()); + + // And the cap refuses to answer rather than answering partially. + let err = store(&addr).max_list_keys(500).list("blobs/").unwrap_err(); + assert!(err.contains("exceeds max_list_keys"), "{err}"); +} + +#[test] +fn keys_needing_encoding_survive_the_round_trip() { + let (addr, _) = stub(); + let s = store(&addr); + for key in [ + "unicode/reçu-café.txt", + "spaces/a b c.txt", + "plus/a+b=c.txt", + "punct/(paren)[brack]&.txt", + "deep/a/b/c/d/e.bin", + ] { + s.put(key, key.as_bytes(), "text/plain") + .unwrap_or_else(|e| panic!("put {key}: {e}")); + assert_eq!( + s.get(key).unwrap().unwrap(), + key.as_bytes(), + "get {key} came back wrong" + ); + } + // The same keys survive the XML listing, which encodes them differently + // than the request path does. + let listed: Vec = s.list("").unwrap().into_iter().map(|o| o.key).collect(); + assert!( + listed.contains(&"unicode/reçu-café.txt".to_string()), + "{listed:?}" + ); + assert!( + listed.contains(&"punct/(paren)[brack]&.txt".to_string()), + "{listed:?}" + ); + assert_eq!(listed.len(), 5); +} + +#[test] +fn store_errors_become_useful_messages() { + let (addr, _) = stub(); + // A bucket the store does not have: the XML ``/`` reach the + // caller, not a byte dump — and a 404 on a *write* is an error, never a + // quietly-successful no-op. + let s = S3Store::new("wrong-bucket", "us-east-1", "AKIAEXAMPLE", "secret-key") + .with_endpoint(&addr) + .insecure_http() + .storage(PlainHttp::new()); + let err = s.put("k.txt", b"x", "text/plain").unwrap_err(); + assert!(err.contains("HTTP 404"), "{err}"); + assert!(err.contains("NoSuchBucket"), "{err}"); + assert!(err.contains("The specified bucket does not exist"), "{err}"); + assert!(err.contains("s3 put 'k.txt'"), "{err}"); + + // A missing *bucket* on a read is indistinguishable from a missing key at + // the HTTP level, so `get` reports absence — the documented consequence of + // 404 meaning both things. + assert_eq!(s.get("k.txt").unwrap(), None); +} + +#[test] +fn a_dead_endpoint_fails_fast_and_clearly() { + // Port 1 on loopback: nothing listens, and the error says so rather than + // hanging or surfacing as a missing object. + let s = S3Store::new(BUCKET, "us-east-1", "AK", "SK") + .with_endpoint("127.0.0.1:1") + .insecure_http() + .storage(PlainHttp::new()); + let err = s.get("k.txt").unwrap_err(); + assert!(err.contains("connect"), "{err}"); +} + +/// The `https` transport's plumbing — config on stdin, temp-file upload, +/// `--head` for HEAD, `-i` header parsing — driven against the same stub over +/// plaintext (`allow_http`). TLS itself is `curl`'s business, not ours; what is +/// ours is everything around it. +#[test] +fn system_curl_transport_moves_bytes() { + let curl_present = std::process::Command::new("curl") + .arg("--version") + .output() + .is_ok_and(|o| o.status.success()); + if !curl_present { + eprintln!("skipping: no curl on PATH"); + return; + } + + let (addr, _) = stub(); + let s = S3Store::new(BUCKET, "us-east-1", "AKIAEXAMPLE", "secret-key") + .with_endpoint(&addr) + .insecure_http() + .storage(sutegi_storage::SystemCurl::new().allow_http()); + + s.put("via/curl.txt", b"through a subprocess", "text/plain") + .unwrap(); + assert_eq!( + s.get("via/curl.txt").unwrap().unwrap(), + b"through a subprocess" + ); + let meta = s.stat("via/curl.txt").unwrap().unwrap(); + assert_eq!((meta.size, meta.content_type.as_str()), (20, "text/plain")); + assert_eq!(s.list("via/").unwrap().len(), 1); + assert!(s.delete("via/curl.txt").unwrap()); + assert_eq!(s.get("via/curl.txt").unwrap(), None); + + // A body big enough to cross curl's Expect: 100-continue threshold, whose + // interim header block the parser has to skip. + let big = vec![b'z'; 2 * 1024 * 1024]; + s.put("via/big.bin", &big, "application/octet-stream") + .unwrap(); + assert_eq!(s.get("via/big.bin").unwrap().unwrap().len(), big.len()); + + // No staged upload survives the call. + let leftovers = std::fs::read_dir(std::env::temp_dir()) + .unwrap() + .filter_map(|e| e.ok()) + .filter(|e| e.file_name().to_string_lossy().starts_with("sutegi-s3-")) + .count(); + assert_eq!(leftovers, 0, "temp upload files leaked"); +} + +#[test] +fn presigning_and_moving_bytes_share_one_credential() { + let (addr, _) = stub(); + let s = store(&addr); + s.put("shared.txt", b"hi", "text/plain").unwrap(); + let url = s.presigner().presign_get("shared.txt", 300).unwrap(); + assert!( + url.starts_with(&format!("http://{addr}/{BUCKET}/shared.txt?")), + "{url}" + ); + assert!(url.contains("X-Amz-Signature=")); +} diff --git a/crates/sutegi/src/lib.rs b/crates/sutegi/src/lib.rs index b73aad9..21ceee3 100644 --- a/crates/sutegi/src/lib.rs +++ b/crates/sutegi/src/lib.rs @@ -18,7 +18,7 @@ //! | `template` | sutegi-template | Blade-style template engine (`{{ }}`, `@if`, `@foreach`, `@include`) | //! | `mail` | sutegi-mail (+ template) | Email builder, themed messages, Transport seam, smtp/sendmail/log drivers | //! | `auth-mail` | + sutegi-auth/mail | email-verification + password-reset flows | -//! | `storage` | sutegi-storage (pure std) | file storage: local fs + S3 presigned URLs | +//! | `storage` | sutegi-storage (pure std) | file storage: local fs, S3/R2 objects, presigned URLs | //! | `storage-db` | + sutegi-orm | blobs in SQLite/Postgres over the `Backend` seam | //! | `ws` | sutegi-ws | WebSockets: `App::ws` on the sharded kqueue/epoll reactor | //! | `pubsub` | sutegi-pubsub | in-process topic fan-out behind the `Broker` seam | @@ -499,7 +499,9 @@ pub mod prelude { pub use sutegi_auth::AuthMail; #[cfg(feature = "storage")] - pub use sutegi_storage::{FsStorage, ObjectMeta, S3Store, Storage}; + pub use sutegi_storage::{ + FsStorage, HttpTransport, ObjectMeta, PlainHttp, S3Storage, S3Store, Storage, SystemCurl, + }; #[cfg(feature = "sqlite")] pub use sutegi_orm::db::Db; diff --git a/examples/storage/src/main.rs b/examples/storage/src/main.rs index 12edcfe..ee7c0d9 100644 --- a/examples/storage/src/main.rs +++ b/examples/storage/src/main.rs @@ -1,7 +1,8 @@ -//! A file server on sutegi's storage layer — [`FsStorage`] for the bytes -//! (single-node, zero-ops: one directory on disk), plus the agent-native S3 -//! shape: `presign_upload` / `presign_download` tools that mint time-limited -//! URLs so the **agent moves the bytes itself**, straight to the object store. +//! A file server on sutegi's storage layer, with **the backend chosen at boot +//! and the routes never told which**: a directory on disk by default, an +//! S3/R2 bucket when `STORAGE=s3`. Plus the agent-native S3 shape — +//! `presign_upload` / `presign_download` tools that mint time-limited URLs so +//! the **agent moves the bytes itself**, straight to the object store. //! //! ```text //! curl -T report.pdf localhost:8080/files/report.pdf @@ -12,9 +13,12 @@ //! -d '{"key":"report.pdf"}' # S3 URL (needs S3_* env) //! ``` //! -//! S3 presigning activates when `S3_BUCKET`, `S3_ACCESS_KEY` and -//! `S3_SECRET_KEY` are set (`S3_REGION` defaults to `us-east-1`; -//! `S3_ENDPOINT` points at R2/MinIO/… and switches to path-style). +//! The S3 credentials come from `S3_BUCKET`, `S3_ACCESS_KEY`, `S3_SECRET_KEY` +//! (`S3_REGION` defaults to `us-east-1`; `S3_ENDPOINT` points at R2/MinIO/… and +//! switches to path-style). They power presigning on their own; add +//! `STORAGE=s3` and the file routes themselves move their bytes through the +//! bucket. `S3_INSECURE=1` selects plaintext + the pure-`std` transport for an +//! in-cluster MinIO; otherwise `https` goes through the system `curl`. use sutegi::prelude::*; @@ -34,6 +38,27 @@ fn s3_from_env() -> Option { Some(s3) } +/// The storage backend as the routes see it: a trait object, so the choice +/// below is the only code in this file that knows where bytes live. +type Store = Box; + +/// `STORAGE=s3` puts the file routes on the bucket; anything else is a +/// directory on disk. This is the whole cost of swapping the backend. +fn store_from_env() -> Store { + if std::env::var("STORAGE").as_deref() == Ok("s3") { + let s3 = s3_from_env().expect("STORAGE=s3 needs S3_BUCKET / S3_ACCESS_KEY / S3_SECRET_KEY"); + return if std::env::var("S3_INSECURE").is_ok() { + // Plaintext to a store on a trusted network path: no TLS, no + // subprocess, pure std. SigV4 still signs the payload hash. + Box::new(s3.insecure_http().storage(PlainHttp::new())) + } else { + Box::new(s3.storage(SystemCurl::new())) + }; + } + let root = std::env::var("STORAGE_ROOT").unwrap_or_else(|_| "files".to_string()); + Box::new(FsStorage::new(root).expect("open storage root")) +} + fn presign(s3: &Option, args: &Json, put: bool) -> Result { let s3 = s3.as_ref().ok_or_else(|| { Error::new( @@ -60,8 +85,7 @@ fn presign(s3: &Option, args: &Json, put: bool) -> Result } fn main() -> std::io::Result<()> { - let root = std::env::var("STORAGE_ROOT").unwrap_or_else(|_| "files".to_string()); - let store = FsStorage::new(root).expect("open storage root"); + let store = store_from_env(); let s3_up = s3_from_env(); let s3_down = s3_up.clone(); // Gate the agent surface: presigning mints URLs for arbitrary object keys, @@ -86,7 +110,7 @@ fn main() -> std::io::Result<()> { .state(store) .get("/", "Health check.", |_| "sutegi storage up") .get("/files", "List stored files.", |c| { - let items = c.state::().list("")?; + let items = c.state::().list("")?; Ok::<_, Error>(json( 200, &Json::arr(items.iter().map(ObjectMeta::to_json).collect()), @@ -94,14 +118,14 @@ fn main() -> std::io::Result<()> { }) .put("/files/:name", "Store the raw request body as a file.", |c| { let ct = c.header("content-type").unwrap_or(""); - c.state::().put(name(c), &c.req.body, ct)?; + c.state::().put(name(c), &c.req.body, ct)?; Ok::<_, Error>(json(201, &Json::obj(vec![("key", Json::str(name(c)))]))) }) .get( "/files/:name", "Download a file with its stored content type.", |c| -> Result { - let store = c.state::(); + let store = c.state::(); match store.stat(name(c))? { Some(meta) => { let bytes = store.get(name(c))?.unwrap_or_default(); @@ -114,7 +138,7 @@ fn main() -> std::io::Result<()> { }, ) .delete("/files/:name", "Delete a file.", |c| { - let removed = c.state::().delete(name(c))?; + let removed = c.state::().delete(name(c))?; Ok::<_, Error>(json( 200, &Json::obj(vec![("deleted", Json::Bool(removed))]), diff --git a/landing/src/Docs.svelte b/landing/src/Docs.svelte index 3540ac6..ce50789 100644 --- a/landing/src/Docs.svelte +++ b/landing/src/Docs.svelte @@ -453,8 +453,10 @@ let msg = MailMessage::new() .line("If you didn't sign up, ignore this email."); mailer.send(msg.build("Acme"))?;`; - const cStorage = `// FsStorage: single-node, one directory on disk. Same Storage trait as DbStorage / S3. -let store = FsStorage::new("files")?; + const cStorage = `// One trait, three backends. Swap the type you construct, not the call sites. +let store = FsStorage::new("files")?; // single-node, one directory +let store = S3Store::r2(&acct, "media", &ak, &sk) // …or a Cloudflare R2 bucket + .storage(SystemCurl::new()); .put("/files/:name", "Upload", |c| { let ct = c.header("content-type").unwrap_or(""); @@ -1133,9 +1135,14 @@ App::new("api") (put/get/stat/delete/list/ get_reader, with traversal-validated keys). FsStorage writes to a local directory (atomic temp-and-rename); DbStorage (the storage-db feature) stores blobs in - SQLite or Postgres for multi-pod files with no new infrastructure; and S3Store is a pure-std - SigV4 presigner for AWS/R2/MinIO. Presigning is the agent-native trick: mint a time-limited URL and let the - client (or agent) move the bytes straight to the object store — they never pass through your server. + SQLite or Postgres for multi-pod files with no new infrastructure; and S3Storage puts the same + trait on a real bucket — AWS S3, Cloudflare R2, MinIO, Garage. It moves the bytes over an injected + HttpTransport, which is how an S3 client exists here with no TLS stack and no dependency: + SystemCurl borrows the system curl for https, PlainHttp is pure std + for a store on a trusted path. Requests are SigV4-signed over the real payload hash, and downloads are + ETag-verified. The same credentials also presign, which is the agent-native trick: mint a time-limited URL + and let the client (or agent) move the bytes straight to the object store — they never pass through your + server.

{@render code(cStorage, 'storage')}