Side B delivers a substantial, coherent feature: it converts entity fetch from an implicit auto-fetch to an explicit user-initiated, event-sourced import (EntityImported event, replay support, configurable API base for testing, dotenv support), plus real tests and an integration test harness with fixtures. Side A is a smaller refactor (Deque->Vec/List) that removes an eager trim optimization in favor of a query-time cap, which is reasonable but much lower impact and scope than B's architectural and testability improvements.
constitution · epochs · watch · epoch 3
c_a896b2dc05d5 (tommy-mor) vs c_e57094c6229a (tommy-mor)
download prompt · raw event · cmp_e3ea7f79519869
council reasoning
B adds durable EntityImported events, stores full API payloads for replay, switches Reddit loading to explicit user fetch with real integration coverage, and rewires the broker/UI around that model—lasting product and architecture value. A is a sound but narrower storage simplification (Deque→List/Vec, trim-on-write → cap-on-read) without comparable feature or durability impact.
Side B introduces a substantive new architecture for importing and persisting external entity data: it adds an `EntityImported` event, stores full API payloads for replay, refactors Reddit fetching to append events instead of mutating state directly, adds explicit user-triggered fetch UI, and includes replay and integration tests. Side A mainly replaces a durable deque with an append-only list, removes write-time trimming in favor of query-time capping, updates the reducer to use `Vec`, and adds a focused test; while useful, it is a narrower storage optimization compared with B's end-to-end feature and persistence design.
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_e57094c6229a (tommy-mor)
message
[40b975bf] nice
diff preview
diff --git a/Cargo.lock b/Cargo.lock
index 266e876bb7ccbe788beb1d5bd53ad5b45ee5825b..2cea973082716e761ef6f5dd5886acc08ff9aac0 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -222,6 +222,12 @@ dependencies = [
"syn",
]
+[[package]]
+name = "dotenvy"
+version = "0.15.7"
+source = "registry+https://github.com/rust-lang/crates.io-index"
+checksum = "1aaf95b3e5c8f23aa320147307562d361db0ae0d51242340f558153b4eb2439b"
+
[[package]]
name = "encoding_rs"
version = "0.8.35"
@@ -1238,6 +1244,7 @@ version = "0.0.1"
dependencies = [
"axum",
"axum-extra",
+ "dotenvy",
"maud",
"reqwest",
"serde",
diff --git a/server/Cargo.toml b/server/Cargo.toml
index 4677fedcb45292eebebe7e9cf6ce2f5738f18ddf..bd600138b613bd0f546bdec217a5334cdcb20aa5 100644
--- a/server/Cargo.toml
+++ b/server/Cargo.toml
@@ -17,6 +17,7 @@ tower-http = { version = "0.5", features = ["trace"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
reqwest = { version = "0.12", features = ["json"] }
+dotenvy = "0.15"
[dev-dependencies]
reqwest = { version = "0.12", features = ["json"] }
diff --git a/server/src/api/ui_html.rs b/server/src/api/ui_html.rs
index d2024bd4582bcc8482b461b2ba4fedbd8bff7c66..b33a84e8bb5e817b26592868d88090e6d664d950 100644
--- a/server/src/api/ui_html.rs
+++ b/server/src/api/ui_html.rs
@@ -6,7 +6,7 @@ use axum::{
use std::collections::HashMap;
use crate::{
- html::{input_panel, js_string_literal, ranking_panel, JsBuilder},
+ html::{entity_section, input_panel, js_string_literal, ranking_panel, JsBuilder},
parser::parse_reddit_url,
path_types::ItemId,
reddit::ensure_partial_tree,
@@ -87,6 +87,20 @@ pub async fn post_ui_html(
.into_response()
}
},
+ HtmlUiAction::FetchEntity { item } => {
+ let id = parse_item_param(&item);
+ if id.is_root() {
+ return ui_js_warn("nothing to fetch for the root").into_response();
+ }
+ state.queue_entity_fetch(id.clone());
+ let tree = state.tree.read().await;
+ let empty = crate::reducer::NodeState::default();
+ let node = tree.get(&id).unwrap_or(&empty);
+ let panel = entity_section(&id, node, true);
+ JsBuilder::new()
+ .morph_selector("#entity-section", panel)
+ .into_response()
+ },
}
}
diff --git a/server/src/events.rs b/server/src/events.rs
index ed5be6b13b9d46e838831d6ce0f96f569b401730..07ce24b5e56cf72b0b442c3c3241efbf6c3b006a 100644
--- a/server/src/events.rs
+++ b/server/src/events.rs
@@ -1,4 +1,5 @@
use serde::{Deserialize, Serialize};
+use serde_json::Value;
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(tag = "type", rename_all = "snake_case")]
@@ -18,4 +19,10 @@ pub enum Event {
},
/// Register a node path in the fractal tree (no external fetch).
NodeEnsured { id: String },
+ /// Full upstream API payload for a node (domain-specific view derived at replay/render time).
+ EntityImported {
+ id: String,
+ ts: i64,
+ payload: Value,
+ },
}
diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs
index df6505021d9f446c2b453e20e3eb3cf696a111f9..caf1309c8d93b47104499c57f9cc35ee7631fbb9 100644
--- a/server/src/html/mod.rs
+++ b/server/src/html/mod.rs
@@ -10,6 +10,7 @@ use crate::{
form_template::template_json_compact,
path_types::ItemId,
ranking::{top_bottom, RankedItem},
+ reddit::is_fetchable,
reducer::{GroupState, NodeState},
state::AppState,
ui_action::UI_RPC_FIELD,
@@ -151,7 +152,7 @@ pub fn breadcrumb_path(item: &ItemId) -> Markup {
fn entity_panel(node: &NodeState) -> Markup {
html! {
@if let Some(data) = &node.data {
- section id="entity-panel" class="demo-panel entity-card" {
+ div id="entity-panel" class="entity-card" {
h2 { (data.title) }
@if let Some(author) = &data.author {
p class="muted small" { "by " (author) }
@@ -164,6 +165,42 @@ fn entity_panel(node: &NodeState) -> Markup {
}
}
+/// Reddit/API import control — only shown on fetchable pages; never auto-fires.
+pub fn fetch_entity_panel(item: &ItemId, has_data: bool, fetching: bool) -> Markup {
+ if !is_fetchable(item) {
+ return html! {};
+ }
+ let label = if fetching {
+ "Fetching…"
+ } else if has_data {
+ "Fetch more"
+ } else {
+ "Fetch from Reddit"
+ };
+ let rpc = template_json_compact(&serde_json::json!({
+ "action": "fetch_entity",
+ "item": item.as_str(),
+ }))
+ .expect("fetch_entity rpc template");
+ html! {
+ form method="post" action="/ui" id="fetch-entity-form" class="fetch-entity-form" {
+ input type="hidden" name=(UI_RPC_FIELD) value=(rpc);
+ button type="submit" class="btn-secondary" disabled=(fetching) { (label) }
+ }
+ }
+}
+
+/// Entity card + explicit fetch control (morphed as `#entity-section`).
+pub fn entity_section(item: &ItemId, node: &NodeState, fetching: bool) -> Markup {
+ let has_data = node.data.is_some();
+ html! {
+ section id="entity-section" class="demo-panel" {
+ (entity_panel(node))
+ (fetch_entity_panel(item, has_data, fetching))
+ }
+ }
+}
+
fn rank_list(label: &str, items: &[RankedItem], start_rank: usize) -> Markup {
html! {
@if !items.is_empty() {
@@ -260,7 +297,7 @@ async fn item_page(state: AppState, uri: Uri, item: ItemId) -> Markup {
h1 { "sorter" }
(input_panel("", None))
(breadcrumb_path(&item))
- (entity_panel(node))
+ (entity_section(&item, node, false))
(ranking_panel(&item, group))
};
layout("sorter2", body, views)
@@ -272,16 +309,5 @@ pub async fn home(State(state): State<AppState>, uri: Uri) -> impl IntoResponse
pub async fn browse(State(state): State<AppState>, uri: Uri) -> impl IntoResponse {
let item = ItemId::from_browse_uri(uri.path()).unwrap_or(ItemId::root());
- if item.as_str().starts_with("reddit.com") {
- let needs_fetch = {
- let tree = state.tree.read().await;
- tree.get(&item)
- .map(|n| n.data.is_none())
- .unwrap_or(true)
- };
- if needs_fetch {
- state.reddit.request_fetch(item.clone());
- }
- }
item_page(state, uri, item).await
}
diff --git a/server/src/main.rs b/server/src/main.rs
index c22ec6c9f5358e5ec99fb83210dc351938505a93..1f0cddc39302b35b0cd6a6219f44c9d59202facf 100644
--- a/server/src/main.rs
+++ b/server/src/main.rs
@@ -2,6 +2,10 @@ use sorter2_server::state::AppConfig;
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error + Send + Sync>> {
+ if std::env::var("SORTER2_SKIP_DOTENV").is_err() {
+ let _ = dotenvy::dotenv();
+ }
+
tracing_subscriber::fmt()
.with_env_filter(
tracing_subscriber::EnvFilter::try_from_default_env()
diff --git a/server/src/reddit.rs b/server/src/reddit.rs
index 90053ad03b1d7c8e94f325dd4ee64c2b4f7da900..ff0f01e57b18af878eb5be3efc47204a7673589d 100644
--- a/server/src/reddit.rs
+++ b/server/src/reddit.rs
@@ -6,11 +6,15 @@ use std::time::{Duration, Instant};
use reqwest::{header, Client, StatusCode};
use serde::Deserialize;
+use serde_json::Value;
use tokio::sync::{mpsc, RwLock};
use crate::{
+ event_log::EventLog,
+ events::Event,
+ html::now_ms,
path_types::ItemId,
- reducer::{EntityData, GlobalTree},
+ reducer::GlobalTree,
};
/// Bootstrap blank nodes along a URL path so breadcrumbs and voting work before fetch.
@@ -20,6 +24,8 @@ pub fn ensure_partial_tree(tree: &mut GlobalTree, id: &ItemId) {
pub struct RedditCommand {
pub id: ItemId,
+ /// User-initiated fetch bypasses the in-memory "recently fetched" cache.
+ pub force: bool,
}
#[derive(Clone)]
@@ -33,19 +39,31 @@ struct RedditCredentials {
client_secret: String,
}
+#[derive(Clone)]
+pub struct RedditApiConfig {
+ pub api_base: String,
+ pub oauth_base: String,
+ pub user_agent: String,
+ creds: Option<RedditCredentials>,
+}
+
struct OAuthToken {
access_token: String,
expires_at: Instant,
}
impl RedditBroker {
- pub fn spawn(tree: Arc<RwLock<GlobalTree>>, user_agent: &str) -> Self {
+ pub fn spawn(
+ tree: Arc<RwLock<GlobalTree>>,
+ event_log: Arc<EventLog>,
+ config: RedditApiConfig,
+ ) -> Self {
let (tx, rx) = mpsc::channel(100);
let mut headers = header::HeaderMap::new();
headers.insert(
header::USER_AGENT,
- header::HeaderValue::from_str(user_agent).expect("valid user agent"),
+ header::HeaderValue::from_str(&config.user_agent).expect("valid user agent"),
);
let client = Client::builder()
@@ -54,22 +72,38 @@ impl RedditBroker {
.build()
.expect("reqwest client");
- let creds = RedditCredentials::from_env();
- tokio::spawn(reddit_worker(rx, tree, client, creds));
+ tokio::spawn(reddit_worker(rx, tree, event_log, client, config));
Self { tx }
}
- /// Fire-and-forget: queue a fetch; worker updates the tree when done.
- pub fn request_fetch(&self, id: ItemId) {
- let _ = self.tx.try_send(RedditCommand { id });
+ /// Queue a fetch; drops when the channel is full (backpressure).
+ pub fn request_fetch(&self, id: ItemId, force: bool) {
+ let _ = self.tx.try_send(RedditCommand { id, force });
+ }
+}
+
+impl RedditApiConfig {
+ pub fn from_env() -> Self {
+ Self {
+ api_base: reddit_api_base(),
+ oauth_base: reddit_oauth_base(),
+ user_agent: default_user_agent(),
+ creds: RedditCredentials::from_env(),
+ }
}
}
impl RedditCredentials {
+ /// Reddit's OAuth docs call these "client id" and "client secret"; the app
+ /// registration UI often labels them "app id" / "app secret" — same values.
fn from_env() -> Option<Self> {
- let client_id = std::env::var("REDDIT_CLIENT_ID").ok()?;
- let client_secret = std::env::var("REDDIT_CLIENT_SECRET").ok()?;
+ let client_id = std::env::var("REDDIT_CLIENT_ID")
+ .or_else(|_| std::env::var("REDDIT_APP_ID"))
+ .ok()?;
+ let client_secret = std::env::var("REDDIT_CLIENT_SECRET")
+ .or_else(|_| std::env::var("REDDIT_APP_SECRET"))
+ .ok()?;
if client_id.is_empty() || client_secret.is_empty() {
return None;
}
@@ -80,29 +114,63 @@ impl RedditCredentials {
}
}
+pub fn reddit_api_base() -> String {
+ std::env::var("REDDIT_API_BASE").unwrap_or_else(|_| "https://www.reddit.com".into())
+}
+
+pub fn reddit_oauth_base() -> String {
+ std::env::var("REDDIT_OAUTH_BASE").unwrap_or_else(|_| "https://www.reddit.com".into())
+}
+
pub fn default_user_agent() -> String {
std::env::var("REDDIT_USER_AGENT").unwrap_or_else(|_| {
"web:sorter2.social:v0.0.1 (by /u/sorter2)".to_string()
})
}
+/// True when this node can be loaded from the Reddit JSON API.
+pub fn is_fetchable(id: &ItemId) -> bool {
+ !map_item_to_reddit_api(id, "https://example.com").is_empty()
+}
+
+/// Derive UI-facing fields from a stored payload (Reddit-specific when under reddit.com).
+pub fn entity_view_from_payload(id: &ItemId, payload: &Value) -> Option<crate::reducer::EntityData> {
+ if id.as_str().starts_with("reddit.com") {
+ return parse_reddit_view(id, payload);
+ }
+ None
+}
+
+/// Apply a full API payload to the in-memory tree (view derived for known domains).
+pub fn apply_entity_import(tree: &mut GlobalTree, id: &ItemId, payload: Value) {
+ let view = entity_view_from_payload(id, &payload);
+ tree.apply_entity_raw(id, payload, view);
+}
+
async fn red
… preview truncated; 22,673 characters omittedHardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.