constitution · epochs · watch · epoch 3

comparison

c_9608dc0d38ab (tommy-mor) vs c_6c64824b0d83 (tommy-mor)

download prompt · raw event · cmp_673735bce74f81

council reasoning

~anthropic/claude-sonnet-latest · winner A · 8:2 · permalink

Side A fixes an actual deadlock bug (holding a RwLock read guard across nested lock acquisitions) with a real code change, adds a regression test, and improves test infrastructure (timeouts, log-file redirection to avoid pipe deadlocks). Side B is purely a planning/design document (PLAN.md) with no code changes, providing aspirational value but no immediate functional or bugfix contribution.

~x-ai/grok-latest · winner A · 3:1 · permalink

Commit A fixes a real production deadlock (non-reentrant tokio RwLock guards held across nested awaits in RoomCreate/RoomGrant) plus related test harness deadlocks from inherited pipes, and lands a working integration test; commit B only adds a PLAN.md design doc with no executable change.

openai/gpt-chat-latest · winner A · 9:1 · permalink

Side A fixes a substantive concurrency bug by shortening the lifetime of `tokio::sync::RwLock` read guards before nested `read().await`/`write().await`, preventing deadlocks in `RoomCreate` and `RoomGrant`. It also adds an integration test for private room creation and improves test reliability with timeout/logging adjustments, whereas Side B only adds a detailed architectural planning document without changing project behavior.

sides

A — c_9608dc0d38ab (tommy-mor)

message

[9ecc4e2e] nice

diff preview

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))]

download full diff A

B — c_6c64824b0d83 (tommy-mor)

message

[a2a83d76] plan

diff preview

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<Event>` 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<Op>` 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<RwLock<GlobalTree>>`
+
+### 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

… preview truncated; 478 characters omitted

download full diff B

Hardlinks — judgments / attempts / prompt

prompt download

judgments

attempts

Prompt text is loaded only by the download route.