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: [9e20d06c] Add sorterc dev tool for offline DSL compile and JSONL lint. Introduce a workspace-only binary that validates .sorter files into ranking JSON and scans events.jsonl for corrupt or unreplayable ingests. Co-authored-by: Cursor Side A — unified diff (full patch): diff --git a/Cargo.lock b/Cargo.lock index bf8153d9c723af97122c9ffdd4a7cfe82e853bb6..a07734f089b466440c3ae6fc1087ce85fc24ce62 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1826,6 +1826,17 @@ dependencies = [ "windows-sys 0.60.2", ] +[[package]] +name = "sorterc" +version = "0.0.1" +dependencies = [ + "anyhow", + "clap", + "serde", + "serde_json", + "slugsocial-server", +] + [[package]] name = "spin" version = "0.9.8" diff --git a/Cargo.toml b/Cargo.toml index 149cbf07901eab57c593184ff8719a75d530f1da..25337acdd61e44b20f354c78fed4a88caf896280 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [workspace] -members = ["server", "cli"] +members = ["server", "cli", "sorterc"] resolver = "2" diff --git a/agents.md b/agents.md index d8b801e454fdf37e7ac6038b91a69f83b0746d59..ce646ed3cd7123be4732dccec4a6800467e651e7 100644 --- a/agents.md +++ b/agents.md @@ -93,6 +93,17 @@ SLUG_GOOGLE_CLIENT_SECRET=mock After OAuth completes, the pending-session poll returns a `slug_…` bearer token for API calls. +### Dev-only offline tooling + +**`sorterc`** — workspace binary, not published via npm. Compiles `.sorter` files and lints `events.jsonl` without a server: + +``` +cargo run -p sorterc -- compile path/to/doc.sorter [--base events.jsonl] [--room public] [--pretty] +cargo run -p sorterc -- scan path/to/events.jsonl [--pretty] +``` + +`compile` validates DSL, simulates ingest against empty (or `--base`) reducer state, and prints JSON rankings. `scan` reports corrupt JSONL lines and ingests that fail DSL replay. + ### Testing - **Rust tests:** `cargo nextest run --workspace` (163 tests; requires `cargo-nextest`) diff --git a/server/src/lib.rs b/server/src/lib.rs index c1d477d21aea03aff00e6f0689b0b4379d0d68d2..ad8e31099c807fb5844acb16cd5086a2f19327a7 100644 --- a/server/src/lib.rs +++ b/server/src/lib.rs @@ -10,6 +10,7 @@ pub mod form_template; pub mod html; pub mod identity; pub mod middleware; +pub mod offline; pub mod path_types; pub mod ranking; pub mod reducer; diff --git a/server/src/offline.rs b/server/src/offline.rs new file mode 100644 index 0000000000000000000000000000000000000000..54ad0ded096a305ef8454ab2cdd1c3af71b14f5d --- /dev/null +++ b/server/src/offline.rs @@ -0,0 +1,333 @@ +//! Offline `.sorter` compilation and JSONL diagnostics (no network, no auth). + +use std::collections::HashSet; +use std::path::Path; + +use serde::Serialize; +use slug_types::{CheckScopeRanking, RankComponent, RankRow, paths::GardenItemUrl}; + +use crate::{ + api::{resolve_item, validate_ingest_document}, + dsl, + events::{Event, Ingest}, + path_types::ItemId, + reducer::{ReducerState, ScopeId, scope_from_room_wire}, + scope_rank::build_children_rankings, +}; + +#[derive(Debug, Clone, Serialize)] +pub struct CompileStats { + pub items: usize, + pub votes: usize, + pub prose_blocks: usize, +} + +#[derive(Debug, Serialize)] +pub struct CompileResult { + pub ok: bool, + pub threads: Vec, + pub rankings: Vec, + pub stats: CompileStats, +} + +#[derive(Debug, Clone, Serialize)] +pub struct CompileError { + pub ok: bool, + pub error: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub hint: Option, +} + +#[derive(Debug, Clone, Serialize)] +pub struct BadJsonLine { + pub line: usize, + pub message: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct MalformedIngest { + pub line: usize, + pub id: String, + pub room_id: String, + pub thread_tag: String, + pub reason: String, +} + +#[derive(Debug, Clone, Serialize)] +pub struct ScanResult { + pub ok: bool, + pub path: String, + pub total_lines: usize, + pub parsed_events: usize, + pub bad_json_lines: Vec, + pub malformed_ingests: Vec, + pub skipped_ingests: usize, +} + +fn document_stats(doc: &dsl::Document) -> CompileStats { + let mut items = 0usize; + let mut votes = 0usize; + let mut prose_blocks = 0usize; + for stmt in &doc.statements { + match stmt { + dsl::Stmt::Item { .. } => items += 1, + dsl::Stmt::Vote { .. } => votes += 1, + dsl::Stmt::Prose { .. } => prose_blocks += 1, + } + } + CompileStats { + items, + votes, + prose_blocks, + } +} + +fn threads_in_document(text: &str) -> Vec { + let mut out = HashSet::new(); + for line in text.lines() { + let trimmed = line.trim(); + if !trimmed.starts_with('#') { + continue; + } + let rest = trimmed.trim_start_matches('#').trim(); + if rest.is_empty() { + continue; + } + let tag = rest.split_whitespace().next().unwrap_or(rest); + let tag = tag.split(':').next().unwrap_or(tag).trim(); + if tag.is_empty() { + continue; + } + out.insert(format!("#{}", crate::canonical_path::canonicalize_tag(tag))); + } + let mut tags: Vec = out.into_iter().collect(); + tags.sort(); + tags +} + +fn voted_parent_scopes(doc: &dsl::Document) -> Vec { + let mut parents = HashSet::new(); + for stmt in &doc.statements { + if let dsl::Stmt::Vote { item1, item2, .. } = stmt { + if let (Ok(a), Ok(b)) = (resolve_item(item1), resolve_item(item2)) { + if let Some(p) = a.parent() { + parents.insert(p); + } + if let Some(p) = b.parent() { + parents.insert(p); + } + } + } + } + let mut out: Vec = parents.into_iter().collect(); + out.sort(); + out +} + +fn rankings_for_simulated( + simulated: &ReducerState, + scope: &ScopeId, + room_wire: &str, + doc: &dsl::Document, +) -> Vec { + voted_parent_scopes(doc) + .iter() + .map(|parent| { + let scoped_content = simulated + .content_for_scope(&scope) + .unwrap_or_else(|| simulated.public()); + let scoped = build_children_rankings(scoped_content, parent); + let components: Vec = scoped + .component_rankings + .into_iter() + .map(|comp| RankComponent { + pairs: comp.pairs, + ranking: comp + .ranked + .into_iter() + .map(|r| RankRow { + item: GardenItemUrl::from_stored(&r.item, room_wire), + score: r.score, + percent: None, + }) + .collect(), + }) + .collect(); + CheckScopeRanking { + parent: GardenItemUrl::from_stored(parent, room_wire).into_inner(), + components, + unranked_items: scoped + .unranked_items + .into_iter() + .map(|it| GardenItemUrl::from_stored(&it, room_wire)) + .collect(), + } + }) + .collect() +} + +/// Validate and simulate one `.sorter` document against optional base reducer state. +pub fn compile_document( + base: &ReducerState, + room: &str, + text: &str, +) -> Result { + let room_key = room.trim(); + let scope = scope_from_room_wire(room_key); + let validated = validate_ingest_document(base, text, &scope).map_err(|(_, message, hint)| { + CompileError { + ok: false, + error: message, + hint, + } + })?; + + let event = Event::Ingest(Ingest { + ts: validated.ts, + id: uuid::Uuid::new_v4().to_string(), + raw: validated.raw_text.clone(), + principal: "offline".to_string(), + delegate: None, + room_id: room_key.to_string(), + thread_tag: "offline".to_string(), + }); + + let mut simulated = base.clone(); + simulated.apply_event(event); + + Ok(CompileResult { + ok: true, + threads: threads_in_document(text), + rankings: rankings_for_simulated(&simulated, &scope, room_key, &validated.doc), + stats: document_stats(&validated.doc), + }) +} + +fn ingest_parse_error(raw: &str) -> Option { + dsl::parse_full(raw).err().map(|e| e.to_string()) +} + +fn load_events_from_jsonl(path: &Path) -> Result<(Vec<(usize, Event)>, Vec), std::io::Error> { + let text = std::fs::read_to_string(path)?; + let mut events = Vec::new(); + let mut bad_json_lines = Vec::new(); + for (idx, line) in text.lines().enumerate() { + let line_no = idx + 1; + let trimmed = line.trim(); + if trimmed.is_empty() { + continue; + } + match serde_json::from_str::(trimmed) { + Ok(ev) => events.push((line_no, ev)), + Err(e) => bad_json_lines.push(BadJsonLine { + line: line_no, + message: e.to_string(), + }), + } + } + Ok((events, bad_json_lines)) +} + +/// Replay a JSONL event log into reducer state (same rules as server boot). +pub fn load_reducer_from_jsonl(path: &Path) -> Result<(ReducerState, Vec), std::io::Error> { + let (events, bad_json_lines) = load_events_from_jsonl(path)?; + let mut state = ReducerState::default(); + for (_line_no, ev) in events { + state.apply_event(ev); + } + Ok((state, bad_json_lines)) +} + +/// Scan an events.jsonl for corrupt JSON lines and ingests that fail DSL replay. +pub fn scan_jsonl(path: &Path) -> Result { + let text = std::fs::read_to_string(path)?; + let total_lines = text.lines().count(); + let (events, bad_json_lines) = load_events_from_jsonl(path)?; + + let mut malformed_ingests = Vec::new(); + let mut skipped_ingests = 0usize; + let mut state = ReducerState::default(); + let parsed_events = events.len(); + + for (line_no, ev) in events { + if let Event::Ingest(ref ing) = ev { + if let Some(reason) = ingest_parse_error(&ing.raw) { + malformed_ingests.push(MalformedIngest { + line: line_no, + id: ing.id.clone(), + room_id: ing.room_id.clone(), + thread_tag: ing.thread_tag.clone(), + reason, + }); + } + let before = state.ingests_by_id.len(); + state.apply_event(ev); + if state.ingests_by_id.len() == before { + skipped_ingests += 1; + } + } else { + state.apply_event(ev); + } + } + + let ok = bad_json_lines.is_empty() && malformed_ingests.is_empty() && skipped_ingests == 0; + + Ok(ScanResult { + ok, + path: path.display().to_string(), + total_lines, + parsed_events, + bad_json_lines, + malformed_ingests, + skipped_ingests, + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + const TUTORIAL: &str = include_str!("../tests/fixtures/tutorial.sorter"); + + #[test] + fn compile_tutorial_fixture_emits_rankings() { + let result = compile_document(&ReducerState::default(), "public", TUTORIAL).unwrap(); + assert!(result.ok); + assert!(!result.threads.is_empty()); + assert!(result.stats.items >= 6); + assert!(result.stats.votes >= 6); + assert!(!result.rankings.is_empty()); + } + + #[test] + fn compile_rejects_vote_on_missing_item() { + let err = compile_document( + &ReducerState::default(), + "public", + "{ reason }\n~/missing/a 2:1 ~/missing/b", + ) + .unwrap_err(); + assert!(!err.ok); + assert!(err.error.contains("undefined")); + } + + #[test] + fn scan_empty_jsonl_is_ok() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("events.jsonl"); + std::fs::write(&path, "").unwrap(); + let report = scan_jsonl(&path).unwrap(); + assert!(report.ok); + assert!(report.bad_json_lines.is_empty()); + } + + #[test] + fn scan_reports_bad_json_line() { + let dir = tempfile::tempdir().unwrap(); + let path = dir.path().join("events.jsonl"); + std::fs::write(&path, "{not json}\n").unwrap(); + let report = scan_jsonl(&path).unwrap(); + assert!(!report.ok); + assert_eq!(report.bad_json_lines.len(), 1); + } +} diff --git a/sorterc/Cargo.toml b/sorterc/Cargo.toml new file mode 100644 index 0000000000000000000000000000000000000000..92d477aff9b53291fb1a266db83065c1c800791d --- /dev/null +++ b/sorterc/Cargo.toml @@ -0,0 +1,18 @@ +[package] +name = "sorterc" +version = "0.0.1" +edition = "2021" +license = "MIT" +publish = false +description = "Offline .sorter compiler and events.jsonl linter (dev only)" + +[[bin]] +name = "sorterc" +path = "src/main.rs" + +[dependencies] +anyhow = "1" +clap = { version = "4", features = ["derive"] } +serde = { version = "1", features = ["derive"] } +serde_json = "1" +slugsocial-server = { path = "../server" } diff --git a/sorterc/readme.md b/sorterc/readme.md new file mode 100644 index 0000000000000000000000000000000000000000..1ebcc3fc935541ea9e47e0458ec67a750fe19fff --- /dev/null +++ b/sorterc/readme.md @@ -0,0 +1,92 @@ +# sorterc + +Dev-only offline tooling for the slug `.sorter` DSL and `events.jsonl` event log. + +`sorterc` is **not** published via npm and does not talk to slug.social. It reuses the same parser, validator, and ranking code as the server, but runs entirely on local files. + +## Build + +From the repo root: + +```bash +cargo build -p sorterc +cargo run -p sorterc -- --help +``` + +## Commands + +### `compile` — evaluate a `.sorter` document + +Reads a `.sorter` file (or `-` for stdin), validates the DSL, simulates one ingest against reducer state, and prints JSON rankings to stdout. + +```bash +cargo run -p sorterc -- compile path/to/doc.sorter +cargo run -p sorterc -- compile path/to/doc.sorter --pretty +cargo run -p sorterc -- compile - --pretty # stdin +cargo run -p sorterc -- compile doc.sorter --base events.jsonl # seed garden from log +cargo run -p sorterc -- compile doc.sorter --room public # default room +``` + +**Flags** + +| Flag | Description | +|------|-------------| +| `--base PATH` | Replay an `events.jsonl` first, then compile against that garden state | +| `--room ID` | Room wire id (`public` or private room id). Default: `public` | +| `--pretty` | Pretty-print JSON | + +**Success output** (shape): + +```json +{ + "ok": true, + "threads": ["#my-thread"], + "rankings": [ … ], + "stats": { "items": 3, "votes": 2, "prose_blocks": 5 } +} +``` + +Rankings use the same structure as the server's dry-run check: parent scope, connected components, scores, unranked items. + +**Error output** exits with code 1: + +```json +{ + "ok": false, + "error": "parse error", + "hint": "…" +} +``` + +### `scan` — lint an `events.jsonl` + +Reads a JSONL event log and reports problems without starting a server. + +```bash +cargo run -p sorterc -- scan events.jsonl +cargo run -p sorterc -- scan events.jsonl --pretty +``` + +Reports: + +- **bad JSON lines** — lines that are not valid JSON +- **malformed ingests** — ingest events whose `raw` DSL fails to parse +- **skipped ingests** — ingests dropped during replay (same behavior as server boot) + +Exits 0 when clean, 1 when any issue is found. + +## Typical uses + +- Iterate on `.sorter` files in an editor and pipe through `compile` to see rankings instantly +- Verify a downloaded or edited `events.jsonl` before uploading to Fly +- Debug "malformed ingest" warnings from production boot logs +- CI or pre-commit checks on fixture docs (no OAuth, no network) + +## What it does not do + +- Post to slug.social or append to a live log +- Authenticate users or bind agents +- Run browser/UI tests +- Replace `slugsocial public check` for operators who want the full RPC path against a running server + +For live server dry-run against current garden state, use `npx slugsocial public check` or `POST /try/check` in the browser. diff --git a/sorterc/src/main.rs b/sorterc/src/main.rs new file mode 100644 index 0000000000000000000000000000000000000000..71382c180085cb0ad71043c852f8db5d3a48a284 --- /dev/null +++ b/sorterc/src/main.rs @@ -0,0 +1,117 @@ +use std::path::{Path, PathBuf}; + +use anyhow::{bail, Context, Result}; +use clap::{Parser, Subcommand}; +use slugsocial_server::{ + offline::{self, CompileError, CompileResult, ScanResult}, + reducer::ReducerState, +}; + +#[derive(Parser)] +#[command( + name = "sorterc", + about = "Offline .sorter compiler and events.jsonl linter (dev only)", + version +)] +struct Cli { + #[command(subcommand)] + cmd: Command, +} + +#[derive(Subcommand)] +enum Command { + /// Parse and simulate a .sorter document; emit ranking JSON to stdout. + Compile { + /// `.sorter` file, or `-` for stdin. + file: PathBuf, + /// Room wire id (`public` or private room id). + #[arg(long, default_value = "public")] + room: String, + /// Optional events.jsonl to replay before compiling (seed garden state). + #[arg(long)] + base: Option, + /// Pretty-print JSON. + #[arg(long)] + pretty: bool, + }, + /// Scan an events.jsonl for corrupt JSON lines and malformed ingests. + Scan { + file: PathBuf, + #[arg(long)] + pretty: bool, + }, +} + +fn read_input(path: &Path) -> Result { + if path.as_os_str() == "-" { + use std::io::Read; + let mut buf = String::new(); + std::io::stdin().read_to_string(&mut buf)?; + Ok(buf) + } else { + std::fs::read_to_string(path) + .with_context(|| format!("read {}", path.display())) + } +} + +fn load_base_state(base: Option<&Path>) -> Result { + let Some(path) = base else { + return Ok(ReducerState::default()); + }; + let (state, bad_lines) = offline::load_reducer_from_jsonl(path) + .with_context(|| format!("load base jsonl {}", path.display()))?; + if !bad_lines.is_empty() { + bail!( + "base jsonl has {} corrupt line(s); fix or omit --base", + bad_lines.len() + ); + } + Ok(state) +} + +fn print_json(value: &T, pretty: bool) -> Result<()> { + if pretty { + println!("{}", serde_json::to_string_pretty(value)?); + } else { + println!("{}", serde_json::to_string(value)?); + } + Ok(()) +} + +fn run_compile(file: PathBuf, room: String, base: Option, pretty: bool) -> Result<()> { + let text = read_input(&file)?; + let base_state = load_base_state(base.as_deref())?; + match offline::compile_document(&base_state, &room, &text) { + Ok(result) => { + print_json::(&result, pretty)?; + Ok(()) + } + Err(err) => { + print_json::(&err, pretty)?; + std::process::exit(1); + } + } +} + +fn run_scan(file: PathBuf, pretty: bool) -> Result<()> { + let report = offline::scan_jsonl(&file) + .with_context(|| format!("scan {}", file.display()))?; + print_json::(&report, pretty)?; + if !report.ok { + std::process::exit(1); + } + Ok(()) +} + +fn main() -> Result<()> { + let cli = Cli::parse(); + match cli.cmd { + Command::Compile { + file, + room, + base, + pretty, + } => run_compile(file, room, base, pretty), + Command::Scan { file, pretty } => run_scan(file, pretty), + } +} 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)