constitution · epochs · watch · epoch 3

comparison

c_a896b2dc05d5 (tommy-mor) vs c_16438843de8f (tommy-mor)

download prompt · raw event · cmp_526d29d81887ba

council reasoning

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

Side A is a coherent, well-reasoned refactor (deque→vec, moving the recent-votes cap to query time) with clear rationale in comments and a new test verifying correctness, directly improving code clarity and correctness. Side B is a disorganized 'seed' commit dumping unrelated deployment configs (Dockerfile, fly.toml, deps.edn for Clojure) and duplicate/misplaced source files (event_log.rs, views.rs at repo root) with no tests or integration, offering little lasting architectural value.

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

A is a deliberate storage/reducer redesign: Deque→List/Vec, write-time trim removed for append-only votes with a query-time cap, schema bump, call-site cleanup, and a new cap test—lasting correctness and simpler projection apply. B mostly dumps seed scaffolding (root event_log/views, deps.edn, Dockerfile/fly.toml) that is foundational but rough/noisy versus A’s integrated design change.

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

Side A introduces a lasting redesign of recent-vote storage by replacing the durable deque with an append-only list, removing post-commit trimming, updating the reducer and loading logic to cap results at query time, bumping the storage schema version, and adding a regression test for the new behavior. Side B adds useful infrastructure (Dockerfile, Fly configuration, and new `EventLog`/`ViewStore` modules), but those additions are more ancillary than the core persistence and data-model changes in Side A.

sides

A — c_a896b2dc05d5 (tommy-mor)

message

[1531154d] dequeue -> vec

diff preview

diff --git a/server/src/projection_apply.rs b/server/src/projection_apply.rs
index 9c8990a8af927f35d3344c8d0872a516aba56b86..ad404bacb8bcdd5ae0e682cff97f974fd44528ea 100644
--- a/server/src/projection_apply.rs
+++ b/server/src/projection_apply.rs
@@ -6,8 +6,6 @@
 //! batch as the (non-idempotent) edge merges guarantees exactly-once application
 //! across replay.
 
-use std::collections::BTreeSet;
-
 use crate::{
     event_log::EventLogError,
     events::{Event, EventRecord},
@@ -44,7 +42,6 @@ pub fn apply_records(
 
     let db = projection_store.db();
     let mut batch = db.batch();
-    let mut vote_parents: BTreeSet<ItemId> = BTreeSet::new();
     let mut last_seq = 0u64;
 
     for record in records {
@@ -70,7 +67,6 @@ pub fn apply_records(
                     *ts,
                 )
                 .map_err(|e| EventLogError::Apply(e.to_string()))?;
-                vote_parents.insert(parent);
             }
             Event::NodeEnsured { id } => {
                 let parsed = parse_event_id(id)?;
@@ -85,11 +81,5 @@ pub fn apply_records(
         .commit_with(durable::Durability::DisableWal)
         .map_err(|e| EventLogError::Apply(e.to_string()))?;
 
-    for parent in vote_parents {
-        projection_store
-            .trim_recent_votes(&parent)
-            .map_err(|e| EventLogError::Apply(e.to_string()))?;
-    }
-
     Ok(())
 }
diff --git a/server/src/projection_store.rs b/server/src/projection_store.rs
index 8576d671f351004426207894ac35594ddb0f70cf..9a8953d010029d3639dc3987687554bab8b7663e 100644
--- a/server/src/projection_store.rs
+++ b/server/src/projection_store.rs
@@ -18,7 +18,7 @@ use crate::{
 
 const PROJECTION_CURSOR_KEY: &str = "cursor";
 const PROJECTION_SCHEMA_KEY: &str = "schema_version";
-const PROJECTION_SCHEMA_VERSION: u64 = 3;
+const PROJECTION_SCHEMA_VERSION: u64 = 4;
 
 #[derive(Debug, thiserror::Error)]
 pub enum ProjectionStoreError {
@@ -142,16 +142,6 @@ impl ProjectionStore {
         Ok(tree)
     }
 
-    /// Cap a node's recent-vote window after applying votes (best-effort, blind).
-    pub(crate) fn trim_recent_votes(&self, parent: &ItemId) -> Result<(), ProjectionStoreError> {
-        node(parent).recent_votes().truncate_back(
-            &self.db,
-            crate::storage_schema::RECENT_VOTES_CAP,
-            Durability::DisableWal,
-        )?;
-        Ok(())
-    }
-
     /// Cache Reddit display content outside the event log (must be evicted per policy).
     pub fn put_ephemeral_content(
         &self,
diff --git a/server/src/reducer.rs b/server/src/reducer.rs
index 0c75c85150bb9e5f578bbadf58b3e43f8a80be4b..759918b8c0eb8f8bf1ed0911d8877adaa55c8ea6 100644
--- a/server/src/reducer.rs
+++ b/server/src/reducer.rs
@@ -1,4 +1,4 @@
-use std::collections::{HashMap, HashSet, VecDeque};
+use std::collections::{HashMap, HashSet};
 
 use serde::{Deserialize, Serialize};
 
@@ -52,7 +52,7 @@ pub struct GroupState {
     pub idx_to_item: Vec<ItemId>,
     pub edges: HashMap<(usize, usize), f64>,
     pub voted_pairs: HashSet<(usize, usize)>,
-    pub recent_votes: VecDeque<VoteData>,
+    pub recent_votes: Vec<VoteData>,
 }
 
 impl GroupState {
@@ -62,7 +62,7 @@ impl GroupState {
             idx_to_item: Vec::new(),
             edges: HashMap::new(),
             voted_pairs: HashSet::new(),
-            recent_votes: VecDeque::with_capacity(200),
+            recent_votes: Vec::new(),
         }
     }
 
@@ -111,10 +111,7 @@ impl GroupState {
         self.add_edge_weight(b_idx, a_idx, w_a);
         self.add_edge_weight(a_idx, b_idx, w_b);
 
-        self.recent_votes.push_front(vote);
-        while self.recent_votes.len() > 200 {
-            self.recent_votes.pop_back();
-        }
+        self.recent_votes.push(vote);
     }
 }
 
diff --git a/server/src/storage_dto.rs b/server/src/storage_dto.rs
index 9dfb13c53efe4389277625a6ab3bfc18f566a453..3fd6db5cb909ac4896bd8a3ecace796de5f08781 100644
--- a/server/src/storage_dto.rs
+++ b/server/src/storage_dto.rs
@@ -39,7 +39,7 @@ pub struct StoredEntityDataV1 {
     pub link_url: Option<String>,
 }
 
-/// One vote stored in a node's `recent_votes` deque.
+/// One vote stored in a node's `recent_votes` list.
 #[derive(Debug, Clone, Serialize, Deserialize)]
 pub struct StoredVoteV1 {
     pub version: u32,
diff --git a/server/src/storage_schema.rs b/server/src/storage_schema.rs
index bd26e665e084b95b10fdfff091c31e8dc84d07b8..5d2bb1d56927fb61c7c6d2d8602bd6882327f862 100644
--- a/server/src/storage_schema.rs
+++ b/server/src/storage_schema.rs
@@ -2,13 +2,13 @@
 //! durable collections instead of one blob per node.
 //!
 //! A vote updates a handful of keys: a few edge-weight merges, a voted-pair flag,
-//! a recent-vote deque push, and child-link set entries. The in-memory
+//! a recent-vote list append, and child-link set entries. The in-memory
 //! [`crate::reducer::GroupState`] is reconstructed from these keys on read for
 //! rank-centrality.
 
 use std::collections::{BTreeSet, HashMap, HashSet};
 
-use durable::{Batch, Db, Deque, Durable, Leaf, Map, Sum};
+use durable::{Batch, Db, Durable, Leaf, List, Map, Sum};
 
 use crate::{
     path_types::ItemId,
@@ -38,8 +38,8 @@ pub struct NodeSchema {
     pub edges: Map<EdgeKey, Sum<f64>>,
     /// Voted pairs `(min, max) -> true`.
     pub voted_pairs: Map<PairKey, Leaf<bool>>,
-    /// Recent votes, newest at the front (capped on write).
-    pub recent_votes: Deque<Leaf<StoredVoteV1>>,
+    /// Recent votes, append-only oldest-first (cap applied on read).
+    pub recent_votes: List<Leaf<StoredVoteV1>>,
     /// When ephemeral Reddit display content was last fetched (ms); absent after eviction.
     pub fetched_at: Leaf<i64>,
 }
@@ -55,7 +55,7 @@ pub struct Store {
     pub view_meta: Map<String, Leaf<u64>>,
 }
 
-/// Cap on the per-node recent-vote window (matches the in-memory reducer).
+/// Max recent votes returned when loading a node (query-time cap only).
 pub const RECENT_VOTES_CAP: u64 = 200;
 
 fn id_key(id: &ItemId) -> String {
@@ -148,11 +148,14 @@ fn build_group_state(
         }
     }
 
-    // Deque is front=newest; in-memory VecDeque is also front=newest.
-    let mut recent_votes = std::collections::VecDeque::new();
-    for stored in np.recent_votes().iter(db)? {
-        recent_votes.push_back(decode_vote(stored).map_err(durable::Error::Deserialize)?);
-    }
+    // List is index order (oldest first); keep the newest RECENT_VOTES_CAP entries.
+    let stored = np.recent_votes().iter(db)?;
+    let cap = RECENT_VOTES_CAP as usize;
+    let start = stored.len().saturating_sub(cap);
+    let recent_votes = stored[start..]
+        .iter()
+        .map(|s| decode_vote(s.clone()).map_err(durable::Error::Deserialize))
+        .collect::<Result<Vec<_>, _>>()?;
 
     Ok(GroupState {
         item_to_idx,
@@ -248,7 +251,7 @@ pub fn vote_writes(
     };
     batch.write(pnode.voted_pairs().key(&(lo, hi)).set(&true));
 
-    // Recent votes (newest at front).
+    // Recent votes (append-only; cap on read).
     let stored = encode_vote(&VoteData {
         ts,
         a: a_id,
@@ -260,7 +263,7 @@ pub fn vote_writes(
         delegate: None,
         thread_tag: "default".to_string(),
     });
-    batch.push_front(&pnode.recent_votes(), &stored)?;
+    batch.push(&pnode.recent_votes(), &stored)?;
     Ok(())
 }
 
@@ -314,6 +317,35 @@ mod tests {
         assert!(load_node_state(&db, &parent).unwrap().is_none());
     }
 
+    #[test]
+    fn load_caps_recent_votes_at_query_time() {
+        let dir = tempfile::tempdir().unwrap();
+        let db = Db::open(dir.path()).unwrap();
+        let parent = ItemId::root();
+
+        let mut batch = db.batch();
+        for i in 0..RECENT_VOTES_CAP + 10 {
+            vote_writes(&mut batch, &parent, "alpha", "beta", 1, 0, i as i64).unwrap();
+        }
+        batch.commit().unwrap();
+
+        assert_eq!(
+            node(&parent).recent_votes().len(&db).unwrap(),
+            RECENT_VOTES_CAP + 10
+        );
+
+        let node_state = load_node_state(&db, &parent).unwrap().unwrap();
+        assert_eq!(node_state.local_ranking.recent_votes.len(), RECENT_VOTES_CAP as usize);
+        assert_eq!(
+            node_state.local_ranking.recent_votes.first().map(|v| v.ts),
+            Some(10)
+        );
+        assert_eq!(
+            node_state.local_ranking.recent_votes.last().map(|v| v.ts),
+            Some(RECENT_VOTES_CAP as i64 + 9)
+        );
+    }
+
     #[test]
     fn missing_node_is_none() {
         let dir = tempfile::tempdir().unwrap();

download full diff A

B — c_16438843de8f (tommy-mor)

message

[4cd0d15d] more seed

diff preview

diff --git a/Dockerfile b/Dockerfile
new file mode 100644
index 0000000000000000000000000000000000000000..9cb07c60cb0da063f747cfbf1b3b876ecb8ba03e
--- /dev/null
+++ b/Dockerfile
@@ -0,0 +1,34 @@
+# time 0.3.47+ requires Rust 1.88 (edition 2024)
+FROM rust:1.88-slim as builder
+
+WORKDIR /build
+
+RUN apt-get update && \
+    apt-get install -y pkg-config libssl-dev && \
+    rm -rf /var/lib/apt/lists/*
+
+# Copy source and build. (Keep it simple to avoid remote build cache oddities.)
+COPY . .
+RUN cargo build --release --package slugsocial-server
+
+FROM debian:bookworm-slim
+
+RUN apt-get update && \
+    apt-get install -y ca-certificates && \
+    rm -rf /var/lib/apt/lists/*
+
+WORKDIR /app
+
+COPY --from=builder /build/target/release/slugsocial-server /app/slugsocial-server
+
+# Create data directory for persistent volume
+RUN mkdir -p /data
+
+ENV SLUG_DATA_DIR=/data
+ENV SLUG_EVENT_LOG=/data/events.jsonl
+ENV PORT=8080
+
+EXPOSE 8080
+
+CMD ["/app/slugsocial-server"]
+
diff --git a/deps.edn b/deps.edn
new file mode 100644
index 0000000000000000000000000000000000000000..0bf892d44f491cb2313e01ae8a942c3097c52948
--- /dev/null
+++ b/deps.edn
@@ -0,0 +1,10 @@
+{:paths ["." "test"]
+ :deps {cheshire/cheshire {:mvn/version "5.13.0"}
+        http-kit/http-kit {:mvn/version "2.8.0"}
+        babashka/fs {:mvn/version "0.5.32"}
+        babashka/process {:mvn/version "0.6.25"}
+        com.blockether/spel {:mvn/version "0.7.11"}}
+ :aliases
+ {:kaocha {:extra-deps {lambdaisland/kaocha {:mvn/version "1.91.1392"}
+                        lambdaisland/kaocha-junit-xml {:mvn/version "1.17.101"}}
+          :main-opts ["-m" "kaocha.runner"]}}}
diff --git a/event_log.rs b/event_log.rs
new file mode 100644
index 0000000000000000000000000000000000000000..eaae0d495e43a45d6590603892265a62cc92906e
--- /dev/null
+++ b/event_log.rs
@@ -0,0 +1,83 @@
+use std::path::{Path, PathBuf};
+
+use tokio::{
+    fs::{self, OpenOptions},
+    io::{AsyncBufReadExt, AsyncWriteExt, BufReader},
+};
+
+use crate::events::Event;
+
+#[derive(Debug, thiserror::Error)]
+pub enum EventLogError {
+    #[error("io error: {0}")]
+    Io(#[from] std::io::Error),
+    #[error("json error: {0}")]
+    Json(#[from] serde_json::Error),
+}
+
+#[derive(Debug, Clone)]
+pub struct EventLog {
+    path: PathBuf,
+}
+
+impl EventLog {
+    pub fn new(path: impl Into<PathBuf>) -> Self {
+        Self { path: path.into() }
+    }
+
+    pub fn path(&self) -> &Path {
+        &self.path
+    }
+
+    pub async fn ensure_parent_dir(&self) -> Result<(), EventLogError> {
+        if let Some(parent) = self.path.parent() {
+            fs::create_dir_all(parent).await?;
+        }
+        Ok(())
+    }
+
+    pub async fn append(&self, event: &Event) -> Result<(), EventLogError> {
+        self.ensure_parent_dir().await?;
+        let mut f: tokio::fs::File = OpenOptions::new()
+            .create(true)
+            .append(true)
+            .open(&self.path)
+            .await?;
+
+        let mut line = serde_json::to_string(event)?;
+        line.push('\n');
+        f.write_all(line.as_bytes()).await?;
+        f.flush().await?;
+        Ok(())
+    }
+
+    /// Load events from JSONL. Corrupt lines are skipped and returned as `(line_no, line)`.
+    pub async fn load_all(&self) -> Result<(Vec<Event>, Vec<(usize, String)>), EventLogError> {
+        if !fs::try_exists(&self.path).await? {
+            return Ok((vec![], vec![]));
+        }
+
+        let f = fs::File::open(&self.path).await?;
+        let mut reader = BufReader::new(f).lines();
+
+        let mut events = Vec::new();
+        let mut bad_lines = Vec::new();
+
+        let mut line_no: usize = 0;
+        while let Some(line) = reader.next_line().await? {
+            line_no += 1;
+            let trimmed = line.trim();
+            if trimmed.is_empty() {
+                continue;
+            }
+            match serde_json::from_str::<Event>(trimmed) {
+                Ok(ev) => events.push(ev),
+                Err(_) => bad_lines.push((line_no, line)),
+            }
+        }
+
+        Ok((events, bad_lines))
+    }
+}
+
+
diff --git a/fly.toml b/fly.toml
new file mode 100644
index 0000000000000000000000000000000000000000..bbb9345e527452db1d87a549213645c195eae5fc
--- /dev/null
+++ b/fly.toml
@@ -0,0 +1,42 @@
+app = "slugsocial"
+primary_region = "iad"
+
+[build]
+  dockerfile = "Dockerfile"
+
+[env]
+  SLUG_DATA_DIR = "/data"
+  SLUG_EVENT_LOG = "/data/events.jsonl"
+  PORT = "8080"
+
+[[services]]
+  internal_port = 8080
+  protocol = "tcp"
+
+  [[services.ports]]
+    port = 80
+    handlers = ["http"]
+    force_https = true
+
+  [[services.ports]]
+    port = 443
+    handlers = ["tls", "http"]
+
+  [services.concurrency]
+    type = "connections"
+    hard_limit = 1000
+    soft_limit = 500
+
+  [[services.http_checks]]
+    interval = "10s"
+    timeout = "2s"
+    grace_period = "5s"
+    method = "GET"
+    path = "/healthz"
+    protocol = "http"
+    tls_skip_verify = false
+
+[[mounts]]
+  source = "slugsocial_data"
+  destination = "/data"
+
diff --git a/views.rs b/views.rs
new file mode 100644
index 0000000000000000000000000000000000000000..d4f0ffc49475f014698b4da0de6f476884430813
--- /dev/null
+++ b/views.rs
@@ -0,0 +1,63 @@
+use std::{
+    collections::HashMap,
+    sync::{Arc, Mutex},
+};
+use tokio::sync::mpsc;
+
+type CountMap = Arc<Mutex<HashMap<String, u64>>>;
+
+#[derive(Clone)]
+pub struct ViewStore {
+    counts: CountMap,
+    flush_tx: mpsc::Sender<()>,
+}
+
+impl ViewStore {
+    pub fn new(json_path: &str) -> Self {
+        // Load existing counts from disk on startup (best-effort)
+        let initial: HashMap<String, u64> = std::fs::read_to_string(json_path)
+            .ok()
+            .and_then(|s| serde_json::from_str(&s).ok())
+            .unwrap_or_default();
+
+        let counts: CountMap = Arc::new(Mutex::new(initial));
+        let (flush_tx, mut flush_rx) = mpsc::channel::<()>(64);
+        let path = json_path.to_string();
+
+        let counts_for_writer = counts.clone();
+        tokio::spawn(async move {
+            while flush_rx.recv().await.is_some() {
+                while flush_rx.try_recv().is_ok() {}
+
+                let snapshot: HashMap<String, u64> = {
+                    counts_for_writer.lock().unwrap().clone()
+                };
+
+                let path = path.clone();
+                let _ = tokio::task::spawn_blocking(move || {
+                    if let Ok(json) = serde_json::to_string(&snapshot) {
+                        let tmp = format!("{path}.tmp");
+                        if std::fs::write(&tmp, &json).is_ok() {
+                            let _ = std::fs::rename(&tmp, &path);
+                        }
+                    }
+                })
+                .await;
+            }
+        });
+
+        Self { counts, flush_tx }
+    }
+
+    pub fn increment(&self, path: String) {
+        {
+            let mut map = self.counts.lock().unwrap();
+            *map.entry(path).or_insert(0) += 1;
+        }
+        let _ = self.flush_tx.try_send(());
+    }
+
+    pub fn get_views(&self, path: &str) -> u64 {
+        self.counts.lock().unwrap().get(path).copied().unwrap_or(0)
+    }
+}

download full diff B

Hardlinks — judgments / attempts / prompt

prompt download

judgments

attempts

Prompt text is loaded only by the download route.