From 3de8590b6bb03b6dd9863e8ceb9de43c6af1ab2d Mon Sep 17 00:00:00 2001 From: "Snow Lee (Sungwon)" Date: Thu, 13 Aug 2026 12:20:36 -0700 Subject: [PATCH] feat(sync): gzip the sync wire, without touching what a hash means (#160) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Nothing on the /store/* wire was compressed, while the corpus it carries is markdown and source. Compression lands as a pure transport concern: content addressing, the storage layout and the journal format all stay over the uncompressed bytes. The two legs are not symmetric. Pull needs no negotiation — net/http already sends Accept-Encoding: gzip and inflates transparently — so devices built before this get it the day the hub ships; a real pre-compression binary receives 19,958 bytes for a 148 KB corpus (7.4x) with no client change. Push is negotiated through sign()'s accept_encoding, because a gzip body posted to an old hub would be stored under the sha256 of its plaintext. The hub inflates ABOVE spool — the sha a key promises, the ops a journal carries and the size that gets billed are all plaintext properties — and the inflate is bounded at 256 MiB, because Content-Encoding severs the one-wire-byte-one-disk-byte relationship that made spool safe unbounded. The presigned direct-to-storage leg stays raw and is asserted to. Known deployment caveat: a compressed push clears ContentLength, so it goes out chunked where every push was sized before. A reverse proxy that buffers or rejects chunked request bodies would fail pushes (degrading to Offline and retrying, not losing data). --- CLAUDE.md | 4 +- architecture/webapp-server.md | 14 +- internal/remote/compress.go | 82 ++++++ internal/remote/compress_test.go | 77 ++++++ internal/remote/http.go | 64 ++++- internal/remote/push_compress_test.go | 187 ++++++++++++++ internal/syncer/compress_e2e_test.go | 189 ++++++++++++++ internal/webapp/compress_e2e_test.go | 222 +++++++++++++++++ internal/webapp/store.go | 118 ++++++++- internal/webapp/store_compress_test.go | 329 +++++++++++++++++++++++++ 10 files changed, 1275 insertions(+), 11 deletions(-) create mode 100644 internal/remote/compress.go create mode 100644 internal/remote/compress_test.go create mode 100644 internal/remote/push_compress_test.go create mode 100644 internal/syncer/compress_e2e_test.go create mode 100644 internal/webapp/compress_e2e_test.go create mode 100644 internal/webapp/store_compress_test.go diff --git a/CLAUDE.md b/CLAUDE.md index c0ca25c..9f8251a 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -40,12 +40,12 @@ Package roles (`internal/`): - **`journal`** — the core data model. Every change is an `Op` (`put`/`delete`) in a per-device append-only JSONL log. `Less` defines the total order `(lamport, time, device, seq)`; `Replay` folds all ops into the volume state, last-writer-wins per path. Everything else is machinery around this. - **`store`** — a volume's local on-disk state: content-addressed blob store (`blobs//`), per-device journal copies, the per-mount materialization cache (`state-.json`, size+mtime fingerprints for cheap change detection), sync state (lamport clock + push cursor), and the exclusive flock that serializes cycles. -- **`remote`** — the `Backend` interface (Put/Get/List/Exists) with `file://`, `s3://`, `gs://`, and `https://` implementations (`https://` syncs through a `bdrive serve` server's `/api/store` API — the client device holds no storage credentials). `PutSigner` is the optional presign capability (S3/GCS). Remote layout: `blobs/` + `journal/.jsonl` under the URL prefix. +- **`remote`** — the `Backend` interface (Put/Get/List/Exists) with `file://`, `s3://`, `gs://`, and `https://` implementations (`https://` syncs through a `bdrive serve` server's `/api/store` API — the client device holds no storage credentials). `PutSigner` is the optional presign capability (S3/GCS). Remote layout: `blobs/` + `journal/.jsonl` under the URL prefix. The `https://` wire is gzipped (`compress.go`) — **transport only**: hashes, storage and the journal format are always over the uncompressed bytes, `Compressible` skips already-compressed content, and the presigned direct-to-storage leg stays raw. Pull needs no negotiation (net/http sends `Accept-Encoding: gzip` and inflates transparently, so `do` must never set that header); push compresses only when `sign()` answered `accept_encoding`, which an older hub does not. - **`syncer`** — the heart: `Session.Cycle()` runs one pass: scan → commit local ops → pull peer journals → preserve conflict copies → materialize merged state → push blobs + own journal. Read the package doc comment in `syncer.go` first. `ignore.go` holds the path filter (`.bdriveignore` rules + the `.bdrive` include list), applied symmetrically in scan and materialize; a newly filtered path is dropped from the cache *without* a delete op so opting out locally never deletes remotely. - **`daemon`** — per-mount background loop (detached process, `daemon.pid`/`daemon.log`/`daemon.lock` in the mount's volume dir). Scans every `--scan-interval` (3s), talks to the remote every `--remote-interval` (10s) or immediately after local edits. Re-reads `.bdrive/config.json` each tick; if it vanishes (folder moved/renamed/deleted) the daemon **exits cleanly without propagating deletes** — the next bdrive command at the new location resumes it (self-heal on next touch). **Liveness is the flock on `daemon.lock`, never the pidfile**: the pidfile outlives its process (it lives in `$BDRIVE_HOME`, which survives reboots), so a recycled pid used to read as a live daemon — making `status` lie and `Start` a silent no-op, which broke the one documented recovery. The kernel drops the flock when the holder dies, including at reboot and on a crash, and holding the lock also makes two daemons on one mount impossible (two writers of one journal). **The pid `stop` signals is written inside the lock file and cleared with it** — `daemon.pid` is display-only, since nothing binds its contents to the lock holder and a `kill -9`'d daemon leaves it behind for a recycled pid to inherit. A mid-run change to the folder's `remote` is not followed: the daemon exits cleanly and the next bdrive command in the folder starts one for whatever the config then says. - **`autostart`** — the login registration that undoes a reboot: one unit per machine, user-level, running `bdrive resume` — macOS `~/Library/LaunchAgents/ai.beardrive.daemon.plist` (`RunAtLoad`, deliberately no `KeepAlive` since the job exits); Linux `$XDG_CONFIG_HOME/systemd/user/beardrive.service` (`Type=oneshot`, no `Restart=`) plus the `default.target.wants` symlink that `systemctl --user enable` would create, since systemd ignores a unit nothing wants. Linux also requires systemd to be booted (`/run/systemd/system`, i.e. `sd_booted`) — otherwise `Install` returns `ErrUnsupported` rather than writing a file nothing would ever read (Alpine/runit, WSL1, slim containers), which starts a daemon for every enrolled, unpaused mount — so mounts added later need no re-registration and `bdrive stop` still means stay stopped. `init` installs it (`--no-autostart` skips). Writing the file is the whole job: no `launchctl` shell-out, so a test or a packaging script can't register a real login item as a side effect, and launchd loads it at the next login anyway. Windows uses a per-user `HKCU\...\CurrentVersion\Run` value (`golang.org/x/sys/windows/registry`) — no admin, no COM (a Startup `.lnk` would need it), no `schtasks`, and it shows up in Task Manager's Startup tab where a user can disable it. Its tests exist but have **never been executed** (written and compile-checked from macOS); they run the first time the suite runs on Windows. **`GOOS=windows go build ./...` does not pass yet** and this package is not the reason: `internal/store`'s `Lock` uses `syscall.Flock` and `internal/daemon` uses `syscall.Kill`/`Setsid`, all unix-only — a Windows port means `LockFileEx` and a kill story for a platform with no SIGTERM, which is a separate change against the sync invariants. `autostart_other.go` (`!darwin && !linux && !windows`) covers the BSDs with `ErrUnsupported`; shared bits (`writeIfDifferent`, `selfPath`) live in the tag-free `autostart.go`. - **`config`** — global state under `$BDRIVE_HOME` (default `~/.bdrive`): device identity (`device.json`), settings (`settings.json`: default server + device token + signed-in account), and the mount registry (`mounts.json`, keyed by **stable mount id**, holding only each mount's last-known path). The per-folder `.bdrive/` directory (`project.go`) holds `config.json` with the mount id + volume/remote/include; **nothing is keyed by the folder path**, so renames/moves are free — `ResolveMount` self-heals the registry path, and the volume store lives at `~/.bdrive/volumes//`. `.bdrive/` is never synced and holds no credentials. -- **`webapp`** — the `bdrive serve` server, in two modes. Single-volume: `Source` is a `DirSource` (plain folder from disk) or `RemoteSource` (folds journals into a file tree with per-file provenance). Hub: `Root` + `Projects` host many projects on one storage root, each under `//` via `remote.Prefixed`; `ProjectDB` (`projects.go`) is a file-backed registry (JSON, loaded at open, rewritten atomically per change) with create-or-join-by-name semantics, name-scoped per organization. Orgs (`orgs.go`, file-backed `orgs.json`) wall projects by membership (email → owner|member): every per-project route — viewer APIs, uploads, history, shares management, the `/store/*` sync proxy — 403s for non-members, `/api/projects` lists only your orgs' projects, owners mint expiring multi-use invite links (`/join/`), and a pre-org hub migrates all projects into a "default" org (all existing accounts join, oldest owns) at startup. `QuotaProvider` (`quota.go`) is the plan-enforcement seam mirroring `AuthProvider` — CheckWrite/RecordUsage on every write path, CheckSeat on invite redemption; OSS ships only `UnlimitedQuota`, managed deployments swap the provider. Renders markdown (goldmark + Obsidian `[[wikilinks]]`). With `--upload` it accepts writes: browser uploads (`upload.go` — direct-to-storage via presigned URLs when the backend implements `remote.PutSigner`, relayed otherwise; ops journaled under the server's own device) and the per-project `/api/p//store/*` proxy (`store.go`) that whole devices sync through — the `https://` remote backend (`remote/http.go`) is its client; journals are never presigned, only immutable blobs. Frontend is a React + TypeScript app (`webapp/frontend/`, Vite + Tailwind v4 + shadcn/ui — Radix-based components copied into `src/components/ui`, themed from the BearDrive tokens in `src/tw.css`; TanStack query/table/virtual, react-hook-form + zod, cmdk, sonner, lucide-react) whose **built output is committed** at `webapp/static/` — the `go:embed static` target — so plain `go build` needs no Node; after any `frontend/src` change run `npm run build` there and commit the new `static/` (`frontend/check-dist.sh` verifies freshness; e2e suite: `npm run e2e` — Playwright against the seeded harness in `e2e_serve_test.go`, port 8993). It learns everything from `/api/config` (+ `/api/projects` in hub mode) and never sees storage info or credentials. It uses native History-API path routing (`//` in hub mode, `/` in volume mode, `/join/` for invites — no `#`, slashes stay literal). **Every user-facing page owns a URL path**: new surfaces are view routes (`//{dashboard|history|install|settings}[/]`, `VIEW_ROUTES` in `router.ts`; renamed segments live in `LEGACY_VIEWS` and are normalized away on arrival) so deep links, reload, and back/forward always work — never URL-less panel state (the org/hub admin panels are the legacy exceptions; don't add more) implemented by the in-repo synchronous router `frontend/src/nav.ts` + `router.ts` (deliberately NOT a router library: react-router v7's startTransition navigation left stale views on screen); `Server.frontend` serves `index.html` as the SPA fallback for any non-asset, non-API/auth/share route so deep links and refreshes resolve (hashed `assets/*` are cached immutable, everything else no-cache), and all client API/asset URLs are root-absolute so a deep path doesn't break relative resolution; the shell also carries `nosniff` + `frame-ancestors 'none'`, and `immutable` is set only after an asset is found (a miss under `assets/` must not pin the shell in a shared cache for a year). Rendered markdown is transformed as a string before mounting and link clicks are delegated on the container — never patch the `dangerouslySetInnerHTML` subtree after commit (React re-applies the markup on unrelated updates and discards DOM patches). **Read heat** (`reads.go`): a `ReadLedger` (hub-only, nil = off, config `reads` block) aggregates read telemetry into daily per-actor buckets, debounced to 10-minute visits, folded into all-time rows past `retention_days` — viewer file/render/download = human (recorded via the project id the `proj()` resolver stashes in the request context), `/s/*` hits = share, device-reported reads (`POST /api/p//reads`) = agent; `/store/*` replication and history `/blob` views are NEVER reads. `GET /api/p//heat?prefix=&days=` returns counts/distinct-readers/last-read only — actor identities (the email/device/token in the buckets) must never appear in an API response, with one stated exception: `?by=device` reports agent **device** ids, which History already shows every project member. That exception holds only because ingest validates the id — `POST /api/p//reads` records a read as an agent actor only when the reported `X-Bdrive-Device` is shaped like a device id at all and no OTHER account has been seen syncing under it (`devices.go:ownsDevice` → `DeviceRegistry.MayActAs`); otherwise the report is accepted and counted for nobody. That route never registers a device: **only `/store/*` traffic does**, and it registers per `(account, id)`, so naming someone else's id claims nothing and cannot lock its real owner out. Emails and share tokens never leave the server in any shape. Recording and flushing degrade silently (log once); telemetry must never fail a request or a sync cycle. The frontend shows heat dots on folder listings and a per-project Dashboard quadrant (reads × staleness, route `//dashboard`) — both visible to every project member, since `/heat` is membership-gated and identity-free. **Hub metadata persistence** (accounts, projects, orgs+invites, shares, devices, read buckets — never blobs or journals) sits behind a pluggable `MetaStore` of typed repos (`db.go`): the service structs (`BuiltinAuth`, `OrgDB`, `ProjectDB`, `ShareDB`, `DeviceRegistry`, `ReadLedger`) keep their in-memory maps + logic and persist each change as one record through a repo (the `ReadRepo` alone is batch-oriented — one flush, one write). Two backends — `db_file.go` (the historical JSON files, still the zero-dep default, reached via the `Open*(path)` constructors) and `db_sql.go` (one `database/sql` impl over pure-Go drivers: `modernc.org/sqlite` locally, `jackc/pgx` for Postgres/Supabase, portable schema + idempotent migrations + transactional multi-row writes). `web.go`'s `database` config (`{driver:file|sqlite|postgres, dsn}`) selects it; file is default and untouched. `db_conformance_test.go` runs the same service ops against every backend. +- **`webapp`** — the `bdrive serve` server, in two modes. Single-volume: `Source` is a `DirSource` (plain folder from disk) or `RemoteSource` (folds journals into a file tree with per-file provenance). Hub: `Root` + `Projects` host many projects on one storage root, each under `//` via `remote.Prefixed`; `ProjectDB` (`projects.go`) is a file-backed registry (JSON, loaded at open, rewritten atomically per change) with create-or-join-by-name semantics, name-scoped per organization. Orgs (`orgs.go`, file-backed `orgs.json`) wall projects by membership (email → owner|member): every per-project route — viewer APIs, uploads, history, shares management, the `/store/*` sync proxy — 403s for non-members, `/api/projects` lists only your orgs' projects, owners mint expiring multi-use invite links (`/join/`), and a pre-org hub migrates all projects into a "default" org (all existing accounts join, oldest owns) at startup. `QuotaProvider` (`quota.go`) is the plan-enforcement seam mirroring `AuthProvider` — CheckWrite/RecordUsage on every write path, CheckSeat on invite redemption; OSS ships only `UnlimitedQuota`, managed deployments swap the provider. Renders markdown (goldmark + Obsidian `[[wikilinks]]`). With `--upload` it accepts writes: browser uploads (`upload.go` — direct-to-storage via presigned URLs when the backend implements `remote.PutSigner`, relayed otherwise; ops journaled under the server's own device) and the per-project `/api/p//store/*` proxy (`store.go`) that whole devices sync through — the `https://` remote backend (`remote/http.go`) is its client; journals are never presigned, only immutable blobs. That proxy gzips what it answers and inflates what it is sent — the inflate sits ABOVE `spool` (the sha, the op count and the billed size are all properties of the plaintext) and is bounded by `maxInflatedPut`, since a compressed body severs the one-wire-byte-one-disk-byte relationship that made `spool` safe unbounded; `RecordUsage` charges uncompressed, `RecordEgress` compressed. Frontend is a React + TypeScript app (`webapp/frontend/`, Vite + Tailwind v4 + shadcn/ui — Radix-based components copied into `src/components/ui`, themed from the BearDrive tokens in `src/tw.css`; TanStack query/table/virtual, react-hook-form + zod, cmdk, sonner, lucide-react) whose **built output is committed** at `webapp/static/` — the `go:embed static` target — so plain `go build` needs no Node; after any `frontend/src` change run `npm run build` there and commit the new `static/` (`frontend/check-dist.sh` verifies freshness; e2e suite: `npm run e2e` — Playwright against the seeded harness in `e2e_serve_test.go`, port 8993). It learns everything from `/api/config` (+ `/api/projects` in hub mode) and never sees storage info or credentials. It uses native History-API path routing (`//` in hub mode, `/` in volume mode, `/join/` for invites — no `#`, slashes stay literal). **Every user-facing page owns a URL path**: new surfaces are view routes (`//{dashboard|history|install|settings}[/]`, `VIEW_ROUTES` in `router.ts`; renamed segments live in `LEGACY_VIEWS` and are normalized away on arrival) so deep links, reload, and back/forward always work — never URL-less panel state (the org/hub admin panels are the legacy exceptions; don't add more) implemented by the in-repo synchronous router `frontend/src/nav.ts` + `router.ts` (deliberately NOT a router library: react-router v7's startTransition navigation left stale views on screen); `Server.frontend` serves `index.html` as the SPA fallback for any non-asset, non-API/auth/share route so deep links and refreshes resolve (hashed `assets/*` are cached immutable, everything else no-cache), and all client API/asset URLs are root-absolute so a deep path doesn't break relative resolution; the shell also carries `nosniff` + `frame-ancestors 'none'`, and `immutable` is set only after an asset is found (a miss under `assets/` must not pin the shell in a shared cache for a year). Rendered markdown is transformed as a string before mounting and link clicks are delegated on the container — never patch the `dangerouslySetInnerHTML` subtree after commit (React re-applies the markup on unrelated updates and discards DOM patches). **Read heat** (`reads.go`): a `ReadLedger` (hub-only, nil = off, config `reads` block) aggregates read telemetry into daily per-actor buckets, debounced to 10-minute visits, folded into all-time rows past `retention_days` — viewer file/render/download = human (recorded via the project id the `proj()` resolver stashes in the request context), `/s/*` hits = share, device-reported reads (`POST /api/p//reads`) = agent; `/store/*` replication and history `/blob` views are NEVER reads. `GET /api/p//heat?prefix=&days=` returns counts/distinct-readers/last-read only — actor identities (the email/device/token in the buckets) must never appear in an API response, with one stated exception: `?by=device` reports agent **device** ids, which History already shows every project member. That exception holds only because ingest validates the id — `POST /api/p//reads` records a read as an agent actor only when the reported `X-Bdrive-Device` is shaped like a device id at all and no OTHER account has been seen syncing under it (`devices.go:ownsDevice` → `DeviceRegistry.MayActAs`); otherwise the report is accepted and counted for nobody. That route never registers a device: **only `/store/*` traffic does**, and it registers per `(account, id)`, so naming someone else's id claims nothing and cannot lock its real owner out. Emails and share tokens never leave the server in any shape. Recording and flushing degrade silently (log once); telemetry must never fail a request or a sync cycle. The frontend shows heat dots on folder listings and a per-project Dashboard quadrant (reads × staleness, route `//dashboard`) — both visible to every project member, since `/heat` is membership-gated and identity-free. **Hub metadata persistence** (accounts, projects, orgs+invites, shares, devices, read buckets — never blobs or journals) sits behind a pluggable `MetaStore` of typed repos (`db.go`): the service structs (`BuiltinAuth`, `OrgDB`, `ProjectDB`, `ShareDB`, `DeviceRegistry`, `ReadLedger`) keep their in-memory maps + logic and persist each change as one record through a repo (the `ReadRepo` alone is batch-oriented — one flush, one write). Two backends — `db_file.go` (the historical JSON files, still the zero-dep default, reached via the `Open*(path)` constructors) and `db_sql.go` (one `database/sql` impl over pure-Go drivers: `modernc.org/sqlite` locally, `jackc/pgx` for Postgres/Supabase, portable schema + idempotent migrations + transactional multi-row writes). `web.go`'s `database` config (`{driver:file|sqlite|postgres, dsn}`) selects it; file is default and untouched. `db_conformance_test.go` runs the same service ops against every backend. `cmd/bdrive/` is a thin cobra CLI over these packages (`login`, `logout`, `init`, `stop`, `sync`, `scope`, `forget`, `status`, `log`, `share`, `export`, `import`, `url`, `hooks`, `resume`, `autostart`, `read-log`, `serve`, `whoami`, `daemon`, `version` — `mnt`/`umnt`/`remote` are gone; `init` is the front door and `stop` pauses). `export`/`import` (`migrate.go`) move a whole project between hubs with full fidelity: the archive is the remote store layout (all devices' journals + all blobs) in a tar.gz, streamed through the existing `remote.Backend` — no server-side support needed, so it works against any hub in either direction (the anti-lock-in story for cloud-hesitant users). `bdrive login` signs the device in (bare form uses the remembered server or `config.DefaultServer` = beardrive.ai; loopback-callback browser flow in `login.go`, `--device` for headless) and stores server+token+account in `settings.json`; `bdrive logout` revokes this device's token on the hub (`DELETE /api/auth/token`, authenticated by the token itself) and then clears the saved token+account, keeping the remembered server unless `--forget`; a revocation it could not reach the hub for is reported, never swallowed. Switching hubs is `bdrive login ` then re-`init` — `init` is the only thing that writes a folder's remote (always a hub, `server + "/p/" + id`); there is no client command to point a folder at a raw bucket. `bdrive init` is interactive on a TTY (survey menus: create-new vs connect-existing with a project list; whole-folder vs only-some-subfolders) with full flag bypass (`--name/--project/--only/--yes`) and never prompts without a TTY; it runs the login flow first when there is no session, writes `.bdrive/config.json`, seeds `.bdriveignore`, registers agent sync hooks in each platform's USER config (`~/.claude/settings.json` and friends — once per machine, never inside a project: platforms read hook config only from the directory a session starts in, so a per-project file covers only sessions that start there and, living in a mount, would sync to the team; `Install` also migrates away hooks older versions wrote into projects), and starts sync via `startSync`; re-running it resumes — including after a folder move. **A mount is always exactly the folder named** — there is no re-rooting flag. Syncing only part of a mount is `--only wiki,docs`, which writes a bdrive-managed block of `.bdriveignore` negation rules (`cmd/bdrive/scopefile.go`; `bdrive scope add/rm` edits the same block) rather than a second scope mechanism: the old `Include` list in `config.json` is legacy — still honored, never written. Because the rules live in the synced `.bdriveignore`, scope is team-wide, which is why `sync --prune` refuses when `!` rules are present (it would strip everything outside the scope from the hub for everyone; `bdrive forget ` is the per-path tool). `init` also refuses a second folder for a project this device already syncs — one device writes one journal per project, so two mounts would overwrite each other's ops. `bdrive serve -c config.json` configures the server from a file, explicit flags winning. diff --git a/architecture/webapp-server.md b/architecture/webapp-server.md index 3a0240b..216f840 100644 --- a/architecture/webapp-server.md +++ b/architecture/webapp-server.md @@ -125,7 +125,7 @@ classDiagram opsNameTheirAuthor(ops) whose name journalKeepsItsOps(ctx, be, key, ops) } - note for journalDoor "store.go — the invariant "each device writes only its own journal" is now ENFORCED here, not assumed. The key must be journal/<canonical device id>.jsonl for the device in the request header, that device must already be owned by the caller (DeviceRegistry.OwnerOf) or the caller must be a project admin (the recovery arm) — the old first-writer-claims arm is gone. Every op must pass journal.SafePath + config.ReservedPath on its Path and journal.SafeText on Note/Author/UserName, must name its own owner's account, and the upload must keep every Seq the stored journal already had: append-only, 409 on truncation. Bodies are spooled first, and a blob PUT must hash to the key it claims" + note for journalDoor "store.go — the invariant "each device writes only its own journal" is now ENFORCED here, not assumed. The key must be journal/<canonical device id>.jsonl for the device in the request header, that device must already be owned by the caller (DeviceRegistry.OwnerOf) or the caller must be a project admin (the recovery arm) — the old first-writer-claims arm is gone. Every op must pass journal.SafePath + config.ReservedPath on its Path and journal.SafeText on Note/Author/UserName, must name its own owner's account, and the upload must keep every Seq the stored journal already had: append-only, 409 on truncation. Bodies are spooled first, and a blob PUT must hash to the key it claims. A Content-Encoding: gzip body is inflated ABOVE the spool — the sha, the op count and the billed size are all properties of the plaintext, so nothing below that line knows compression happened — and the inflate is bounded (maxInflatedPut, 256 MiB, only when an encoding was declared), because compression severs the one-wire-byte-one-disk-byte relationship that made spool safe unbounded" note for journalDoor "Delta sync grew the key space: validStoreKey also accepts chunks/<sha256> (content-addressed, PUT must hash to its key, presigned like blobs incl. refuse-existing) and manifests/<sha256> (keyed by the whole FILE's sha — not its own content hash — so it is never presigned and gets two ingest gates instead: every chunk it names must already EXIST in the store, and the key is WRITE-ONCE — an identical re-put is a 200 no-op so an interrupted push can retry, a different body 409s. Together these make "a manifest exists ⟹ its chunks exist" an invariant every consumer can lean on: the client's push skip-proof, reassemble, and bdrive import)" class Backend { @@ -139,6 +139,14 @@ classDiagram note for Backend "internal/remote — impls: localBackend (file://), s3Backend, gcsBackend, httpBackend (https:// hub), Prefixed wrapper" note for Backend "Key handling is fallible now: Prefixed.key and localBackend.path RETURN AN ERROR (safeKey / store.UnderRoot) rather than concatenating, so a `..` key cannot walk out of a project's prefix or out of a file:// root — and Prefixed.List re-checks the STRIPPED key on the way out, since the prefix it removes is the only thing that was ever validated. The httpBackend client is origin-bound: the device token is keyed to settings.Server, SameOrigin is the one rule, refuseOffOriginRedirect is its CheckRedirect, a presign target must be https on a trusted origin (directTargetOK), and List drops keys failing journal.SafePath and clamps a negative Size. gcs SignPut now signs Content-Length too. Object carries Modified (S3 LastModified, GCS Updated, file mtime; zero where the backend has none) — RemoteSource.verify reads it to decide when a blob can no longer be rewritten by a presigned URL" + class wireCodec { + <> + +Compressible(r) (rejoined, worth, err) + +AcceptsGzip(req) bool + putPlan.AcceptEncoding []string + } + note for wireCodec "Transport compression, and TRANSPORT ONLY: content addressing, the storage layout and the journal format are all over the UNCOMPRESSED bytes. One probe helper serves both legs — it gzips the first 64 KiB and keeps the compressed form only if the sample shrank, so already-compressed content (JPEG, zip, weights) passes through untouched; the reader it returns is the stream REJOINED, since a probe that eats bytes corrupts every transfer. The two legs are not symmetric. PULL needs no negotiation: net/http sends Accept-Encoding: gzip itself and inflates transparently, so devices built before this get it free — which is why httpBackend.do must never set that header. PUSH cannot be unilateral, because an old hub would store the gzip bytes under the sha256 of the plaintext, so the client compresses only when sign() advertised accept_encoding (absent on an old hub → raw). putDirect stays raw: a presigned upload has no hub in the path to inflate it" + class AuthProvider { <> +CLILoginPath() @@ -360,7 +368,7 @@ classDiagram +Write(p) n +n int64 } - note for countingWriter "Bills what actually reached the client. FileInfo.Size and the journal's Size are claims made BEFORE the write; a reader who abandons a download halfway must not be charged for the whole file" + note for countingWriter "Bills what actually reached the client. FileInfo.Size and the journal's Size are claims made BEFORE the write; a reader who abandons a download halfway must not be charged for the whole file. It wraps the SOCKET and gzip writes INTO it, never the reverse — RecordEgress is a bandwidth meter, so a compressed response must report its compressed size, and getting that backwards bills plaintext while nothing fails" class grant { +project +org +key @@ -399,6 +407,8 @@ classDiagram Server ..> secretScan : handleShareCreate scans the first 1 MiB unless confirmed or alreadyPublic secretScan ..> secretFinding Server ..> countingWriter : every bytes-out route that bills + Server ..> wireCodec : gzip on /store/ GET+list, inflate above spool on PUT + Backend ..> wireCodec : httpBackend gzips a relayed PUT when sign() allows reservations ..> Backend : reconcile — did the blob land Server *-- journalDoor : /store/* is the only way a device writes journalDoor ..> DeviceRegistry : OwnerOf gates the journal key diff --git a/internal/remote/compress.go b/internal/remote/compress.go new file mode 100644 index 0000000..fbd965d --- /dev/null +++ b/internal/remote/compress.go @@ -0,0 +1,82 @@ +package remote + +import ( + "bytes" + "compress/gzip" + "io" + "net/http" + "strings" +) + +// Transport compression for the sync wire. Nothing on it was compressed: +// blobs and journals crossed as raw application/octet-stream in both +// directions, while the corpus they carry is 5–10 KB markdown and source +// files that gzip ~3.4x (TestCompressionTextCorpusRatio measures it). +// +// This is a TRANSPORT concern only, and the whole feature depends on that +// staying true: content addressing is over the UNCOMPRESSED bytes, storage +// holds uncompressed objects, and the journal format is untouched. A hub that +// stored a gzip body under the sha256 of its plaintext would have broken +// every device's blob check, so the hub inflates before it hashes +// (handleStorePut) and never after. +// +// The codec is gzip rather than zstd for one decisive reason: net/http already +// sends `Accept-Encoding: gzip` on every request whose caller did not set that +// header itself, and transparently inflates the response. httpBackend.do does +// not set it, so the entire pull leg compresses for binaries built before this +// existed, with no client change at all. zstd would forfeit that. + +// probeWindow is how much of a stream is sampled to decide whether the rest is +// worth compressing. Big enough to see past a file header, small enough that +// the sample is a buffer rather than a spool. +const probeWindow = 64 << 10 + +// probeMargin is the share of the sample compression has to save before it is +// worth paying for. Already-compressed content (JPEG, zip, model weights) gets +// ~0.1% BIGGER under gzip, so the margin is really a sign test with slack. +const probeMargin = 0.9 + +// Compressible reports whether a stream is worth gzipping, by compressing its +// first probeWindow bytes and checking that the sample actually shrank. +// +// It returns the stream REJOINED — the sampled bytes followed by whatever is +// left — because the probe has to consume the bytes it judges. That identity is +// the property this helper lives or dies on: a probe that eats bytes silently +// corrupts every push and every pull that runs through it, and the failure +// shows up as a sha mismatch far from here. compress_test.go asserts it for a +// stream longer than the window, one shorter, and an empty one. +func Compressible(r io.Reader) (io.Reader, bool, error) { + sample, err := io.ReadAll(io.LimitReader(r, probeWindow)) + rejoined := io.MultiReader(bytes.NewReader(sample), r) + if err != nil { + return rejoined, false, err + } + if len(sample) == 0 { + return rejoined, false, nil + } + var n countingSink + gz := gzip.NewWriter(&n) + if _, err := gz.Write(sample); err != nil { + return rejoined, false, err + } + if err := gz.Close(); err != nil { + return rejoined, false, err + } + return rejoined, float64(n) < float64(len(sample))*probeMargin, nil +} + +type countingSink int + +func (c *countingSink) Write(p []byte) (int, error) { *c += countingSink(len(p)); return len(p), nil } + +// AcceptsGzip reports whether a request's Accept-Encoding allows a gzipped +// answer. Go's own transport sets that header on every request this package +// makes, which is why old devices get the compressed pull leg for free. +func AcceptsGzip(r *http.Request) bool { + for _, enc := range strings.Split(r.Header.Get("Accept-Encoding"), ",") { + if name, _, _ := strings.Cut(enc, ";"); strings.EqualFold(strings.TrimSpace(name), "gzip") { + return true + } + } + return false +} diff --git a/internal/remote/compress_test.go b/internal/remote/compress_test.go new file mode 100644 index 0000000..fbbc5f0 --- /dev/null +++ b/internal/remote/compress_test.go @@ -0,0 +1,77 @@ +package remote + +import ( + "bytes" + "io" + "math/rand" + "net/http" + "net/http/httptest" + "strings" + "testing" +) + +// The probe consumes the bytes it judges, so the stream it hands back must be +// byte-identical to the one it was given — for a stream longer than the probe +// window, one shorter, and an empty one. A probe that eats bytes corrupts every +// push and pull that runs through it, and the damage surfaces as a sha +// mismatch nowhere near this file. +func TestCompressibleRejoinsTheStream(t *testing.T) { + cases := []struct { + name string + in []byte + want bool + }{ + {"text past the window", []byte(strings.Repeat("package main // hello hello\n", 5000)), true}, + {"text under the window", []byte(strings.Repeat("hello beardrive\n", 100)), true}, + {"already compressed", randomBytes(200 << 10), false}, + {"tiny", []byte("hi"), false}, + {"empty", nil, false}, + } + for _, c := range cases { + t.Run(c.name, func(t *testing.T) { + got, worth, err := Compressible(bytes.NewReader(c.in)) + if err != nil { + t.Fatal(err) + } + if worth != c.want { + t.Errorf("worth = %v, want %v", worth, c.want) + } + rejoined, err := io.ReadAll(got) + if err != nil { + t.Fatal(err) + } + if !bytes.Equal(rejoined, c.in) { + t.Fatalf("rejoined stream is %d bytes, want the original %d", len(rejoined), len(c.in)) + } + }) + } +} + +func TestAcceptsGzip(t *testing.T) { + cases := map[string]bool{ + "": false, + "identity": false, + "gzip": true, + "deflate, gzip;q=1.0, *;q=0": true, + "GZIP": true, + "x-gzip": false, // a different token, not a prefix match + } + for hdr, want := range cases { + r := httptest.NewRequest(http.MethodGet, "/", nil) + if hdr != "" { + r.Header.Set("Accept-Encoding", hdr) + } + if got := AcceptsGzip(r); got != want { + t.Errorf("AcceptsGzip(%q) = %v, want %v", hdr, got, want) + } + } +} + +// randomBytes stands in for already-compressed content (JPEG, zip, model +// weights): incompressible by construction, which is the whole point. +func randomBytes(n int) []byte { + b := make([]byte, n) + rng := rand.New(rand.NewSource(1)) + rng.Read(b) + return b +} diff --git a/internal/remote/http.go b/internal/remote/http.go index 6206f28..437d280 100644 --- a/internal/remote/http.go +++ b/internal/remote/http.go @@ -2,6 +2,7 @@ package remote import ( "bytes" + "compress/gzip" "context" "encoding/json" "fmt" @@ -165,6 +166,13 @@ func refuseOffOriginRedirect(req *http.Request, via []*http.Request) error { // do sends the request with this device's credential attached, plus the // identity headers the server's device registry records for history (name, // OS; the server observes the IP itself). +// +// It deliberately does NOT set Accept-Encoding. net/http adds `gzip` itself +// whenever the caller has not, and transparently inflates the response — which +// is the entire pull half of transport compression, free and backward +// compatible. Setting the header here turns that off silently: the hub would +// still answer `Content-Encoding: gzip`, nothing would inflate it, and every +// blob would fail its sha check while looking like a corrupt hub. func (b *httpBackend) do(req *http.Request) (*http.Response, error) { if b.token != "" { req.Header.Set("Authorization", "Bearer "+b.token) @@ -354,7 +362,7 @@ func (b *httpBackend) Put(ctx context.Context, key string, r io.Reader, size int // No usable destination: relay through the hub, which already holds // this device's credential and is the party it chose to trust. } - return b.putViaServer(ctx, key, r, size) + return b.putViaServer(ctx, plan, key, r, size) } // directTargetOK decides whether this device will hand a file's bytes to the @@ -389,6 +397,23 @@ type putPlan struct { URL string `json:"url"` Method string `json:"method"` Headers map[string]string `json:"headers"` + // AcceptEncoding is what the hub will accept on the relayed PUT body. + // Push cannot be unilateral the way pull is: a gzipped body posted to a + // hub that does not inflate is stored verbatim under the sha256 of its + // PLAINTEXT — a 400 for a blob, and a silently mis-stored journal. So the + // client compresses only when the hub says so, and an older hub says + // nothing at all (absent field → nil → raw). sign() runs before every + // single put, so this costs no extra round trip and needs no config flag. + AcceptEncoding []string `json:"accept_encoding"` +} + +func (p putPlan) acceptsGzip() bool { + for _, enc := range p.AcceptEncoding { + if strings.EqualFold(strings.TrimSpace(enc), "gzip") { + return true + } + } + return false } func (b *httpBackend) sign(ctx context.Context, key string, size int64) (putPlan, error) { @@ -444,7 +469,34 @@ func (b *httpBackend) putDirect(ctx context.Context, plan putPlan, r io.Reader, return nil } -func (b *httpBackend) putViaServer(ctx context.Context, key string, r io.Reader, size int64) error { +// putViaServer relays the bytes through the hub, gzipping them when the hub +// advertised that it inflates (plan.AcceptEncoding) and the content is worth +// compressing. putDirect deliberately stays raw: a presigned upload lands in +// the object store under the sha256 of the plaintext, with no hub in the path +// to inflate it, so compressing that leg would corrupt content addressing at +// rest. +func (b *httpBackend) putViaServer(ctx context.Context, plan putPlan, key string, r io.Reader, size int64) error { + gzipped := false + if plan.acceptsGzip() { + probed, worth, err := Compressible(r) + if err != nil { + return err + } + r, gzipped = probed, worth + } + if gzipped { + pr, pw := io.Pipe() + src := r + go func() { + gz := gzip.NewWriter(pw) + _, err := io.Copy(gz, src) + if cerr := gz.Close(); err == nil { + err = cerr + } + pw.CloseWithError(err) + }() + r = pr + } req, err := http.NewRequestWithContext(ctx, http.MethodPut, b.endpoint("object", url.Values{"key": {key}}), r) if err != nil { @@ -452,6 +504,14 @@ func (b *httpBackend) putViaServer(ctx context.Context, key string, r io.Reader, } nameJournalDevice(req, key) req.ContentLength = size + if gzipped { + req.Header.Set("Content-Encoding", "gzip") + // The compressed length is not knowable without compressing twice, so + // the request goes out chunked. The hub's spool() already treats a -1 + // length as the normal case — it is why it measures the body instead of + // believing a header. + req.ContentLength = -1 + } resp, err := b.do(req) if err != nil { return err diff --git a/internal/remote/push_compress_test.go b/internal/remote/push_compress_test.go new file mode 100644 index 0000000..d7d19cf --- /dev/null +++ b/internal/remote/push_compress_test.go @@ -0,0 +1,187 @@ +package remote + +import ( + "bytes" + "compress/gzip" + "context" + "io" + "net/http" + "net/http/httptest" + "strings" + "testing" +) + +// The mixed-fleet contract for the push leg, both directions, at the one place +// it is decided: what sign() advertised. +// +// New client ↔ OLD hub is the dangerous half. An old hub does not inflate, so a +// gzipped body would be stored verbatim under the sha256 of its plaintext — +// rejected outright for a blob, silently mis-stored for a journal. The old hub +// says nothing about encodings, so the client must send raw. +func TestPushCompressesOnlyWhenTheHubAdvertisesIt(t *testing.T) { + payload := strings.Repeat("# notes\nthe corpus is markdown and source, which gzips well\n", 500) + + for _, tc := range []struct { + name string + signAnswer string + wantGzip bool + }{ + {"old hub says nothing", `{"mode":"server"}`, false}, + {"old hub with an empty list", `{"mode":"server","accept_encoding":[]}`, false}, + {"hub speaks another codec", `{"mode":"server","accept_encoding":["zstd"]}`, false}, + {"new hub", `{"mode":"server","accept_encoding":["gzip"]}`, true}, + } { + t.Run(tc.name, func(t *testing.T) { + var gotBody []byte + var gotEncoding string + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/store/sign") { + w.Header().Set("Content-Type", "application/json") + w.Write([]byte(tc.signAnswer)) + return + } + gotEncoding = r.Header.Get("Content-Encoding") + gotBody, _ = io.ReadAll(r.Body) + w.Write([]byte(`{"ok":true}`)) + })) + defer ts.Close() + + be, err := Open(context.Background(), ts.URL+"/p/p-0123abcd") + if err != nil { + t.Fatal(err) + } + defer be.Close() + if err := be.Put(context.Background(), "blobs/"+strings.Repeat("a", 64), + strings.NewReader(payload), int64(len(payload))); err != nil { + t.Fatal(err) + } + + if !tc.wantGzip { + if gotEncoding != "" { + t.Fatalf("Content-Encoding = %q, want none", gotEncoding) + } + if string(gotBody) != payload { + t.Fatal("body is not the plaintext it was handed") + } + return + } + if gotEncoding != "gzip" { + t.Fatalf("Content-Encoding = %q, want gzip", gotEncoding) + } + if len(gotBody) >= len(payload) { + t.Fatalf("compressed body is %d bytes, larger than the %d it started as", len(gotBody), len(payload)) + } + gz, err := gzip.NewReader(bytes.NewReader(gotBody)) + if err != nil { + t.Fatal(err) + } + plain, err := io.ReadAll(gz) + if err != nil { + t.Fatal(err) + } + // Content addressing is over the UNCOMPRESSED bytes: what the hub + // inflates has to be exactly what the key names. + if string(plain) != payload { + t.Fatal("the body does not inflate to what was pushed") + } + }) + } +} + +// Incompressible content must cross untouched even against a hub that offers +// gzip — chunked large files are mostly already-compressed binary, where gzip +// pays CPU to make the payload ~0.1% bigger. +func TestPushSkipsIncompressibleContent(t *testing.T) { + payload := randomBytes(300 << 10) + var gotEncoding string + var gotLen int + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if strings.HasSuffix(r.URL.Path, "/store/sign") { + w.Header().Set("Content-Type", "application/json") + w.Write([]byte(`{"mode":"server","accept_encoding":["gzip"]}`)) + return + } + gotEncoding = r.Header.Get("Content-Encoding") + body, _ := io.ReadAll(r.Body) + gotLen = len(body) + w.Write([]byte(`{"ok":true}`)) + })) + defer ts.Close() + + be, err := Open(context.Background(), ts.URL+"/p/p-0123abcd") + if err != nil { + t.Fatal(err) + } + defer be.Close() + if err := be.Put(context.Background(), "blobs/"+strings.Repeat("b", 64), + bytes.NewReader(payload), int64(len(payload))); err != nil { + t.Fatal(err) + } + if gotEncoding != "" { + t.Fatalf("Content-Encoding = %q on incompressible content, want none", gotEncoding) + } + if gotLen != len(payload) { + t.Fatalf("wire carried %d bytes for a %d-byte payload", gotLen, len(payload)) + } +} + +// The presigned leg must stay raw even when the hub advertises gzip on the +// same sign() response. A direct upload goes to the object store under the +// sha256 of the PLAINTEXT with no hub in the path to inflate it, so a stray +// compression here would corrupt content addressing at rest — the one failure +// in this change that storage would keep forever rather than reject. +// +// Not reachable through any hub fixture in the tree: they all run on file:// +// storage, which implements no PutSigner, so every plan is mode:"server". +// Managed hubs are S3/GCS-backed, which makes this the production path. +func TestPresignedUploadIsNeverCompressed(t *testing.T) { + payload := strings.Repeat("# highly compressible markdown\n", 2000) + + var gotBody []byte + var gotEncoding string + var relayed bool + // One origin serving both roles: directTargetOK refuses a presign target + // that is neither https nor the hub's own origin, so a second httptest + // server would be declined and silently relayed instead — which would pass + // this test while proving nothing about putDirect. + var hub *httptest.Server + hub = httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + switch { + case strings.HasSuffix(r.URL.Path, "/store/sign"): + w.Header().Set("Content-Type", "application/json") + // The hub advertises gzip — it does on every sign answer — AND + // hands back a presigned destination. The client must honor the + // second and ignore the first. + w.Write([]byte(`{"mode":"direct","exists":false,"accept_encoding":["gzip"],"url":"` + + hub.URL + `/presigned-blob","method":"PUT"}`)) + case r.URL.Path == "/presigned-blob": + gotEncoding = r.Header.Get("Content-Encoding") + gotBody, _ = io.ReadAll(r.Body) + w.WriteHeader(http.StatusOK) + default: + relayed = true + http.Error(w, "relayed through the hub", http.StatusNotFound) + } + })) + defer hub.Close() + + be, err := Open(context.Background(), hub.URL+"/p/p-0123abcd") + if err != nil { + t.Fatal(err) + } + defer be.Close() + if err := be.Put(context.Background(), "blobs/"+strings.Repeat("c", 64), + strings.NewReader(payload), int64(len(payload))); err != nil { + t.Fatal(err) + } + if relayed { + t.Fatal("the upload was relayed through the hub; putDirect was never exercised") + } + if gotEncoding != "" { + t.Fatalf("presigned upload carried Content-Encoding %q — the object store has no hub to inflate it", gotEncoding) + } + if string(gotBody) != payload { + t.Fatalf("presigned upload sent %d bytes, want the %d-byte plaintext the key is the hash of", + len(gotBody), len(payload)) + } +} diff --git a/internal/syncer/compress_e2e_test.go b/internal/syncer/compress_e2e_test.go new file mode 100644 index 0000000..b774e6b --- /dev/null +++ b/internal/syncer/compress_e2e_test.go @@ -0,0 +1,189 @@ +package syncer + +import ( + "context" + "fmt" + "net/http" + "net/http/httptest" + "path/filepath" + "sync/atomic" + "testing" + + "github.com/runbear-io/beardrive/internal/remote" + "github.com/runbear-io/beardrive/internal/webapp" +) + +// The measurement that justifies the feature, taken where it is real: HTTP +// bodies in and out of the hub. +// +// It has to be measured HERE and not one layer up. A counter wrapping a +// remote.Backend sits ABOVE httpBackend and can only ever see the plaintext it +// hands down — it cannot observe transport encoding at all, and the file:// +// backend those counters usually wrap is never compressed by this change. +// countingHub proxies the hub's own handler, so what it counts is the wire. + +// countingHub is a hub whose HTTP bodies are counted in both directions. +func countingHub(t *testing.T, storage remote.Backend) (*httptest.Server, webapp.Project, *wireCount) { + t.Helper() + db, err := webapp.OpenProjectDB(filepath.Join(t.TempDir(), "projects.json")) + if err != nil { + t.Fatal(err) + } + p, _, err := db.GetOrCreate("vol", "") + if err != nil { + t.Fatal(err) + } + srv := &webapp.Server{ + Root: storage, Projects: db, Refresh: 0, + Device: webapp.Identity{ID: "hubdev", Name: "hub", Author: "hub@test"}, + Upload: webapp.UploadConfig{Enabled: true}, + } + n := &wireCount{} + h := srv.Handler() + ts := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + r.Body = &countingReadCloser{rc: r.Body, n: &n.up} + h.ServeHTTP(&countingRW{ResponseWriter: w, n: &n.down}, r) + })) + t.Cleanup(ts.Close) + return ts, p, n +} + +type wireCount struct{ up, down atomic.Int64 } + +func (w *wireCount) reset() { w.up.Store(0); w.down.Store(0) } + +type countingReadCloser struct { + rc interface { + Read([]byte) (int, error) + Close() error + } + n *atomic.Int64 +} + +func (c *countingReadCloser) Read(p []byte) (int, error) { + n, err := c.rc.Read(p) + c.n.Add(int64(n)) + return n, err +} +func (c *countingReadCloser) Close() error { return c.rc.Close() } + +type countingRW struct { + http.ResponseWriter + n *atomic.Int64 +} + +func (c *countingRW) Write(p []byte) (int, error) { + n, err := c.ResponseWriter.Write(p) + c.n.Add(int64(n)) + return n, err +} + +// A text-heavy project must cross the wire at ≥2.5x reduction in BOTH +// directions, and the same project must still converge byte for byte. +func TestCompressionWireRatio(t *testing.T) { + storage := sharedRemote(t) + ts, p, wire := countingHub(t, storage) + viaServer, err := remote.Open(context.Background(), ts.URL+"/p/"+p.ID) + if err != nil { + t.Fatal(err) + } + defer viaServer.Close() + + a := newDevice(t, "deva", viaServer) + raw := 0 + for i := 0; i < 20; i++ { + body := fmt.Sprintf("# note %d\n\n%s", i, textish(i)) + write(t, a.Folder, fmt.Sprintf("notes/note-%02d.md", i), body) + raw += len(body) + } + wire.reset() + cycle(t, a) + up := wire.up.Load() + t.Logf("push: %d bytes on the wire for a %d-byte corpus (%.2fx)", up, raw, float64(raw)/float64(up)) + if ratio := float64(raw) / float64(up); ratio < 2.5 { + t.Fatalf("push carried %d bytes for a %d-byte corpus (%.2fx), want at least 2.5x", up, raw, ratio) + } + + // A second device pulls the same corpus back down through the same hub. + b, err := remote.Open(context.Background(), ts.URL+"/p/"+p.ID) + if err != nil { + t.Fatal(err) + } + defer b.Close() + dev := newDevice(t, "devb", b) + wire.reset() + if res := cycle(t, dev); res.PulledOps != 20 { + t.Fatalf("pulled %d ops, want 20", res.PulledOps) + } + down := wire.down.Load() + t.Logf("pull: %d bytes on the wire for a %d-byte corpus (%.2fx)", down, raw, float64(raw)/float64(down)) + if ratio := float64(raw) / float64(down); ratio < 2.5 { + t.Fatalf("pull carried %d bytes for a %d-byte corpus (%.2fx), want at least 2.5x", down, raw, ratio) + } + // Compression is a transport concern: the bytes on disk are the bytes that + // were written, or none of the above matters. + for i := 0; i < 20; i++ { + want := fmt.Sprintf("# note %d\n\n%s", i, textish(i)) + if got := read(t, dev.Folder, fmt.Sprintf("notes/note-%02d.md", i)); got != want { + t.Fatalf("note %d did not converge", i) + } + } +} + +// Already-compressed content must cross within ~1% of its raw size, in both +// directions: the skip path exists so gzip does not pay CPU to grow a JPEG. +func TestCompressionSkipsIncompressiblePayload(t *testing.T) { + storage := sharedRemote(t) + ts, p, wire := countingHub(t, storage) + viaServer, err := remote.Open(context.Background(), ts.URL+"/p/"+p.ID) + if err != nil { + t.Fatal(err) + } + defer viaServer.Close() + + a := newDevice(t, "deva", viaServer) + payload := string(pseudoJPEG(512 << 10)) + write(t, a.Folder, "photo.jpg", payload) + wire.reset() + cycle(t, a) + if up := wire.up.Load(); float64(up) > float64(len(payload))*1.01+4096 { + t.Fatalf("push carried %d bytes for a %d-byte incompressible payload", up, len(payload)) + } + + b, err := remote.Open(context.Background(), ts.URL+"/p/"+p.ID) + if err != nil { + t.Fatal(err) + } + defer b.Close() + dev := newDevice(t, "devb", b) + wire.reset() + cycle(t, dev) + if down := wire.down.Load(); float64(down) > float64(len(payload))*1.01+4096 { + t.Fatalf("pull carried %d bytes for a %d-byte incompressible payload", down, len(payload)) + } + if read(t, dev.Folder, "photo.jpg") != payload { + t.Fatal("the incompressible payload did not converge") + } +} + +// textish is a stand-in for the real corpus: markdown and source, 5–10 KB. +func textish(seed int) string { + var b []byte + for i := 0; i < 120; i++ { + b = append(b, fmt.Sprintf("- item %d/%d: the sync wire carries markdown and source files\n", seed, i)...) + } + return string(b) +} + +// pseudoJPEG stands in for content that is already compressed. +func pseudoJPEG(n int) []byte { + b := make([]byte, n) + x := uint32(2463534242) + for i := range b { + x ^= x << 13 + x ^= x >> 17 + x ^= x << 5 + b[i] = byte(x) + } + return b +} diff --git a/internal/webapp/compress_e2e_test.go b/internal/webapp/compress_e2e_test.go new file mode 100644 index 0000000..f1b0ed0 --- /dev/null +++ b/internal/webapp/compress_e2e_test.go @@ -0,0 +1,222 @@ +package webapp + +import ( + "bytes" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" +) + +// spyHub is a proxy in front of an existing hub that watches ONE client: it +// counts the bodies both ways and remembers which content encodings crossed. +// +// Per-client, and deliberately not countingHub with a mark taken between +// phases. `bdrive init` starts a daemon, so a mark placed after it measures +// whatever the daemon had not already done — which is how the first draft of +// this test reported a 282x pull ratio and proved nothing. Giving the client +// its own front door makes every counted byte and every observed header that +// client's, whenever it happened. +type spy struct { + up, down atomic.Int64 + mu sync.Mutex + sentEnc map[string]bool // Content-Encoding the client PUT with + gotEnc map[string]bool // Content-Encoding the hub answered it with +} + +func (s *spy) note(m map[string]bool, enc string) { + s.mu.Lock() + defer s.mu.Unlock() + m[enc] = true +} + +func (s *spy) saw(m map[string]bool, enc string) bool { + s.mu.Lock() + defer s.mu.Unlock() + return m[enc] +} + +func spyHub(t *testing.T, inner *httptest.Server) (*httptest.Server, *spy) { + t.Helper() + sp := &spy{sentEnc: map[string]bool{}, gotEnc: map[string]bool{}} + proxy := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method == http.MethodPut && strings.Contains(r.URL.Path, "/store/") { + sp.note(sp.sentEnc, r.Header.Get("Content-Encoding")) + } + r.Body = &countingReader{r: r.Body, n: &sp.up} + rec := &encSpyRW{ResponseWriter: w, sp: sp} + inner.Config.Handler.ServeHTTP(rec, r) + })) + t.Cleanup(proxy.Close) + return proxy, sp +} + +type encSpyRW struct { + http.ResponseWriter + sp *spy + noted bool +} + +func (e *encSpyRW) WriteHeader(code int) { + e.note() + e.ResponseWriter.WriteHeader(code) +} + +func (e *encSpyRW) Write(p []byte) (int, error) { + e.note() + n, err := e.ResponseWriter.Write(p) + e.sp.down.Add(int64(n)) + return n, err +} + +func (e *encSpyRW) note() { + if !e.noted { + e.noted = true + e.sp.note(e.sp.gotEnc, e.Header().Get("Content-Encoding")) + } +} + +// Transport compression against the REAL binary that predates it (oldBinRef, +// pinned at the commit before delta sync — which is also the commit before +// this). Both halves of the mixed-fleet claim, driven end to end rather than +// argued: +// +// - the old binary's PULL is compressed with no client change, because +// net/http asks for gzip and inflates on its own. This is the whole reason +// the read leg needed no negotiation. +// - the old binary's PUSH stays raw and still converges, because it never +// reads the accept_encoding sign() now answers — the same field an old HUB +// omits, which is what keeps a new client from posting gzip bytes to a hub +// that would store them under the sha256 of the plaintext. +func TestCompressionE2E_OldBinaryPullsCompressedAndPushesRaw(t *testing.T) { + inner := startTestHub(t) + fresh := newCLIEnvOn(t, inner) + oldDoor, sp := spyHub(t, inner) + old := newCLIEnvBin(t, oldDoor, buildOldBinary(t)) + + dirA := filepath.Join(t.TempDir(), "proj") + initProject(t, fresh, dirA, "compress-e2e", false) + raw := 0 + for i := 0; i < 40; i++ { + body := []byte(fmt.Sprintf("# note %d\n\n%s", i, markdownish(i))) + if err := os.WriteFile(filepath.Join(dirA, fmt.Sprintf("note-%02d.md", i)), body, 0o644); err != nil { + t.Fatal(err) + } + raw += len(body) + } + syncNow(t, fresh, dirA) + + // The old binary pulls the corpus it has never seen. Everything it has + // ever received is behind its own door, so no window has to be guessed. + dirOld := filepath.Join(t.TempDir(), "proj") + initProject(t, old, dirOld, "compress-e2e", true) + syncNow(t, old, dirOld) + for i := 0; i < 40; i++ { + want := fmt.Sprintf("# note %d\n\n%s", i, markdownish(i)) + got, err := os.ReadFile(filepath.Join(dirOld, fmt.Sprintf("note-%02d.md", i))) + if err != nil || string(got) != want { + t.Fatalf("old binary did not converge on note %d: %v", i, err) + } + } + pulled := sp.down.Load() + t.Logf("old binary received %d bytes total for a %d-byte corpus (%.2fx)", pulled, raw, float64(raw)/float64(pulled)) + if !sp.saw(sp.gotEnc, "gzip") { + t.Fatal("the hub never answered the old binary with Content-Encoding: gzip — the free win is not happening") + } + // Everything that client has ever been sent, sign-in and listings included, + // against the corpus alone. + if ratio := float64(raw) / float64(pulled); ratio < 2.5 { + t.Fatalf("old binary received %d bytes for a %d-byte corpus (%.2fx), want at least 2.5x", pulled, raw, ratio) + } + + // The old binary pushes. It never learned to read accept_encoding, so every + // body it PUT must be raw — asserted on the headers rather than on a byte + // count, so a daemon tick cannot change the answer. + back := []byte("# from the old client\n\n" + markdownish(99)) + if err := os.WriteFile(filepath.Join(dirOld, "from-old.md"), back, 0o644); err != nil { + t.Fatal(err) + } + syncNow(t, old, dirOld) + if sp.saw(sp.sentEnc, "gzip") { + t.Fatal("the old binary sent a compressed body — it cannot know how, so the hub must have been asked to inflate one it never sent") + } + if !sp.saw(sp.sentEnc, "") { + t.Fatal("the old binary never PUT anything; the push half was not exercised") + } + syncNow(t, fresh, dirA) + if got, err := os.ReadFile(filepath.Join(dirA, "from-old.md")); err != nil || !bytes.Equal(got, back) { + t.Fatalf("the new client did not converge on the old binary's raw push: %v", err) + } +} + +// Delta sync grew the key space by two classes after this feature was written. +// Both ride the same compressed wire, and both keep the gates that make them +// safe — because those gates read the spooled PLAINTEXT, below the inflate. +// +// The manifest is the one that matters: it is never presigned, so in production +// it ALWAYS goes through the relay path that compresses. A large file that +// syncs as chunks + manifest is the case where a gzip body would otherwise have +// been the thing the write-once compare and the chunks-exist gate saw. +func TestCompressionE2E_ChunksAndManifestsOverGzip(t *testing.T) { + inner := startTestHub(t) + door, sp := spyHub(t, inner) + a := newCLIEnvBin(t, door, "") + b := newCLIEnvOn(t, inner) + + dirA := filepath.Join(t.TempDir(), "proj") + initProject(t, a, dirA, "compress-chunks", false) + // Past the chunking threshold, and compressible — so the manifest, the + // chunks and the journal all take the gzip path rather than the skip path. + var big []byte + for len(big) < 12<<20 { + big = append(big, fmt.Sprintf("line %d: %s\n", len(big), markdownish(len(big)%7))...) + } + if err := os.WriteFile(filepath.Join(dirA, "big.md"), big, 0o644); err != nil { + t.Fatal(err) + } + syncNow(t, a, dirA) + // The push really did compress — otherwise the rest of this test is only + // re-testing delta sync. + if !sp.saw(sp.sentEnc, "gzip") { + t.Fatal("the chunked push sent nothing compressed") + } + + dirB := filepath.Join(t.TempDir(), "proj") + initProject(t, b, dirB, "compress-chunks", true) + syncNow(t, b, dirB) + got, err := os.ReadFile(filepath.Join(dirB, "big.md")) + if err != nil || !bytes.Equal(got, big) { + t.Fatalf("chunked file did not converge over the compressed wire: %v, %d bytes", err, len(got)) + } + + // Edit it: the second push re-puts a manifest under a NEW key and re-sends + // only the changed chunks, all compressed. If the write-once compare or the + // chunks-exist gate had been reading gzip bytes, this is where it breaks. + copy(big[6<<20:], []byte("EDITED")) + if err := os.WriteFile(filepath.Join(dirA, "big.md"), big, 0o644); err != nil { + t.Fatal(err) + } + syncNow(t, a, dirA) + syncNow(t, b, dirB) + if got, err := os.ReadFile(filepath.Join(dirB, "big.md")); err != nil || !bytes.Equal(got, big) { + t.Fatalf("the edit did not converge: %v", err) + } + // A repeat sync must be a no-op, not a manifest conflict: an identical + // re-put is the retry path, and it is compared on the plaintext. + syncNow(t, a, dirA) +} + +// markdownish is a stand-in for the real corpus — markdown and source, the +// content that made compression worth doing. +func markdownish(seed int) string { + var b []byte + for i := 0; i < 60; i++ { + b = append(b, fmt.Sprintf("- item %d/%d: the sync wire carries markdown and source files\n", seed, i)...) + } + return string(b) +} diff --git a/internal/webapp/store.go b/internal/webapp/store.go index 0c848cc..78f803c 100644 --- a/internal/webapp/store.go +++ b/internal/webapp/store.go @@ -2,6 +2,7 @@ package webapp import ( "bytes" + "compress/gzip" "context" "encoding/json" "errors" @@ -45,6 +46,26 @@ func validStoreKey(key string) bool { chunkKeyRe.MatchString(key) || manifestKeyRe.MatchString(key) } +// storeAcceptEncoding is what handleStoreSign tells a client this hub will +// inflate on a relayed PUT. A hub older than this answers without the field, +// and the client sends raw — which is the whole mixed-fleet story for the push +// leg (the pull leg needs no negotiation at all, see remote/compress.go). +var storeAcceptEncoding = []string{"gzip"} + +// maxInflatedPut bounds what one compressed PUT may write to the hub's disk. +// +// spool() is unbounded, which is safe exactly as long as a byte on the wire +// costs a byte on disk. Content-Encoding breaks that: a 1 MB body can inflate +// to an arbitrary write, and it lands BEFORE CheckWrite can refuse it, because +// nothing knows the size until the inflate is done. The bound applies only when +// the body declares an encoding, so no honest raw push that works today can +// start failing. +// +// ponytail: 256 MiB, mirroring maxImportBlob on the archive path — a precedent, +// not a measurement. A real workload that hits it wants a server-side knob, not +// a bigger constant. +const maxInflatedPut = 256 << 20 + // storeSource returns the volume's RemoteSource; only real beardrive // remotes have a store to expose. func storeSource(v *volume, w http.ResponseWriter) *RemoteSource { @@ -197,7 +218,25 @@ func (s *Server) handleStoreList(v *volume, w http.ResponseWriter, r *http.Reque storageErr(w, http.StatusBadGateway, "storage is temporarily unavailable", err) return } - writeJSON(w, map[string]any{"objects": objs}) + writeStoreJSON(w, r, map[string]any{"objects": objs}) +} + +// writeStoreJSON is writeJSON that compresses when the caller accepts it. The +// listing is the first call of every sync cycle on every device, it is JSON, +// and it is highly repetitive — one key per blob — so it is compressed without +// probing. Only this route uses it: writeJSON serves the whole browser API and +// is out of this change's scope, and handleStoreExists answers one boolean, +// which is smaller than a gzip header. +func writeStoreJSON(w http.ResponseWriter, r *http.Request, v any) { + w.Header().Set("Content-Type", "application/json") + if !remote.AcceptsGzip(r) { + json.NewEncoder(w).Encode(v) + return + } + w.Header().Set("Content-Encoding", "gzip") + gz := gzip.NewWriter(w) + defer gz.Close() + json.NewEncoder(gz).Encode(v) } func (s *Server) handleStoreGet(v *volume, w http.ResponseWriter, r *http.Request) { @@ -226,6 +265,17 @@ func (s *Server) handleStoreGet(v *volume, w http.ResponseWriter, r *http.Reques return } defer rc.Close() + // Compression is decided before a single header is written, because + // Content-Encoding cannot be added after the first Write. + var src io.Reader = rc + gzipOK := false + if remote.AcceptsGzip(r) { + src, gzipOK, err = remote.Compressible(rc) + if err != nil { + storageErr(w, http.StatusBadGateway, "could not read the object", err) + return + } + } w.Header().Set("Content-Type", "application/octet-stream") // The sync proxy is a stored-bytes door like the other two: a // cookie-authenticated GET whose URL one member can hand another, answering @@ -234,8 +284,20 @@ func (s *Server) handleStoreGet(v *volume, w http.ResponseWriter, r *http.Reques // Recorded, never checked. This is a device syncing: refusing it here // surfaces as ErrForbidden, which the syncer reads as "access is gone — // pause and touch nothing". Sync must not break over a bill. + // + // The counter wraps the SOCKET and gzip writes into it, never the other way + // round: RecordEgress is a bandwidth meter, so it has to report what left + // the machine. Inverted, it silently bills plaintext for every compressed + // response and no test fails. cw := &countingWriter{w: w} - io.Copy(cw, rc) + if gzipOK { + w.Header().Set("Content-Encoding", "gzip") + gz := gzip.NewWriter(cw) + io.Copy(gz, src) + gz.Close() + } else { + io.Copy(cw, src) + } s.quota().RecordEgress(s.orgOf(r.PathValue("project")), cw.n) } @@ -316,7 +378,7 @@ func (s *Server) handleStoreSign(v *volume, w http.ResponseWriter, r *http.Reque return } if exists, err := rs.Backend.Exists(r.Context(), req.Key); err == nil && exists { - writeJSON(w, map[string]any{"mode": "direct", "exists": true}) + writeJSON(w, map[string]any{"mode": "direct", "exists": true, "accept_encoding": storeAcceptEncoding}) return } if signer, ok := rs.Backend.(remote.PutSigner); ok { @@ -335,13 +397,14 @@ func (s *Server) handleStoreSign(v *volume, w http.ResponseWriter, r *http.Reque writeJSON(w, map[string]any{ "mode": "direct", "url": signed.URL, "method": signed.Method, "headers": signed.Headers, "expires": signed.Expires.UTC(), + "accept_encoding": storeAcceptEncoding, }) return } s.claimGrant(project, req.Key) // nothing was granted: give it back } } - writeJSON(w, map[string]any{"mode": "server"}) + writeJSON(w, map[string]any{"mode": "server", "accept_encoding": storeAcceptEncoding}) } // journalOps reads the operations a spooled journal body carries, exactly the @@ -517,6 +580,32 @@ func journalKeepsItsOps(ctx context.Context, be remote.Backend, key string, ops return true, nil } +// storePutBody hands back the plaintext of a PUT body, inflating it when the +// client declared an encoding, and reports whether it did. The reader it +// returns is capped one byte past maxInflatedPut so a bomb can never write more +// than that to disk before the caller measures it and refuses. +// +// It answers the client itself on the two ways this can be the client's fault: +// an encoding this hub does not implement, and a body that says gzip and is +// not one (gzip.NewReader reads the header eagerly, so that is caught here +// rather than halfway through a spool). +func storePutBody(w http.ResponseWriter, r *http.Request) (io.Reader, bool, bool) { + enc := strings.TrimSpace(r.Header.Get("Content-Encoding")) + if enc == "" { + return r.Body, false, true + } + if !strings.EqualFold(enc, "gzip") { + http.Error(w, "unsupported Content-Encoding "+enc, http.StatusUnsupportedMediaType) + return nil, false, false + } + gz, err := gzip.NewReader(r.Body) + if err != nil { + http.Error(w, "body declares Content-Encoding: gzip but is not gzip", http.StatusBadRequest) + return nil, false, false + } + return io.LimitReader(gz, maxInflatedPut+1), true, true +} + func (s *Server) handleStorePut(v *volume, w http.ResponseWriter, r *http.Request) { rs := storeSource(v, w) if rs == nil { @@ -536,13 +625,32 @@ func (s *Server) handleStorePut(v *volume, w http.ResponseWriter, r *http.Reques // (Content-Length is -1 on any chunked request, which made every unsized // put free), and how many ops a journal write actually authors. // Cost: one temp file per put on the hub's busiest write path. - tmp, size, sum, err := spool(r.Body) + // + // Inflating sits ABOVE the spool, because every single thing this handler + // goes on to be sure of is a property of the plaintext: the sha a blob key + // promises, the ops journalOps counts, the size CheckWrite bills, and the + // append-only check. Nothing below this line knows compression happened. + body, inflated, ok := storePutBody(w, r) + if !ok { + return + } + tmp, size, sum, err := spool(body) if err != nil { + if inflated { + // A gzip stream that truncates or fails its CRC is the client's + // body, not the hub's storage; 502 would blame the wrong machine. + http.Error(w, "could not decompress the body", http.StatusBadRequest) + return + } storageErr(w, http.StatusBadGateway, "could not store the object", err) return } defer os.Remove(tmp.Name()) defer tmp.Close() + if inflated && size > maxInflatedPut { + http.Error(w, "compressed body inflates past this hub's limit", http.StatusRequestEntityTooLarge) + return + } // Blobs and chunks are content-addressed: the key IS the content's hash. // Manifests are not — their key is the whole FILE's sha, which the hub // cannot check without reading every chunk; readers verify by reassembly. diff --git a/internal/webapp/store_compress_test.go b/internal/webapp/store_compress_test.go new file mode 100644 index 0000000..f10e6a6 --- /dev/null +++ b/internal/webapp/store_compress_test.go @@ -0,0 +1,329 @@ +package webapp + +import ( + "bytes" + "compress/gzip" + "encoding/json" + "io" + "math/rand" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "sync" + "testing" + + "github.com/runbear-io/beardrive/internal/journal" + "github.com/runbear-io/beardrive/internal/remote" +) + +// Transport compression on /store/*. Everything here is about one invariant +// holding while the wire changes underneath it: content addressing is over the +// UNCOMPRESSED bytes. The hub inflates before it hashes, stores plaintext, and +// bills plaintext — the only thing that is ever compressed is the transfer. + +const textCorpus = "# release notes\n\nthe corpus is markdown and source files, " + + "five to ten kilobytes each, which is exactly what gzip is good at.\n" + +// meterQuota records what the hub bills, separately per meter: usage is a +// storage bill (uncompressed), egress is a bandwidth bill (compressed). +type meterQuota struct { + UnlimitedQuota + mu sync.Mutex + usage int64 + egress int64 +} + +func (q *meterQuota) RecordUsage(_ string, b int64) { + q.mu.Lock() + defer q.mu.Unlock() + q.usage += b +} + +func (q *meterQuota) RecordEgress(_ string, b int64) { + q.mu.Lock() + defer q.mu.Unlock() + q.egress += b +} + +// doRaw sends an exact body with exact headers — reads_test.go's doHdr +// JSON-marshals what it is given, and here the bytes on the wire are the +// subject. +func doRaw(t *testing.T, h http.Handler, method, url string, body []byte, hdr map[string]string) *httptest.ResponseRecorder { + t.Helper() + req := httptest.NewRequest(method, url, bytes.NewReader(body)) + for k, v := range hdr { + req.Header.Set(k, v) + } + rec := httptest.NewRecorder() + h.ServeHTTP(rec, req) + return rec +} + +func gzipBytes(t *testing.T, b []byte) []byte { + t.Helper() + var buf bytes.Buffer + gz := gzip.NewWriter(&buf) + if _, err := gz.Write(b); err != nil { + t.Fatal(err) + } + if err := gz.Close(); err != nil { + t.Fatal(err) + } + return buf.Bytes() +} + +func gunzipBytes(t *testing.T, b []byte) []byte { + t.Helper() + gz, err := gzip.NewReader(bytes.NewReader(b)) + if err != nil { + t.Fatal(err) + } + out, err := io.ReadAll(gz) + if err != nil { + t.Fatal(err) + } + return out +} + +func incompressible(n int) []byte { + b := make([]byte, n) + rand.New(rand.NewSource(7)).Read(b) + return b +} + +// The pull leg. It needs no client change at all — net/http sends +// Accept-Encoding: gzip on its own and inflates transparently — so this is what +// a device built before this feature existed gets for free, and it must still +// arrive as the exact stored bytes. +func TestStoreGetCompresses(t *testing.T) { + srv, p, root := newHub(t, false, nil) + q := &meterQuota{} + srv.Quota = q + f := newFakeRemoteAt(t, filepath.Join(root, p.ID)) + text := strings.Repeat(textCorpus, 200) + f.put("deva", "notes.md", text) + h := srv.Handler() + key := "blobs/" + shaOf(text) + url := "/api/p/" + p.ID + "/store/object?key=" + key + + rec := doRaw(t, h, "GET", url, nil, map[string]string{"Accept-Encoding": "gzip"}) + if rec.Code != 200 || rec.Header().Get("Content-Encoding") != "gzip" { + t.Fatalf("gzip GET: %d %q", rec.Code, rec.Header().Get("Content-Encoding")) + } + wire := rec.Body.Bytes() + if got := string(gunzipBytes(t, wire)); got != text { + t.Fatal("the compressed response does not inflate to the stored bytes") + } + if ratio := float64(len(text)) / float64(len(wire)); ratio < 2.5 { + t.Fatalf("wire ratio %.2fx, want at least 2.5x", ratio) + } + // The bandwidth meter reports what left the socket, not what was stored. + // Inverted (counter inside the gzip writer) this reads len(text) and + // nothing else fails. + if q.egress != int64(len(wire)) { + t.Fatalf("RecordEgress = %d, want the compressed %d", q.egress, len(wire)) + } + + // A caller that did not ask gets exactly what it always got. + rec = doRaw(t, h, "GET", url, nil, nil) + if rec.Code != 200 || rec.Header().Get("Content-Encoding") != "" || rec.Body.String() != text { + t.Fatalf("plain GET: %d %q", rec.Code, rec.Header().Get("Content-Encoding")) + } +} + +// Already-compressed content must cross untouched: gzip on a JPEG pays CPU to +// make the payload ~0.1% bigger. +func TestStoreGetSkipsIncompressible(t *testing.T) { + srv, p, root := newHub(t, false, nil) + f := newFakeRemoteAt(t, filepath.Join(root, p.ID)) + raw := incompressible(300 << 10) + f.put("deva", "photo.jpg", string(raw)) + h := srv.Handler() + url := "/api/p/" + p.ID + "/store/object?key=blobs/" + shaOf(string(raw)) + + for _, ae := range []string{"gzip", ""} { + rec := doRaw(t, h, "GET", url, nil, map[string]string{"Accept-Encoding": ae}) + if rec.Code != 200 || rec.Header().Get("Content-Encoding") != "" { + t.Fatalf("Accept-Encoding %q: %d %q", ae, rec.Code, rec.Header().Get("Content-Encoding")) + } + if n := rec.Body.Len(); float64(n) > float64(len(raw))*1.01 { + t.Fatalf("wire carried %d bytes for a %d-byte payload", n, len(raw)) + } + if !bytes.Equal(rec.Body.Bytes(), raw) { + t.Fatal("body is not the stored bytes") + } + } +} + +// The listing is the first call of every cycle on every device, and it is +// repetitive JSON. +func TestStoreListCompresses(t *testing.T) { + srv, p, root := newHub(t, false, nil) + f := newFakeRemoteAt(t, filepath.Join(root, p.ID)) + for i := 0; i < 50; i++ { + f.put("deva", "notes/"+strings.Repeat("x", i+1)+".md", textCorpus+strings.Repeat("y", i+1)) + } + h := srv.Handler() + url := "/api/p/" + p.ID + "/store/list?prefix=blobs/" + + rec := doRaw(t, h, "GET", url, nil, map[string]string{"Accept-Encoding": "gzip"}) + if rec.Code != 200 || rec.Header().Get("Content-Encoding") != "gzip" { + t.Fatalf("gzip list: %d %q", rec.Code, rec.Header().Get("Content-Encoding")) + } + var list struct { + Objects []remote.Object `json:"objects"` + } + if err := json.Unmarshal(gunzipBytes(t, rec.Body.Bytes()), &list); err != nil { + t.Fatal(err) + } + if len(list.Objects) != 50 { + t.Fatalf("objects = %d, want 50", len(list.Objects)) + } + plain := doRaw(t, h, "GET", url, nil, nil) + if plain.Header().Get("Content-Encoding") != "" { + t.Fatal("a caller that did not ask for gzip got it anyway") + } + if rec.Body.Len() >= plain.Body.Len() { + t.Fatalf("compressed listing is %d bytes, raw is %d", rec.Body.Len(), plain.Body.Len()) + } +} + +// The push leg is negotiated, and this is the advertisement the client reads. +func TestStoreSignAdvertisesGzip(t *testing.T) { + srv, p, _ := newHub(t, true, nil) + h := srv.Handler() + rec := do(t, h, "POST", "/api/p/"+p.ID+"/store/sign", + map[string]any{"key": "blobs/" + shaOf("hi"), "size": 2}) + var plan struct { + Mode string `json:"mode"` + AcceptEncoding []string `json:"accept_encoding"` + } + mustJSON(t, rec, &plan) + if plan.Mode != "server" || len(plan.AcceptEncoding) != 1 || plan.AcceptEncoding[0] != "gzip" { + t.Fatalf("plan = %+v", plan) + } +} + +// A compressed PUT is stored as plaintext under the plaintext's hash, and +// billed for the plaintext — the object stored is uncompressed, so that is what +// the storage bill is for. +func TestStorePutInflatesBeforeItHashes(t *testing.T) { + srv, p, root := newHub(t, true, nil) + q := &meterQuota{} + srv.Quota = q + h := srv.Handler() + text := strings.Repeat(textCorpus, 200) + key := "blobs/" + shaOf(text) + body := gzipBytes(t, []byte(text)) + if len(body) >= len(text) { + t.Fatal("fixture is not actually compressed") + } + + rec := doRaw(t, h, "PUT", "/api/p/"+p.ID+"/store/object?key="+key, body, + map[string]string{"Content-Encoding": "gzip"}) + if rec.Code != 200 { + t.Fatalf("gzip put: %d %s", rec.Code, rec.Body) + } + stored, err := os.ReadFile(filepath.Join(root, p.ID, key)) + if err != nil { + t.Fatal(err) + } + if string(stored) != text { + t.Fatal("storage does not hold the plaintext") + } + if q.usage != int64(len(text)) { + t.Fatalf("RecordUsage = %d, want the uncompressed %d", q.usage, len(text)) + } +} + +// Content addressing is unchanged: a compressed body that decodes to bytes +// which are not what the key names is refused, exactly like a raw one. +func TestStorePutGzipMustStillHashToItsKey(t *testing.T) { + srv, p, _ := newHub(t, true, nil) + h := srv.Handler() + body := gzipBytes(t, []byte("the plaintext nobody asked for")) + rec := doRaw(t, h, "PUT", "/api/p/"+p.ID+"/store/object?key=blobs/"+shaOf("something else"), body, + map[string]string{"Content-Encoding": "gzip"}) + if rec.Code != http.StatusBadRequest || !strings.Contains(rec.Body.String(), "does not hash to its key") { + t.Fatalf("mismatched gzip put: %d %s", rec.Code, rec.Body) + } +} + +// Content-Encoding severs the "a byte on the wire costs a byte on disk" +// relationship that made spool() safe unbounded, and the inflate lands before +// CheckWrite can refuse anything. A small body must not become an arbitrary +// hub-side write. +func TestStorePutRefusesAGzipBomb(t *testing.T) { + srv, p, root := newHub(t, true, nil) + h := srv.Handler() + // ~256 MiB of zeros compresses to a few hundred KB — the shape of the bomb, + // small enough to keep the test fast. + bomb := gzipBytes(t, make([]byte, maxInflatedPut+1)) + if len(bomb) > 1<<20 { + t.Fatalf("bomb fixture is %d bytes, expected it to compress much harder", len(bomb)) + } + rec := doRaw(t, h, "PUT", "/api/p/"+p.ID+"/store/object?key=blobs/"+strings.Repeat("a", 64), bomb, + map[string]string{"Content-Encoding": "gzip"}) + if rec.Code != http.StatusRequestEntityTooLarge { + t.Fatalf("bomb: %d %s, want 413", rec.Code, rec.Body) + } + blobs, _ := os.ReadDir(filepath.Join(root, p.ID, "blobs")) + if len(blobs) != 0 { + t.Fatalf("the bomb stored %d objects", len(blobs)) + } +} + +// The two ways a declared encoding can be the client's fault, kept apart from +// "the hub's storage broke". +func TestStorePutRejectsBadEncodings(t *testing.T) { + srv, p, _ := newHub(t, true, nil) + h := srv.Handler() + url := "/api/p/" + p.ID + "/store/object?key=blobs/" + shaOf("hi") + + rec := doRaw(t, h, "PUT", url, []byte("not gzip at all"), map[string]string{"Content-Encoding": "gzip"}) + if rec.Code != http.StatusBadRequest { + t.Fatalf("lying Content-Encoding: %d %s, want 400", rec.Code, rec.Body) + } + rec = doRaw(t, h, "PUT", url, []byte("hi"), map[string]string{"Content-Encoding": "br"}) + if rec.Code != http.StatusUnsupportedMediaType { + t.Fatalf("unsupported codec: %d %s, want 415", rec.Code, rec.Body) + } + // Truncated gzip: the header parses, the stream does not finish. + full := gzipBytes(t, []byte(strings.Repeat(textCorpus, 50))) + rec = doRaw(t, h, "PUT", url, full[:len(full)/2], map[string]string{"Content-Encoding": "gzip"}) + if rec.Code != http.StatusBadRequest { + t.Fatalf("truncated gzip: %d %s, want 400", rec.Code, rec.Body) + } +} + +// A journal body is read for its ops after the inflate, so every rule that +// protects the log still applies to a compressed push — including the +// append-only one, which is the invariant a compressed body could otherwise +// have smuggled past. +func TestStorePutGzippedJournalKeepsItsOps(t *testing.T) { + srv, p, _ := newHub(t, true, nil) + h := srv.Handler() + url := "/api/p/" + p.ID + "/store/object?key=journal/deva.jsonl" + hdr := map[string]string{"Content-Encoding": "gzip", "X-Bdrive-Device": "deva"} + + line := func(seq int64, path string) []byte { + b, err := json.Marshal(journal.Op{ + Kind: journal.KindPut, Path: path, Seq: seq, Device: "deva", + Blob: shaOf(path), Size: 1, Lamport: seq, + }) + if err != nil { + t.Fatal(err) + } + return append(b, '\n') + } + two := append(line(1, "a.md"), line(2, "b.md")...) + if rec := doRaw(t, h, "PUT", url, gzipBytes(t, two), hdr); rec.Code != 200 { + t.Fatalf("gzipped journal put: %d %s", rec.Code, rec.Body) + } + // Same journal, one op short: refused, compressed or not. + if rec := doRaw(t, h, "PUT", url, gzipBytes(t, line(1, "a.md")), hdr); rec.Code != http.StatusConflict { + t.Fatalf("truncating journal put: %d %s, want 409", rec.Code, rec.Body) + } +}