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: [1531154d] dequeue -> vec Side A — unified diff (full patch): 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 = 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, pub edges: HashMap<(usize, usize), f64>, pub voted_pairs: HashSet<(usize, usize)>, - pub recent_votes: VecDeque, + pub recent_votes: Vec, } 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, } -/// 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>, /// Voted pairs `(min, max) -> true`. pub voted_pairs: Map>, - /// Recent votes, newest at the front (capped on write). - pub recent_votes: Deque>, + /// Recent votes, append-only oldest-first (cap applied on read). + pub recent_votes: List>, /// When ephemeral Reddit display content was last fetched (ms); absent after eviction. pub fetched_at: Leaf, } @@ -55,7 +55,7 @@ pub struct Store { pub view_meta: Map>, } -/// 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::, _>>()?; 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(); Side B — contributor: tommy-mor Side B — commit message: [4cd0d15d] more seed Side B — unified diff (full patch): 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) -> 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, 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::(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>>; + +#[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 = 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 = { + 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) + } +}