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: [5db58b98] Improve vote pair selection for spanning trees and rank refinement. Prefer attaching unranked items to established components before comparing isolates, then zip down adjacent rank-centrality pairs once the pool is fully connected, skipping pairs that already have votes. Co-authored-by: Cursor Side A — unified diff (full patch): diff --git a/server/src/pair.rs b/server/src/pair.rs index 54b5d2417e9dba04ed8df422156e274c2b2f76b2..c14de4b0502c8b5a17cddf3746077739d56e03e0 100644 --- a/server/src/pair.rs +++ b/server/src/pair.rs @@ -3,13 +3,21 @@ //! Pair selection prefers **bridge** votes — comparisons between items in //! different connected components of the voted-pairs graph — so the pool //! merges into one ranking group before refining within it. +//! +//! Among unvoted bridges, prefer merging established voted components, then +//! attaching a never-voted child to an established component, and only then +//! comparing two never-voted children (so the voted graph grows as one tree). +//! +//! Once every pool child sits in one voted component, refinement **zips** down +//! the rank-centrality order: prefer 1 vs 2, then 2 vs 3, and so on, skipping +//! pairs that already have a vote. use rand::seq::SliceRandom; use std::collections::{HashMap, HashSet}; use crate::{ path_types::ItemId, - ranking::connected_components_from_voted_pairs, + ranking::{connected_components_from_voted_pairs, ranked_items}, reducer::{GlobalTree, GroupState}, }; @@ -28,36 +36,77 @@ fn pair_is_voted(group: &GroupState, a: &ItemId, b: &ItemId) -> bool { group.voted_pairs.contains(&(i, j)) } -/// Component id per pool item: voted-pairs graph components plus one id per -/// never-voted child. -fn component_ids(group: &GroupState, pool: &[ItemId]) -> HashMap { +/// Voted-pairs layout for pool items: component id per item plus which ids are +/// multi-node voted components (ranked groups in the UI). +struct ComponentLayout { + ids: HashMap, + established: HashSet, +} + +fn component_layout(group: &GroupState, pool: &[ItemId]) -> ComponentLayout { let n = group.idx_to_item.len(); let (comps, isolates) = connected_components_from_voted_pairs(n, group.voted_pairs.iter().copied()); - let mut out: HashMap = HashMap::new(); + let mut established = HashSet::new(); + let mut ids: HashMap = HashMap::new(); for (comp_idx, comp) in comps.iter().enumerate() { + if comp.len() >= 2 { + established.insert(comp_idx); + } for &idx in comp { if idx < n { - out.insert(group.idx_to_item[idx].clone(), comp_idx); + ids.insert(group.idx_to_item[idx].clone(), comp_idx); } } } let mut next = comps.len(); for &idx in &isolates { if idx < n { - out.insert(group.idx_to_item[idx].clone(), next); + ids.insert(group.idx_to_item[idx].clone(), next); next += 1; } } for item in pool { - out.entry(item.clone()).or_insert_with(|| { + ids.entry(item.clone()).or_insert_with(|| { let id = next; next += 1; id }); } - out + ComponentLayout { ids, established } +} + +/// Every pool child shares one multi-node voted component (spanning tree phase done). +fn pool_fully_connected(layout: &ComponentLayout, pool: &[ItemId]) -> bool { + if pool.len() < 2 { + return false; + } + let mut comp_id = None; + for item in pool { + let Some(id) = layout.ids.get(item) else { + return false; + }; + if !layout.established.contains(id) { + return false; + } + match comp_id { + None => comp_id = Some(*id), + Some(expected) if expected == *id => {} + _ => return false, + } + } + comp_id.is_some() +} + +/// Pool children that appear in `group`, sorted best rank first. +fn ranked_pool_order(group: &GroupState, pool: &[ItemId]) -> Vec { + let pool_set: HashSet<_> = pool.iter().collect(); + ranked_items(group) + .into_iter() + .map(|r| r.item) + .filter(|id| pool_set.contains(id)) + .collect() } #[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] @@ -72,19 +121,108 @@ enum PairPriority { WithinVoted = 3, } -fn pair_priority( +/// Tie-break among unvoted bridge pairs. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +enum BridgeSubPriority { + /// Both endpoints lie in established (multi-node) voted components. + MergeEstablished = 0, + /// One established component member and one never-voted child. + AttachIsolate = 1, + /// Two never-voted children (separate singleton components). + IsolatePair = 2, +} + +/// Tie-break among within-component pairs once the pool is one connected group. +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +struct WithinSubPriority { + /// 1 = adjacent ranks (i vs i+1); larger = farther apart in the order. + rank_gap: usize, + /// min rank index of the two — zip from the top (1 vs 2 before 2 vs 3). + zip_index: usize, +} + +const WITHIN_SUB_WORST: WithinSubPriority = WithinSubPriority { + rank_gap: usize::MAX, + zip_index: usize::MAX, +}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)] +struct PairSortKey { + priority: PairPriority, + bridge_sub: BridgeSubPriority, + within_sub: WithinSubPriority, +} + +fn item_in_established(layout: &ComponentLayout, item: &ItemId) -> bool { + layout + .ids + .get(item) + .is_some_and(|id| layout.established.contains(id)) +} + +fn bridge_sub_priority(layout: &ComponentLayout, a: &ItemId, b: &ItemId) -> BridgeSubPriority { + let a_est = item_in_established(layout, a); + let b_est = item_in_established(layout, b); + match (a_est, b_est) { + (true, true) => BridgeSubPriority::MergeEstablished, + (true, false) | (false, true) => BridgeSubPriority::AttachIsolate, + (false, false) => BridgeSubPriority::IsolatePair, + } +} + +fn within_sub_priority( + group: &GroupState, + pool: &[ItemId], + layout: &ComponentLayout, + a: &ItemId, + b: &ItemId, +) -> WithinSubPriority { + if !pool_fully_connected(layout, pool) { + return WITHIN_SUB_WORST; + } + let order = ranked_pool_order(group, pool); + let (Some(i), Some(j)) = (order.iter().position(|x| x == a), order.iter().position(|x| x == b)) + else { + return WITHIN_SUB_WORST; + }; + WithinSubPriority { + rank_gap: i.abs_diff(j), + zip_index: i.min(j), + } +} + +fn pair_sort_key( group: &GroupState, - components: &HashMap, + pool: &[ItemId], + layout: &ComponentLayout, a: &ItemId, b: &ItemId, -) -> PairPriority { +) -> PairSortKey { let voted = pair_is_voted(group, a, b); - let bridge = components.get(a) != components.get(b); - match (bridge, voted) { + let bridge = layout.ids.get(a) != layout.ids.get(b); + let priority = match (bridge, voted) { (true, false) => PairPriority::BridgeUnvoted, (false, false) => PairPriority::WithinUnvoted, (true, true) => PairPriority::BridgeVoted, (false, true) => PairPriority::WithinVoted, + }; + let bridge_sub = if priority == PairPriority::BridgeUnvoted { + bridge_sub_priority(layout, a, b) + } else { + BridgeSubPriority::MergeEstablished + }; + let within_sub = if matches!( + priority, + PairPriority::WithinUnvoted | PairPriority::WithinVoted + ) { + within_sub_priority(group, pool, layout, a, b) + } else { + WITHIN_SUB_WORST + }; + PairSortKey { + priority, + bridge_sub, + within_sub, } } @@ -109,9 +247,12 @@ fn candidate_pairs(pool: &[ItemId], exclude: Option<(&ItemId, &ItemId)>) -> Vec< /// Pick the next pair to vote on within `pool`. /// -/// 1. Prefer unvoted **bridge** pairs (connect separate ranking components). -/// 2. Then unvoted within-component pairs (refinement). -/// 3. Then already-voted pairs (re-compare). +/// 1. Prefer unvoted **bridge** pairs (connect separate ranking components), +/// with sub-priority: merge established components, attach an isolate to +/// established, then compare two isolates. +/// 2. Then unvoted within-component pairs; when the pool is one connected group, +/// prefer adjacent ranks (1 vs 2, 2 vs 3, …) in order, skipping voted pairs. +/// 3. Then already-voted pairs (re-compare), with the same zip ordering. pub fn suggest_next_pair_in_pool( group: &GroupState, pool: &[ItemId], @@ -121,15 +262,15 @@ pub fn suggest_next_pair_in_pool( if candidates.is_empty() { return None; } - let components = component_ids(group, pool); + let layout = component_layout(group, pool); let best = candidates .iter() - .map(|(a, b)| (pair_priority(group, &components, a, b), (a, b))) - .min_by_key(|(p, _)| *p)? + .map(|(a, b)| (pair_sort_key(group, pool, &layout, a, b), (a, b))) + .min_by_key(|(k, _)| *k)? .0; let best_pairs: Vec<(ItemId, ItemId)> = candidates .into_iter() - .filter(|(a, b)| pair_priority(group, &components, a, b) == best) + .filter(|(a, b)| pair_sort_key(group, pool, &layout, a, b) == best) .collect(); best_pairs.choose(&mut rand::thread_rng()).cloned() } @@ -303,6 +444,38 @@ mod tests { assert!(from_ab && from_cd, "expected bridge pair, got {:?}", chosen); } + #[test] + fn suggest_prefers_attach_over_isolate_pair_among_many_unranked() { + let parent = ItemId::parse("reddit.com/r/rust").unwrap(); + let mut tree = seed_children( + &parent, + &[ + "reddit.com/r/rust/a", + "reddit.com/r/rust/b", + "reddit.com/r/rust/c", + "reddit.com/r/rust/d", + "reddit.com/r/rust/e", + ], + ); + let ab = + VoteData::from_recorded(1, "reddit.com/r/rust/a", "reddit.com/r/rust/b", 2, 1).unwrap(); + tree.apply_vote(&parent, ab); + let group = tree.get(&parent).unwrap().local_ranking.clone(); + let pool = children_of(&tree, &parent); + let pair = suggest_next_pair_in_pool(&group, &pool, None).unwrap(); + let chosen = pair_set(&pair); + let from_ab = + chosen.contains("reddit.com/r/rust/a") || chosen.contains("reddit.com/r/rust/b"); + let from_cde = chosen.contains("reddit.com/r/rust/c") + || chosen.contains("reddit.com/r/rust/d") + || chosen.contains("reddit.com/r/rust/e"); + assert!( + from_ab && from_cde, + "expected ranked+unranked attach, got {:?}", + chosen + ); + } + #[test] fn suggest_connects_isolate_to_existing_component() { let parent = ItemId::parse("reddit.com/r/rust").unwrap(); @@ -325,6 +498,65 @@ mod tests { assert!(chosen.contains("reddit.com/r/rust/a") || chosen.contains("reddit.com/r/rust/b")); } + #[test] + fn suggest_zips_adjacent_ranks_when_tree_complete() { + let parent = ItemId::parse("reddit.com/r/rust").unwrap(); + let mut tree = seed_children( + &parent, + &[ + "reddit.com/r/rust/a", + "reddit.com/r/rust/b", + "reddit.com/r/rust/c", + ], + ); + // Star at a connects all three; b-c is the only unvoted adjacent pair left. + for (a, b, l, r) in [ + ("reddit.com/r/rust/a", "reddit.com/r/rust/b", 3, 1), + ("reddit.com/r/rust/a", "reddit.com/r/rust/c", 2, 1), + ] { + let v = VoteData::from_recorded(1, a, b, l, r).unwrap(); + tree.apply_vote(&parent, v); + } + let group = tree.get(&parent).unwrap().local_ranking.clone(); + let pool = children_of(&tree, &parent); + let pair = suggest_next_pair_in_pool(&group, &pool, None).unwrap(); + let chosen = pair_set(&pair); + // a-b and a-c voted; b-c is the only unvoted adjacent pair in rank order. + assert!(chosen.contains("reddit.com/r/rust/b")); + assert!(chosen.contains("reddit.com/r/rust/c")); + } + + #[test] + fn suggest_zip_prefers_1v2_before_2v3_when_both_unvoted() { + let parent = ItemId::parse("reddit.com/r/rust").unwrap(); + let mut tree = seed_children( + &parent, + &[ + "reddit.com/r/rust/a", + "reddit.com/r/rust/b", + "reddit.com/r/rust/c", + "reddit.com/r/rust/d", + ], + ); + // Hub at c connects all four; leave rank-adjacent a-b and b-c unvoted. + for (a, b, l, r) in [ + ("reddit.com/r/rust/c", "reddit.com/r/rust/d", 3, 1), + ("reddit.com/r/rust/b", "reddit.com/r/rust/c", 2, 1), + ("reddit.com/r/rust/a", "reddit.com/r/rust/c", 2, 1), + ] { + let v = VoteData::from_recorded(1, a, b, l, r).unwrap(); + tree.apply_vote(&parent, v); + } + let group = tree.get(&parent).unwrap().local_ranking.clone(); + let pool = children_of(&tree, &parent); + assert!(pool_fully_connected(&component_layout(&group, &pool), &pool)); + let pair = suggest_next_pair_in_pool(&group, &pool, None).unwrap(); + let chosen = pair_set(&pair); + // Top adjacent unvoted edge should be a-b (zip index 0), not b-c (index 1). + assert!(chosen.contains("reddit.com/r/rust/a")); + assert!(chosen.contains("reddit.com/r/rust/b")); + } + #[test] fn resolve_pair_picks_from_pool() { let parent = ItemId::parse("reddit.com/r/rust").unwrap(); Side B — contributor: tommy-mor Side B — commit message: [8d8230d1] reddit Side B — unified diff (full patch): diff --git a/.gitignore b/.gitignore index 4c7073f9fac0c30fd2050d79a60ef447af58ebeb..ada462e900d24a3a6d08165d158c80f79c35a5a5 100644 --- a/.gitignore +++ b/.gitignore @@ -7,3 +7,4 @@ data/ repomix-output.xml dev-data/ +.env diff --git a/server/Cargo.toml b/server/Cargo.toml index 7906a8547d56b8e6a48ef59c37aa82a8510fdee9..4677fedcb45292eebebe7e9cf6ce2f5738f18ddf 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -16,6 +16,7 @@ tower = "0.5" tower-http = { version = "0.5", features = ["trace"] } tracing = "0.1" tracing-subscriber = { version = "0.3", features = ["env-filter"] } +reqwest = { version = "0.12", features = ["json"] } [dev-dependencies] reqwest = { version = "0.12", features = ["json"] } diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs index 2864407ed6e8ec284a1dc663acf1805534566b62..df6505021d9f446c2b453e20e3eb3cf696a111f9 100644 --- a/server/src/html/mod.rs +++ b/server/src/html/mod.rs @@ -272,5 +272,16 @@ pub async fn home(State(state): State, uri: Uri) -> impl IntoResponse pub async fn browse(State(state): State, 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/reddit.rs b/server/src/reddit.rs index d203dca09245daf869b3aa942898447700ae69fb..90053ad03b1d7c8e94f325dd4ee64c2b4f7da900 100644 --- a/server/src/reddit.rs +++ b/server/src/reddit.rs @@ -1,4 +1,12 @@ -//! Reddit API import (async, decoupled from UI request path). +//! Reddit API import via a single background worker (rate limits, dedup, backoff). + +use std::collections::{HashMap, HashSet}; +use std::sync::Arc; +use std::time::{Duration, Instant}; + +use reqwest::{header, Client, StatusCode}; +use serde::Deserialize; +use tokio::sync::{mpsc, RwLock}; use crate::{ path_types::ItemId, @@ -10,12 +18,401 @@ pub fn ensure_partial_tree(tree: &mut GlobalTree, id: &ItemId) { tree.ensure_path(id); } -/// Placeholder for Reddit JSON import. Returns entity data when implemented. -pub async fn fetch_reddit_entity(_id: &ItemId) -> Option { - None +pub struct RedditCommand { + pub id: ItemId, +} + +#[derive(Clone)] +pub struct RedditBroker { + tx: mpsc::Sender, +} + +#[derive(Clone)] +struct RedditCredentials { + client_id: String, + client_secret: String, +} + +struct OAuthToken { + access_token: String, + expires_at: Instant, +} + +impl RedditBroker { + pub fn spawn(tree: Arc>, user_agent: &str) -> 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"), + ); + + let client = Client::builder() + .default_headers(headers) + .timeout(Duration::from_secs(15)) + .build() + .expect("reqwest client"); + + let creds = RedditCredentials::from_env(); + tokio::spawn(reddit_worker(rx, tree, client, creds)); + + 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 }); + } +} + +impl RedditCredentials { + fn from_env() -> Option { + let client_id = std::env::var("REDDIT_CLIENT_ID").ok()?; + let client_secret = std::env::var("REDDIT_CLIENT_SECRET").ok()?; + if client_id.is_empty() || client_secret.is_empty() { + return None; + } + Some(Self { + client_id, + client_secret, + }) + } +} + +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() + }) } -/// Apply fetched entity data to a node (called from async worker). -pub fn apply_entity(tree: &mut GlobalTree, id: &ItemId, data: EntityData) { - tree.set_entity_data(id, data); +async fn reddit_worker( + mut rx: mpsc::Receiver, + tree: Arc>, + client: Client, + creds: Option, +) { + let mut in_flight = HashSet::new(); + let mut recently_fetched: HashMap = HashMap::new(); + let mut current_delay = Duration::from_secs(1); + let mut oauth: Option = None; + let cache_ttl = Duration::from_secs(300); + + while let Some(cmd) = rx.recv().await { + let now = Instant::now(); + recently_fetched.retain(|_, t| now.duration_since(*t) < cache_ttl); + + if in_flight.contains(&cmd.id) || recently_fetched.contains_key(&cmd.id) { + continue; + } + + in_flight.insert(cmd.id.clone()); + let fetch_id = cmd.id.clone(); + + tokio::time::sleep(current_delay).await; + + if let Some(c) = &creds { + oauth = ensure_oauth_token(&client, c, oauth.take()).await; + } + + let token = oauth.as_ref().map(|t| t.access_token.as_str()); + let use_oauth = token.is_some(); + + match do_fetch(&client, &fetch_id, use_oauth, token).await { + Ok(FetchOutcome::Entity(data)) => { + let mut w = tree.write().await; + w.set_entity_data(&fetch_id, data); + recently_fetched.insert(fetch_id.clone(), Instant::now()); + current_delay = Duration::from_millis(600); + } + Ok(FetchOutcome::NotFound) => { + recently_fetched.insert(fetch_id.clone(), Instant::now()); + } + Ok(FetchOutcome::RateLimited { reset_secs }) => { + let wait = Duration::from_secs(reset_secs.max(1)); + tracing::warn!( + "Reddit rate limit for {}; sleeping {}s", + fetch_id, + wait.as_secs() + ); + tokio::time::sleep(wait).await; + current_delay = (current_delay * 2).min(Duration::from_secs(60)); + } + Err(e) => { + tracing::warn!("Reddit fetch failed for {}: {}", fetch_id, e); + current_delay = (current_delay * 2).min(Duration::from_secs(60)); + } + } + + in_flight.remove(&fetch_id); + } +} + +enum FetchOutcome { + Entity(EntityData), + NotFound, + RateLimited { reset_secs: u64 }, +} + +async fn ensure_oauth_token( + client: &Client, + creds: &RedditCredentials, + existing: Option, +) -> Option { + if let Some(t) = existing { + if Instant::now() < t.expires_at - Duration::from_secs(60) { + return Some(t); + } + } + + let resp = client + .post("https://www.reddit.com/api/v1/access_token") + .basic_auth(&creds.client_id, Some(&creds.client_secret)) + .form(&[("grant_type", "client_credentials")]) + .send() + .await; + + let resp = match resp { + Ok(r) => r, + Err(e) => { + tracing::warn!("Reddit OAuth token request failed: {e}"); + return None; + } + }; + + if !resp.status().is_success() { + tracing::warn!("Reddit OAuth token HTTP {}", resp.status()); + return None; + } + + #[derive(Deserialize)] + struct TokenResponse { + access_token: String, + expires_in: u64, + } + + let body: TokenResponse = match resp.json().await { + Ok(b) => b, + Err(e) => { + tracing::warn!("Reddit OAuth token parse failed: {e}"); + return None; + } + }; + + Some(OAuthToken { + access_token: body.access_token, + expires_at: Instant::now() + Duration::from_secs(body.expires_in), + }) +} + +async fn do_fetch( + client: &Client, + id: &ItemId, + use_oauth: bool, + bearer: Option<&str>, +) -> Result { + let url = map_item_to_reddit_api(id, use_oauth); + if url.is_empty() { + return Ok(FetchOutcome::NotFound); + } + + let mut req = client.get(&url); + if let Some(token) = bearer { + req = req.bearer_auth(token); + } + + let resp = req.send().await.map_err(|e| e.to_string())?; + + if resp.status() == StatusCode::TOO_MANY_REQUESTS { + let reset = rate_limit_reset_secs(&resp); + return Ok(FetchOutcome::RateLimited { reset_secs: reset }); + } + + if resp.status() == StatusCode::SERVICE_UNAVAILABLE { + return Err("Reddit unavailable (503)".to_string()); + } + + if !resp.status().is_success() { + return Ok(FetchOutcome::NotFound); + } + + if rate_limit_remaining(&resp) == Some(0) { + let reset = rate_limit_reset_secs(&resp); + return Ok(FetchOutcome::RateLimited { reset_secs: reset }); + } + + let bytes = resp.bytes().await.map_err(|e| e.to_string())?; + Ok(parse_reddit_json(id, &bytes) + .map(FetchOutcome::Entity) + .unwrap_or(FetchOutcome::NotFound)) +} + +fn rate_limit_remaining(resp: &reqwest::Response) -> Option { + resp.headers() + .get("x-ratelimit-remaining") + .and_then(|v| v.to_str().ok()) + .and_then(|s| s.parse::().ok()) + .map(|f| f.floor() as u64) +} + +fn rate_limit_reset_secs(resp: &reqwest::Response) -> u64 { + resp.headers() + .get("x-ratelimit-reset") + .and_then(|v| v.to_str().ok()) + .and_then(|s| s.parse::().ok()) + .map(|f| f.ceil() as u64) + .unwrap_or(5) +} + +/// Map canonical item id to Reddit JSON API URL. +pub fn map_item_to_reddit_api(id: &ItemId, oauth: bool) -> String { + let path = id.as_str(); + if !path.starts_with("reddit.com/") && path != "reddit.com" { + return String::new(); + } + + let base = if oauth { + "https://oauth.reddit.com" + } else { + "https://www.reddit.com" + }; + + let segments: Vec<&str> = path.split('/').collect(); + + if let Some(i) = segments.iter().position(|&p| p == "comments") { + if segments.len() > i + 1 { + let api_path = segments[1..=i + 1].join("/"); + return format!("{base}/{api_path}.json?raw_json=1"); + } + } + + if segments.len() == 3 && segments[1] == "r" { + return format!("{base}/r/{}/about.json?raw_json=1", segments[2]); + } + + String::new() +} + +fn parse_reddit_json(id: &ItemId, bytes: &[u8]) -> Option { + let v: serde_json::Value = serde_json::from_slice(bytes).ok()?; + let segments: Vec<&str> = id.as_str().split('/').collect(); + + if segments.iter().any(|&p| p == "comments") { + parse_post_listing(&v) + } else { + parse_subreddit_about(&v) + } +} + +fn parse_subreddit_about(v: &serde_json::Value) -> Option { + let data = v.get("data")?; + let title = data + .get("title") + .or_else(|| data.get("display_name")) + .and_then(|t| t.as_str())? + .to_string(); + let body_html = data + .get("public_description_html") + .or_else(|| data.get("public_description")) + .and_then(|t| t.as_str()) + .map(|s| s.to_string()); + let thumb_url = data + .get("icon_img") + .or_else(|| data.get("community_icon")) + .and_then(|t| t.as_str()) + .filter(|s| !s.is_empty()) + .map(|s| s.to_string()); + + Some(EntityData { + title, + author: None, + body_html, + thumb_url, + }) +} + +fn parse_post_listing(v: &serde_json::Value) -> Option { + let listing = v.as_array()?.first()?; + let child = listing + .pointer("/data/children/0/data")?; + let title = child.get("title")?.as_str()?.to_string(); + let author = child + .get("author") + .and_then(|a| a.as_str()) + .filter(|a| *a != "[deleted]") + .map(|s| s.to_string()); + let body_html = child + .get("selftext_html") + .and_then(|t| t.as_str()) + .filter(|s| !s.is_empty()) + .map(|s| s.to_string()); + let thumb_url = child + .get("thumbnail") + .and_then(|t| t.as_str()) + .filter(|s| s.starts_with("http")) + .map(|s| s.to_string()); + + Some(EntityData { + title, + author, + body_html, + thumb_url, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn map_subreddit_about_url() { + let id = ItemId::parse("reddit.com/r/rust").unwrap(); + assert_eq!( + map_item_to_reddit_api(&id, false), + "https://www.reddit.com/r/rust/about.json?raw_json=1" + ); + assert_eq!( + map_item_to_reddit_api(&id, true), + "https://oauth.reddit.com/r/rust/about.json?raw_json=1" + ); + } + + #[test] + fn map_post_url() { + let id = + ItemId::parse("reddit.com/r/amitheasshole/comments/1trnvdl").unwrap(); + assert_eq!( + map_item_to_reddit_api(&id, false), + "https://www.reddit.com/r/amitheasshole/comments/1trnvdl.json?raw_json=1" + ); + } + + #[test] + fn map_non_reddit_empty() { + let id = ItemId::opaque("example.com/foo"); + assert!(map_item_to_reddit_api(&id, false).is_empty()); + } + + #[test] + fn parse_subreddit_fixture() { + let json = r#"{"kind":"t5","data":{"title":"Rust","display_name":"rust","public_description":"systems"}}"#; + let entity = parse_reddit_json( + &ItemId::parse("reddit.com/r/rust").unwrap(), + json.as_bytes(), + ) + .unwrap(); + assert_eq!(entity.title, "Rust"); + } + + #[test] + fn parse_post_fixture() { + let json = r#"[{"kind":"Listing","data":{"children":[{"kind":"t3","data":{"title":"AITA","author":"op","selftext_html":"<p>hi</p>","thumbnail":"https://b.thumbs.redditmedia.com/x.jpg"}}]}}]"#; + let entity = parse_reddit_json( + &ItemId::parse("reddit.com/r/x/comments/abc").unwrap(), + json.as_bytes(), + ) + .unwrap(); + assert_eq!(entity.title, "AITA"); + assert_eq!(entity.author.as_deref(), Some("op")); + } } diff --git a/server/src/state.rs b/server/src/state.rs index cc1722f5a5bf4d415f2327ea585c488a15a75592..8c03aa60c15aee803a534439400b69935b1a3d84 100644 --- a/server/src/state.rs +++ b/server/src/state.rs @@ -5,9 +5,10 @@ use tokio::sync::RwLock; use crate::{ event_log::EventLog, events::Event, + journal::JournalClient, path_types::ItemId, + reddit::{default_user_agent, RedditBroker}, reducer::{GlobalTree, VoteData}, - journal::JournalClient, views::ViewStore, }; @@ -73,6 +74,7 @@ pub struct AppState { pub views: ViewStore, pub tree: Arc>, journal: JournalClient, + pub reddit: RedditBroker, } impl AppState { @@ -112,6 +114,7 @@ impl AppState { let tree = Arc::new(RwLock::new(tree)); let journal = JournalClient::spawn(tree.clone(), event_log.clone()); + let reddit = RedditBroker::spawn(tree.clone(), &default_user_agent()); Self { cfg: Arc::new(cfg), @@ -119,6 +122,7 @@ impl AppState { views, tree, journal, + reddit, } } @@ -127,8 +131,11 @@ impl AppState { id: id.as_str().to_string(), }; self.event_log.append(&event).await.map_err(|e| e.to_string())?; - let mut w = self.tree.write().await; - w.ensure_path(id); + { + let mut w = self.tree.write().await; + w.ensure_path(id); + } + self.reddit.request_fetch(id.clone()); Ok(()) } diff --git a/todo b/todo new file mode 100644 index 0000000000000000000000000000000000000000..d196e8cb4cc80ccb95eeff01c73607d520e13212 --- /dev/null +++ b/todo @@ -0,0 +1,7 @@ +reddit import (only on explicit request) +reddit rendering +vote redering +pair chosing +nsfw gate + +logins (uuid user, two sides, oauths, and pseudonyms)