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: [993d359c] Fix Reddit children import wiring and unranked child labels. Listing imports attach posts directly under the subreddit without ensure_path pulling comment-path segments in, and the ranking panel shows imported titles. Update integration tests for JS SSE morphs and children fetch. Co-authored-by: Cursor Side A — unified diff (full patch): diff --git a/server/src/api/ui_html.rs b/server/src/api/ui_html.rs index d1defd28242fd2ca3b886adc91a7070bef75e653..e649a7d192feade465e19ce6187a829f6ec74372 100644 --- a/server/src/api/ui_html.rs +++ b/server/src/api/ui_html.rs @@ -58,7 +58,7 @@ pub async fn post_ui_html( let tree = state.tree.read().await; let empty = crate::reducer::NodeState::default(); let node = tree.get(&parent).unwrap_or(&empty); - let panel = ranking_panel(&parent, node); + let panel = ranking_panel(&parent, node, &tree); JsBuilder::new() .morph_selector("#ranking-panel", panel) .into_response() diff --git a/server/src/fetch/mod.rs b/server/src/fetch/mod.rs index 0177bb161cea1b72a100b52efdfc5e710271c9eb..35f968c21c69ca557bab7951413e3cfbbccfebfd 100644 --- a/server/src/fetch/mod.rs +++ b/server/src/fetch/mod.rs @@ -107,7 +107,7 @@ pub fn fetch_entity_stream( let mut b = JsBuilder::new() .morph_selector("#entity-section", html::entity_section(&id, node, false)); if kind == FetchKind::Children { - b = b.morph_selector("#ranking-panel", ranking_panel(&id, node)); + b = b.morph_selector("#ranking-panel", ranking_panel(&id, node, &tree)); } yield Ok(js_event(b.build())); } diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs index 27ce9118c73ec5643e04363ab1a36cf5da6101bf..e88cc43ddc9d8100f7994be6f5960ec4d8f22c55 100644 --- a/server/src/html/mod.rs +++ b/server/src/html/mod.rs @@ -15,7 +15,7 @@ use crate::{ ranking::{ connected_components_from_voted_pairs, ranked_items_subset, RankedItem, MAX_ITERS, TOL, }, - reducer::NodeState, + reducer::{GlobalTree, NodeState}, state::AppState, ui_action::UI_RPC_FIELD, }; @@ -182,8 +182,15 @@ fn display_label(id: &ItemId) -> String { .to_string() } +fn child_label(tree: &GlobalTree, id: &ItemId) -> String { + tree.get(id) + .and_then(|n| n.data.as_ref()) + .map(|d| d.title.clone()) + .unwrap_or_else(|| display_label(id)) +} + /// Plain (unscored) list of children that have no votes yet. -fn unranked_list(label: &str, items: &[ItemId]) -> Markup { +fn unranked_list(label: &str, items: &[ItemId], tree: &GlobalTree) -> Markup { html! { @if !items.is_empty() { h3 class="rank-heading muted small" { (label) } @@ -191,7 +198,7 @@ fn unranked_list(label: &str, items: &[ItemId]) -> Markup { @for it in items { li { a href=(item_href(it)) { - strong { (display_label(it)) } + strong { (child_label(tree, it)) } } } } @@ -200,7 +207,7 @@ fn unranked_list(label: &str, items: &[ItemId]) -> Markup { } } -pub fn ranking_panel(item: &ItemId, node: &NodeState) -> Markup { +pub fn ranking_panel(item: &ItemId, node: &NodeState, tree: &GlobalTree) -> Markup { let group = &node.local_ranking; let n = group.idx_to_item.len(); let (comps, _isolates) = @@ -248,7 +255,7 @@ pub fn ranking_panel(item: &ItemId, node: &NodeState) -> Markup { @let label = if multi { format!("Ranking group {}", gi + 1) } else { "Ranking".to_string() }; (rank_list(&label, ranked, 1)) } - (unranked_list("Unranked", &unranked)) + (unranked_list("Unranked", &unranked, tree)) } } } @@ -297,7 +304,7 @@ async fn item_page(state: AppState, uri: Uri, item: ItemId) -> Markup { (input_panel("", None)) (breadcrumb_path(&item)) (entity_section(&item, node, false)) - (ranking_panel(&item, node)) + (ranking_panel(&item, node, &tree)) }; layout("sorter2", body, views) } diff --git a/server/src/reddit.rs b/server/src/reddit.rs index a0eb688709478ee0185b953b41a5d26cc354764d..69d979bc7e4a1cb078f70114dc539bdc1986b574 100644 --- a/server/src/reddit.rs +++ b/server/src/reddit.rs @@ -300,9 +300,16 @@ async fn reddit_worker( } { let mut tree = tree.write().await; - apply_entity_import(&mut tree, &child_id, child_payload); if kind == FetchKind::Children { - tree.link_child(&fetch_id, &child_id); + let view = entity_view_from_payload(&child_id, &child_payload); + tree.apply_entity_under_parent( + &fetch_id, + &child_id, + child_payload, + view, + ); + } else { + apply_entity_import(&mut tree, &child_id, child_payload); } } written += 1; diff --git a/server/src/reducer.rs b/server/src/reducer.rs index a36cd5c9287d61536f9f4a3f6b0df2342388857a..4e42d0369dab50bb2f8ca664aa69b628292f6c07 100644 --- a/server/src/reducer.rs +++ b/server/src/reducer.rs @@ -209,14 +209,24 @@ impl GlobalTree { } } - /// Directly attach `child` under `parent`, bypassing path-based nesting. - /// Used for imported listings (e.g. a subreddit's posts) so they show up - /// as children of the subreddit rather than a deep `…/comments/` path. - pub fn link_child(&mut self, parent: &ItemId, child: &ItemId) { + /// Import entity data for `id` and attach it as a direct child of `parent` + /// without running [`Self::ensure_path`] on `id` (avoids Reddit `/comments/` + /// parent rules pulling intermediate path segments into the subreddit). + pub fn apply_entity_under_parent( + &mut self, + parent: &ItemId, + id: &ItemId, + payload: Value, + view: Option, + ) { self.ensure_path(parent); - self.ensure_path(child); + self.ensure_node(id); + if let Some(node) = self.nodes.get_mut(id) { + node.entity_raw = Some(payload); + node.data = view; + } if let Some(p) = self.nodes.get_mut(parent) { - p.children.insert(child.clone()); + p.children.insert(id.clone()); } } } diff --git a/test/reddit_import.clj b/test/reddit_import.clj index 84cbdcf7965e50290313cbce2243a16b583d2097..45a2a19f20799d77e84d8aa64735ab5e7e45f97c 100644 --- a/test/reddit_import.clj +++ b/test/reddit_import.clj @@ -67,6 +67,36 @@ (do (Thread/sleep 200) (recur)) false))))) +(defn- run-reddit-fetch-assertions [app-base data-dir] + (let [browse-url (str app-base "/~/https://reddit.com/r/rust") + log-path (str data-dir "/events.jsonl") + before (:out (process/shell {:out :string :err :string} + "curl" "-sf" browse-url))] + (is (str/includes? before "Fetch from Reddit")) + (is (not (str/includes? before "The Rust Programming Language"))) + (let [sse (curl-fetch-ui-sse app-base "reddit.com/r/rust" "self")] + (is (zero? (:exit sse)) "POST /ui fetch_entity (self) SSE succeeds") + (is (str/includes? (:out sse) "Idiomorph.morph")) + (is (str/includes? (:out sse) "The Rust Programming Language")) + (is (wait-event-log log-path 2000) "event log written")) + (let [after (:out (process/shell {:out :string :err :string} + "curl" "-sf" browse-url)) + log (slurp (io/file log-path))] + (is (str/includes? after "The Rust Programming Language")) + (is (str/includes? log "\"type\":\"entity_imported\"")) + (is (str/includes? log "\"subscribers\":350000")) + (is (str/includes? log "\"display_name\":\"rust\""))) + (let [children-sse (curl-fetch-ui-sse app-base "reddit.com/r/rust" "children")] + (is (zero? (:exit children-sse)) "POST /ui fetch_entity (children) SSE succeeds") + (is (str/includes? (:out children-sse) "Idiomorph.morph")) + (is (str/includes? (:out children-sse) "Announcing Rust 1.99"))) + (let [after-children (:out (process/shell {:out :string :err :string} + "curl" "-sf" browse-url)) + log2 (slurp (io/file log-path))] + (is (str/includes? after-children "Announcing Rust 1.99")) + (is (str/includes? after-children "Unranked")) + (is (str/includes? log2 "announcing_rust_199"))))) + (deftest reddit-fetch-via-mock-api (testing "Fetch more queues import; event log stores full payload; page shows title" (let [root (repo-root) @@ -102,24 +132,7 @@ bin)] (try (is (wait-health app-base 20000) "app healthz") - (let [browse-url (str app-base "/~/https://reddit.com/r/rust") - before (:out (process/shell {:out :string :err :string} - "curl" "-sf" browse-url))] - (is (str/includes? before "Fetch from Reddit")) - (is (not (str/includes? before "The Rust Programming Language"))) - (let [log-path (str data-dir "/events.jsonl") - sse (curl-fetch-ui-sse app-base "reddit.com/r/rust")] - (is (zero? (:exit sse)) "POST /ui fetch_entity SSE succeeds") - (is (str/includes? (:out sse) "event: complete")) - (is (str/includes? (:out sse) "The Rust Programming Language")) - (is (wait-event-log log-path 2000) "event log written") - (let [after (:out (process/shell {:out :string :err :string} - "curl" "-sf" browse-url)) - log (slurp (io/file log-path))] - (is (str/includes? after "The Rust Programming Language")) - (is (str/includes? log "\"type\":\"entity_imported\"")) - (is (str/includes? log "\"subscribers\":350000")) - (is (str/includes? log "\"display_name\":\"rust\""))))) + (run-reddit-fetch-assertions app-base data-dir) (finally (process/destroy proc)))) (finally diff --git a/test/smoke.clj b/test/smoke.clj index 11887c48282088e140d823a88ba616f6325835b3..ce9f958c84b9a89ae55e215ab519f1df6435e24b 100644 --- a/test/smoke.clj +++ b/test/smoke.clj @@ -49,7 +49,7 @@ (is (wait-health base 15000) "server responds to /healthz") (let [home (:out (process/shell {:out :string :err :string} "curl" "-sf" (str base "/")))] - (is (str/includes? home "vote-panel")) + (is (str/includes? home "entity-section")) (is (str/includes? home "ranking-panel")) (is (str/includes? home "parser-panel")) (is (str/includes? home "__rpc__"))) 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)