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: [3f35edab] progress Side A — unified diff (full patch): diff --git a/server/src/api/mod.rs b/server/src/api/mod.rs index a10ce662105cff8fad949c6b83f7035ce79bed18..a986f706ea4b261cbaf004c02b4cf84184b41371 100644 --- a/server/src/api/mod.rs +++ b/server/src/api/mod.rs @@ -3,6 +3,7 @@ mod helpers; mod rpc; mod stream; mod validate; +mod ui_html; mod web_post; pub use auth::{ @@ -33,6 +34,7 @@ pub use stream::{get_html_stream, get_stream}; pub use validate::{normalize_room_and_thread, validate_ingest_document, ValidatedIngest}; +pub use ui_html::post_ui_html; pub use web_post::{check_web_ingest, post_web_ingest, post_web_redact}; #[cfg(test)] diff --git a/server/src/api/ui_html.rs b/server/src/api/ui_html.rs new file mode 100644 index 0000000000000000000000000000000000000000..2b40a72059981d558768f73d189b991f3448c257 --- /dev/null +++ b/server/src/api/ui_html.rs @@ -0,0 +1,139 @@ +//! Single `POST /ui` entry for browser [`crate::html::ui_action::HtmlUiAction`] (JSON in `__rpc__` + holes). + +use axum::{ + body::Body, + extract::State, + http::{header, HeaderMap, StatusCode}, + response::{IntoResponse, Response}, + Form, +}; +use axum_extra::extract::cookie::CookieJar; +use std::collections::HashMap; + +use crate::{ + api::{ + auth::optional_principal, + web_post::{run_check_web_ingest, run_post_web_ingest, run_post_web_redact, WebPostForm, WebRedactForm}, + }, + html::{ + fragment_public_new_thread_form, fragment_room_new_thread_form, login_to_post_hint_markup, + parse_html_ui_from_form, user_can_post_room, user_can_view_room, HtmlUiAction, JsBuilder, + ThreadNav, + }, + state::AppState, +}; + +pub async fn post_ui_html( + State(state): State, + headers: HeaderMap, + jar: CookieJar, + Form(form): Form>, +) -> impl IntoResponse { + let action = match parse_html_ui_from_form(&form) { + Ok(a) => a, + Err(e) => return ui_js_warn(&e.to_string()).into_response(), + }; + + match action { + HtmlUiAction::PostIngest { + room, + thread_tag, + text, + error_target, + form_id, + } => { + run_post_web_ingest( + &state, + &headers, + &jar, + WebPostForm { + room, + thread_tag, + text, + error_target, + form_id, + }, + ) + .await + } + HtmlUiAction::CheckIngest { + room, + thread_tag, + text, + error_target, + form_id, + } => { + run_check_web_ingest( + &state, + &headers, + &jar, + WebPostForm { + room, + thread_tag, + text, + error_target, + form_id, + }, + ) + .await + } + HtmlUiAction::RedactPost { post_id } => { + run_post_web_redact(&state, &headers, &jar, WebRedactForm { post_id }).await + } + HtmlUiAction::ExpandPublicNewThreadForm => { + let reduced = state.reduced.read().await; + let user = optional_principal(&headers, &jar, &reduced); + drop(reduced); + let markup = if user.is_some() { + fragment_public_new_thread_form(true) + } else { + login_to_post_hint_markup() + }; + JsBuilder::new() + .morph_selector("#public-new-thread-ui-slot", markup) + .into_response() + } + HtmlUiAction::ExpandRoomNewThreadForm { room_wire } => { + let room_wire = room_wire.trim().to_string(); + if room_wire.is_empty() { + return ui_js_warn("missing room").into_response(); + } + let reduced = state.reduced.read().await; + let user = optional_principal(&headers, &jar, &reduced); + if !reduced.rooms.contains(&room_wire) { + drop(reduced); + return ui_js_warn("room not found").into_response(); + } + if !user_can_view_room(&reduced, &room_wire, user.as_deref()) { + drop(reduced); + return ui_js_warn("forbidden").into_response(); + } + let can_post = user + .as_ref() + .map(|u| user_can_post_room(&reduced, &room_wire, u)) + .unwrap_or(false); + drop(reduced); + let Some(nav) = ThreadNav::from_room_id(&room_wire) else { + return ui_js_warn("bad room").into_response(); + }; + let markup = if can_post { + fragment_room_new_thread_form(&nav, true) + } else { + login_to_post_hint_markup() + }; + JsBuilder::new() + .morph_selector("#room-new-thread-ui-slot", markup) + .into_response() + } + } +} + +fn ui_js_warn(msg: &str) -> Response { + use crate::html::js_string_literal; + let js = format!("console.warn({});", js_string_literal(msg)); + Response::builder() + .status(StatusCode::OK) + .header(header::CONTENT_TYPE, "text/javascript; charset=utf-8") + .body(Body::from(js)) + .unwrap() +} diff --git a/server/src/api/web_post.rs b/server/src/api/web_post.rs index a64010e382d3039821c836a5529adad0fe67cce5..265025f41ff1548d05b2d2d5d84202245388053f 100644 --- a/server/src/api/web_post.rs +++ b/server/src/api/web_post.rs @@ -222,8 +222,18 @@ pub async fn post_web_redact( jar: CookieJar, Form(form): Form, ) -> impl IntoResponse { + run_post_web_redact(&state, &headers, &jar, form).await +} + +/// Shared with [`crate::api::ui_html::post_ui_html`]. +pub(crate) async fn run_post_web_redact( + state: &AppState, + headers: &HeaderMap, + jar: &CookieJar, + form: WebRedactForm, +) -> Response { let reduced = state.reduced.read().await; - let Some(_username) = optional_principal(&headers, &jar, &reduced) else { + let Some(_username) = optional_principal(headers, jar, &reduced) else { drop(reduced); return js_redirect("/login").into_response(); }; @@ -239,8 +249,8 @@ pub async fn post_web_redact( return js_redirect("/login").into_response(); }; - match rpc_post_redact(&state, &headers, form.post_id).await { - Ok(RpcResult::RedactPostOk {}) => redact_success_response(&state).await.into_response(), + match rpc_post_redact(state, headers, form.post_id).await { + Ok(RpcResult::RedactPostOk {}) => redact_success_response(state).await.into_response(), Ok(_) => (StatusCode::BAD_REQUEST, "unexpected response").into_response(), Err((msg, hint)) => { let detail = hint.as_deref().unwrap_or(""); @@ -255,8 +265,18 @@ pub async fn post_web_ingest( jar: CookieJar, Form(form): Form, ) -> impl IntoResponse { + run_post_web_ingest(&state, &headers, &jar, form).await +} + +/// Shared with [`crate::api::ui_html::post_ui_html`] (`POST /ui`). +pub(crate) async fn run_post_web_ingest( + state: &AppState, + headers: &HeaderMap, + jar: &CookieJar, + form: WebPostForm, +) -> Response { let reduced = state.reduced.read().await; - let Some(_username) = optional_principal(&headers, &jar, &reduced) else { + let Some(_username) = optional_principal(headers, jar, &reduced) else { drop(reduced); return js_redirect("/login").into_response(); }; @@ -282,8 +302,8 @@ pub async fn post_web_ingest( .into_response(); } - match rpc_post_with_bearer(&state, &bearer, room.clone(), thread_tag.clone(), text).await { - Ok(RpcResult::PostOk { .. }) => post_success_response(&state, &form, &headers, &jar) + match rpc_post_with_bearer(state, &bearer, room.clone(), thread_tag.clone(), text).await { + Ok(RpcResult::PostOk { .. }) => post_success_response(state, &form, headers, jar) .await .into_response(), Ok(_) => form_js_error(&form, "unexpected response", "Post did not return PostOk.").into_response(), @@ -297,8 +317,18 @@ pub async fn check_web_ingest( jar: CookieJar, Form(form): Form, ) -> impl IntoResponse { + run_check_web_ingest(&state, &headers, &jar, form).await +} + +/// Shared with [`crate::api::ui_html::post_ui_html`] (`POST /ui`). +pub(crate) async fn run_check_web_ingest( + state: &AppState, + headers: &HeaderMap, + jar: &CookieJar, + form: WebPostForm, +) -> Response { let reduced = state.reduced.read().await; - let Some(_username) = optional_principal(&headers, &jar, &reduced) else { + let Some(_username) = optional_principal(headers, jar, &reduced) else { drop(reduced); return js_redirect("/login").into_response(); }; @@ -324,7 +354,7 @@ pub async fn check_web_ingest( return js_clear_errors(&form_error_target(&form)).into_response(); } - match rpc_check_with_bearer(&state, &bearer, room, form.text.clone()).await { + match rpc_check_with_bearer(state, &bearer, room, form.text.clone()).await { Ok(RpcResult::CheckOk { .. }) => js_clear_errors(&form_error_target(&form)).into_response(), Ok(_) => form_js_error(&form, "unexpected response", "Check did not return CheckOk.").into_response(), Err((msg, hint)) => form_js_error(&form, &msg, hint.as_deref().unwrap_or("")).into_response(), diff --git a/server/src/form_template.rs b/server/src/form_template.rs new file mode 100644 index 0000000000000000000000000000000000000000..3709c2c09a859da006e4af173413d5d235bc19be --- /dev/null +++ b/server/src/form_template.rs @@ -0,0 +1,142 @@ +//! Plan2-style JSON templates with `{"$form": "field_name"}` holes, filled from +//! `application/x-www-form-urlencoded` (or any `String` → `String` map) **before** +//! deserializing into a typed struct. +//! +//! # Wire format +//! +//! Templates are **compact JSON** (`serde_json::to_string`): one line, no pretty +//! printing, strings escaped per JSON rules (`\"`, `\n`, etc.). Embed that string +//! in HTML attributes or text nodes with normal HTML escaping (e.g. maud), not +//! bespoke encodings. +//! +//! # Power vs flat hidden fields +//! +//! A form is always a string→string map. You can fake depth with dotted keys (`a.b.c`), +//! but one structured blob (`__rpc__` = compact JSON) gives you nested objects, +//! arrays, and optional fields without inventing a new naming scheme each time. +//! +//! # Security +//! +//! Substitution runs **before** `serde` into your command type. It does not fix +//! authorization: if the client can replace the hidden `__rpc__` value, they can +//! change the command shape unless you validate (signed blob, server-side session +//! context, or treat the blob as hints only). Same threat model as any hidden field. + +use serde::Serialize; +use serde_json::Value; +use std::collections::HashMap; + +/// Serialize a value to compact JSON for a hidden `__rpc__` (or similar) field. +pub fn template_json_compact(v: &T) -> serde_json::Result { + serde_json::to_string(v) +} + +/// Recursively walk the JSON AST and replace `{"$form": "key"}` with the submitted +/// string for `key` (empty if missing). Other keys are unchanged. +pub fn substitute_form_vars(val: &mut Value, form_data: &HashMap) { + match val { + Value::Object(map) => { + if map.len() == 1 { + if let Some(Value::String(field_name)) = map.get("$form") { + let submitted = form_data + .get(field_name.as_str()) + .map(|s| s.as_str()) + .unwrap_or(""); + *val = Value::String(submitted.to_string()); + return; + } + } + for v in map.values_mut() { + substitute_form_vars(v, form_data); + } + } + Value::Array(arr) => { + for v in arr.iter_mut() { + substitute_form_vars(v, form_data); + } + } + _ => {} + } +} + +/// Parse JSON, apply [`substitute_form_vars`], return the mutated value. +pub fn fill_template_from_form( + template_json: &str, + form_data: &HashMap, +) -> Result { + let mut v: Value = serde_json::from_str(template_json)?; + substitute_form_vars(&mut v, form_data); + Ok(v) +} + +#[cfg(test)] +mod tests { + use super::*; + use serde::Deserialize; + + #[derive(Debug, Deserialize, PartialEq, Eq)] + struct Demo { + room: String, + thread_tag: String, + nested: Nested, + } + + #[derive(Debug, Deserialize, PartialEq, Eq)] + struct Nested { + text: String, + } + + #[test] + fn holes_become_strings() { + let json = r#"{ + "room": "public", + "thread_tag": {"$form": "tag"}, + "nested": {"text": {"$form": "body"}} + }"#; + let mut form = HashMap::new(); + form.insert("tag".into(), "foo".into()); + form.insert("body".into(), "hello\nworld".into()); + + let v = fill_template_from_form(json, &form).unwrap(); + let d: Demo = serde_json::from_value(v).unwrap(); + assert_eq!( + d, + Demo { + room: "public".into(), + thread_tag: "foo".into(), + nested: Nested { + text: "hello\nworld".into(), + }, + } + ); + } + + #[test] + fn missing_form_key_is_empty_string() { + let json = r#"{"x": {"$form": "nope"}}"#; + let mut form = HashMap::new(); + form.insert("other".into(), "y".into()); + let v = fill_template_from_form(json, &form).unwrap(); + assert_eq!(v["x"], ""); + } + + #[test] + fn array_of_holes() { + let json = r#"{"items": [{"$form": "a"}, {"$form": "b"}]}"#; + let mut form = HashMap::new(); + form.insert("a".into(), "1".into()); + form.insert("b".into(), "2".into()); + let v = fill_template_from_form(json, &form).unwrap(); + assert_eq!(v["items"], serde_json::json!(["1", "2"])); + } + + #[test] + fn template_json_compact_escapes_and_single_line() { + let s = template_json_compact(&serde_json::json!({ + "x": "quote\"and\nnewline" + })) + .unwrap(); + assert!(!s.contains('\n')); + assert!(s.contains("\\\"") || s.contains("\\n")); + } +} diff --git a/server/src/html/forum.rs b/server/src/html/forum.rs index f789e347dc136f8ffd356f4492f0bd76cea04bf0..b1c037f3dd93c7bb96d6624cab6019228432c501 100644 --- a/server/src/html/forum.rs +++ b/server/src/html/forum.rs @@ -11,12 +11,15 @@ use crate::{ api::optional_principal, canonical_path::{canonicalize_item, canonicalize_tag}, events::ThreadCapability, + form_template::template_json_compact, identity::parse_username, reducer::{scope_from_room_wire, ReducerState, ScopeId}, state::AppState, timeago, }; +use super::ui_action::{HtmlUiAction, UI_RPC_FIELD}; + use super::{ bc_segment, bc_threads, cli_panel, layout, now_ms, profile_href, recency_class, render_linkified_with_embeds_in_scope, theme_from_jar, theme_next_from_uri, JsBuilder, @@ -289,7 +292,7 @@ fn rooms_for_user(reduced: &ReducerState, username: &str) -> Vec { v } -fn user_can_view_room(reduced: &ReducerState, room_id: &str, username: Option<&str>) -> bool { +pub(crate) fn user_can_view_room(reduced: &ReducerState, room_id: &str, username: Option<&str>) -> bool { if !reduced.rooms.contains(room_id) { return false; } @@ -299,7 +302,7 @@ fn user_can_view_room(reduced: &ReducerState, room_id: &str, username: Option<&s reduced.user_has_cap(room_id, u, ThreadCapability::View) } -fn user_can_post_room(reduced: &ReducerState, room_id: &str, username: &str) -> bool { +pub(crate) fn user_can_post_room(reduced: &ReducerState, room_id: &str, username: &str) -> bool { reduced.user_has_cap(room_id, username, ThreadCapability::Post) } @@ -591,7 +594,6 @@ pub async fn home( let nav = ThreadNav::public(); let reduced_read = state.reduced.read().await; let strip = auth_strip(&headers, &jar, &reduced_read); - let show_forms = user.is_some(); drop(reduced_read); let page = layout( @@ -617,8 +619,14 @@ pub async fn home( } } p class="muted" { "dark = time-ordered · light = vote-ranked" } + div class="thread-feed-toolbar" { + form method="POST" action="/ui" { + input type="hidden" name=(UI_RPC_FIELD) value=(expand_public_new_thread_rpc_value()); + button type="submit" class="section-add-btn" { "+" } + } + } + div id="public-new-thread-ui-slot" {} (render_thread_feed(Some(&nav), "thread-feed", &public_rows, now)) - (new_thread_form_public(show_forms)) (cli_panel("npx slugsocial public forum list")) }, None, @@ -901,7 +909,13 @@ pub async fn room_page( h3 { "threads" } (render_thread_feed(Some(&nav), "room-thread-feed", &rows, now)) @if show_new { - (new_thread_form_for_room(&nav, show_new)) + div class="thread-feed-toolbar" { + form method="POST" action="/ui" { + input type="hidden" name=(UI_RPC_FIELD) value=(expand_room_new_thread_rpc_value(&nav)); + button type="submit" class="section-add-btn" { "+" } + } + } + div id="room-new-thread-ui-slot" {} } (cli_panel(&forum_cli)) (cli_panel(&garden_cli)) @@ -945,6 +959,31 @@ fn new_thread_form_for_room(nav: &ThreadNav, show: bool) -> Markup { } } +pub(crate) fn expand_public_new_thread_rpc_value() -> String { + template_json_compact(&HtmlUiAction::ExpandPublicNewThreadForm).expect("static json") +} + +pub(crate) fn expand_room_new_thread_rpc_value(nav: &ThreadNav) -> String { + template_json_compact(&HtmlUiAction::ExpandRoomNewThreadForm { + room_wire: nav.room_wire.clone(), + }) + .expect("static json") +} + +pub(crate) fn login_to_post_hint_markup() -> Markup { + html! { + p class="muted" { "log in to post" } + } +} + +pub(crate) fn fragment_public_new_thread_form(show: bool) -> Markup { + new_thread_form_public(show) +} + +pub(crate) fn fragment_room_new_thread_form(nav: &ThreadNav, show: bool) -> Markup { + new_thread_form_for_room(nav, show) +} + async fn thread_post_view_inner( state: AppState, tag: String, diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs index 53041a0cf0b30e6f20749c92ecc15a4e9ac56e09..a4fd2dd60ac283f7eb9b2e64522501b21d1e27e9 100644 --- a/server/src/html/mod.rs +++ b/server/src/html/mod.rs @@ -16,6 +16,7 @@ mod editor; mod forum; mod garden; mod search; +pub mod ui_action; use breadcrumb_path::OntologyPath; pub use auth::{auth_complete_page, auth_signed_in_fragment, choose_username_error_fragment, choose_username_page}; @@ -27,9 +28,15 @@ pub use forum::{ thread_post_collapse_deleted, thread_post_expand, thread_post_expand_deleted, thread_post_view, thread_view, ThreadNav, }; + +pub(crate) use forum::{ + fragment_public_new_thread_form, fragment_room_new_thread_form, login_to_post_hint_markup, + user_can_post_room, user_can_view_room, +}; pub use garden::{garden_index, ontology_path, room_garden_index, room_ontology_path}; pub use search::{search_page, search_results_fragment}; pub use forum::user_profile_page; +pub use ui_action::{parse_html_ui_from_form, HtmlUiAction, HtmlUiParseError, UI_RPC_FIELD}; /// Public profile URL path for a stored username (no `@`). pub(crate) fn profile_href(username: &str) -> String { diff --git a/server/src/html/ui_action.rs b/server/src/html/ui_action.rs new file mode 100644 index 0000000000000000000000000000000000000000..3bc69ecbb7e16a59d2d393b2d59f9cd3fedb12b1 --- /dev/null +++ b/server/src/html/ui_action.rs @@ -0,0 +1,116 @@ +//! Browser-only UI commands: JSON in hidden `__rpc__` plus hole fill ([`crate::form_template`]). +//! Not part of [`slug_types::RpcCommand`] (CLI / JSON API). + +use crate::form_template::fill_template_from_form; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::collections::HashMap; +use thiserror::Error; + +/// Form field name for the compact JSON template (possibly with `{"$form":"…"}` holes). +pub const UI_RPC_FIELD: &str = "__rpc__"; + +/// HTML form / fetch `POST /ui` payload after template fill and deserialization. +#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)] +#[serde(tag = "action", rename_all = "snake_case")] +pub enum HtmlUiAction { + /// Same semantics as `POST /post` (forum ingest). + PostIngest { + room: String, + thread_tag: String, + text: String, + #[serde(default)] + error_target: Option, + #[serde(default)] + form_id: Option, + }, + /// Same as `POST /post/check`. + CheckIngest { + room: String, + thread_tag: String, + text: String, + #[serde(default)] + error_target: Option, + #[serde(default)] + form_id: Option, + }, + /// Same as `POST /post/redact`. + RedactPost { + post_id: String, + }, + /// Morph `#public-new-thread-ui-slot` to the new-thread form (or login hint). + ExpandPublicNewThreadForm, + /// Morph `#room-new-thread-ui-slot` for the given room wire id. + ExpandRoomNewThreadForm { + room_wire: String, + }, +} + +#[derive(Debug, Error)] +pub enum HtmlUiParseError { + #[error("missing __rpc__ field")] + MissingRpc, + #[error("invalid template json: {0}")] + Template(serde_json::Error), + #[error("invalid ui action: {0}")] + Action(serde_json::Error), +} + +/// Parse `__rpc__` JSON, apply `$form` holes from the rest of the form map, deserialize. +pub fn parse_html_ui_from_form(form: &HashMap) -> Result { + let template = form + .get(UI_RPC_FIELD) + .ok_or(HtmlUiParseError::MissingRpc)?; + let mut hole_map = form.clone(); + hole_map.remove(UI_RPC_FIELD); + let v: Value = fill_template_from_form(template, &hole_map).map_err(HtmlUiParseError::Template)?; + serde_json::from_value(v).map_err(HtmlUiParseError::Action) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn round_trip_post_ingest_with_holes() { + let template = serde_json::json!({ + "action": "post_ingest", + "room": "public", + "thread_tag": {"$form": "thread_tag"}, + "text": {"$form": "text"}, + "error_target": "e", + "form_id": "f", + }); + let mut form = HashMap::new(); + form.insert( + UI_RPC_FIELD.to_string(), + serde_json::to_string(&template).unwrap(), + ); + form.insert("thread_tag".into(), "x".into()); + form.insert("text".into(), "body".into()); + + let a = parse_html_ui_from_form(&form).unwrap(); + assert_eq!( + a, + HtmlUiAction::PostIngest { + room: "public".into(), + thread_tag: "x".into(), + text: "body".into(), + error_target: Some("e".into()), + form_id: Some("f".into()), + } + ); + } + + #[test] + fn expand_public_unit_variant() { + let template = serde_json::json!({ "action": "expand_public_new_thread_form" }); + let mut form = HashMap::new(); + form.insert( + UI_RPC_FIELD.to_string(), + serde_json::to_string(&template).unwrap(), + ); + let a = parse_html_ui_from_form(&form).unwrap(); + assert_eq!(a, HtmlUiAction::ExpandPublicNewThreadForm); + } +} diff --git a/server/src/lib.rs b/server/src/lib.rs index b9d4b79791ba8dd3a3c26ff777943a2548092de5..5f5b6b0f01b29350c7415e5801ec5caa74f0b452 100644 --- a/server/src/lib.rs +++ b/server/src/lib.rs @@ -3,6 +3,7 @@ pub mod paths; pub mod api; pub mod canonical_path; pub mod dsl; +pub mod form_template; pub mod html; pub mod event_log; pub mod events; @@ -33,6 +34,7 @@ pub fn create_app(state: AppState) -> Router { .route("/post", post(api::post_web_ingest)) .route("/post/redact", post(api::post_web_redact)) .route("/post/check", post(api::check_web_ingest)) + .route("/ui", post(api::post_ui_html)) .route("/theme", post(crate::html::post_theme)) .route("/sse", get(api::get_html_stream)) .route("/stream", get(api::get_stream)) 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)