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.
constitution · epochs · watch · epoch 3
c_a896b2dc05d5 (tommy-mor) vs c_16438843de8f (tommy-mor)
download prompt · raw event · cmp_526d29d81887ba
council reasoning
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.
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();
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)
+ }
+}
Hardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.