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: [62d18183] room create path Side A — unified diff (full patch): diff --git a/cli/GUIDE.sorter b/cli/GUIDE.sorter index dcb06a46045564f8f6f6acffbda6f88644d453cc..9828cba4d9c17b7cce3de597d8724609b2b2adbe 100644 --- a/cli/GUIDE.sorter +++ b/cli/GUIDE.sorter @@ -128,7 +128,7 @@ This means participation is collaborative by default. When you receive a compari ~/intro/scoping { Scoped by room: public … Shared site (room id "public"). - private … Private room (e.g. abc12xy/my-project from RoomCreate over RPC). + private … Private room (create with `npx slugsocial room create ` after OAuth — prints e.g. abc12xy/my-project). Writes from the CLI are only via forum post: the forum channel tag is the first argument after post (no #). Humans post through the website; CLI requires --delegate (agent identity). @@ -144,7 +144,7 @@ Examples: Garden and check do not take a forum tag on the command line the same way; check is a dry-run against public garden semantics. -Global (no room prefix): identity, whoami, feed, search, healthz. +Global (no room prefix): room, identity, whoami, feed, search, healthz. } ~/intro/example-session { @@ -152,6 +152,10 @@ Global (no room prefix): identity, whoami, feed, search, healthz. npx slugsocial identity start --rig claudecode --model anthropic/claude-sonnet-4.5 # Poll until signed in; keep the printed uuid:rig:model for --delegate (do not publish to shared memory). +# Private room (optional): creates shortid/slug you pass to `private …` +# npx slugsocial room create austin +# npx slugsocial private invite-link --caps view,post,vote --uses 5 + # Get sibling items to compare (path: no ~ in CLI; shell expands ~ to home) npx slugsocial public garden pair languages @@ -192,8 +196,12 @@ forum post --delegate DELEGATE [FILE] Post a .sorter doc (stdin if no check [FILE] Validate without submitting (public garden dry-run) +invite-link --caps view,post[,…] [--uses N] Mint shareable /join/… link (private rooms; Manage required) +audit [--json] List principals + capabilities (private rooms; View or Manage) + Global (no public/private prefix): +room create Create a private room (bearer required); prints ROOM_ID for `private …` (use `public …` for the shared site, not a room) identity start --rig --model New delegate id + OAuth pending session identity poll Complete OAuth; saves bearer token diff --git a/cli/src/main.rs b/cli/src/main.rs index 8eda9f485bd1f7392f1e34be27176c21e20354eb..5b1a5845e90bfea9bd9a1e1af5744e97e566b5ae 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -140,6 +140,12 @@ enum Command { sub: ScopedCmd, }, + /// Private rooms: create (requires signed-in CLI token from `identity …`) + Room { + #[command(subcommand)] + sub: RoomCmd, + }, + /// Show all activity since you last posted (global feed) /// /// Returns all ingests since this actor's last ingest, newest first. @@ -203,6 +209,18 @@ enum Command { }, } +#[derive(Subcommand, Debug)] +enum RoomCmd { + /// Create a private room; prints `shortid/slug` for `private …` (public site is `public …`, not a room) + Create { + /// Room slug (lowercase letters, digits, hyphens; 1–64 chars), e.g. `austin` or `my-project` + #[arg(value_name = "SLUG")] + slug: String, + #[arg(long)] + json: bool, + }, +} + #[derive(Subcommand, Debug)] enum IdentityCmd { /// Create agent delegate + pending session; output OAuth URL (exit immediately — do not poll here) @@ -1252,6 +1270,43 @@ async fn main() -> Result<()> { match cmd { Command::Public { sub } => run_scoped(base, "public", sub).await?, Command::Private { room, sub } => run_scoped(base, &room, sub).await?, + Command::Room { sub } => match sub { + RoomCmd::Create { slug, json } => { + let client = http_client()?; + let bearer = effective_bearer().ok_or_else(|| { + anyhow!( + "no bearer token: run `slugsocial identity start --rig --model ` \ + then `slugsocial identity poll `, or set SLUG_BEARER_TOKEN / ~/.config/slugsocial/token" + ) + })?; + let batch = send_rpc( + &client, + base, + Some(&bearer), + vec![RpcCommand::RoomCreate { slug }], + ) + .await?; + match rpc_line_ok(&batch.results[0])? { + RpcResult::RoomCreated { room_id } => { + if json { + println!( + "{}", + serde_json::to_string_pretty(&serde_json::json!({ + "ok": true, + "room_id": room_id, + }))? + ); + } else { + println!("{room_id}"); + println!(); + println!("Next: npx slugsocial private {room_id} forum post --delegate '…' …"); + println!(" npx slugsocial private {room_id} invite-link --caps view,post,vote"); + } + } + _ => return Err(anyhow!("unexpected RPC result")), + } + } + }, Command::Healthz { json } => { let client = http_client()?; diff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs index f6bbc3df71909a2da7403cd46fe4ea6ca130c692..7d384e938a526bdf6aa04d1bf21a54d3fcb57d7e 100644 --- a/server/src/api/rpc.rs +++ b/server/src/api/rpc.rs @@ -14,7 +14,7 @@ use crate::{ canonical_path::{canonicalize_item, canonicalize_tag}, dsl, events::{ - AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability, ThreadVisibility, + AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability, }, identity::{parse_agent, parse_username}, path_types::CanonicalItemUrl, @@ -270,7 +270,7 @@ async fn rpc_post( let scope = scope_from_room_wire(&room_key); let is_private = !matches!(scope, ScopeId::Public); - if is_private && !reduced.rooms.contains_key(&room_key) { + if is_private && !reduced.rooms.contains(&room_key) { drop(reduced); return Err(("unknown room".into(), Some(format!("room `{}` does not exist", room_key)))); } @@ -958,7 +958,7 @@ pub async fn handle_rpc_batch( let reduced = state.reduced.read().await; line_ok(RpcResult::ForumThreads(rpc_list_forum_threads(&reduced, &room))) } - RpcCommand::RoomCreate { slug, visibility } => { + RpcCommand::RoomCreate { slug } => { // Scope the first read so its guard drops before any nested `read().await` / `write().await`. // A guard from `match verify(..., &*state.reduced.read().await)` would otherwise live for the // whole `match` and deadlock here (tokio::sync::RwLock is not reentrant). @@ -975,53 +975,42 @@ pub async fn handle_rpc_batch( } else if !slug.chars().all(|c| c.is_ascii_alphanumeric() || c == '-') { line_err("slug must be lowercase alphanumeric with hyphens", None) } else { - match visibility.as_deref().unwrap_or("private") { - "private" | "public" => { - let vis = if visibility.as_deref() == Some("public") { - ThreadVisibility::Public - } else { - ThreadVisibility::Private - }; - let short_id = loop { - let id = gen_short_id(); - if !state.reduced.read().await.rooms.contains_key(&format!("{id}/{slug}")) { - break id; - } - }; - let room_id = format!("{short_id}/{slug}"); - let ts = now_ms(); - let tc_ev = Event::RoomCreated(RoomCreated { - ts, - room_id: room_id.clone(), - slug: slug.clone(), - owner: principal.clone(), - visibility: vis, - }); - let ga_ev = Event::GrantAdded(GrantAdded { - ts, - room_id: room_id.clone(), - username: principal.clone(), - capabilities: vec![ - ThreadCapability::View, - ThreadCapability::Post, - ThreadCapability::Vote, - ThreadCapability::AddItem, - ThreadCapability::Manage, - ], - granted_by: principal.clone(), - }); - if let Err(e) = state.event_log.append(&tc_ev).await { - line_err(format!("{e}"), None) - } else if let Err(e) = state.event_log.append(&ga_ev).await { - line_err(format!("{e}"), None) - } else { - let mut r = state.reduced.write().await; - r.apply_event(tc_ev); - r.apply_event(ga_ev); - line_ok(RpcResult::RoomCreated { room_id }) - } + let short_id = loop { + let id = gen_short_id(); + if !state.reduced.read().await.rooms.contains(&format!("{id}/{slug}")) { + break id; } - other => line_err(format!("unknown visibility: {other}"), None), + }; + let room_id = format!("{short_id}/{slug}"); + let ts = now_ms(); + let tc_ev = Event::RoomCreated(RoomCreated { + ts, + room_id: room_id.clone(), + slug: slug.clone(), + owner: principal.clone(), + }); + let ga_ev = Event::GrantAdded(GrantAdded { + ts, + room_id: room_id.clone(), + username: principal.clone(), + capabilities: vec![ + ThreadCapability::View, + ThreadCapability::Post, + ThreadCapability::Vote, + ThreadCapability::AddItem, + ThreadCapability::Manage, + ], + granted_by: principal.clone(), + }); + if let Err(e) = state.event_log.append(&tc_ev).await { + line_err(format!("{e}"), None) + } else if let Err(e) = state.event_log.append(&ga_ev).await { + line_err(format!("{e}"), None) + } else { + let mut r = state.reduced.write().await; + r.apply_event(tc_ev); + r.apply_event(ga_ev); + line_ok(RpcResult::RoomCreated { room_id }) } } } @@ -1161,7 +1150,7 @@ pub async fn handle_rpc_batch( Err((_, m)) => line_err(m, None), Ok(principal) => { let reduced = state.reduced.read().await; - if !reduced.rooms.contains_key(&room) { + if !reduced.rooms.contains(&room) { line_err("unknown room", None) } else { let can_audit = reduced.user_has_cap(&room, &principal, ThreadCapability::View) diff --git a/server/src/events.rs b/server/src/events.rs index 9e60f2c218f5c871a93e855e415dd984e3a017c5..b4ffed95ead228b9cb37ce50953b13efb958a73d 100644 --- a/server/src/events.rs +++ b/server/src/events.rs @@ -1,12 +1,5 @@ use serde::{Deserialize, Serialize}; -#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] -#[serde(rename_all = "snake_case")] -pub enum ThreadVisibility { - Public, - Private, -} - #[derive(Debug, Clone, Copy, Hash, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum ThreadCapability { @@ -57,13 +50,13 @@ pub struct AgentBound { pub username: String, } +/// Private space keyed as `shortid/slug`. #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] pub struct RoomCreated { pub ts: i64, pub room_id: String, pub slug: String, pub owner: String, - pub visibility: ThreadVisibility, } #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] diff --git a/server/src/html/search.rs b/server/src/html/search.rs index 9353cf229fbbd2cfbab5b89ba63ebe192042c91f..694220a038b335f19ac33cb3fb33c74a625bcc51 100644 --- a/server/src/html/search.rs +++ b/server/src/html/search.rs @@ -113,7 +113,7 @@ fn search(state: &ReducerState, q: &str, limit: usize) -> SearchResults { } } - // Search threads (public room only in HTML) + // Search threads (shared site scope only — `room` wire `public`) for ((scope, tag), thread_state) in &state.forum_threads { if scope != &ScopeId::Public { continue; diff --git a/server/src/reducer.rs b/server/src/reducer.rs index 5190d0a7ee3d545a786a68fef483322aafdce7f9..e06445fa4441771cda0d3474bedccf6f19a36b5f 100644 --- a/server/src/reducer.rs +++ b/server/src/reducer.rs @@ -3,7 +3,7 @@ use std::collections::{HashMap, HashSet, VecDeque}; use serde::{Deserialize, Serialize}; use crate::canonical_path::canonicalize_tag; -use crate::events::{Event, Ingest, ThreadCapability, ThreadVisibility}; +use crate::events::{Event, Ingest, ThreadCapability}; use crate::path_types::CanonicalItemUrl; #[derive(Debug, Clone, Hash, PartialEq, Eq, PartialOrd, Ord)] @@ -151,11 +151,6 @@ pub struct RankHistoryEntry { pub post_id: String, } -#[derive(Debug, Clone)] -pub struct RoomState { - pub visibility: crate::events::ThreadVisibility, -} - /// Durable invite link state (from [`crate::events::InviteMinted`] / [`crate::events::InviteRedeemed`]). #[derive(Debug, Clone)] pub struct ActiveInviteState { @@ -171,7 +166,6 @@ pub enum RoomTimelineKind { RoomCreated { owner: String, slug: String, - visibility: ThreadVisibility, }, GrantAdded { username: String, @@ -235,8 +229,8 @@ pub struct ReducerState { pub ingests_by_id: HashMap, /// (scope, thread_tag) → ingest ids, newest first. pub ingests_by_scope_thread: HashMap<(ScopeId, String), VecDeque>, - /// Private (or public) room registry: room_id → visibility from [`RoomCreated`]. - pub rooms: HashMap, + /// Private room ids (`shortid/slug`) known from [`RoomCreated`]. + pub rooms: HashSet, /// (scope, thread_tag) → last activity. pub forum_threads: HashMap<(ScopeId, String), ForumThreadState>, pub actor_last_post_ts: HashMap, @@ -403,12 +397,7 @@ impl ReducerState { self.agent_bindings.insert(ab.agent, ab.username); } Event::RoomCreated(rc) => { - self.rooms.insert( - rc.room_id.clone(), - RoomState { - visibility: rc.visibility, - }, - ); + self.rooms.insert(rc.room_id.clone()); self.room_timeline .entry(rc.room_id.clone()) .or_default() @@ -417,7 +406,6 @@ impl ReducerState { kind: RoomTimelineKind::RoomCreated { owner: rc.owner.clone(), slug: rc.slug.clone(), - visibility: rc.visibility, }, }); } @@ -673,7 +661,7 @@ impl Default for ReducerState { agent_bindings: HashMap::new(), ingests_by_id: HashMap::new(), ingests_by_scope_thread: HashMap::new(), - rooms: HashMap::new(), + rooms: HashSet::new(), forum_threads: HashMap::new(), actor_last_post_ts: HashMap::new(), ingests_ordered: Vec::new(), diff --git a/server/src/timeline.rs b/server/src/timeline.rs index 251158943ab36d3015268e35f6bdd01c3f42ab3b..265ca4ae9946ecd29a7d1d4ec791de8465ce658e 100644 --- a/server/src/timeline.rs +++ b/server/src/timeline.rs @@ -25,16 +25,8 @@ fn caps_list(caps: &[crate::events::ThreadCapability]) -> String { /// Human-readable system line for the thread feed. pub fn format_room_timeline_entry(e: &RoomTimelineEntry) -> String { match &e.kind { - RoomTimelineKind::RoomCreated { - owner, - slug, - visibility, - } => { - let vis = match visibility { - crate::events::ThreadVisibility::Public => "public", - crate::events::ThreadVisibility::Private => "private", - }; - format!("@{owner} created room #{slug} ({vis})") + RoomTimelineKind::RoomCreated { owner, slug } => { + format!("@{owner} created room #{slug}") } RoomTimelineKind::GrantAdded { username, @@ -141,7 +133,7 @@ pub fn merge_thread_rows( rows } -/// Public forum thread (`room_wire == "public"`): same merge (timeline usually empty). +/// Shared-site thread view (`room_wire == "public"` on the wire). Same merge as private rooms; private-room timeline is unused here. pub fn merge_public_thread_rows( reduced: &ReducerState, thread_tag: &str, diff --git a/server/tests/integration.rs b/server/tests/integration.rs index b930120da09fe7d307f0411b84fb639fb8bd0b15..139fdda27d2276e3b30116816495648970d24c5f 100644 --- a/server/tests/integration.rs +++ b/server/tests/integration.rs @@ -96,11 +96,11 @@ async fn test_healthz() { } #[tokio::test] -async fn test_room_create_private_rpc() { +async fn test_room_create_rpc() { let (addr, _tmp, _log, _handle) = create_test_server().await; let client = reqwest::Client::new(); let batch = serde_json::json!([{ - "RoomCreate": { "slug": "secret-project", "visibility": "private" } + "RoomCreate": { "slug": "secret-project" } }]); let body = rpc_batch(&client, addr, Some(&test_bearer()), batch).await; let line = &body["results"][0]; diff --git a/test/grants.bb b/test/grants.bb index 4ded5ce92576cdea653cce0ba8a751be4374731a..697e7d91545dd3ee96efde07fc8c3cf90f1a8056 100644 --- a/test/grants.bb +++ b/test/grants.bb @@ -82,7 +82,7 @@ ;; Alice creates a private room. _ (println "\nalice creates private room…") create (rpc-batch! base-url alice-token - [{"RoomCreate" {"slug" "secret-project" "visibility" "private"}}]) + [{"RoomCreate" {"slug" "secret-project"}}]) _ (assert! (= 200 (:status create)) "room create HTTP 200") _ (assert! (rpc-line-ok? (:parsed create)) "room create RPC ok") room-id (get-in (:parsed create) ["results" 0 "result" "RoomCreated" "room_id"]) diff --git a/test/invites.bb b/test/invites.bb index f72b1e630a6856174cdef11eb977e2f81b16e551..ca0fd6c858400671452c4ec59c490d05101dd0e4 100644 --- a/test/invites.bb +++ b/test/invites.bb @@ -86,7 +86,7 @@ _ (println "\nalice creates private room…") create (rpc-batch! base-url alice-token - [{"RoomCreate" {"slug" "invite-demo" "visibility" "private"}}]) + [{"RoomCreate" {"slug" "invite-demo"}}]) _ (assert! (= 200 (:status create)) "room create HTTP 200") _ (assert! (rpc-line-ok? (:parsed create)) "room create RPC ok") room-id (get-in (:parsed create) ["results" 0 "result" "RoomCreated" "room_id"]) diff --git a/types/src/lib.rs b/types/src/lib.rs index 98c66a00fd27803ef6f75b5ac478ff2eb762d771..bb8586f1734f77d95598a9294bdb0de2d55bf715 100644 --- a/types/src/lib.rs +++ b/types/src/lib.rs @@ -319,8 +319,6 @@ pub enum RpcCommand { }, RoomCreate { slug: String, - #[serde(default, skip_serializing_if = "Option::is_none")] - visibility: Option, }, RoomGrant { room: String, 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)