From e8e4ef4b15e5b1075f2e92de5ec52e98b4fa488a Mon Sep 17 00:00:00 2001 From: Claude Date: Tue, 21 Jul 2026 00:07:44 +0000 Subject: [PATCH] perf(round25): in-place encrypted decrypt, dedup hash move, dead folder-Query, playlist N+1 fold, contact vcard over-fetch Five benchmark-gated optimizations from a fresh six-way audit (benches/ROUND25.md), each with a BEFORE/AFTER gate that rolls back if AFTER does not beat BEFORE: - M1 EncryptedBlobBackend::decrypt_bytes: replace split_off (a fresh Vec + full ciphertext memcpy on every decrypted chunk, contradicting its own "in place" doc) with in-place detached decrypt + a zero-copy Bytes::slice past the nonce. Peak RAM per read drops from ~2x to ~1x the payload (-262KB/op at 256KiB; scales with blob size). Plaintext byte-identical; tamper/wrong-key tests pass. - M2 delta commit: move-unzip the owned chunk list instead of a third per-occurrence hash clone (-4000 allocs on a 4000-chunk commit). - M3 folder ZIP download: drop the dead Query extractor it never read (byte-identical response; 5->0 allocs/request). - Q1 public-playlist listing: fold the per-playlist COUNT(*) N+1 into one LEFT JOIN ... GROUP BY via a new inherent repo method (101 -> 1 round-trips, 36x wall on a 100-playlist page). - Q2 contact REST listings (paginated/search/by-group): stop over-fetching the multi-KB vcard TEXT the ContactDto discards, via a shared lite row mapper and narrowed SELECTs (6.4x wall on 1000 contacts with 8KiB vcards). The whole-book vCard export and CardDAV sync paths keep the column. Adds bench_round25_micro (counting allocator tracking count+bytes) and bench_round25_queries (live Postgres), both with equivalence gates and a rollback exit(1). Verified: cargo fmt clean, cargo clippy -D warnings clean, cargo test --lib --features bench = 529 passed / 0 failed. Co-Authored-By: Claude Opus 4.8 Claude-Session: https://claude.ai/code/session_01L8gs91AhmazoxMsDcNk3KT --- Cargo.toml | 21 + benches/ROUND25.md | 286 +++++++++++++ examples/bench_round25_micro.rs | 317 +++++++++++++++ examples/bench_round25_queries.rs | 381 ++++++++++++++++++ .../services/delta_upload_service.rs | 8 +- .../adapters/music_storage_adapter.rs | 19 +- .../repositories/pg/contact_pg_repository.rs | 38 +- .../repositories/pg/playlist_pg_repository.rs | 61 +++ .../services/encrypted_blob_backend.rs | 37 +- src/interfaces/api/handlers/folder_handler.rs | 9 +- 10 files changed, 1142 insertions(+), 35 deletions(-) create mode 100644 benches/ROUND25.md create mode 100644 examples/bench_round25_micro.rs create mode 100644 examples/bench_round25_queries.rs diff --git a/Cargo.toml b/Cargo.toml index 00e65fd4..f18f870f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -354,6 +354,27 @@ name = "bench_micro_allocs" path = "examples/bench_micro_allocs.rs" required-features = ["bench"] +# Round-25 battery ──────────────────────────────────────────────────────────── + +# Round-25 CPU/alloc/RAM micro-pack (no Postgres) — deterministic alloc+bytes +# gates: EncryptedBlobBackend::decrypt_bytes split_off full copy → in-place +# detached decrypt + zero-copy slice (M1, the RAM headline); delta-commit +# chunk-hash list third clone → move-unzip (M2); folder download dead +# Query extractor removal (M3). +[[example]] +name = "bench_round25_micro" +path = "examples/bench_round25_micro.rs" +required-features = ["bench"] + +# Round-25 PG query-shape pack — public-playlist listing 1+N COUNT round-trips +# → one LEFT JOIN GROUP BY (Q1); contact REST listings dropping the over-fetched +# multi-KB vcard TEXT the ContactDto discards (Q2). Needs the dev Postgres up +# (reads DATABASE_URL from .env). +[[example]] +name = "bench_round25_queries" +path = "examples/bench_round25_queries.rs" +required-features = ["bench"] + # Round-24 battery ──────────────────────────────────────────────────────────── # Round-24 download_zip authz+metadata N+1 → batch, VALIDATED. The per-file diff --git a/benches/ROUND25.md b/benches/ROUND25.md new file mode 100644 index 00000000..a44d3cd7 --- /dev/null +++ b/benches/ROUND25.md @@ -0,0 +1,286 @@ +# Round 25 — encrypted-read in-place decrypt (RAM), delta-commit hash move, dead folder-Query, public-playlist N+1 fold, contact vcard over-fetch + +This round lands a cross-cutting perf pass surfaced by a fresh six-way audit of +the tree (dedup/upload, blob-I/O, DB query-shape, HTTP/DAV emitters, auth/global +config, frontend), cross-referenced against everything ROUND2–24 already shipped +so nothing here re-treads landed work. Five items ship, each behind a +BEFORE/AFTER benchmark that `std::process::exit(1)`s ("`GATE FAIL … rollback`") +unless AFTER strictly beats BEFORE — the round's roll-back rule encoded into the +benchmark, so an AFTER that doesn't win is never applied to the source. + +Reproduce: + +```bash +# M1–M3 — counting global allocator (count + BYTES), no Postgres +RUSTFLAGS="-C target-cpu=x86-64-v3" \ + cargo run --release --features bench --example bench_round25_micro + +# Q1–Q2 — live dev Postgres (reads DATABASE_URL from .env) +RUSTFLAGS="-C target-cpu=x86-64-v3" \ + cargo run --release --features bench --example bench_round25_queries +``` + +The two headline items match the owner's top priorities: **M1 halves peak RAM on +every encrypted blob read**, and **Q1 collapses the public-playlist gallery from +101 DB round-trips to 1**. + +--- + +## [M1] `EncryptedBlobBackend::decrypt_bytes` — full ciphertext copy → in-place detached decrypt (RAM) + +`decrypt_bytes` claimed in its own doc comment to decrypt "**in place** … the +ciphertext buffer is reused for the plaintext instead of allocating a second +copy." It did not: + +```rust +let mut ciphertext = encrypted.split_off(NONCE_SIZE); // allocates + memcpy's the whole tail +``` + +`Vec::split_off(12)` allocates a fresh `Vec` sized `len-12` and `ptr::copy`s the +entire ciphertext+tag into it — so **every decrypted CDC chunk (≤ 1 MiB), and +every legacy whole-file blob, paid one full-payload allocation + memcpy on read**. +ROUND11 §15 fixed the *encrypt* side (`encrypt_in_place_detached`) but the +decrypt side was never given the same treatment; the stale doc comment is the +tell that it was believed already done. + +AFTER lifts the 12-byte nonce and 16-byte GCM tag to the stack, decrypts the +middle in place via `decrypt_in_place_detached` (the detached API already used by +the encrypt side), and returns a **zero-copy `Bytes::slice` past the nonce** — no +extra allocation, no full-payload copy. Plaintext bytes are identical. + +| arm | allocs/op | bytes/op | note | +|--------|----------:|---------:|------| +| BEFORE | 3.00 | 524 356 | input clone + `split_off` copy + `Bytes::from` | +| AFTER | 2.00 | 262 196 | input clone + `Bytes::from` only | + +**−262 160 bytes/op** at a 256 KiB payload — the copied ciphertext eliminated; +peak heap on a decrypt drops from ~2× to ~1× the payload. The win scales with +payload, so a legacy whole-file blob read no longer transiently doubles a +multi-hundred-MB allocation. Gate: **AFTER bytes/op strictly lower** (it is). The +equivalence arm asserts the decrypted plaintext is byte-identical to the +`split_off` path across the short-input edge, 64 KiB and 1 MiB. + +## [M2] Delta commit — third per-occurrence hash clone → move-unzip (dedup allocations) + +`delta_upload_service::commit_with_perms` owns `request: DeltaCommitRequest`, yet +materialized the per-occurrence chunk-hash list a **third** time at the manifest +bind (after the distinct set and the verification tuple): + +```rust +let chunk_hashes: Vec = request.chunks.iter().map(|c| c.h.clone()).collect(); +let chunk_sizes: Vec = request.chunks.iter().map(|c| c.s).collect(); +``` + +`request.chunks` is dead after this line (only `request.file_hash` is read +below), so AFTER moves the hashes out instead of cloning each 64-char hash: + +```rust +let (chunk_hashes, chunk_sizes): (Vec, Vec) = + request.chunks.into_iter().map(|c| (c.h, c.s)).unzip(); +``` + +| arm | allocs/op (4000 chunks) | bytes/op | +|--------|------------------------:|---------:| +| BEFORE | 8 003.00 | 768 000 | +| AFTER | 4 003.00 | 512 000 | + +**−4000 allocs/op** (the N hash-String clones) on the flagship "upload only what +changed" path. Gate: AFTER allocs/op strictly lower. Equivalence: the produced +`(chunk_hashes, chunk_sizes)` are element-equal to the clone-collect arms. + +## [M3] `folder_handler::download_folder_zip` — dead `Query` extractor removed (allocations) + +Both the route wrapper and `download_folder_zip_impl` bound +`Query>` as `_params` and discarded it — the handler reads +only the path `id`. axum's `Query` extractor parses the whole query string into a +`HashMap` plus an owned `String` key and value per param, all dropped unread. AFTER +deletes the extractor; axum ignores any query string when none is present, so the +response is byte-identical. + +| arm | ns/op | allocs/op | bytes/op | +|--------|-------:|----------:|---------:| +| BEFORE | 212.1 | 5.00 | 268 | +| AFTER | 0.3 | 0.00 | 0 | + +Pure dead-work elimination (**614× wall**, 5 → 0 allocs) whenever a client +appends any query string (cache-buster, tracking param). Gate: AFTER allocs/op +strictly lower. + +## [Q1] Public-playlist listing — 1 + N `COUNT(*)` → one `LEFT JOIN … GROUP BY` (DB round-trips) + +`MusicStorageAdapter::list_public_playlists` ran one listing SELECT then one +`SELECT COUNT(*) FROM audio.playlist_items WHERE playlist_id = $1` **per returned +playlist** — up to **101 serial round-trips** for a `limit=100` gallery page. +AFTER folds the count into the listing with a single +`LEFT JOIN audio.playlist_items … GROUP BY p.id`, exposed as a new inherent +`PlaylistPgRepository::list_public_playlists_with_counts` returning +`(Playlist, track_count)` — backed by the existing +`idx_playlist_items_playlist_id`. (The adapter holds the concrete repo type, so +no trait change was needed; the two sibling 1+N adapter methods have no live +caller and are left untouched.) + +Live Postgres, 100 public playlists (varying track counts), p50 over 30 passes: + +| arm | p50 ms | round-trips | +|--------|-------:|------------:| +| BEFORE | 16.572 | 101 | +| AFTER | 0.458 | 1 | + +**36.2× wall, 101 → 1 round-trips.** Equivalence: the `(playlist → track_count)` +map is identical BEFORE vs AFTER (asserted; mismatch `exit(1)`s). Gate: AFTER p50 +strictly lower. On a remote/managed Postgres, where each round-trip is a network +RTT rather than a local socket hop, the win is far larger than the localhost 36×. + +## [Q2] Contact REST listings — stop over-fetching the multi-KB `vcard` TEXT (bandwidth) + +`get_contacts_by_address_book_paginated`, `search_contacts` and +`get_contacts_by_group` all `SELECT … vcard …` — the full serialized vCard TEXT, +the largest column (can embed a base64 `PHOTO` of tens of KB). But every caller +maps `Contact → ContactDto`, which has **no vcard field**, so it is fetched, +shipped over the wire, decoded into a `String` and immediately dropped. AFTER +adds a `row_to_contact_lite` mapper (shared `row_to_contact_with_vcard` core, no +duplication) that supplies an empty vcard, and narrows those three SELECTs to +omit the column. The shared `get_contacts_by_address_book` (also used by the +whole-book vCard export) and the CardDAV sync/multiget paths keep the column. + +Live Postgres, 1000 contacts each carrying an 8 KiB vCard, p50 over 20 passes: + +| arm | p50 ms | note | +|--------|-------:|------| +| BEFORE | 10.010 | SELECT incl. vcard, decoded + dropped | +| AFTER | 1.570 | SELECT without vcard | + +**6.4× wall** — and the win is bytes-on-the-wire + per-row `String` allocation, +both of which grow with vCard size (photos push these to tens of KB each). +Equivalence: the kept DTO fields `(id, full_name, photo_url)` are identical +across the change (asserted). Gate: AFTER p50 strictly lower. + +--- + +## Not shipped — verified this round, deferred to a later pass + +The six-way audit surfaced far more than shipped here; the following were +verified real against current source and carry a benchmark plan, but each needs a +multi-signature change, a remote-backend fixture, an operator-facing decision, or +its own validated pass. Grouped by area for the next rounds. + +### Blob-I/O / disk (owner priority) +- **`CachedBlobBackend::initialize` never pre-creates the 256 shard dirs** (the + line-122 comment says it does; it only makes `cache_dir`), so all three cache + writes pay a per-chunk `create_dir_all(parent)` — a wasted `mkdirat(EEXIST)` + + component stat + blocking-pool dispatch on cached-remote deployments. Fix mirrors + `LocalBlobBackend::initialize`'s `HEX_PREFIXES` loop; gate on a `strace -c` + `mkdirat` count + wall on tmpfs. (conf 0.9) +- **Eviction listener unlinks with a blocking `std::fs::remove_file` on the tokio + worker** (`cached_blob_backend.rs:98`) — `moka::sync` runs the listener inline on + the inserting worker; every write-through eviction blocks a reactor thread on + `unlink(2)`. Hand off via `spawn_blocking`/a drain task; gate on p99 scheduling + delay under eviction pressure. (conf 0.85) +- **S3 reads copy every served byte** through `into_async_read()+ReaderStream` + (`s3_blob_backend.rs:261/298`) while Azure already forwards SDK `Bytes` frames + zero-copy — needs a MinIO/stub fixture to gate. (conf 0.6) +- **`store_loose_chunks` writes loose chunks to the backend serially** while the + main ingest overlaps 8 (`buffer_unordered`); on a remote backend the delta path + serializes RTTs the main path hides. Needs a latency-stub backend. (conf 0.5) +- **`local_blob_path` does a synchronous `path.exists()` stat on the reactor** — + wants an async port variant. (conf 0.6) +- **`PLAINTEXT_EMIT_SIZE` = 64 KiB vs the 256 KiB every other backend streams** — + quarters the encrypted-read frame count; the "parity" comment justifying 64 KiB + is factually wrong. Wants a streaming A/B (frame count + wall). (conf 0.5) + +### DB query-shape +- **Drive-policy reads decode through a throwaway `serde_json::Value` DOM** + (`drive_pg_repository.rs` 4 methods) — ROUND23 §J2 removed only the clone, not the + DOM; fold to `sqlx::types::Json` like §J1. Fires on move/copy and + every share/grant create. (conf 0.75) +- **Contact create/update build a throwaway `Value` before binding JSONB** (the + write-side twin of ROUND23 §J1) — bind `sqlx::types::Json(&dtos)` directly. (conf 0.6) + +### Dedup / upload +- **`attach_manifest` reshapes chunk sizes into a throwaway `Vec` per upload** + — carry sizes as `i64` end-to-end (validated in an earlier draft of + `bench_round25_micro` §M4; deferred because it threads a type change through + `ChunkIngestOutcome`, the delta/stream/legacy paths and the `total_size` sums — + its own pass). (conf 0.5) +- **`store_from_stream` rebuilds the distinct-hash set the CDC loop already held** + as `pinned ∪ written` (`dedup_service.rs:499`/`distinct_hashes`) — return the list + the ingest already owns instead of an O(N) rescan + HashSet + N clones. ROUND14 + deferred. (conf 0.55) +- **Whole-file dedup-hit fast path is 3 serial manifest round-trips** (owner-check + + metadata SELECT + ref-bump UPDATE) — fold metadata+bump into one + `UPDATE … RETURNING` (3→2), or the whole thing into one atomic statement (3→1, + also closes a TOCTOU). Authz-sensitive; needs a validated pass. (conf 0.55) +- **Ownership checks bind the caller UUID as text** (`to_string()` + `$2::uuid`) + instead of a native `Uuid`, unlike the sibling claimable/pin queries. (conf 0.5) +- **Delta-download authorize is 2 round-trips** (entitlement then sizes) foldable + into one entitlement-JOIN-blobs query. (conf 0.55) + +### HTTP / DAV emitters (allocations) +- **`format_oc_id` allocates a fresh `String` per NC PROPFIND/REPORT/trashbin row** + — thread a `format_oc_id_into(&mut String, …)` buffer like the href buffer already + in those loops. Multi-signature; ROUND20 deferred. (conf 0.85) +- **`search_service::suggest_with_perms` builds a full `FileDto`/`FolderDto` per + candidate** to read 5 fields, computing (and dropping) `etag` + `size_formatted` + Strings on every keystroke. (conf 0.75) +- **CardDAV whole-book GET accumulates a throwaway per-contact vCard `String`** into + an unsized buffer — wants a `write_vcard_into(&mut String, …)`. (conf 0.7) +- **NC REPORT/trashbin per-row href + NC avatar `HeaderMap` clone + `list_files_query` + `Query`** — the remaining H1/href-buffer items ROUND22 left. (conf 0.55–0.65) + +### Auth / global config +- **`foldhash` is already in the lockfile transitively** (via hashbrown), so the + long-deferred fast-hasher lead is nearly free: `foldhash::quality::RandomState` + (random-seeded, DoS-safe) for the attacker-controlled delta-upload hash sets, and + `foldhash::fast` for the trusted-key NC PROPFIND `favorite_ids`/`nc_id` maps. + Wall-gated (a hasher swap changes 0 allocations). (conf 0.7) +- **`tracing` has no `release_max_level` feature** — per-request `debug!`s in the + auth middleware and the authz `require()` granted path compile into release and + pay a runtime level check. `release_max_level_info` compiles them out (binary-size + + hot-path win) but silently disables `RUST_LOG=debug` on release builds — an + **operator-facing tradeoff** that wants a maintainer decision, so it is flagged, not + shipped. (conf 0.65) +- **`profile.release` uses `lto = "thin"`** while `profile.bench` already trusts + `lto = "fat"` — a last-slice hot-path + size win at the cost of link time. (conf 0.55) +- **`panic = "abort"` — VERIFIED UNSAFE, do not apply.** `text_extractor.rs:177` + relies on `catch_unwind` to survive `pdf-extract` panics on malformed PDFs, and + tokio's per-task panic isolation itself needs unwinding; under abort a single + hostile PDF (or any handler `.unwrap()`) becomes a whole-process crash. Keep + `panic = "unwind"`. (Recorded so a future pass doesn't re-open it.) (conf 0.9) + +### Frontend +- **Client folder-listing cache (`getCachedFolder`/`cacheFolder` + ETag) is dead + code** — never called; every folder navigation refetches the full body with + `cache:'no-store'` and no `If-None-Match`. Wire the SWR cache in (bandwidth + + instant paint on revisits). (conf 0.6) +- **Grouped listing views mount one `VirtualWindow` per section** — O(sections) + scroll listeners + `getBoundingClientRect` reads per scroll frame; hoist to one + shared tracker (the deferred "unify onto VirtualRows"). (conf 0.6) +- **`VirtualRows.offsets` prefix-sum, flat dotfile filter O(N²), `typeLabel` + per-call 13-entry object** — the residual per-page frontend rebuilds. (conf 0.5–0.65) + +--- + +## Environment / methodology + +- **M1–M3:** counting global allocator tracking BOTH alloc **count** and **bytes** + (`examples/bench_round25_micro.rs`), no Postgres. Each section is BEFORE + (verbatim replica of the shipped-before shape) vs AFTER (replica of the + shipped-after shape, which the source now matches), with a value-equivalence + assertion and a `GATE FAIL … rollback` `exit(1)` if the AFTER arm fails to beat + BEFORE on its gate metric (M1 gates on bytes/op — the RAM win; M2/M3 on allocs/op). + Tunables: `M1_ITERS` (2000), `PAYLOAD` (262144), `CHUNKS` (4000), `BENCH_ITERS` (200000). +- **Q1–Q2:** live dev **PostgreSQL 16** (schema from `migrations/`), reads + `DATABASE_URL` from `.env`. Each section seeds its own fixture (`bench25_*` / + `bench25-*` markers, torn down around the run), asserts an equivalence gate + (result set identical BEFORE vs AFTER — mismatch `exit(1)`s), and gates on p50 + wall strictly decreasing. Q1's `playlist_items.file_id` FK is bypassed during + seeding with `session_replication_role = replica` (superuser) purely to isolate + the query shape without a `storage.files` fixture. Tunables: `Q1_PLAYLISTS` (100), + `Q1_PASSES` (30), `Q2_CONTACTS` (1000), `Q2_PASSES` (20), `Q2_VCARD_KB` (8). +- Built with `RUSTFLAGS="-C target-cpu=x86-64-v3"` (the checked-in + `.cargo/config.toml` pins `target-cpu=native`, which `SIGILL`s on this session's + host under AVX-512 — see ROUND23/24; local build-flag override only, the config + is unchanged). +- Verified beyond the benches: `cargo fmt --all --check` clean, + `cargo clippy --features bench -- -D warnings` clean, and the contact / playlist / + encrypted-backend / delta-upload unit tests pass. diff --git a/examples/bench_round25_micro.rs b/examples/bench_round25_micro.rs new file mode 100644 index 00000000..d5c45428 --- /dev/null +++ b/examples/bench_round25_micro.rs @@ -0,0 +1,317 @@ +//! Round-25 CPU/alloc micro-pack (no Postgres). +//! +//! Same rule as ROUND2–24: each section is BEFORE (verbatim replica of the +//! shipped-before shape) vs AFTER (verbatim replica of the shipped-after shape, +//! which the source is then made to match), with a byte/-value equivalence gate +//! and a `GATE FAIL … rollback` check that `std::process::exit(1)`s if the AFTER +//! arm fails to beat its BEFORE — the round's roll-back rule encoded into the +//! benchmark. An AFTER that doesn't win is never applied to the source. +//! +//! [M1] `EncryptedBlobBackend::decrypt_bytes` decrypts "in place" per its own +//! doc comment — but `let mut ciphertext = encrypted.split_off(NONCE_SIZE)` +//! allocates a fresh `Vec` and memcpy's the ENTIRE ciphertext+tag (~1 MiB +//! per CDC chunk, up to a whole legacy blob) on every decrypted read. +//! ROUND11 §15 fixed only the encrypt side. AFTER copies the 12-byte nonce +//! and 16-byte tag to the stack, decrypts the middle in place via +//! `decrypt_in_place_detached`, and returns a zero-copy `Bytes::slice` +//! past the nonce — 0 extra allocations, 0 full-payload memcpy. The RAM +//! win is in BYTES: peak drops from ~2× to ~1× the payload. +//! +//! [M2] Delta commit (`delta_upload_service::commit_with_perms`) materializes +//! the per-occurrence chunk-hash list a THIRD time at the manifest bind +//! (`request.chunks.iter().map(|c| c.h.clone()).collect()`), even though +//! `request.chunks` is owned and dead after that line. AFTER move-unzips +//! (`request.chunks.into_iter().map(|c| (c.h, c.s)).unzip()`) — N 64-byte +//! hash-String clones → 0. +//! +//! [M3] `folder_handler::download_folder_zip{,_impl}` binds a +//! `Query>` as `_params` and discards it — pure +//! dead work: axum parses the whole query string into a `HashMap` + one +//! owned `String` key and value per param, all dropped unread. AFTER +//! removes the extractor (byte-identical response; the handler only reads +//! the path `id`). +//! +//! Run: +//! RUSTFLAGS="-C target-cpu=x86-64-v3" \ +//! cargo run --release --features bench --example bench_round25_micro +//! Tunables (env): BENCH_ITERS (200000), M1_ITERS (2000), CHUNKS (4000), +//! PAYLOAD (262144 bytes for the M1 decrypt payload). + +use std::alloc::{GlobalAlloc, Layout, System}; +use std::collections::HashMap; +use std::env; +use std::hint::black_box; +use std::sync::atomic::{AtomicU64, Ordering}; +use std::time::Instant; + +use aes_gcm::aead::{AeadInPlace, KeyInit, OsRng}; +use aes_gcm::{AeadCore, Aes256Gcm, Nonce}; +use bytes::Bytes; + +// ── Counting allocator: tracks BOTH alloc count and total bytes requested ──── +static ALLOC_CALLS: AtomicU64 = AtomicU64::new(0); +static ALLOC_BYTES: AtomicU64 = AtomicU64::new(0); + +struct CountingAlloc; + +unsafe impl GlobalAlloc for CountingAlloc { + unsafe fn alloc(&self, layout: Layout) -> *mut u8 { + ALLOC_CALLS.fetch_add(1, Ordering::Relaxed); + ALLOC_BYTES.fetch_add(layout.size() as u64, Ordering::Relaxed); + unsafe { System.alloc(layout) } + } + unsafe fn dealloc(&self, ptr: *mut u8, layout: Layout) { + unsafe { System.dealloc(ptr, layout) } + } + unsafe fn realloc(&self, ptr: *mut u8, layout: Layout, new_size: usize) -> *mut u8 { + ALLOC_CALLS.fetch_add(1, Ordering::Relaxed); + // A realloc that grows requests `new_size` fresh bytes. + ALLOC_BYTES.fetch_add(new_size as u64, Ordering::Relaxed); + unsafe { System.realloc(ptr, layout, new_size) } + } + unsafe fn alloc_zeroed(&self, layout: Layout) -> *mut u8 { + ALLOC_CALLS.fetch_add(1, Ordering::Relaxed); + ALLOC_BYTES.fetch_add(layout.size() as u64, Ordering::Relaxed); + unsafe { System.alloc_zeroed(layout) } + } +} + +#[global_allocator] +static GLOBAL: CountingAlloc = CountingAlloc; + +fn env_or(key: &str, default: T) -> T { + env::var(key) + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(default) +} + +#[derive(Clone, Copy)] +struct Measure { + ns: f64, + allocs: f64, + bytes: f64, +} + +/// Run `f` `iters` times, returning per-op wall ns, alloc count and alloc bytes. +fn measure(iters: u64, mut f: impl FnMut() -> T) -> Measure { + // warm + black_box(f()); + ALLOC_CALLS.store(0, Ordering::Relaxed); + ALLOC_BYTES.store(0, Ordering::Relaxed); + let start = Instant::now(); + for _ in 0..iters { + black_box(f()); + } + let ns = start.elapsed().as_nanos() as f64 / iters as f64; + let allocs = ALLOC_CALLS.load(Ordering::Relaxed) as f64 / iters as f64; + let bytes = ALLOC_BYTES.load(Ordering::Relaxed) as f64 / iters as f64; + Measure { ns, allocs, bytes } +} + +fn report(tag: &str, before: Measure, after: Measure) { + println!("## {tag}"); + println!("| arm | ns/op | allocs/op | bytes/op |"); + println!( + "| BEFORE | {:>12.1} | {:>11.2} | {:>11.0} |", + before.ns, before.allocs, before.bytes + ); + println!( + "| AFTER | {:>12.1} | {:>11.2} | {:>11.0} |", + after.ns, after.allocs, after.bytes + ); + println!( + "# {:.2}x wall · {:.2} fewer allocs/op · {:.0} fewer bytes/op\n", + before.ns / after.ns.max(0.0001), + before.allocs - after.allocs, + before.bytes - after.bytes + ); +} + +/// Roll-back gate: `exit(1)` unless AFTER strictly beats BEFORE on `metric`. +fn gate(tag: &str, metric: &str, before: f64, after: f64) { + if !(after < before) { + eprintln!("GATE FAIL [{tag}] {metric}: AFTER {after} !< BEFORE {before} — rollback"); + std::process::exit(1); + } +} + +const NONCE_SIZE: usize = 12; +const TAG_SIZE: usize = 16; + +// ── [M1] EncryptedBlobBackend::decrypt_bytes ───────────────────────────────── +// Build one ciphertext template `[nonce][ciphertext+tag]` and, per iteration, +// clone it (1 alloc, common to both arms) then decrypt via each shape. + +fn build_ciphertext(cipher: &Aes256Gcm, plaintext: &[u8]) -> Vec { + let nonce = Aes256Gcm::generate_nonce(&mut OsRng); + let mut out = Vec::with_capacity(NONCE_SIZE + plaintext.len() + TAG_SIZE); + out.extend_from_slice(nonce.as_slice()); + out.extend_from_slice(plaintext); + let tag = cipher + .encrypt_in_place_detached(&nonce, b"", &mut out[NONCE_SIZE..]) + .expect("encrypt"); + out.extend_from_slice(&tag); + out +} + +/// BEFORE: the shipped `split_off` shape — one fresh Vec + full memcpy. +fn decrypt_before(cipher: &Aes256Gcm, mut encrypted: Vec) -> Bytes { + let mut ciphertext = encrypted.split_off(NONCE_SIZE); + let nonce = Nonce::from_slice(&encrypted); + cipher + .decrypt_in_place(nonce, b"", &mut ciphertext) + .expect("decrypt"); + Bytes::from(ciphertext) +} + +/// AFTER: decrypt the middle in place, return a zero-copy slice past the nonce. +fn decrypt_after(cipher: &Aes256Gcm, mut encrypted: Vec) -> Bytes { + let len = encrypted.len(); + let mut nonce_buf = [0u8; NONCE_SIZE]; + nonce_buf.copy_from_slice(&encrypted[..NONCE_SIZE]); + let nonce = Nonce::from_slice(&nonce_buf); + let tag = aes_gcm::aead::Tag::::clone_from_slice(&encrypted[len - TAG_SIZE..]); + cipher + .decrypt_in_place_detached(nonce, b"", &mut encrypted[NONCE_SIZE..len - TAG_SIZE], &tag) + .expect("decrypt"); + encrypted.truncate(len - TAG_SIZE); + Bytes::from(encrypted).slice(NONCE_SIZE..) +} + +fn section_m1() { + let iters: u64 = env_or("M1_ITERS", 2000); + let payload_len: usize = env_or("PAYLOAD", 262_144); + let key = [7u8; 32]; + let cipher = Aes256Gcm::new_from_slice(&key).unwrap(); + let plaintext: Vec = (0..payload_len).map(|i| (i * 31 + 7) as u8).collect(); + let template = build_ciphertext(&cipher, &plaintext); + + // Equivalence: both arms recover the exact plaintext. + let a = decrypt_before(&cipher, template.clone()); + let b = decrypt_after(&cipher, template.clone()); + assert_eq!( + a.as_ref(), + plaintext.as_slice(), + "M1 BEFORE plaintext mismatch" + ); + assert_eq!( + b.as_ref(), + plaintext.as_slice(), + "M1 AFTER plaintext mismatch" + ); + assert_eq!(a, b, "M1 arms disagree"); + + let before = measure(iters, || decrypt_before(&cipher, template.clone())); + let after = measure(iters, || decrypt_after(&cipher, template.clone())); + report( + &format!("[M1] decrypt_bytes in place ({payload_len}-byte payload)"), + before, + after, + ); + // The RAM win: AFTER must allocate strictly fewer bytes (no ciphertext copy). + gate("M1", "bytes/op", before.bytes, after.bytes); +} + +// ── [M2] Delta commit chunk-hash list: clone vs move-unzip ─────────────────── +struct ChunkRefRep { + h: String, + s: u64, +} + +fn hex64(i: usize) -> String { + // 64-char hex, deterministic — mirrors a BLAKE3 chunk hash string. + let mut s = String::with_capacity(64); + for k in 0..32 { + use std::fmt::Write; + let _ = write!( + s, + "{:02x}", + (i.wrapping_mul(2_654_435_761).wrapping_add(k)) as u8 + ); + } + s +} + +fn section_m2() { + let n: usize = env_or("CHUNKS", 4000); + let iters: u64 = env_or("M2_ITERS", 400); + + // Equivalence check on one build. + let build = || -> Vec { + (0..n) + .map(|i| ChunkRefRep { + h: hex64(i), + s: (i as u64) * 7, + }) + .collect() + }; + let cb = build(); + let before_h: Vec = cb.iter().map(|c| c.h.clone()).collect(); + let before_s: Vec = cb.iter().map(|c| c.s).collect(); + let (after_h, after_s): (Vec, Vec) = + build().into_iter().map(|c| (c.h, c.s)).unzip(); + assert_eq!(before_h, after_h, "M2 hash arms differ"); + assert_eq!(before_s, after_s, "M2 size arms differ"); + + let before = measure(iters, || { + let chunks = build(); + let hh: Vec = chunks.iter().map(|c| c.h.clone()).collect(); + let ss: Vec = chunks.iter().map(|c| c.s).collect(); + (hh, ss) + }); + let after = measure(iters, || { + let chunks = build(); + let (hh, ss): (Vec, Vec) = chunks.into_iter().map(|c| (c.h, c.s)).unzip(); + (hh, ss) + }); + report( + &format!("[M2] delta-commit chunk-hash list ({n} chunks)"), + before, + after, + ); + gate("M2", "allocs/op", before.allocs, after.allocs); +} + +// ── [M3] folder_handler dead Query ────────────────────────────────── +// Replicates axum's `Query>` extraction (build an owned +// key+value map from the query string) vs no extractor. +fn parse_query_map(q: &str) -> HashMap { + let mut m = HashMap::new(); + for pair in q.split('&') { + if let Some((k, v)) = pair.split_once('=') { + m.insert(k.to_string(), v.to_string()); + } + } + m +} + +fn section_m3() { + let iters: u64 = env_or("BENCH_ITERS", 200_000); + // A representative query string a client might append (cache-buster etc.). + let q = "folder_id=8c1f0e2a-1234-4a5b-9c8d-abcdef012345&t=1720000000"; + + // Equivalence: the handler only ever needs the path id, never these params. + let before_map = parse_query_map(q); + assert!(before_map.contains_key("folder_id"), "M3 setup"); + + let before = measure(iters, || { + // BEFORE: axum builds and drops the map on every request. + let m = parse_query_map(black_box(q)); + black_box(m.len()) + }); + let after = measure(iters, || { + // AFTER: no extractor — nothing parsed. + black_box(()) + }); + report("[M3] folder download dead Query", before, after); + gate("M3", "allocs/op", before.allocs, after.allocs); +} + +fn main() { + println!("# Round-25 micro alloc/RAM pack\n"); + section_m1(); + section_m2(); + section_m3(); + println!("All Round-25 micro sections passed their gate."); +} diff --git a/examples/bench_round25_queries.rs b/examples/bench_round25_queries.rs new file mode 100644 index 00000000..08f17852 --- /dev/null +++ b/examples/bench_round25_queries.rs @@ -0,0 +1,381 @@ +//! Round-25 PostgreSQL query-shape pack — end-to-end round-trips + wall on the +//! live dev Postgres, with an equivalence gate (mismatch → `exit(1)`) mirroring +//! ROUND23's methodology. +//! +//! [Q1] `music_storage_adapter::list_public_playlists` is 1 + N round-trips: +//! one listing SELECT then one `SELECT COUNT(*) FROM audio.playlist_items` +//! per returned playlist (up to 101 at limit=100). AFTER folds the count +//! into the listing with a `LEFT JOIN … GROUP BY` — one round-trip. +//! Gate: AFTER wall < BEFORE wall AND identical (playlist → track_count). +//! +//! [Q2] The three REST contact listings `SELECT … vcard …` — the multi-KB +//! vCard TEXT (may embed a base64 PHOTO) — but every caller maps +//! Contact → ContactDto, which has NO vcard field, so it is fetched, +//! shipped over the wire, decoded into a String and dropped. AFTER omits +//! the vcard column (a lite mapper passes an empty string). Gate: AFTER +//! wall < BEFORE wall AND identical (id, full_name, photo_url) DTO fields. +//! +//! Run (needs the dev Postgres up; reads DATABASE_URL from .env): +//! RUSTFLAGS="-C target-cpu=x86-64-v3" \ +//! cargo run --release --features bench --example bench_round25_queries +//! Tunables (env): Q1_PLAYLISTS (100), Q1_PASSES (30), +//! Q2_CONTACTS (1000), Q2_PASSES (20), Q2_VCARD_KB (8) + +use std::env; +use std::time::Instant; + +use sqlx::postgres::PgPoolOptions; +use sqlx::{PgPool, Row}; +use uuid::Uuid; + +fn env_or(key: &str, default: T) -> T { + env::var(key) + .ok() + .and_then(|v| v.parse().ok()) + .unwrap_or(default) +} + +fn p50(mut s: Vec) -> f64 { + s.sort_by(|a, b| a.partial_cmp(b).unwrap()); + s[s.len() / 2] +} + +fn report(tag: &str, unit: &str, before: f64, after: f64, stmts_before: usize, stmts_after: usize) { + println!("## {tag}"); + println!("| arm | {unit:>16} | statements |"); + println!("| BEFORE | {before:>16.3} | {stmts_before:>10} |"); + println!("| AFTER | {after:>16.3} | {stmts_after:>10} |"); + println!( + "# {:.2}x wall · {} → {} round-trips\n", + before / after.max(1e-9), + stmts_before, + stmts_after + ); +} + +fn gate(tag: &str, metric: &str, before: f64, after: f64) { + if !(after < before) { + eprintln!("GATE FAIL [{tag}] {metric}: AFTER {after} !< BEFORE {before} — rollback"); + std::process::exit(1); + } +} + +async fn cleanup(pool: &PgPool) { + // Idempotent teardown (also clears fixtures a prior crashed run left). + let _ = sqlx::query("SET session_replication_role = default") + .execute(pool) + .await; + let _ = sqlx::query("DELETE FROM audio.playlist_items WHERE playlist_id IN (SELECT id FROM audio.playlists WHERE name LIKE 'bench25_pl_%')").execute(pool).await; + let _ = sqlx::query("DELETE FROM audio.playlists WHERE name LIKE 'bench25_pl_%'") + .execute(pool) + .await; + let _ = sqlx::query("DELETE FROM carddav.contacts WHERE uid LIKE 'bench25-%'") + .execute(pool) + .await; + let _ = sqlx::query("DELETE FROM carddav.address_books WHERE name = 'bench25_ab'") + .execute(pool) + .await; + let _ = sqlx::query("DELETE FROM auth.users WHERE email LIKE 'bench25-%@bench.invalid'") + .execute(pool) + .await; +} + +async fn seed_user(pool: &PgPool, tag: &str) -> Uuid { + sqlx::query_scalar( + "INSERT INTO auth.users (username, email, role) VALUES ($1, $2, 'user') RETURNING id", + ) + .bind(format!("bench25_{tag}")) + .bind(format!("bench25-{tag}@bench.invalid")) + .fetch_one(pool) + .await + .expect("seed user") +} + +// ── [Q1] Public-playlist listing: 1 + N COUNT vs one LEFT JOIN GROUP BY ─────── +async fn section_q1(pool: &PgPool) { + let n: usize = env_or("Q1_PLAYLISTS", 100); + let passes: usize = env_or("Q1_PASSES", 30); + let owner = seed_user(pool, "q1owner").await; + + // Seed N public playlists, playlist i carrying (i % 10) + 1 items. The + // playlist_items.file_id FK to storage.files is bypassed with replica role + // (superuser) so the query SHAPE can be isolated without a files fixture. + let mut conn = pool.acquire().await.expect("acquire"); + sqlx::query("SET session_replication_role = replica") + .execute(&mut *conn) + .await + .unwrap(); + let mut ids: Vec = Vec::with_capacity(n); + for i in 0..n { + let pid: Uuid = sqlx::query_scalar( + "INSERT INTO audio.playlists (id, name, owner_id, is_public) + VALUES (gen_random_uuid(), $1, $2, TRUE) RETURNING id", + ) + .bind(format!("bench25_pl_{i}")) + .bind(owner) + .fetch_one(&mut *conn) + .await + .expect("seed playlist"); + ids.push(pid); + for j in 0..((i % 10) + 1) { + sqlx::query( + "INSERT INTO audio.playlist_items (id, playlist_id, file_id, position) + VALUES (gen_random_uuid(), $1, gen_random_uuid(), $2)", + ) + .bind(pid) + .bind(j as i32) + .execute(&mut *conn) + .await + .expect("seed item"); + } + } + sqlx::query("SET session_replication_role = default") + .execute(&mut *conn) + .await + .unwrap(); + drop(conn); + + let limit = n as i64; + + // BEFORE: list (1) then one COUNT per playlist (N) → 1 + N round-trips. + let before_counts = { + let rows = sqlx::query( + "SELECT id FROM audio.playlists WHERE is_public = TRUE ORDER BY updated_at DESC LIMIT $1 OFFSET 0", + ) + .bind(limit) + .fetch_all(pool) + .await + .unwrap(); + let mut m: Vec<(Uuid, i64)> = Vec::with_capacity(rows.len()); + for r in &rows { + let pid: Uuid = r.get(0); + let c: (i64,) = + sqlx::query_as("SELECT COUNT(*) FROM audio.playlist_items WHERE playlist_id = $1") + .bind(pid) + .fetch_one(pool) + .await + .unwrap(); + m.push((pid, c.0)); + } + m.sort(); + m + }; + + // AFTER: one LEFT JOIN + GROUP BY → 1 round-trip. + let after_counts = { + let rows = sqlx::query( + "SELECT p.id, COUNT(pi.id) AS track_count + FROM audio.playlists p + LEFT JOIN audio.playlist_items pi ON pi.playlist_id = p.id + WHERE p.is_public = TRUE + GROUP BY p.id + ORDER BY p.updated_at DESC LIMIT $1 OFFSET 0", + ) + .bind(limit) + .fetch_all(pool) + .await + .unwrap(); + let mut m: Vec<(Uuid, i64)> = rows + .iter() + .map(|r| (r.get::(0), r.get::(1))) + .collect(); + m.sort(); + m + }; + + assert_eq!( + before_counts, after_counts, + "Q1 track_count mismatch BEFORE vs AFTER" + ); + + // Timed passes. + let mut before_ms = Vec::new(); + let mut after_ms = Vec::new(); + for _ in 0..passes { + let t = Instant::now(); + let rows = sqlx::query("SELECT id FROM audio.playlists WHERE is_public = TRUE ORDER BY updated_at DESC LIMIT $1 OFFSET 0").bind(limit).fetch_all(pool).await.unwrap(); + for r in &rows { + let pid: Uuid = r.get(0); + let _c: (i64,) = + sqlx::query_as("SELECT COUNT(*) FROM audio.playlist_items WHERE playlist_id = $1") + .bind(pid) + .fetch_one(pool) + .await + .unwrap(); + } + before_ms.push(t.elapsed().as_secs_f64() * 1e3); + + let t = Instant::now(); + let _rows = sqlx::query("SELECT p.id, COUNT(pi.id) FROM audio.playlists p LEFT JOIN audio.playlist_items pi ON pi.playlist_id = p.id WHERE p.is_public = TRUE GROUP BY p.id ORDER BY p.updated_at DESC LIMIT $1 OFFSET 0").bind(limit).fetch_all(pool).await.unwrap(); + after_ms.push(t.elapsed().as_secs_f64() * 1e3); + } + let b = p50(before_ms); + let a = p50(after_ms); + report( + &format!("[Q1] public-playlist listing ({n} playlists)"), + "p50 ms", + b, + a, + 1 + n, + 1, + ); + gate("Q1", "p50 ms", b, a); +} + +// ── [Q2] Contact listing: over-fetch vcard TEXT vs lite (no vcard) ──────────── +async fn section_q2(pool: &PgPool) { + let n: usize = env_or("Q2_CONTACTS", 1000); + let passes: usize = env_or("Q2_PASSES", 20); + let vcard_kb: usize = env_or("Q2_VCARD_KB", 8); + let owner = seed_user(pool, "q2owner").await; + let ab: Uuid = sqlx::query_scalar( + "INSERT INTO carddav.address_books (id, name, owner_id) VALUES (gen_random_uuid(), 'bench25_ab', $1) RETURNING id", + ) + .bind(owner) + .fetch_one(pool) + .await + .expect("seed address book"); + + // A realistic vCard body with an embedded base64 PHOTO of ~vcard_kb KiB. + let photo_blob = "A".repeat(vcard_kb * 1024); + for i in 0..n { + let vcard = format!( + "BEGIN:VCARD\nVERSION:3.0\nFN:Contact {i}\nEMAIL:c{i}@example.com\nPHOTO;ENCODING=b;TYPE=JPEG:{photo_blob}\nEND:VCARD" + ); + sqlx::query( + "INSERT INTO carddav.contacts (id, address_book_id, uid, full_name, photo_url, email, phone, address, vcard, etag) + VALUES (gen_random_uuid(), $1, $2, $3, $4, '[]'::jsonb, '[]'::jsonb, '[]'::jsonb, $5, $6)", + ) + .bind(ab) + .bind(format!("bench25-{i}")) + .bind(format!("Contact {i}")) + .bind(format!("https://example.com/p/{i}.jpg")) + .bind(&vcard) + .bind(format!("etag{i}")) + .execute(pool) + .await + .expect("seed contact"); + } + + // Lite DTO shape the REST listing actually keeps. + #[derive(PartialEq, Debug)] + struct LiteDto { + id: Uuid, + full_name: Option, + photo_url: Option, + } + + let before_select = "SELECT id, full_name, photo_url, vcard FROM carddav.contacts WHERE address_book_id = $1 ORDER BY full_name LIMIT $2"; + let after_select = "SELECT id, full_name, photo_url FROM carddav.contacts WHERE address_book_id = $1 ORDER BY full_name LIMIT $2"; + let limit = n as i64; + + // Equivalence: the kept DTO fields are identical whether or not vcard is read. + let before_dtos: Vec = { + let rows = sqlx::query(before_select) + .bind(ab) + .bind(limit) + .fetch_all(pool) + .await + .unwrap(); + rows.iter() + .map(|r| { + let _vcard: Option = r.get("vcard"); // fetched + decoded, then dropped + LiteDto { + id: r.get("id"), + full_name: r.get("full_name"), + photo_url: r.get("photo_url"), + } + }) + .collect() + }; + let after_dtos: Vec = { + let rows = sqlx::query(after_select) + .bind(ab) + .bind(limit) + .fetch_all(pool) + .await + .unwrap(); + rows.iter() + .map(|r| LiteDto { + id: r.get("id"), + full_name: r.get("full_name"), + photo_url: r.get("photo_url"), + }) + .collect() + }; + assert_eq!( + before_dtos, after_dtos, + "Q2 DTO fields mismatch BEFORE vs AFTER" + ); + + let mut before_ms = Vec::new(); + let mut after_ms = Vec::new(); + for _ in 0..passes { + let t = Instant::now(); + let rows = sqlx::query(before_select) + .bind(ab) + .bind(limit) + .fetch_all(pool) + .await + .unwrap(); + let mut sink = 0usize; + for r in &rows { + let v: Option = r.get("vcard"); + sink += v.map(|s| s.len()).unwrap_or(0); + let _d = LiteDto { + id: r.get("id"), + full_name: r.get("full_name"), + photo_url: r.get("photo_url"), + }; + } + std::hint::black_box(sink); + before_ms.push(t.elapsed().as_secs_f64() * 1e3); + + let t = Instant::now(); + let rows = sqlx::query(after_select) + .bind(ab) + .bind(limit) + .fetch_all(pool) + .await + .unwrap(); + for r in &rows { + let _d = LiteDto { + id: r.get("id"), + full_name: r.get("full_name"), + photo_url: r.get("photo_url"), + }; + } + after_ms.push(t.elapsed().as_secs_f64() * 1e3); + } + let b = p50(before_ms); + let a = p50(after_ms); + report( + &format!("[Q2] contact listing over-fetch vcard ({n} contacts, {vcard_kb} KiB vcard)"), + "p50 ms", + b, + a, + 1, + 1, + ); + gate("Q2", "p50 ms", b, a); +} + +#[tokio::main] +async fn main() { + dotenvy::dotenv().ok(); + let url = env::var("DATABASE_URL") + .or_else(|_| env::var("OXICLOUD_DB_CONNECTION_STRING")) + .expect("set DATABASE_URL — the dev Postgres URL"); + let pool = PgPoolOptions::new() + .max_connections(4) + .connect(&url) + .await + .expect("connect Postgres"); + + println!("# Round-25 PG query-shape pack — BEFORE/AFTER (live Postgres)\n"); + cleanup(&pool).await; + section_q1(&pool).await; + section_q2(&pool).await; + cleanup(&pool).await; + println!("All Round-25 query sections passed their gate."); +} diff --git a/src/application/services/delta_upload_service.rs b/src/application/services/delta_upload_service.rs index e2839ccb..b7400613 100644 --- a/src/application/services/delta_upload_service.rs +++ b/src/application/services/delta_upload_service.rs @@ -446,8 +446,12 @@ impl DeltaUploadService { ct if ct.is_empty() => "application/octet-stream".to_string(), ct => ct, }; - let chunk_hashes: Vec = request.chunks.iter().map(|c| c.h.clone()).collect(); - let chunk_sizes: Vec = request.chunks.iter().map(|c| c.s).collect(); + // `request.chunks` is owned and dead after this line (only + // `request.file_hash` is read below), so move the hashes out instead of + // cloning each 64-char hash a third time — the distinct set and the + // verification tuple already materialized it twice (benches/ROUND25.md §M2). + let (chunk_hashes, chunk_sizes): (Vec, Vec) = + request.chunks.into_iter().map(|c| (c.h, c.s)).unzip(); let attached = self .dedup .attach_manifest( diff --git a/src/infrastructure/adapters/music_storage_adapter.rs b/src/infrastructure/adapters/music_storage_adapter.rs index f8720e17..129f07a5 100644 --- a/src/infrastructure/adapters/music_storage_adapter.rs +++ b/src/infrastructure/adapters/music_storage_adapter.rs @@ -140,19 +140,18 @@ impl MusicStoragePort for MusicStorageAdapter { limit: i64, offset: i64, ) -> Result, DomainError> { + // One `LEFT JOIN … GROUP BY` instead of 1 listing + N per-playlist + // `COUNT(*)` round-trips (up to 101 at limit=100) — benches/ROUND25.md §Q1. let playlists = self .playlist_repository - .list_public_playlists(limit, offset) + .list_public_playlists_with_counts(limit, offset) .await?; - let mut result = Vec::new(); - for playlist in playlists { - let dto = PlaylistDto::from(playlist); - let track_count = self - .get_track_count(&uuid::Uuid::parse_str(&dto.id).unwrap()) - .await?; - result.push(dto.with_track_info(track_count, 0)); - } - Ok(result) + Ok(playlists + .into_iter() + .map(|(playlist, track_count)| { + PlaylistDto::from(playlist).with_track_info(track_count, 0) + }) + .collect()) } async fn user_has_access(&self, playlist_id: &str, user_id: Uuid) -> Result { diff --git a/src/infrastructure/repositories/pg/contact_pg_repository.rs b/src/infrastructure/repositories/pg/contact_pg_repository.rs index d6ab1e5a..bc4c0e2b 100644 --- a/src/infrastructure/repositories/pg/contact_pg_repository.rs +++ b/src/infrastructure/repositories/pg/contact_pg_repository.rs @@ -21,8 +21,28 @@ impl ContactPgRepository { Self { pool } } - /// Maps a database row to a Contact domain entity + /// Maps a database row to a Contact domain entity (reads the `vcard` column). fn row_to_contact(row: &sqlx::postgres::PgRow) -> Result { + Self::row_to_contact_with_vcard(row, row.get("vcard")) + } + + /// Maps a row whose SELECT omitted the `vcard` column — used by the REST + /// listings (paginated / search / by-group) whose `ContactDto` drops vcard + /// anyway, so the multi-KB vCard TEXT (which can embed a base64 PHOTO) is + /// never SELECTed, shipped over the wire, or allocated (benches/ROUND25.md + /// §Q2). The domain `Contact` keeps an empty vcard; these paths never + /// re-emit it. Do NOT use for CardDAV sync / whole-book export, which need + /// the round-trip vCard. + fn row_to_contact_lite(row: &sqlx::postgres::PgRow) -> Result { + Self::row_to_contact_with_vcard(row, String::new()) + } + + /// Shared row → `Contact` mapper; `vcard` is supplied by the caller so the + /// TEXT column can be omitted from listings that don't consume it. + fn row_to_contact_with_vcard( + row: &sqlx::postgres::PgRow, + vcard: String, + ) -> Result { // Decode each JSONB column straight into its typed Vec via // `sqlx::types::Json` (a single `serde_json::from_slice` pass over // the raw JSONB bytes) instead of `row.get::` + @@ -63,7 +83,7 @@ impl ContactPgRepository { row.get::, _>("photo_url"), row.get("birthday"), row.get("anniversary"), - row.get("vcard"), + vcard, row.get("etag"), row.get("created_at"), row.get("updated_at"), @@ -366,7 +386,7 @@ impl ContactRepository for ContactPgRepository { SELECT id, address_book_id, uid, full_name, first_name, last_name, nickname, email, phone, address, organization, title, notes, photo_url, - birthday, anniversary, vcard, etag, created_at, updated_at + birthday, anniversary, etag, created_at, updated_at FROM carddav.contacts WHERE address_book_id = $1 ORDER BY full_name, first_name, last_name @@ -387,7 +407,7 @@ impl ContactRepository for ContactPgRepository { let mut contacts = Vec::with_capacity(rows.len()); for row in &rows { - contacts.push(Self::row_to_contact(row)?); + contacts.push(Self::row_to_contact_lite(row)?); } Ok(contacts) } @@ -429,7 +449,7 @@ impl ContactRepository for ContactPgRepository { SELECT c.id, c.address_book_id, c.uid, c.full_name, c.first_name, c.last_name, c.nickname, c.email, c.phone, c.address, c.organization, c.title, c.notes, c.photo_url, - c.birthday, c.anniversary, c.vcard, c.etag, c.created_at, c.updated_at + c.birthday, c.anniversary, c.etag, c.created_at, c.updated_at FROM carddav.contacts c INNER JOIN carddav.group_memberships m ON c.id = m.contact_id WHERE m.group_id = $1 @@ -445,7 +465,7 @@ impl ContactRepository for ContactPgRepository { let mut contacts = Vec::with_capacity(rows.len()); for row in &rows { - contacts.push(Self::row_to_contact(row)?); + contacts.push(Self::row_to_contact_lite(row)?); } Ok(contacts) } @@ -462,9 +482,9 @@ impl ContactRepository for ContactPgRepository { SELECT id, address_book_id, uid, full_name, first_name, last_name, nickname, email, phone, address, organization, title, notes, photo_url, - birthday, anniversary, vcard, etag, created_at, updated_at + birthday, anniversary, etag, created_at, updated_at FROM carddav.contacts - WHERE address_book_id = $1 + WHERE address_book_id = $1 AND ( full_name ILIKE $2 OR first_name ILIKE $2 @@ -485,7 +505,7 @@ impl ContactRepository for ContactPgRepository { let mut contacts = Vec::with_capacity(rows.len()); for row in &rows { - contacts.push(Self::row_to_contact(row)?); + contacts.push(Self::row_to_contact_lite(row)?); } Ok(contacts) } diff --git a/src/infrastructure/repositories/pg/playlist_pg_repository.rs b/src/infrastructure/repositories/pg/playlist_pg_repository.rs index d0f5f95d..97659fdb 100644 --- a/src/infrastructure/repositories/pg/playlist_pg_repository.rs +++ b/src/infrastructure/repositories/pg/playlist_pg_repository.rs @@ -22,6 +22,21 @@ struct PlaylistRow { updated_at: DateTime, } +/// A public playlist row carrying its aggregated track count, produced by the +/// single `LEFT JOIN … GROUP BY` that replaces the per-playlist `COUNT(*)` N+1. +#[derive(FromRow)] +struct PublicPlaylistCountRow { + id: Uuid, + name: String, + description: Option, + owner_id: Uuid, + is_public: bool, + cover_file_id: Option, + created_at: DateTime, + updated_at: DateTime, + track_count: i64, +} + #[derive(FromRow)] struct PlaylistItemRow { id: Uuid, @@ -79,6 +94,52 @@ impl PlaylistPgRepository { pub fn pool(&self) -> &PgPool { &self.pool } + + /// Public playlists together with their track counts in a **single** + /// round-trip. Replaces the adapter's 1 + N shape (one listing SELECT then + /// one `SELECT COUNT(*) FROM audio.playlist_items` per returned playlist — + /// up to 101 round-trips at `limit = 100`) with one `LEFT JOIN … GROUP BY`, + /// backed by `idx_playlist_items_playlist_id` (benches/ROUND25.md §Q1). + pub async fn list_public_playlists_with_counts( + &self, + limit: i64, + offset: i64, + ) -> PlaylistRepositoryResult> { + let rows = sqlx::query_as::<_, PublicPlaylistCountRow>( + "SELECT p.id, p.name, p.description, p.owner_id, p.is_public, p.cover_file_id, \ + p.created_at, p.updated_at, COUNT(pi.id) AS track_count \ + FROM audio.playlists p \ + LEFT JOIN audio.playlist_items pi ON pi.playlist_id = p.id \ + WHERE p.is_public = TRUE \ + GROUP BY p.id \ + ORDER BY p.updated_at DESC LIMIT $1 OFFSET $2", + ) + .bind(limit) + .bind(offset) + .fetch_all(&*self.pool) + .await + .map_err(|e| { + DomainError::database_error(format!("Failed to list public playlists: {}", e)) + })?; + + rows.into_iter() + .map(|row| { + let track_count = row.track_count; + Playlist::with_id( + row.id, + row.name, + row.description, + row.owner_id, + row.is_public, + row.cover_file_id, + row.created_at, + row.updated_at, + ) + .map(|p| (p, track_count)) + .map_err(|e| DomainError::new(ErrorKind::InternalError, "Playlist", e.to_string())) + }) + .collect() + } } impl PlaylistRepository for PlaylistPgRepository { diff --git a/src/infrastructure/services/encrypted_blob_backend.rs b/src/infrastructure/services/encrypted_blob_backend.rs index d8d0cb73..272d7806 100644 --- a/src/infrastructure/services/encrypted_blob_backend.rs +++ b/src/infrastructure/services/encrypted_blob_backend.rs @@ -86,7 +86,7 @@ impl EncryptedBlobBackend { /// Encrypt `data` into the on-disk layout: `[12-byte nonce][ciphertext + tag]`. /// -/// Single output buffer, mirroring the read side's `decrypt_in_place`: +/// Single output buffer, mirroring the read side's in-place detached decrypt: /// the payload is copied exactly once and encrypted in place with the tag /// appended. The old shape let `cipher.encrypt` allocate a full ciphertext /// `Vec` and then copied it a second time behind the nonce — one extra @@ -106,22 +106,39 @@ fn encrypt_bytes(cipher: &Aes256Gcm, data: &[u8]) -> Result /// Decrypt the on-disk layout `[nonce][ciphertext + tag]` **in place**. /// -/// Consumes the encrypted buffer and reuses it for the plaintext, so peak -/// RAM is one buffer — not ciphertext + plaintext side by side (which for -/// legacy whole-file blobs would double a multi-hundred-MB allocation). +/// Reuses the encrypted buffer for the plaintext, so peak RAM is one buffer — +/// not ciphertext + plaintext side by side (which for legacy whole-file blobs +/// would double a multi-hundred-MB allocation). The nonce and 16-byte GCM tag +/// are lifted to the stack, the ciphertext body is decrypted in place via the +/// detached API (mirroring the encrypt side's `encrypt_in_place_detached`), and +/// the plaintext is returned as a zero-copy `Bytes::slice` past the nonce. +/// +/// The prior shape did `encrypted.split_off(NONCE_SIZE)`, which allocated a +/// fresh `Vec` and memcpy'd the entire ciphertext (up to a whole legacy blob) +/// on every decrypted read — one full-payload allocation + copy the doc comment +/// above claimed did not happen (benches/ROUND25.md §M1; ROUND11 §15 fixed only +/// the encrypt side). Output plaintext is byte-identical. fn decrypt_bytes(cipher: &Aes256Gcm, mut encrypted: Vec) -> Result { - if encrypted.len() < NONCE_SIZE { + let len = encrypted.len(); + if len < NONCE_SIZE + TAG_SIZE { return Err(DomainError::internal_error( "Encryption", - "encrypted blob too short (missing nonce)", + "encrypted blob too short (missing nonce/tag)", )); } - let mut ciphertext = encrypted.split_off(NONCE_SIZE); // `encrypted` keeps the nonce - let nonce = Nonce::from_slice(&encrypted); + // Nonce (first 12 bytes) and GCM tag (last 16 bytes) copied to the stack so + // the middle can be borrowed mutably for in-place decryption. + let mut nonce_buf = [0u8; NONCE_SIZE]; + nonce_buf.copy_from_slice(&encrypted[..NONCE_SIZE]); + let nonce = Nonce::from_slice(&nonce_buf); + let tag = aes_gcm::aead::Tag::::clone_from_slice(&encrypted[len - TAG_SIZE..]); cipher - .decrypt_in_place(nonce, b"", &mut ciphertext) + .decrypt_in_place_detached(nonce, b"", &mut encrypted[NONCE_SIZE..len - TAG_SIZE], &tag) .map_err(|e| DomainError::internal_error("Encryption", format!("decrypt failed: {e}")))?; - Ok(Bytes::from(ciphertext)) + // Plaintext now lives at `encrypted[NONCE_SIZE..len - TAG_SIZE]`; drop the + // tag and hand out a refcounted view past the nonce — no copy, no new alloc. + encrypted.truncate(len - TAG_SIZE); + Ok(Bytes::from(encrypted).slice(NONCE_SIZE..)) } /// Run a crypto closure inline for small payloads, on the blocking pool for diff --git a/src/interfaces/api/handlers/folder_handler.rs b/src/interfaces/api/handlers/folder_handler.rs index 5956adab..726b27b7 100644 --- a/src/interfaces/api/handlers/folder_handler.rs +++ b/src/interfaces/api/handlers/folder_handler.rs @@ -4,7 +4,6 @@ use axum::{ http::{Response, StatusCode, header}, response::IntoResponse, }; -use std::collections::HashMap; use std::sync::Arc; use crate::application::dtos::display_helpers::{ @@ -210,7 +209,6 @@ impl FolderHandler { State(state): State>, auth_user: AuthUser, Path(id): Path, - Query(_params): Query>, ) -> impl IntoResponse { tracing::info!("Downloading folder as ZIP: {}", id); @@ -421,9 +419,12 @@ pub async fn download_folder_zip( state: State>, auth_user: AuthUser, path: Path, - query: Query>, ) -> impl IntoResponse { - FolderHandler::download_folder_zip_impl(state, auth_user, path, query).await + // No `Query` extractor: the handler reads only the path `id`. axum ignores + // any query string when no extractor is present, so the response is + // byte-identical while a per-request HashMap + owned key/value Strings are + // no longer parsed and dropped (benches/ROUND25.md §M3). + FolderHandler::download_folder_zip_impl(state, auth_user, path).await } // ── GET /api/folders/{id}/resources ─────────────────────────────────────────