You are a constitutional council ranking individual git commits for ownership allocation. Compare these two commits. Decide which contributed more lasting value to the project. Judge substance, not spectacle: - Prefer correct, lasting design and real bugfixes over churn, formatting, renames, or generated noise. - Prefer clarity and necessity over sheer line count. A small precise change can beat a large diffuse one. - Do not favor a side merely because its patch is longer or noisier. - Weight what the change does for the project, not the contributor's name. Return ONLY a JSON object: {"winner": "A" or "B", "ratio": "N:M", "explanation": "..."} The explanation must cite concrete differences in the patches (1-3 sentences). Side A — contributor: tommy-mor Side A — commit message: [9ecc4e2e] nice Side A — unified diff (full patch): diff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs index 5675dcd2e28062acbe3c037b94b77c33adf32474..a18241f3ed7b4a248748fcefc799f9801ca62461 100644 --- a/server/src/api/rpc.rs +++ b/server/src/api/rpc.rs @@ -936,7 +936,15 @@ pub async fn handle_rpc_batch( line_ok(RpcResult::ForumThreads(rpc_list_forum_threads(&reduced, &room))) } RpcCommand::RoomCreate { slug, visibility } => { - match verify_bearer_principal(&headers, &*state.reduced.read().await) { + // Scope the first read so its guard drops before any nested `read().await` / `write().await`. + // A guard from `match verify(..., &*state.reduced.read().await)` would otherwise live for the + // whole `match` and deadlock here (tokio::sync::RwLock is not reentrant). + let principal = { + let reduced = state.reduced.read().await; + verify_bearer_principal(&headers, &*reduced) + }; + match principal { + Err((_, m)) => line_err(m, None), Ok(principal) => { let slug = slug.trim().to_lowercase(); if slug.is_empty() || slug.len() > 64 { @@ -994,7 +1002,6 @@ pub async fn handle_rpc_batch( } } } - Err((_, m)) => line_err(m, None), } } RpcCommand::RoomGrant { @@ -1002,17 +1009,28 @@ pub async fn handle_rpc_batch( username, capability, } => { - match verify_bearer_principal(&headers, &*state.reduced.read().await) { + let principal = { + let reduced = state.reduced.read().await; + verify_bearer_principal(&headers, &*reduced) + }; + match principal { Err((_, m)) => line_err(m, None), Ok(principal) => { - let reduced = state.reduced.read().await; - if !reduced.user_has_cap(&room, &principal, ThreadCapability::Manage) { + let can_manage = { + let reduced = state.reduced.read().await; + reduced.user_has_cap(&room, &principal, ThreadCapability::Manage) + }; + if !can_manage { line_err("requires Manage capability", None) } else { match parse_username(&username) { Err(msg) => line_err("invalid username", Some(msg)), Ok(target) => { - if !reduced.users_by_provider.values().any(|u| u == &target) { + let user_exists = { + let reduced = state.reduced.read().await; + reduced.users_by_provider.values().any(|u| u == &target) + }; + if !user_exists { line_err(format!("user @{target} not found"), None) } else { match parse_capability(&capability) { diff --git a/server/tests/integration.rs b/server/tests/integration.rs index 9162bd5bee84035e3908bfb9c8e201b7878ca339..b930120da09fe7d307f0411b84fb639fb8bd0b15 100644 --- a/server/tests/integration.rs +++ b/server/tests/integration.rs @@ -95,6 +95,23 @@ async fn test_healthz() { assert_eq!(response.text().await.unwrap(), "ok"); } +#[tokio::test] +async fn test_room_create_private_rpc() { + let (addr, _tmp, _log, _handle) = create_test_server().await; + let client = reqwest::Client::new(); + let batch = serde_json::json!([{ + "RoomCreate": { "slug": "secret-project", "visibility": "private" } + }]); + let body = rpc_batch(&client, addr, Some(&test_bearer()), batch).await; + let line = &body["results"][0]; + assert_eq!(line["ok"], true, "room create: {:?}", line); + let room_id = line["result"]["RoomCreated"]["room_id"].as_str().unwrap(); + assert!( + room_id.contains("/secret-project"), + "expected room_id to contain slug, got {room_id}" + ); +} + #[tokio::test] async fn test_index_page() { // HTML routes are offline during the auth-v3 refactor. diff --git a/test/auth.bb b/test/auth.bb index 611e04f1806ef81d678a91f499597fe55f691dd7..a67cc1de763434176d2f68e419dd37a8cd328395 100644 --- a/test/auth.bb +++ b/test/auth.bb @@ -176,7 +176,7 @@ (assert! (= 200 (:status poll)) "pending-session poll returns 200") (let [poll-json (json/parse-string (:body poll) true)] (assert! (:complete poll-json) "pending session complete=true") - (assert! (= "@bbuser" (:user poll-json)) "poll returns @bbuser") + (assert! (= "bbuser" (:user poll-json)) "poll returns stored username bbuser") (assert! (clojure.string/starts-with? (:token poll-json) "slug_") "poll returns bearer token") (println "\nwhoami…") @@ -184,7 +184,7 @@ :headers {"Authorization" (str "Bearer " (:token poll-json))})] (assert! (= 200 (:status who)) "whoami returns 200") (let [who-json (json/parse-string (:body who) true)] - (assert! (= "@bbuser" (:user who-json)) "whoami user is @bbuser")))))) + (assert! (= "bbuser" (:user who-json)) "whoami user is bbuser (stored form)")))))) (println "\nCLI: identity start → OAuth → identity poll → whoami…") (let [cli-home (str tmp-dir "/cli-home") @@ -208,7 +208,7 @@ (str "identity poll exits 0 (stderr: " (:err poll-proc) ")")) (let [poll-cli (json/parse-string (:out poll-proc) true)] (assert! (= "complete" (:phase poll-cli)) "identity poll --json phase") - (assert! (= "@cliuser" (:user poll-cli)) "CLI poll user") + (assert! (= "cliuser" (:user poll-cli)) "CLI poll user (stored form)") (assert! (clojure.string/starts-with? (:token poll-cli) "slug_") "CLI poll token") (let [token-path (str cli-home "/.config/slugsocial/token")] (assert! (fs/exists? token-path) "token written under isolated HOME") @@ -219,7 +219,7 @@ (assert! (zero? (:exit who-proc)) (str "whoami exits 0 (stderr: " (:err who-proc) ")")) (let [who-cli (json/parse-string (:out who-proc) true)] - (assert! (= "@cliuser" (:user who-cli)) "CLI whoami uses saved token")))))))) + (assert! (= "cliuser" (:user who-cli)) "CLI whoami uses saved token")))))))) (finally (when-some [s @!server] (common/kill-server s)) diff --git a/test/common.bb b/test/common.bb index ebb16c615a8b40e8830b5d5765d1e935a45538d0..4ed450377cdbb020f5ff164a812dfbbfc23c0f47 100644 --- a/test/common.bb +++ b/test/common.bb @@ -78,9 +78,20 @@ (defn start-server "Start the slugsocial-server binary with the given env map. - Returns the babashka.process map." - [server-bin env-map] - (p/process [server-bin] {:out :inherit :err :inherit :env env-map})) + Returns the babashka.process map. + + When `log-file` (string path) is provided, stdout and stderr are appended there + instead of inheriting the parent descriptors. Inheriting shared pipes while the + parent blocks on HTTP I/O can fill the pipe buffer and deadlock the server on log writes." + ([server-bin env-map] + (start-server server-bin env-map nil)) + ([server-bin env-map log-file] + (p/process [server-bin] + (if log-file + ;; Two string paths (same file): babashka.process can deref the process cleanly. + ;; :err :out + ProcessBuilder$Redirect breaks stream copying in deref/kill-server. + {:env env-map :out log-file :err log-file} + {:out :inherit :err :inherit :env env-map})))) (defn kill-server "Forcibly kill a server process (babashka.process map) and wait for it to exit." diff --git a/test/grants.bb b/test/grants.bb index 7066c1f2a57f39a362052f8524206dccf37bb7b1..793ba216f2de76a77eb20a76d08403db07ff411b 100644 --- a/test/grants.bb +++ b/test/grants.bb @@ -35,13 +35,14 @@ (defn- http-client [] (-> (java.net.http.HttpClient/newBuilder) (.followRedirects java.net.http.HttpClient$Redirect/ALWAYS) + (.connectTimeout (java.time.Duration/ofSeconds 15)) (.build))) (defn- http-get [url & {:keys [headers]}] (let [b (java.net.http.HttpRequest/newBuilder (java.net.URI/create url))] (doseq [[k v] (or headers {})] (.header b k v)) - (let [req (-> b (.GET) (.build)) + (let [req (-> b (.timeout (java.time.Duration/ofSeconds 60)) (.GET) (.build)) resp (.send (http-client) req (java.net.http.HttpResponse$BodyHandlers/ofString))] {:status (.statusCode resp) :body (.body resp)}))) @@ -52,6 +53,7 @@ (doseq [[k v] (or headers {})] (.header b k v)) (let [req (-> b + (.timeout (java.time.Duration/ofSeconds 60)) (.POST (java.net.http.HttpRequest$BodyPublishers/ofString body)) (.build)) resp (.send (http-client) req (java.net.http.HttpResponse$BodyHandlers/ofString))] @@ -67,6 +69,7 @@ b (java.net.http.HttpRequest/newBuilder (java.net.URI/create url))] (.header b "Content-Type" "application/x-www-form-urlencoded") (let [req (-> b + (.timeout (java.time.Duration/ofSeconds 60)) (.POST (java.net.http.HttpRequest$BodyPublishers/ofString pairs)) (.build)) resp (.send (http-client) req (java.net.http.HttpResponse$BodyHandlers/ofString))] Side B — contributor: tommy-mor Side B — commit message: [a2a83d76] plan Side B — unified diff (full patch): diff --git a/PLAN.md b/PLAN.md new file mode 100644 index 0000000000000000000000000000000000000000..38291545caebde2632e3fe409ea54c63abc30f5a --- /dev/null +++ b/PLAN.md @@ -0,0 +1,337 @@ +# sorter2 storage & memory plan + +## Goal + +Fit a **decent chunk of Reddit** into a **256MB** Fly VM while keeping the product simple: one Rust binary, no external database service. + +**RAM should be bounded by query shape**, not dataset size — ideally one rank-centrality graph in memory at a time, plus runtime overhead. + +**`events.jsonl` remains the main database.** Everything on disk elsewhere is a **rebuildable projection**. + +--- + +## Architecture (target) + +``` + ┌─────────────────────────────────┐ + │ events.jsonl (source of truth) │ + └───────────────┬─────────────────┘ + │ + append on every mutation + │ + ▼ + ┌──────────────────────────────────────────────┐ + │ apply event → durable projection (on disk) │ + │ (same semantics as today's in-memory reducer) │ + └──────────────────────────────────────────────┘ + │ │ + ▼ ▼ + ┌──────────────────┐ ┌──────────────────────────┐ + │ entity_payloads │ │ reducer state per scope │ + │ (fat Reddit JSON)│ │ nodes, children, edges, … │ + └──────────────────┘ └──────────────────────────┘ + │ │ + └────────┬───────────┘ + ▼ + ┌──────────────────────────────────────────────┐ + │ RAM per request (or small LRU cache) │ + │ • one scope's GroupState for rank-centrality │ + │ • children + EntityData (small views) │ + │ • RC scratch allocations │ + │ → compute → render → drop / evict │ + └──────────────────────────────────────────────┘ +``` + +This is **event sourcing / CQRS**: + +| Layer | Role | +|-------|------| +| **JSONL** | Canonical write log; audit; disaster recovery | +| **Durable (RocksDB)** | Materialized read model + payload store; rebuildable from JSONL | +| **RAM** | One (or few) hot scopes for ranking and render | + +Durable is **not** a second source of truth. If projection and log diverge, **stream JSONL and rebuild durable**. + +--- + +## What we have today (baseline) + +| Piece | Status | +|-------|--------| +| `events.jsonl` append-only log | ✓ source of truth | +| `EventLog::replay` streaming one event at a time | ✓ no full `Vec` at startup | +| `entity_store` / `entity_db` (RocksDB via `durable`) | ✓ fat payloads off-heap | +| `EntityData` in `GlobalTree` | ✓ small derived views in RAM | +| Full `GlobalTree` replayed at boot | ✗ all nodes, all scopes' `GroupState` in RAM | +| Rank-centrality | ✓ already scoped per parent; reads in-memory `GroupState` | + +**Measured RSS (release, ~395 imports, 2.8MB JSONL):** + +| Scenario | RSS | +|----------|-----| +| Empty data dir | ~14 MB (+ RocksDB baseline) | +| After boot with data | ~21–22 MB | +| Pre-offload (in-tree payloads + vec replay) | ~25–34 MB | + +Payload offload + streaming replay helped startup peak, but **the full in-memory reducer** is still the scaling ceiling. + +--- + +## What lives where (target) + +### JSONL (`events.jsonl`) + +All mutations, append-only: + +- `VoteRecorded` — scope, pair, ratios +- `EntityImported` — id, full upstream payload +- `NodeEnsured` — register path +- (legacy / other event types as present in log) + +### Durable / RocksDB (`{data_dir}/…`) + +Single embedded DB directory. Collections (names tentative): + +| Collection | Contents | Notes | +|------------|----------|-------| +| `entity_payloads` | `ItemId → JSON string` | **Done.** Fat Reddit API blobs | +| `nodes` | `ItemId → { data: EntityData, children: … }` | Small; no raw payload | +| `scopes/{parent}/…` | `GroupState` materialization | edges, voted_pairs, item_to_idx, recent_votes (capped) | + +Nested layout can follow durable's `Map → Map → Vec` patterns (see `durable/docs/001.md` Sorter sketch). + +### RAM + +| Resident | When | +|----------|------| +| Tokio, Axum, reqwest, RocksDB block cache (tuned) | always | +| **One scope slice** | per request (or LRU of few scopes, byte-capped) | +| Rank-centrality temporaries | during render for that scope | + +**Not** in RAM at steady state: all subreddits, all vote graphs, all payloads. + +--- + +## Write path + +Order matters: + +1. Append event to `events.jsonl` (must durably succeed first) +2. Apply event to durable projection (same logic as today's `apply_event`) +3. Invalidate / update in-memory scope cache if that scope is hot + +Journal worker already serializes votes disk → tree; extend to **disk → durable** instead of (eventually) **disk → full GlobalTree**. + +```rust +// conceptual +append(jsonl, event)?; +apply_to_durable(event)?; +scope_cache.invalidate(scope_for(event)); +``` + +On failure after (1): replay from log repairs projection on next boot or via `replay-index` command. + +--- + +## Read path + +For a page under parent scope `P` (e.g. `reddit.com/r/rust`): + +1. **Load scope** from durable (or scope LRU hit) + - `EntityData` + children for listing + - `GroupState` for ranking and pair selection +2. **Run rank-centrality** on that `GroupState` (requires RAM — that's fine) +3. **Render** +4. **Drop** scope from RAM or return to LRU + +Payload fetch (rare): `entity_store.get(id)` only when render needs fields not in `EntityData`. + +--- + +## Startup & recovery + +### Normal startup + +``` +open entity_db (RocksDB) +open / validate scope indexes in same DB +do NOT replay JSONL into RAM +serve requests (cold scopes loaded on demand) +``` + +### Rebuild projection + +``` +stream events.jsonl → apply_event → durable +(one line at a time; same as EventLog::replay today) +``` + +Run when: + +- First deploy of projection layer +- Detected corruption / missing durable dir +- Manual `replay-index` after restoring JSONL from backup + +JSONL is the only file you need to trust for recovery. + +--- + +## Scope cache (RAM bound) + +**Strict mode:** one scope in RAM at a time — simplest, lowest RAM. + +**Practical mode:** LRU cache with **byte budget** (e.g. 64–128MB for scopes on a 256MB VM): + +- Evict least-recently-used scope's `GroupState` + child views +- Reload from durable on next visit + +Eviction policy is independent of storage engine. + +--- + +## Rank-centrality + +No change to the algorithm. It already assumes a whole `GroupState` for one parent scope. + +Moving reducer to durable does **not** remove RC memory cost — it removes **holding every scope's graph at once**. + +Optional later: materialized score vectors on disk, invalidated on vote. Not required for v1 of this plan. + +--- + +## Durable mutations (future API) + +Separate **intent** from **apply** for batching and testability: + +```rust +let m = rankings.path().key("rust").key(day).push_end(score); +db.apply(m)?; // or batch.apply(&[m1, m2, m3]) +``` + +Benefits: + +- One RocksDB `WriteBatch` / one WAL flush per vote or import batch +- Serializable ops for tests +- Aligns with JSONL events at the app layer and storage ops at the durable layer + +Keep chained `entry().push()` as sugar over `path().…; apply()`. + +Type safety: use type-state path builders if we want compile-time nesting; erased `Vec` only if we accept runtime errors at apply. + +--- + +## Storage engine: RocksDB vs alternatives + +**Current choice:** RocksDB via vendored `durable/` workspace crate. + +| Engine | Verdict for sorter2 | +|--------|---------------------| +| **RocksDB** | Good default for LSM, prefix scans, write-heavy votes + bulk imports. C++ dep, tune block cache for 256MB. | +| **sled** | Pure Rust appeal; production reliability history gives pause. Not a priority switch. | +| **fjall** | Pure Rust LSM; evaluate with benchmarks if leaving RocksDB. | +| **redb** | Lighter pure Rust; fine for payload-only store; less ideal for heavy scattered writes across scopes. | +| **SQLite** | Wrong shape for nested fractal tree; OK for a single KV table only. | + +**Switching engines matters less than:** + +1. Batched durable writes +2. Not materializing full reducer in RAM +3. Scope-local load/evict + +RocksDB stays **narrow**: blob attic + materialized reducer projection. Not a replacement for JSONL. + +--- + +## Scaling story (before vs after) + +| | In-memory reducer (today) | Target (JSONL + durable projection) | +|--|---------------------------|-------------------------------------| +| **RAM grows with** | Total nodes + all scopes + (was) payloads | Hot scope count × scope size + runtime | +| **Disk grows with** | JSONL (+ entity_db payloads today) | JSONL + full durable projection | +| **Startup** | O(events) replay into RAM | O(1) open DB | +| **Fails when** | RSS > VM limit | Scope too large for one RC graph, or disk full | +| **Recovery** | Replay JSONL | Replay JSONL → rebuild durable | + +**Rough RAM per active subreddit (~500 posts, moderate votes):** + +| Component | Order of magnitude | +|-----------|-------------------| +| `EntityData` × 500 | 0.5–2 MB | +| `GroupState` | 1–5 MB | +| RC scratch | 1–5 MB | +| Runtime + tuned RocksDB | 15–25 MB | +| **Total one hot scope** | **~20–40 MB** | + +Multiple subreddits fit on 256MB with LRU eviction, not all resident at once. + +--- + +## Implementation phases + +### Phase 0 — Done + +- [x] Vend `durable` as workspace crate (`durable/`) +- [x] `entity_store`: payloads in `{data_dir}/entity_db` +- [x] Remove `entity_raw` from `NodeState` +- [x] `EventLog::replay`: stream JSONL, one `Event` at a time +- [x] `apply_event` in `state.rs` for replay semantics + +### Phase 1 — Durable projection (write path) + +- [ ] Single `Db` under `{data_dir}/store` (payloads + reducer) +- [ ] `apply_event` writes to durable collections (nodes, scopes) in addition to or instead of `GlobalTree` +- [ ] Journal + Reddit import paths use same apply +- [ ] Batched writes where possible (one flush per vote / per import batch) +- [ ] Tests: apply event → read back from durable + +### Phase 2 — Stop full tree at boot + +- [ ] Startup: open durable only; no `GlobalTree::new()` + full replay into RAM +- [ ] `replay-index` command / flag: stream JSONL → durable (offline rebuild) +- [ ] Integration tests use replay-index fixture or temp DB + +### Phase 3 — Scope load on read + +- [ ] `ScopeView { parent, children, group }` loaded from durable +- [ ] HTML / vote / pair / fetch handlers take `ScopeView` instead of `&GlobalTree` +- [ ] Remove or shrink `Arc>` + +### Phase 4 — Scope cache + +- [ ] LRU with byte budget for hot scopes +- [ ] Invalidate on write to that scope +- [ ] Metrics: cache hit/miss, evictions, scope load time + +### Phase 5 — Durable API polish (optional) + +- [ ] Reified path/mutation API + `WriteBatch` integration in `durable` +- [ ] RocksDB tuning preset for 256MB Fly (`block_cache`, write buffers) +- [ ] Document `replay-index` in AGENTS.md + +--- + +## Non-goals (for now) + +- Distributed replication or multi-writer +- Replacing JSONL as canonical store +- SQL query layer over state +- Materialized rank scores on disk (unless RC latency forces it) +- Switching from RocksDB to sled without benchmarks + +--- + +## Open questions + +1. **One DB or two?** `{data_dir}/entity_db` today vs single `{data_dir}/store` — merge on Phase 1? +2. **Scope key encoding** — string `ItemId` paths vs hashed; must match event `scope` field. +3. **Child scope wiring** — Reddit `apply_entity_under_parent` creates children without full path ensure; durable schema must preserve this. +4. **Clojure smoke tests** — still read `events.jsonl`; durable is internal. No change expected. +5. **Fly volume** — `/data` holds JSONL + RocksDB; monitor disk alongside RAM. + +--- + +## Summary + +**JSONL = write log. Durable = full reducer on disk + fat payloads. RAM = one rank-centrality graph (or a small LRU of scopes).** + +Option B from design discussions is not separate from “reducer in durable” — it **is** reducer in durable, with JSONL still owning writes and recovery. The incremental work is: projection on apply, drop full tree at boot, scope load on read, then cache with eviction.