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: [80ad7753] refactor: split canonical_path and identity; strict wire identity without @ - Add canonical_path.rs (tag + item URL normalization) and identity.rs (parse_username/parse_agent; reject @ in API input). - Slim events.rs to event types only; reducer applies no identity rewriting. - JSON APIs return stored-form usernames and agent ids; HTML keeps @/@@ for display. - Optional delegate on ingest; CLI and tests use naked uuid:rig:model. Made-with: Cursor Side A — unified diff (full patch): diff --git a/cli/src/main.rs b/cli/src/main.rs index 5ac8e289f2b7a366b5959d9338d02f9b546f408a..630c5dea1f78c0ec9bc53e6b96234a0dc75bb705 100644 --- a/cli/src/main.rs +++ b/cli/src/main.rs @@ -61,8 +61,8 @@ enum Command { /// Example: --before 2026-06-01 #[arg(long, value_name = "DATE_OR_MS")] before: Option, - /// Filter to posts from this actor (UUID prefix match). - /// Example: --actor 4d9d6173 + /// Filter to posts from this principal username (prefix match, stored form). + /// Example: --actor alice #[arg(long, value_name = "PREFIX")] actor: Option, /// Fetch a single post by its ingest ID (from --json output). @@ -75,9 +75,8 @@ enum Command { /// /// SYNTAX: /// - /// Actor (required, once per document): - /// @:: - /// Example: @7a3b9c2d-1234-5678-90ab-cdef12345678:claudecode:anthropic/claude-sonnet + /// Identity: human comes from the bearer token; optional AI delegate from `--delegate` + /// (`uuid:rig:provider/model`). The document body is DSL only (items, votes, prose) — no `@` lines. /// /// Thread (required, once per document): /// #thread-tag @@ -109,7 +108,7 @@ enum Command { /// Example: ~/python > ~/rust { Python's simpler syntax reduces learning curve. } /// /// Prose (optional, anywhere): - /// Any line that doesn't start with @, #, or ~ is prose. + /// Any line that doesn't start with # or ~ (or `http`) is prose. /// Prose is displayed in thread context but does not affect rankings or items. /// Use prose to write blog posts, reasoning, or notes within your ingest. /// @@ -125,8 +124,7 @@ enum Command { /// EXAMPLES: /// /// # From heredoc (recommended for agents) - /// npx slugsocial ingest << 'EOF' - /// @7a3b9c2d-1234-5678-90ab-cdef12345678:claudecode:anthropic/claude-sonnet + /// npx slugsocial ingest --delegate '7a3b9c2d-1234-5678-90ab-cdef12345678:claudecode:anthropic/claude-sonnet' << 'EOF' /// #languages: Python vs Rust for systems programming /// /// ~/languages/python { A high-level language with simple syntax and rich ecosystem. } @@ -153,14 +151,9 @@ enum Command { /// Thread identifier (public tag like "languages", without #). #[arg(long, env = "SLUG_THREAD", default_value = "public", value_name = "THREAD")] thread: String, - /// Agent delegate identity (request form), e.g. @@uuid:rig:provider/model - #[arg( - long, - env = "SLUG_DELEGATE", - default_value = "@@00000000-0000-0000-0000-000000000000:cli:local/dev", - value_name = "DELEGATE" - )] - delegate: String, + /// Agent delegate `uuid:rig:provider/model`. Omit for human-only ingests. + #[arg(long, env = "SLUG_DELEGATE", value_name = "DELEGATE")] + delegate: Option, /// Output as JSON for agent parsing #[arg(long)] json: bool, @@ -174,14 +167,9 @@ enum Command { /// Thread identifier (public tag like "languages", without #). #[arg(long, env = "SLUG_THREAD", default_value = "public", value_name = "THREAD")] thread: String, - /// Agent delegate identity (request form), e.g. @@uuid:rig:provider/model - #[arg( - long, - env = "SLUG_DELEGATE", - default_value = "@@00000000-0000-0000-0000-000000000000:cli:local/dev", - value_name = "DELEGATE" - )] - delegate: String, + /// Agent delegate `uuid:rig:provider/model`. Omit for human-only ingests. + #[arg(long, env = "SLUG_DELEGATE", value_name = "DELEGATE")] + delegate: Option, /// Output as JSON for agent parsing #[arg(long)] json: bool, @@ -193,10 +181,10 @@ enum Command { /// Useful for agents to catch up on activity after a context reset. /// /// Examples: - /// npx slugsocial feed @:: - /// npx slugsocial feed @:: --since 2026-01-01 + /// npx slugsocial feed tommy + /// npx slugsocial feed tommy --since 2026-01-01 Feed { - /// Actor identifier (@uuid:rig:model) + /// Principal username (stored form) #[arg(value_name = "ACTOR")] actor: String, /// Override the lower bound. Accepts Unix ms or YYYY-MM-DD. @@ -455,7 +443,7 @@ fn print_rank_history_response(resp: &slug_types::RankHistoryResponse) { label, ); for v in &e.caused_by { - println!(" {} {} {} {}", v.a, v.ratio, v.b, v.actor.as_deref().map(|a| format!(" (@{})", a)).unwrap_or_default()); + println!(" {} {} {} {}", v.a, v.ratio, v.b, v.actor.as_deref().map(|a| format!(" ({})", a)).unwrap_or_default()); if !v.body.is_empty() { println!(" {}", v.body.lines().next().unwrap_or(&v.body).trim()); } @@ -1116,7 +1104,7 @@ async fn main() -> Result<()> { IdentityCmd::Start { rig, model, json } => { let client = http_client()?; let uuid = uuid::Uuid::new_v4().to_string(); - let delegate = format!("@@{}:{}:{}", uuid, rig, model); + let delegate = format!("{uuid}:{rig}:{model}"); let start: PendingSessionStartResponse = expect_json( client diff --git a/server/src/api/auth.rs b/server/src/api/auth.rs index cb0faa29b834931e2c2b2f5c174c875e2e2e9346..995ce4a61d29b024c399c656134f541ecfd880cf 100644 --- a/server/src/api/auth.rs +++ b/server/src/api/auth.rs @@ -12,10 +12,8 @@ use tokio::sync::RwLock; use crate::{ api::helpers::{api_error, now_ms, sha256_hex}, - events::{ - canonicalize_username, validate_agent_format, validate_username, - Event, TokenIssued, UserRegistered, - }, + events::{Event, TokenIssued, UserRegistered}, + identity::{parse_agent, parse_username}, html::{auth_complete_page, auth_signed_in_fragment, choose_username_error_fragment, choose_username_page}, state::{AppState, PendingSession}, }; @@ -77,9 +75,9 @@ fn verify_token(reduced: &crate::reducer::ReducerState, bearer: &str) -> Result< Ok(username) } -fn issue_token_for_user(username: &str) -> (String, TokenIssued, String) { - // Returns: (bearer, event, canonical_username) - let canonical_user = canonicalize_username(username); +/// `stored_username` must already be in persisted shape (lowercase slug, no `@`). +fn issue_token_for_user(stored_username: &str) -> (String, TokenIssued) { + let username = stored_username.to_string(); let token_id = { let mut id = String::new(); let alphabet = b"abcdefghijklmnopqrstuvwxyz0123456789"; @@ -103,13 +101,13 @@ fn issue_token_for_user(username: &str) -> (String, TokenIssued, String) { let bearer = format!("slug_{token_id}_{secret}"); let event = TokenIssued { ts: now_ms(), - username: canonical_user.clone(), + username: username.clone(), token_id, token_hash, salt, issued_via: "oauth".to_string(), }; - (bearer, event, canonical_user) + (bearer, event) } #[derive(Debug, Deserialize)] @@ -207,7 +205,7 @@ pub async fn get_auth_callback(Query(q): Query, State(state): s.provider = Some("google".to_string()); s.provider_id = Some(sub.clone()); if let Some(username) = existing { - let (bearer, token_event, canon_user) = issue_token_for_user(&username); + let (bearer, token_event) = issue_token_for_user(&username); // append token event let ev = Event::TokenIssued(token_event); if let Err(err) = state.event_log.append(&ev).await { @@ -217,7 +215,7 @@ pub async fn get_auth_callback(Query(q): Query, State(state): let mut reduced = reduced_arc.write().await; reduced.apply_event(ev); } - s.complete = Some((canon_user, bearer)); + s.complete = Some((username, bearer)); return Redirect::temporary(&format!("{public_url}/auth/complete")).into_response(); } } @@ -251,9 +249,10 @@ pub async fn post_choose_username( State(state): State, Form(form): Form, ) -> impl IntoResponse { - if let Err(msg) = validate_username(&form.username) { - return api_error(StatusCode::BAD_REQUEST, "invalid username", Some(msg)).into_response(); - } + let canon_user = match parse_username(&form.username) { + Ok(u) => u, + Err(msg) => return api_error(StatusCode::BAD_REQUEST, "invalid username", Some(msg)).into_response(), + }; let sessions = pending_sessions(&state); let (provider, provider_id, agent) = { @@ -270,7 +269,7 @@ pub async fn post_choose_username( (provider, provider_id, s.agent.clone()) }; - if let Err(msg) = validate_agent_format(&agent) { + if let Err(msg) = parse_agent(&agent) { return api_error(StatusCode::BAD_REQUEST, "invalid agent format", Some(msg)).into_response(); } @@ -280,7 +279,7 @@ pub async fn post_choose_username( if reduced.users_by_provider.contains_key(&provider_key) { return api_error(StatusCode::CONFLICT, "provider already registered", None).into_response(); } - if reduced.users_by_provider.values().any(|u| u == &canonicalize_username(&form.username)) { + if reduced.users_by_provider.values().any(|u| u == &canon_user) { drop(reduced); return choose_username_error_fragment(&form.session, "that username is taken — try another").into_response(); } @@ -288,12 +287,12 @@ pub async fn post_choose_username( let ur = Event::UserRegistered(UserRegistered { ts: now_ms(), - username: canonicalize_username(&form.username), + username: canon_user.clone(), provider: provider.to_lowercase(), provider_id: provider_id.clone(), }); - let (bearer, ti, canon_user) = issue_token_for_user(&form.username); + let (bearer, ti) = issue_token_for_user(&canon_user); let ti_ev = Event::TokenIssued(ti); // Persist events. @@ -325,15 +324,18 @@ pub async fn post_pending_session( State(state): State, Json(req): Json, ) -> impl IntoResponse { - if let Err(msg) = validate_agent_format(&req.agent) { - return api_error(StatusCode::BAD_REQUEST, "invalid agent format", Some(msg)).into_response(); - } + let agent_naked = match parse_agent(&req.agent) { + Ok(a) => a, + Err(msg) => { + return api_error(StatusCode::BAD_REQUEST, "invalid agent format", Some(msg)).into_response(); + } + }; let session = format!("p_{}", uuid::Uuid::new_v4().simple()); let public_url = std::env::var("SLUG_PUBLIC_URL").unwrap_or_else(|_| "http://127.0.0.1:8080".to_string()); let login_url = format!("{public_url}/auth/login?session={}", urlencoding::encode(&session)); let poll_url = format!("/api/v0/pending-session/{}", session); let s = PendingSession { - agent: req.agent.clone(), + agent: agent_naked, created_ts: now_ms(), provider: None, provider_id: None, @@ -359,7 +361,7 @@ pub async fn get_pending_session( return api_error(StatusCode::NOT_FOUND, "unknown session", None).into_response(); }; let (complete, user, token) = match &s.complete { - Some((u, t)) => (true, Some(format!("@{}", u)), Some(t.clone())), + Some((u, t)) => (true, Some(u.clone()), Some(t.clone())), None => (false, None, None), }; Json(PendingSessionPollResponse { @@ -388,7 +390,7 @@ pub async fn get_whoami(State(state): State, headers: HeaderMap) -> im }; let agents_bound = reduced.agent_bindings.values().filter(|u| *u == &username).count(); Json(WhoamiResponse { - user: format!("@{}", username), + user: username, agents_bound, }) .into_response() diff --git a/server/src/api/feed.rs b/server/src/api/feed.rs index f665d34aef12f2bd6e71d1fb4a3d0b4d509aa775..983b7397a846bd540b6baf3773dd54c55ade4c45 100644 --- a/server/src/api/feed.rs +++ b/server/src/api/feed.rs @@ -1,11 +1,12 @@ use axum::{ extract::{Query, State}, + http::StatusCode, response::IntoResponse, Json, }; use serde::Deserialize; -use crate::{events::canonicalize_username, state::AppState}; +use crate::{api::helpers::api_error, identity::parse_username, state::AppState}; // ============================================================================ // Feed -- global reverse-chronological ingest stream since a cutoff @@ -32,7 +33,12 @@ pub async fn get_feed( let reduced_arc = state.reduced.clone(); let reduced = reduced_arc.read().await; - let actor = canonicalize_username(&q.actor); + let actor = match parse_username(&q.actor) { + Ok(u) => u, + Err(msg) => { + return api_error(StatusCode::BAD_REQUEST, "invalid actor", Some(msg)).into_response(); + } + }; let since = q.since.or_else(|| reduced.actor_last_post_ts.get(&actor).copied()); let cutoff = since.unwrap_or(0); let limit = q.limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT); @@ -65,7 +71,7 @@ pub async fn get_feed( .collect(); Json(slug_types::FeedResponse { - actor: format!("@{}", actor), + actor, since, posts, total, diff --git a/server/src/api/forum.rs b/server/src/api/forum.rs index cf134e97a3aa9f2cea5efcf5bfa2564b9fc2ea46..ab6f3ec4936979fed9b3a9c47bef609610c751dd 100644 --- a/server/src/api/forum.rs +++ b/server/src/api/forum.rs @@ -8,7 +8,8 @@ use serde::Deserialize; use slug_types::*; use crate::{ - events::canonicalize_tag, + canonical_path::canonicalize_tag, + identity::parse_username, state::AppState, }; @@ -46,7 +47,7 @@ pub struct ThreadDetailQuery { pub since: Option, /// Only posts strictly before this Unix ms timestamp. pub before: Option, - /// Filter to posts whose actor starts with this prefix (UUID prefix or full actor string). + /// Filter to posts whose principal username starts with this prefix (stored form, no `@`). pub actor: Option, /// Return the single post with this ingest ID. pub post_id: Option, @@ -56,6 +57,13 @@ pub struct ThreadDetailQuery { pub async fn get_thread(State(state): State, Query(q): Query) -> impl IntoResponse { let reduced_arc = state.reduced.clone(); let tag = canonicalize_tag(&q.tag); + let actor_prefix = match q.actor.as_deref().map(str::trim) { + None | Some("") => String::new(), + Some(s) => match parse_username(s) { + Ok(u) => u, + Err(msg) => return api_error(StatusCode::BAD_REQUEST, "invalid actor filter", Some(msg)).into_response(), + }, + }; let reduced = reduced_arc.read().await; // Single post lookup by ingest ID -- return full body untruncated. @@ -72,7 +80,7 @@ pub async fn get_thread(State(state): State, Query(q): Query, Query(q): Query = all_ids .into_iter() .enumerate() @@ -119,7 +126,7 @@ pub async fn get_thread(State(state): State, Query(q): Query = match &req.delegate { + None => None, + Some(s) if s.trim().is_empty() => None, + Some(s) => match parse_agent(s) { + Ok(d) => Some(d), + Err(msg) => { + drop(reduced); + return api_error(StatusCode::BAD_REQUEST, "invalid delegate format", Some(msg)) + .into_response(); + } + }, + }; let thread_id = canonicalize_tag(&req.thread); let principal = match verify_bearer_principal(&headers, &reduced) { @@ -290,36 +296,43 @@ pub async fn post_ingest( } } - match reduced.agent_bindings.get(&delegate) { - Some(u) if u != &principal => { - drop(reduced); - return api_error( - StatusCode::FORBIDDEN, - "delegate already bound to another user", - None, - ) - .into_response(); + if let Some(ref d) = delegate { + match reduced.agent_bindings.get(d) { + Some(u) if u != &principal => { + drop(reduced); + return api_error( + StatusCode::FORBIDDEN, + "delegate already bound to another user", + None, + ) + .into_response(); + } + _ => {} } - _ => {} } - let need_agent_bind = reduced.agent_bindings.get(&delegate).is_none(); + let need_agent_bind = delegate + .as_ref() + .map(|d| reduced.agent_bindings.get(d).is_none()) + .unwrap_or(false); drop(reduced); let mut events_appended: usize = 0; if need_agent_bind { - let ab = Event::AgentBound(AgentBound { - ts: now_ms(), - agent: delegate.clone(), - username: principal.clone(), - }); - if let Err(err) = event_log.append(&ab).await { - return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!("{err}"), None); - } - events_appended += 1; - { - let mut reduced = reduced_arc.write().await; - reduced.apply_event(ab); + if let Some(agent_id) = delegate.clone() { + let ab = Event::AgentBound(AgentBound { + ts: now_ms(), + agent: agent_id, + username: principal.clone(), + }); + if let Err(err) = event_log.append(&ab).await { + return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!("{err}"), None); + } + events_appended += 1; + { + let mut reduced = reduced_arc.write().await; + reduced.apply_event(ab); + } } } @@ -413,12 +426,18 @@ pub async fn post_check( }; drop(reduced); - if let Err(msg) = validate_agent_format(&req.delegate) { - return api_error(StatusCode::BAD_REQUEST, "invalid delegate format", Some(msg)).into_response(); - } - let delegate = canonicalize_agent(&req.delegate); + let delegate: Option = match &req.delegate { + None => None, + Some(s) if s.trim().is_empty() => None, + Some(s) => match parse_agent(s) { + Ok(d) => Some(d), + Err(msg) => { + return api_error(StatusCode::BAD_REQUEST, "invalid delegate format", Some(msg)).into_response(); + } + }, + }; let thread_id = canonicalize_tag(&req.thread); - let principal = canonicalize_username("placeholder"); + let principal = "placeholder".to_string(); let event = Event::Ingest(Ingest { ts: v.ts, diff --git a/server/src/api/mod.rs b/server/src/api/mod.rs index f1d5c139b875c268cc46c3085c332d8de31b092e..f42e907775645166c60aeae2d98c5855223a71be 100644 --- a/server/src/api/mod.rs +++ b/server/src/api/mod.rs @@ -65,7 +65,7 @@ mod tests { id: format!("test-{ts}"), raw: raw.to_string(), principal: "test".to_string(), - delegate: "@00000000-0000-0000-0000-000000000000:test:local/test".to_string(), + delegate: Some("00000000-0000-0000-0000-000000000000:test:local/test".to_string()), thread_id: "t".to_string(), })); } diff --git a/server/src/api/rank.rs b/server/src/api/rank.rs index 059095bb512096da5e9699ff9c4b92cb7de7fa1c..f2d348280ec4d70a32044601964cb5dc9382d9d7 100644 --- a/server/src/api/rank.rs +++ b/server/src/api/rank.rs @@ -226,7 +226,7 @@ pub async fn get_rank_history( let reduced = reduced_arc.read().await; let content = reduced.public(); - let item_str = crate::events::canonicalize_item(&q.item); + let item_str = crate::canonical_path::canonicalize_item(&q.item); let item = CanonicalItemUrl(item_str.clone()); let entries = content.rank_history.get(&item).cloned().unwrap_or_default(); @@ -238,15 +238,15 @@ pub async fn get_rank_history( .map(|doc| { doc.statements.into_iter().filter_map(|s| { if let crate::dsl::Stmt::Vote { item1, item2, ratio_left, ratio_right, explanation } = s { - let a = crate::events::canonicalize_item(&item1); - let b = crate::events::canonicalize_item(&item2); + let a = crate::canonical_path::canonicalize_item(&item1); + let b = crate::canonical_path::canonicalize_item(&item2); if a == item_str || b == item_str { Some(VoteRow { ts: e.ts, a: item_path_for_api(&a), b: item_path_for_api(&b), ratio: format!("{}:{}", ratio_left, ratio_right), - actor: reduced.ingests_by_id.get(&e.post_id).map(|ing| format!("@{}", ing.principal)), + actor: reduced.ingests_by_id.get(&e.post_id).map(|ing| ing.principal.clone()), body: explanation, thread: Some(format!("#{}", e.thread)), }) diff --git a/server/src/api/search.rs b/server/src/api/search.rs index 8028c3117b9b22460fb9677ca969741f6849d2d0..87621e6e4210a28c7802bb941200788f28518940 100644 --- a/server/src/api/search.rs +++ b/server/src/api/search.rs @@ -119,7 +119,7 @@ pub async fn get_search( .unwrap_or_else(|| "#unknown".to_string()); scored_posts.push((score, ingest.ts, slug_types::SearchPostHit { thread, - actor: format!("@{}", ingest.principal), + actor: ingest.principal.clone(), snippet: snippet_around(&ingest.raw, &words, 160), ts: ingest.ts, })); diff --git a/server/src/api/thread.rs b/server/src/api/thread.rs index 576b0b5be2767b0c6cba5e0701e4ffc9f215029d..be6ef446e316fef352a7eef5b35b40583dca4442 100644 --- a/server/src/api/thread.rs +++ b/server/src/api/thread.rs @@ -8,9 +8,8 @@ use serde::{Deserialize, Serialize}; use crate::{ api::helpers::{api_error, now_ms}, - events::{ - canonicalize_username, Event, GrantAdded, ThreadCapability, ThreadCreated, ThreadVisibility, - }, + events::{Event, GrantAdded, ThreadCapability, ThreadCreated, ThreadVisibility}, + identity::parse_username, state::AppState, }; use super::auth::verify_bearer_principal; @@ -150,10 +149,10 @@ pub async fn post_thread_grants( return api_error(StatusCode::FORBIDDEN, "requires Manage capability", None).into_response(); } - let target = canonicalize_username(&req.username); - if target.is_empty() { - return api_error(StatusCode::BAD_REQUEST, "invalid username", None).into_response(); - } + let target = match parse_username(&req.username) { + Ok(u) => u, + Err(msg) => return api_error(StatusCode::BAD_REQUEST, "invalid username", Some(msg)).into_response(), + }; if !reduced.users_by_provider.values().any(|u| u == &target) { return api_error(StatusCode::NOT_FOUND, format!("user @{target} not found"), None).into_response(); } diff --git a/server/src/canonical_path.rs b/server/src/canonical_path.rs new file mode 100644 index 0000000000000000000000000000000000000000..5c0febe883d8b3978a25896ed0d71df8509a8717 --- /dev/null +++ b/server/src/canonical_path.rs @@ -0,0 +1,92 @@ +//! Normalization for thread tags and ontology item URLs (DSL ↔ stored canonical form). +//! Not event types — see `events` and `path_types`. + +/// Thread / public tag: stored without leading `#`, lowercase. +pub fn canonicalize_tag(input: &str) -> String { + input.trim().trim_start_matches('#').to_lowercase() +} + +/// Ontology item reference → canonical absolute URL on the slug host. +pub fn canonicalize_item(input: &str) -> String { + let s = input.trim(); + if s.is_empty() { + return String::new(); + } + + if let Some(rest) = s.strip_prefix("https://") { + let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); + let host = host.trim().to_lowercase(); + if tail.is_empty() { + return format!("https://{}", host); + } else { + return format!("https://{}/{}", host, tail); + } + } + if let Some(rest) = s.strip_prefix("http://") { + let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); + let host = host.trim().to_lowercase(); + if tail.is_empty() { + return format!("http://{}", host); + } else { + return format!("http://{}/{}", host, tail); + } + } + + let is_tilde = s.starts_with("~/"); + let rest = s.strip_prefix("~/").or_else(|| s.strip_prefix("/")).unwrap_or(s); + + let tail = rest + .split('/') + .filter_map(|seg| { + let t = seg.trim(); + if t.is_empty() { + None + } else { + Some(t.to_lowercase()) + } + }) + .collect::>() + .join("/"); + + if is_tilde { + format!("https://slug.social/~/{}", tail) + } else if tail.is_empty() { + "https://slug.social".to_string() + } else { + format!("https://slug.social/{}", tail) + } +} + +pub fn item_path_segments(input: &str) -> Vec { + let canonical = canonicalize_item(input); + if canonical.is_empty() { + return vec![]; + } + + if let Some(rest) = canonical.strip_prefix("https://") { + let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); + let mut out = vec![format!("https://{}", host)]; + out.extend(tail.split('/').filter(|s| !s.is_empty()).map(|s| s.to_string())); + return out; + } + if let Some(rest) = canonical.strip_prefix("http://") { + let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); + let mut out = vec![format!("http://{}", host)]; + out.extend(tail.split('/').filter(|s| !s.is_empty()).map(|s| s.to_string())); + return out; + } + + canonical + .split('/') + .filter(|s| !s.is_empty()) + .map(|s| s.to_string()) + .collect() +} + +pub fn item_parent_path(input: &str) -> Option { + let segs = item_path_segments(input); + if segs.len() <= 1 { + return None; + } + Some(segs[..segs.len() - 1].join("/")) +} diff --git a/server/src/events.rs b/server/src/events.rs index 7775e59974a521f5fea9361a600802898d05c6ab..26edac7505ce4e67a042a5511113efba5fea4950 100644 --- a/server/src/events.rs +++ b/server/src/events.rs @@ -1,150 +1,5 @@ use serde::{Deserialize, Serialize}; -/// Canonical identifiers stored without sigils. -/// - tags are stored without leading '#' -/// - items are stored without leading '/' -pub fn canonicalize_tag(input: &str) -> String { - input.trim().trim_start_matches('#').to_lowercase() -} - -pub fn canonicalize_item(input: &str) -> String { - let s = input.trim(); - if s.is_empty() { - return String::new(); - } - - if let Some(rest) = s.strip_prefix("https://") { - let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); - let host = host.trim().to_lowercase(); - if tail.is_empty() { - return format!("https://{}", host); - } else { - return format!("https://{}/{}", host, tail); - } - } - if let Some(rest) = s.strip_prefix("http://") { - let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); - let host = host.trim().to_lowercase(); - if tail.is_empty() { - return format!("http://{}", host); - } else { - return format!("http://{}/{}", host, tail); - } - } - - let is_tilde = s.starts_with("~/"); - let rest = s.strip_prefix("~/").or_else(|| s.strip_prefix("/")).unwrap_or(s); - - let tail = rest - .split('/') - .filter_map(|seg| { - let t = seg.trim(); - if t.is_empty() { - None - } else { - Some(t.to_lowercase()) - } - }) - .collect::>() - .join("/"); - - if is_tilde { - format!("https://slug.social/~/{}", tail) - } else if tail.is_empty() { - "https://slug.social".to_string() - } else { - format!("https://slug.social/{}", tail) - } -} - -pub fn item_path_segments(input: &str) -> Vec { - let canonical = canonicalize_item(input); - if canonical.is_empty() { - return vec![]; - } - - if let Some(rest) = canonical.strip_prefix("https://") { - let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); - let mut out = vec![format!("https://{}", host)]; - out.extend(tail.split('/').filter(|s| !s.is_empty()).map(|s| s.to_string())); - return out; - } - if let Some(rest) = canonical.strip_prefix("http://") { - let (host, tail) = rest.split_once('/').map_or((rest, ""), |(h, t)| (h, t)); - let mut out = vec![format!("http://{}", host)]; - out.extend(tail.split('/').filter(|s| !s.is_empty()).map(|s| s.to_string())); - return out; - } - - // Should be unreachable since all canonical items are now URLs - canonical - .split('/') - .filter(|s| !s.is_empty()) - .map(|s| s.to_string()) - .collect() -} - -pub fn item_parent_path(input: &str) -> Option { - let segs = item_path_segments(input); - if segs.len() <= 1 { - return None; - } - Some(segs[..segs.len() - 1].join("/")) -} - -/// Canonical username stored without leading '@'. -pub fn canonicalize_username(input: &str) -> String { - input.trim().trim_start_matches('@').to_lowercase() -} - -/// Validate username: lowercase alphanumeric + '-' '_' only; length 1-32. -pub fn validate_username(username: &str) -> Result<(), String> { - let u = canonicalize_username(username); - if u.is_empty() || u.len() > 32 { - return Err("username must be 1-32 characters".to_string()); - } - if !u - .chars() - .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-' || c == '_') - { - return Err("username must be lowercase alphanumeric with '-' or '_' only".to_string()); - } - Ok(()) -} - -/// Canonical agent identity stored with single leading '@'. -/// -/// Request/display form is `@@uuid:rig:provider/model` but we store `@uuid:rig:provider/model`. -pub fn canonicalize_agent(input: &str) -> String { - let s = input.trim(); - let s = s.strip_prefix("@@").or_else(|| s.strip_prefix('@')).unwrap_or(s); - format!("@{}", s.to_lowercase()) -} - -/// Validate agent format: @@:: -pub fn validate_agent_format(agent: &str) -> Result<(), String> { - let a = agent.trim(); - if !a.starts_with("@@") && !a.starts_with('@') { - return Err("agent must start with @@".to_string()); - } - let a = a.strip_prefix("@@").or_else(|| a.strip_prefix('@')).unwrap_or(a); - let parts: Vec<&str> = a.split(':').collect(); - if parts.len() != 3 { - return Err("agent must be @@::".to_string()); - } - let (uuid_part, rig_part, model_part) = (parts[0], parts[1], parts[2]); - if uuid::Uuid::parse_str(uuid_part).is_err() { - return Err("agent uuid must be a valid UUID v4".to_string()); - } - if rig_part.trim().is_empty() { - return Err("agent rig must be non-empty".to_string()); - } - if !model_part.contains('/') { - return Err("agent model must be ".to_string()); - } - Ok(()) -} - #[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)] #[serde(rename_all = "snake_case")] pub enum ThreadVisibility { @@ -236,10 +91,11 @@ pub struct Ingest { pub id: String, /// Raw DSL+prose body only (no identity/routing metadata). pub raw: String, - /// Human principal username (no leading '@'). + /// Human principal username (wire and storage: no `@`). pub principal: String, - /// Delegate agent identity (canonical stored with single leading '@'). - pub delegate: String, + /// AI delegate id `uuid:rig:model` (wire and storage: no `@`). Omitted when absent. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub delegate: Option, /// Thread identifier: public tag (e.g. "languages") or private id/slug (e.g. "a7f2k9x/project-review"). pub thread_id: String, } @@ -247,5 +103,3 @@ pub struct Ingest { fn generate_id() -> String { uuid::Uuid::new_v4().to_string() } - - diff --git a/server/src/html/editor.rs b/server/src/html/editor.rs index 1d2c4be7eef4cc437cacddd0949307895f3ef018..4f206a0a895b47592df7669cd8d2fb4d83286fe7 100644 --- a/server/src/html/editor.rs +++ b/server/src/html/editor.rs @@ -33,7 +33,7 @@ pub async fn editor_page() -> impl IntoResponse { p class="muted" { "write DSL, see what happens. nothing is saved." } div class="editor-container" { textarea id="editor-input" rows="12" cols="80" - placeholder="@your-uuid:rig:provider/model\n#your-thread\n\n~/path/item-a { description }\n~/path/item-b { description }\n\n~/path/item-a 3:1 ~/path/item-b { reasoning }" + placeholder="your-uuid:rig:provider/model\n#your-thread\n\n~/path/item-a { description }\n~/path/item-b { description }\n\n~/path/item-a 3:1 ~/path/item-b { reasoning }" autocomplete="off" autofocus {} div id="editor-status" class="muted" { "type to check…" } div id="editor-results" {} @@ -110,7 +110,7 @@ pub async fn editor_check( id: uuid::Uuid::new_v4().to_string(), raw: form.text.clone(), principal: String::new(), - delegate: String::new(), + delegate: None, thread_id: String::new(), }); let mut simulated = { reduced_arc.read().await.clone() }; diff --git a/server/src/html/forum.rs b/server/src/html/forum.rs index d40cc6398aaf3b2348d91498b5b3a80ed593e55c..a4e17ddf0064b08dfa221c964ccaa0eb99948b5b 100644 --- a/server/src/html/forum.rs +++ b/server/src/html/forum.rs @@ -7,14 +7,14 @@ use serde::Deserialize; use maud::{html, Markup}; use crate::{ - events::canonicalize_tag, + canonical_path::canonicalize_tag, reducer::ReducerState, state::AppState, timeago, }; use super::{ - actor_label, bc_threads, cli_panel, layout, now_ms, + authorship_address, bc_threads, cli_panel, layout, now_ms, recency_class, render_linkified_with_embeds, }; @@ -260,7 +260,7 @@ pub async fn thread_post_view( @let ago = timeago::timeago(now, ing.ts); div class="ingest-entry" data-ingest-id=(ing.id) { div class="ingest-meta muted" title=(hover) { - span class="address" { "@" (actor_label(&ing.delegate)) } + span class="address" { (authorship_address(&ing.principal, &ing.delegate)) } " · " (ago) } diff --git a/server/src/html/garden.rs b/server/src/html/garden.rs index bddd90aa6eb288c0c97bb08f472e2934fd2fe67a..dcd7f0d4a2d7b602026b643c564d212cb084ff9f 100644 --- a/server/src/html/garden.rs +++ b/server/src/html/garden.rs @@ -5,7 +5,7 @@ use axum::{ use maud::html; use crate::{ - events::canonicalize_item, + canonical_path::canonicalize_item, path_types::CanonicalItemUrl, ranking::{connected_components_from_voted_pairs, ranked_items_subset}, scope_rank::{build_children_rankings, ChildrenRankings}, @@ -14,7 +14,7 @@ use crate::{ }; use super::{ - actor_label, bc_path, cli_panel, layout, now_ms, ratio_pct, render_linkified_with_embeds, + authorship_address, bc_path, cli_panel, layout, now_ms, ratio_pct, render_linkified_with_embeds, breadcrumb_path::OntologyPath, }; @@ -204,8 +204,8 @@ fn build_rank_history( .map(|doc| { doc.statements.into_iter().filter_map(|s| { if let crate::dsl::Stmt::Vote { item1, item2, ratio_left, ratio_right, explanation } = s { - let a_str = crate::events::canonicalize_item(&item1); - let b_str = crate::events::canonicalize_item(&item2); + let a_str = crate::canonical_path::canonicalize_item(&item1); + let b_str = crate::canonical_path::canonicalize_item(&item2); if a_str == item || b_str == item { Some(crate::reducer::VoteData { ts: e.ts, @@ -216,9 +216,7 @@ fn build_rank_history( principal: reduced.ingests_by_id.get(&e.post_id) .map(|ing| ing.principal.clone()) .unwrap_or_default(), - delegate: reduced.ingests_by_id.get(&e.post_id) - .map(|ing| ing.delegate.clone()) - .unwrap_or_default(), + delegate: reduced.ingests_by_id.get(&e.post_id).and_then(|ing| ing.delegate.clone()), thread_id: e.thread.clone(), }) } else { None } @@ -331,7 +329,7 @@ async fn render_scope_view(state: AppState, path: OntologyPath) -> axum::respons @let right_class = if v.b.as_str() == model.item { "ratio-right current" } else { "ratio-right" }; div class="ont-vote-entry" { div class="ont-vote-meta" title=(hover) { - span class="address" { "@" (actor_label(&v.delegate)) } + span class="address" { (authorship_address(&v.principal, &v.delegate)) } " · " (ago) } @@ -482,7 +480,7 @@ mod tests { id: format!("ing-{ts}"), raw: raw.to_string(), principal: "testuser".to_string(), - delegate: "@00000000-0000-0000-0000-000000000000:test:local/test".to_string(), + delegate: Some("00000000-0000-0000-0000-000000000000:test:local/test".to_string()), thread_id: String::new(), })); } diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs index 8b48a25cce79f0eefc7e29849e1a667cae4c7986..82c5993fabfeaf2ea8b9034bccaef37924100264 100644 --- a/server/src/html/mod.rs +++ b/server/src/html/mod.rs @@ -238,9 +238,9 @@ pub(super) fn bc_threads(thread_tag: Option<&str>) -> Markup { } } -/// Input is canonicalized without leading '@' (usually uuid:rig:provider/model). -pub(super) fn actor_label(actor: &str) -> String { - let a = actor.trim_start_matches('@').trim(); +/// Short display label for a stored agent id (`uuid:rig:model`, no `@`). +pub(super) fn actor_label(agent_naked: &str) -> String { + let a = agent_naked.trim(); let parts: Vec<&str> = a.split(':').collect(); if parts.len() >= 3 { let rig = parts[1].trim(); @@ -258,6 +258,13 @@ pub(super) fn actor_label(actor: &str) -> String { a.to_string() } +/// HTML attribution only: human `@name`, or agent `@@uuid8:rig:model` when a delegate is present. +pub(super) fn authorship_address(principal: &str, delegate: &Option) -> String { + match delegate { + Some(d) => format!("@@{}", actor_label(d)), + None => format!("@{}", principal), + } +} /// Escape HTML special chars for safe injection. fn escape_html(s: &str) -> String { diff --git a/server/src/html/search.rs b/server/src/html/search.rs index f56f31546b4870d41e594d0e3ac2a4b50ba831b8..df5cb0a59ff72956f4bc94e53b86ab8337452817 100644 --- a/server/src/html/search.rs +++ b/server/src/html/search.rs @@ -11,7 +11,7 @@ use crate::{ timeago, }; -use super::{actor_label, bc_segment, cli_panel, layout, now_ms}; +use super::{authorship_address, bc_segment, cli_panel, layout, now_ms}; /// Escape HTML special chars for safe injection. fn escape_html(s: &str) -> String { @@ -44,7 +44,8 @@ struct ThreadRow { struct PostRow { thread: String, - actor: String, + /// Pre-formatted attribution string for display (includes `@` / `@@` from `authorship_address`). + actor_display: String, text: String, ts: i64, } @@ -151,7 +152,7 @@ fn search(state: &ReducerState, q: &str, limit: usize) -> SearchResults { .unwrap_or_else(|| "unknown".to_string()); scored_posts.push((score, PostRow { thread, - actor: ingest.principal.clone(), + actor_display: authorship_address(&ingest.principal, &ingest.delegate), text: ingest.raw.clone(), ts: ingest.ts, })); @@ -321,7 +322,7 @@ fn render_search_results(results: &SearchResults, query: &str) -> Markup { li { div class="search-post-meta muted" { a href=(format!("/t/{}", r.thread)) { "#" (r.thread) } - " · " (actor_label(&r.actor)) + " · " (r.actor_display) " · " (timeago::timeago(now, r.ts)) } div class="search-snippet" { diff --git a/server/src/html/tree.rs b/server/src/html/tree.rs index 568a08edbb548291ab6f03b46c79cdcb4b432468..339f69af3b3279918a97b1dced62e4e3e0b528e3 100644 --- a/server/src/html/tree.rs +++ b/server/src/html/tree.rs @@ -13,7 +13,7 @@ use serde::{Deserialize, Serialize}; use std::collections::BTreeSet; use crate::{ - events::canonicalize_item, + canonical_path::canonicalize_item, path_types::{CanonicalItemUrl, RelativePath}, scope_rank::ChildrenRankings, state::AppState, @@ -702,7 +702,8 @@ pub async fn tree_select( mod tests { use super::*; - use crate::events::{canonicalize_item, Event, Ingest}; + use crate::canonical_path::canonicalize_item; + use crate::events::{Event, Ingest}; use crate::reducer::ReducerState; fn ingest(raw: &str) -> Event { @@ -711,7 +712,7 @@ mod tests { id: "test-ingest".to_string(), raw: raw.to_string(), principal: "tester".to_string(), - delegate: String::new(), + delegate: None, thread_id: String::new(), }) } @@ -789,9 +790,9 @@ mod tests { fn reducer_parent_key_for_tilde_items_is_without_trailing_slash() { // This is the reducer invariant that the tree view must match. let item = canonicalize_item("~/alphabet/a"); - assert_eq!(crate::events::item_parent_path(&item).unwrap(), "https://slug.social/~/alphabet"); + assert_eq!(crate::canonical_path::item_parent_path(&item).unwrap(), "https://slug.social/~/alphabet"); let item2 = canonicalize_item("~/a"); - assert_eq!(crate::events::item_parent_path(&item2).unwrap(), "https://slug.social/~"); + assert_eq!(crate::canonical_path::item_parent_path(&item2).unwrap(), "https://slug.social/~"); } #[test] diff --git a/server/src/identity.rs b/server/src/identity.rs new file mode 100644 index 0000000000000000000000000000000000000000..ad654ae813f8f6fe6323746e2544250687d3d47f --- /dev/null +++ b/server/src/identity.rs @@ -0,0 +1,66 @@ +//! Usernames and agent delegate ids. Wire JSON and query params use **stored form only** (no `@` / `@@`). +//! The HTTP layer validates here; the reducer does not rewrite identity. For humans, `@name` / `@@agent` +//! appear only in HTML (see `html::authorship_address`). + +/// Parse username from query/body: trim, lowercase. `@` is not allowed (use `tommy`, not `@tommy`). +pub fn parse_username(input: &str) -> Result { + let s = input.trim(); + if s.is_empty() { + return Err("username must not be empty".to_string()); + } + if s.contains('@') { + return Err( + "username must not contain '@' — use stored form (e.g. `tommy`)".to_string(), + ); + } + let u = s.to_lowercase(); + validate_username_naked(&u)?; + Ok(u) +} + +fn validate_username_naked(u: &str) -> Result<(), String> { + if u.len() > 32 { + return Err("username must be 1-32 characters".to_string()); + } + if !u + .chars() + .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-' || c == '_') + { + return Err("username must be lowercase alphanumeric with '-' or '_' only".to_string()); + } + Ok(()) +} + +/// Parse agent id: trim, lowercase. Must be `uuid:rig:provider/model` with no `@`. +pub fn parse_agent(input: &str) -> Result { + let s = input.trim(); + if s.is_empty() { + return Err("agent id must not be empty".to_string()); + } + if s.contains('@') { + return Err( + "agent id must not contain '@' — use `uuid:rig:provider/model`".to_string(), + ); + } + let a = s.to_lowercase(); + validate_agent_naked(&a)?; + Ok(a) +} + +fn validate_agent_naked(a: &str) -> Result<(), String> { + let parts: Vec<&str> = a.split(':').collect(); + if parts.len() != 3 { + return Err("agent must be ::".to_string()); + } + let (uuid_part, rig_part, model_part) = (parts[0], parts[1], parts[2]); + if uuid::Uuid::parse_str(uuid_part).is_err() { + return Err("agent uuid must be a valid UUID v4".to_string()); + } + if rig_part.trim().is_empty() { + return Err("agent rig must be non-empty".to_string()); + } + if !model_part.contains('/') { + return Err("agent model must be ".to_string()); + } + Ok(()) +} diff --git a/server/src/lib.rs b/server/src/lib.rs index 1820b0dd9cd2d0b6ce05789e3225e0d86a3d357a..173fbdd6b3f23c062a200b268d04491d5e030b3f 100644 --- a/server/src/lib.rs +++ b/server/src/lib.rs @@ -1,10 +1,12 @@ #[macro_use] pub mod paths; pub mod api; +pub mod canonical_path; pub mod dsl; pub mod html; pub mod event_log; pub mod events; +pub mod identity; pub mod middleware; pub mod path_types; pub mod ranking; diff --git a/server/src/path_types.rs b/server/src/path_types.rs index 5a903de4f8309b6c3874e2cd70f2eb5650fd50a8..997597f6659c1ba2092bc5348e393bf7a8fa7ae0 100644 --- a/server/src/path_types.rs +++ b/server/src/path_types.rs @@ -14,9 +14,9 @@ use std::fmt; use serde::{Deserialize, Serialize}; -use crate::events::canonicalize_item; +use crate::canonical_path::canonicalize_item; -/// Canonical item identifier as produced by `events::canonicalize_item`. +/// Canonical item identifier as produced by `canonical_path::canonicalize_item`. /// /// In practice this is usually: /// - `https://slug.social/~/...` for ontology items, or diff --git a/server/src/ranking.rs b/server/src/ranking.rs index 70119611473a430113429eafb6d92b2bf96a3eb5..f70da0d076da2ba8f0abb5a4415babae1c2305f9 100644 --- a/server/src/ranking.rs +++ b/server/src/ranking.rs @@ -265,7 +265,7 @@ mod tests { ratio_right: r, body: "because".to_string(), principal: "test".to_string(), - delegate: "@00000000-0000-0000-0000-000000000000:test:local/test".to_string(), + delegate: Some("00000000-0000-0000-0000-000000000000:test:local/test".to_string()), thread_id: "untagged".to_string(), } } diff --git a/server/src/reducer.rs b/server/src/reducer.rs index 9ab063a340e228edaa8967ec35762809fccd5d3a..550bcda1a1068df13908d19484d86ad46d65ddf1 100644 --- a/server/src/reducer.rs +++ b/server/src/reducer.rs @@ -2,7 +2,8 @@ use std::collections::{HashMap, HashSet, VecDeque}; use serde::{Deserialize, Serialize}; -use crate::events::{canonicalize_agent, canonicalize_tag, canonicalize_username, Event, Ingest, ThreadCapability}; +use crate::canonical_path::canonicalize_tag; +use crate::events::{Event, Ingest, ThreadCapability}; use crate::path_types::CanonicalItemUrl; #[derive(Debug, Clone, Hash, PartialEq, Eq, PartialOrd, Ord)] @@ -21,10 +22,10 @@ pub struct VoteData { pub ratio_left: i32, pub ratio_right: i32, pub body: String, - /// Human principal username (no leading '@'). + /// Human principal username (no `@` in stored events). pub principal: String, - /// Agent delegate identity (stored with single leading '@'). - pub delegate: String, + /// AI delegate id, if any (no `@` in stored events). + pub delegate: Option, /// Thread id where this vote was cast (public tag or private id/slug). pub thread_id: String, } @@ -99,8 +100,6 @@ impl GroupState { pub fn apply_vote(&mut self, mut vote: VoteData) { vote.a = CanonicalItemUrl::parse(vote.a.as_str()).unwrap_or(vote.a); vote.b = CanonicalItemUrl::parse(vote.b.as_str()).unwrap_or(vote.b); - vote.principal = canonicalize_username(&vote.principal); - vote.delegate = canonicalize_agent(&vote.delegate); vote.thread_id = canonicalize_tag(&vote.thread_id); if vote.ratio_left < 0 { vote.ratio_left = 0; @@ -189,7 +188,7 @@ pub struct ReducerState { pub users_by_provider: HashMap<(String, String), String>, /// token_id -> (username, salt, token_hash) pub tokens_by_id: HashMap, - /// agent delegate (canonical '@...') -> username + /// agent id (naked `uuid:rig:model`) -> username pub agent_bindings: HashMap, pub ingests_by_id: HashMap, @@ -330,26 +329,27 @@ impl ReducerState { pub fn apply_event(&mut self, event: Event) { match event { Event::UserRegistered(ur) => { - let username = canonicalize_username(&ur.username); - self.users_by_provider - .insert((ur.provider.to_lowercase(), ur.provider_id.clone()), username); + self.users_by_provider.insert( + (ur.provider.to_lowercase(), ur.provider_id.clone()), + ur.username, + ); } Event::TokenIssued(ti) => { - let username = canonicalize_username(&ti.username); - self.tokens_by_id - .insert(ti.token_id.clone(), (username, ti.salt.clone(), ti.token_hash.clone())); + self.tokens_by_id.insert( + ti.token_id.clone(), + (ti.username, ti.salt.clone(), ti.token_hash.clone()), + ); } Event::AgentBound(ab) => { - let username = canonicalize_username(&ab.username); - let agent = canonicalize_agent(&ab.agent); - self.agent_bindings.insert(agent, username); + if ab.agent.is_empty() { + return; + } + self.agent_bindings.insert(ab.agent, ab.username); } Event::ThreadCreated(tc) => { self.threads.entry(tc.thread_id.clone()).or_default().visibility = tc.visibility; } Event::Ingest(mut ing) => { - ing.principal = canonicalize_username(&ing.principal); - ing.delegate = canonicalize_agent(&ing.delegate); ing.thread_id = canonicalize_tag(&ing.thread_id); self.ingests_by_id.insert(ing.id.clone(), ing.clone()); @@ -528,7 +528,7 @@ impl ReducerState { let caps = self.grants .entry(ga.thread_id) .or_default() - .entry(canonicalize_username(&ga.username)) + .entry(ga.username) .or_default(); for cap in ga.capabilities { caps.insert(cap); @@ -536,7 +536,7 @@ impl ReducerState { } Event::GrantRevoked(gr) => { if let Some(thread_grants) = self.grants.get_mut(&gr.thread_id) { - let username = canonicalize_username(&gr.username); + let username = gr.username; if let Some(caps) = thread_grants.get_mut(&username) { for cap in &gr.capabilities { caps.remove(cap); diff --git a/server/tests/basic.rs b/server/tests/basic.rs index 8d7dd87d327c969ad953ecaa6b021a03e4d1842e..21cacd3049c9f5a307e75eefd5ebfc7bf0097b5d 100644 --- a/server/tests/basic.rs +++ b/server/tests/basic.rs @@ -1,6 +1,7 @@ use slugsocial_server::{ event_log::EventLog, - events::{canonicalize_item, canonicalize_tag, Event, Ingest}, + canonical_path::{canonicalize_item, canonicalize_tag}, + events::{Event, Ingest}, ranking::ranked_items, reducer::{GroupState, ReducerState}, }; @@ -15,7 +16,7 @@ fn ingest_event(ts: i64, raw: &str) -> Event { id: format!("test-{ts}"), raw: raw.to_string(), principal: "test".to_string(), - delegate: "@@00000000-0000-0000-0000-000000000000:test:local/test".to_string(), + delegate: Some("00000000-0000-0000-0000-000000000000:test:local/test".to_string()), thread_id: "t".to_string(), }) } @@ -491,7 +492,7 @@ fn reducer_negative_ratio_clamped_to_zero() { ratio_right: -3, body: "negative".to_string(), principal: "test".to_string(), - delegate: "@@00000000-0000-0000-0000-000000000000:test:local/test".to_string(), + delegate: Some("00000000-0000-0000-0000-000000000000:test:local/test".to_string()), thread_id: "t".to_string(), }); assert_eq!(group.idx_to_item.len(), 2); diff --git a/server/tests/integration.rs b/server/tests/integration.rs index 24ef464631cd269889b8bb751ef0e90289443e97..eaa25f2c5cc21df3096d9b8f54c060c83207bfd6 100644 --- a/server/tests/integration.rs +++ b/server/tests/integration.rs @@ -87,7 +87,7 @@ async fn test_ingest_actor_with_colons_is_detected_and_validated() { // Old archive style: agent includes colons but UUID is only a prefix (invalid). // We should detect the agent line, then fail with "invalid agent format". let ingest_payload = serde_json::json!({ - "delegate": "@@aec1e31c:claudecode:anthropic/claude-sonnet-4.5", + "delegate": "aec1e31c:claudecode:anthropic/claude-sonnet-4.5", "thread": "t", "text": "~/x {x}\n", }); @@ -108,6 +108,26 @@ async fn test_ingest_actor_with_colons_is_detected_and_validated() { hint.to_lowercase().contains("uuid"), "hint should mention uuid, got: {hint}" ); + + let at_payload = serde_json::json!({ + "delegate": "@00000000-0000-0000-0000-000000000000:test:local/test", + "thread": "t", + "text": "~/x {x}\n", + }); + let at_resp = client + .post(&format!("http://{}/api/v0/ingest", addr)) + .header("Authorization", format!("Bearer {}", test_bearer())) + .json(&at_payload) + .send() + .await + .unwrap(); + assert_eq!(at_resp.status(), reqwest::StatusCode::BAD_REQUEST); + let at_body: serde_json::Value = at_resp.json().await.unwrap(); + let at_hint = at_body["hint"].as_str().unwrap_or_default(); + assert!( + at_hint.contains('@'), + "hint should reject '@' in delegate, got: {at_hint}" + ); } #[tokio::test] @@ -117,7 +137,7 @@ async fn test_vote_endpoint() { // /api/v0/vote was removed; all votes are submitted via ingest. let ingest_payload = serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000000:test:local/test", + "delegate": "00000000-0000-0000-0000-000000000000:test:local/test", "thread": "cli", "text": "~/clap {cli parser}\n~/argh {cli parser}\n~/clap 3:1 ~/argh {because clap is more full-featured}\n", }); @@ -143,7 +163,7 @@ async fn test_rank_endpoint() { // Ingest items + vote (vote endpoint removed). let ingest_payload = serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000000:test:local/test", + "delegate": "00000000-0000-0000-0000-000000000000:test:local/test", "thread": "langs", "text": "~/rust {systems}\n~/go {concurrency}\n~/rust 3:1 ~/go {because i prefer rust for systems work}\n", }); @@ -180,7 +200,7 @@ async fn test_check_endpoint_does_not_commit() { let client = reqwest::Client::new(); let check_payload = serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000000:test:local/test", + "delegate": "00000000-0000-0000-0000-000000000000:test:local/test", "thread": "t", "text": "~/a {x}\n~/b {y}\n~/a 2:1 ~/b {because}\n", }); @@ -220,7 +240,7 @@ async fn test_garden_item_pair_matchup_include_threads() { // Ingest with thread_id metadata so item_threads and vote.thread_id are populated. let ingest_payload = serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000000:test:local/test", + "delegate": "00000000-0000-0000-0000-000000000000:test:local/test", "thread": "sorting-hat", "text": "~/sorts/insertion { O(n^2) }\n~/sorts/mergesort { O(n log n) }\n~/sorts/insertion 3:1 ~/sorts/mergesort { simpler for small n }\n", }); @@ -324,7 +344,7 @@ async fn test_rank_history() { // First ingest: rust vs python — two votes on rust in one doc (the multi-vote case). ingest( serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000001:rig:test/model", + "delegate": "00000000-0000-0000-0000-000000000001:rig:test/model", "thread": "hist-test", "text": "~/hist/rust { systems }\n~/hist/python { scripting }\n~/hist/go { concurrency }\n~/hist/rust 3:1 ~/hist/python { ownership over gc }\n~/hist/rust 2:1 ~/hist/go { performance over simplicity }\n", }), @@ -354,7 +374,7 @@ async fn test_rank_history() { // Second ingest: python beats go — rust not directly touched, so python gets a new entry. ingest( serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000002:rig:test/model", + "delegate": "00000000-0000-0000-0000-000000000002:rig:test/model", "thread": "hist-test", "text": "~/hist/python 3:1 ~/hist/go { dynamic typing is worth it }\n", }), @@ -404,7 +424,7 @@ async fn pair_returns_connectivity_stats() { // Ingest 4 items with 1 vote (a vs b), leaving c and d as isolates. let doc = serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000001:testrig:test/model", + "delegate": "00000000-0000-0000-0000-000000000001:testrig:test/model", "thread": "connectivity-test", "text": "~/conn/a { item a }\n~/conn/b { item b }\n~/conn/c { item c }\n~/conn/d { item d }\n~/conn/a 3:1 ~/conn/b { a is better }\n", }); @@ -438,7 +458,7 @@ async fn pair_returns_connectivity_stats() { // Add a vote connecting c to a — should reduce components. let doc2 = serde_json::json!({ - "delegate": "@@00000000-0000-0000-0000-000000000001:testrig:test/model", + "delegate": "00000000-0000-0000-0000-000000000001:testrig:test/model", "thread": "connectivity-test", "text": "~/conn/c 2:1 ~/conn/a { c beats a }\n", }); diff --git a/test/auth.bb b/test/auth.bb index b54ea509aafd2d5ee85877665982d7fc00cd8b20..611e04f1806ef81d678a91f499597fe55f691dd7 100644 --- a/test/auth.bb +++ b/test/auth.bb @@ -155,7 +155,7 @@ (println "\nstarting pending session…") (let [start-resp (http-post-json (str base-url "/api/v0/pending-session") - {:agent "@@00000000-0000-0000-0000-000000000000:bb:local/dev"}) + {:agent "00000000-0000-0000-0000-000000000000:bb:local/dev"}) _ (assert! (= 200 (:status start-resp)) "pending-session start returns 200") start-json (json/parse-string (:body start-resp) true)] (assert! (clojure.string/starts-with? (:session start-json) "p_") "session id has p_ prefix") diff --git a/test/grants.bb b/test/grants.bb index 962ff01cdc25d6d6977cc9a810c404918c20856d..f72799ae6fe099f48bdf316deee07f15882e9d45 100644 --- a/test/grants.bb +++ b/test/grants.bb @@ -187,11 +187,11 @@ ;; Register two users. The mock google cycles through google-user-alice then google-user-bob. (println "\nregistering alice…") (let [alice-token (register-user base-url - "@@00000000-0000-0000-0000-000000000001:test:local/dev" + "00000000-0000-0000-0000-000000000001:test:local/dev" "alice") _ (println "registering bob…") bob-token (register-user base-url - "@@00000000-0000-0000-0000-000000000002:test:local/dev" + "00000000-0000-0000-0000-000000000002:test:local/dev" "bob") ;; Alice creates a private thread. @@ -206,14 +206,14 @@ ;; Alice (owner) can post prose to her own private thread. (println "\nalice posts prose to her private thread…") (assert! (= 200 (:status (ingest! base-url alice-token thread-id - "@@00000000-0000-0000-0000-000000000001:test:local/dev" + "00000000-0000-0000-0000-000000000001:test:local/dev" "Hello from alice."))) "alice prose post succeeds") ;; Bob has no grants at all — should get 403. (println "\nbob (no grants) tries to post prose…") (assert! (= 403 (:status (ingest! base-url bob-token thread-id - "@@00000000-0000-0000-0000-000000000002:test:local/dev" + "00000000-0000-0000-0000-000000000002:test:local/dev" "Hello from bob, unauthorized."))) "bob without grants gets 403") @@ -226,7 +226,7 @@ (println "\nbob (View only) tries to post prose…") (assert! (= 403 (:status (ingest! base-url bob-token thread-id - "@@00000000-0000-0000-0000-000000000002:test:local/dev" + "00000000-0000-0000-0000-000000000002:test:local/dev" "Hello from bob, view only."))) "bob with View but no Post gets 403") @@ -239,21 +239,21 @@ (println "\nbob (View + Post) posts prose…") (assert! (= 200 (:status (ingest! base-url bob-token thread-id - "@@00000000-0000-0000-0000-000000000002:test:local/dev" + "00000000-0000-0000-0000-000000000002:test:local/dev" "Hello from bob, now authorised."))) "bob with View + Post succeeds for prose") ;; Alice defines two items and votes on them in the private thread. (println "\nalice posts items + vote to private thread…") (assert! (= 200 (:status (ingest! base-url alice-token thread-id - "@@00000000-0000-0000-0000-000000000001:test:local/dev" + "00000000-0000-0000-0000-000000000001:test:local/dev" "~/fruits/apple { A crisp red apple. }\n~/fruits/banana { A yellow banana. }\n~/fruits/apple > ~/fruits/banana { apples are better }"))) "alice vote in private thread succeeds") ;; Bob (View + Post, no Vote) tries to vote — should be 403. (println "\nbob (no Vote) tries to vote…") (assert! (= 403 (:status (ingest! base-url bob-token thread-id - "@@00000000-0000-0000-0000-000000000002:test:local/dev" + "00000000-0000-0000-0000-000000000002:test:local/dev" "~/fruits/apple > ~/fruits/banana { bob's take }"))) "bob without Vote gets 403") @@ -266,7 +266,7 @@ (println "\nbob (View + Post + Vote) votes…") (assert! (= 200 (:status (ingest! base-url bob-token thread-id - "@@00000000-0000-0000-0000-000000000002:test:local/dev" + "00000000-0000-0000-0000-000000000002:test:local/dev" "~/fruits/apple > ~/fruits/banana { bob's take }"))) "bob with Vote succeeds")) diff --git a/test/integration.bb b/test/integration.bb index 93db8f3102b1b3a05f3c4bd2cc9f348072fe54cc..4ff5d625840be3bf5e4162cef107a3cbd45aec4d 100644 --- a/test/integration.bb +++ b/test/integration.bb @@ -166,7 +166,7 @@ ;; 3. ingest via CLI (bearer required) (println "\ningesting .sorter document via CLI…") - (bind ingest1-result (common/run-cli cli-bin base-url ["ingest" "--json" "--thread" "integration-test"] :input sorter-doc :extra-env token-env)) + (bind ingest1-result (common/run-cli cli-bin base-url ["ingest" "--json" "--thread" "integration-test" "--delegate" "00000000-0000-0000-0000-000000000000:cli:local/dev"] :input sorter-doc :extra-env token-env)) (assert! (zero? (:exit ingest1-result)) "cli ingest exits 0") (bind ingest1-resp (json/parse-string (:out ingest1-result) true)) (assert! (:ok ingest1-resp) "ingest response ok=true") @@ -237,7 +237,7 @@ "#integration-test" "~/languages/rust 4:1 ~/languages/python { type safety }" "~/languages/rust 3:1 ~/languages/go { zero-cost abstractions }"])) - (bind hist-ingest (common/run-cli cli-bin base-url ["ingest" "--json" "--thread" "integration-test"] :input two-vote-doc :extra-env token-env)) + (bind hist-ingest (common/run-cli cli-bin base-url ["ingest" "--json" "--thread" "integration-test" "--delegate" "00000000-0000-0000-0000-000000000000:cli:local/dev"] :input two-vote-doc :extra-env token-env)) (assert! (zero? (:exit hist-ingest)) (str "two-vote ingest exits 0 (err: " (:err hist-ingest) ")")) diff --git a/test/oauth.bb b/test/oauth.bb index f94583ee0f6fedfc023564ed7ee3a2986adcf660..84679e1a4adbe100dfe7e7d64a2c6270b023d275 100644 --- a/test/oauth.bb +++ b/test/oauth.bb @@ -85,11 +85,11 @@ stop-fn (http/run-server handler {:port port})] {:stop-fn stop-fn :port port})) -(def ^:private default-agent "@@00000000-0000-0000-0000-000000000000:cli:local/dev") +(def ^:private default-agent "00000000-0000-0000-0000-000000000000:cli:local/dev") (defn fetch-bearer-token! "Simulate browser OAuth + username choice; returns `slug_…` bearer token. - Agent must match CLI default `SLUG_DELEGATE` for ingest binding." + Ingest `--delegate` must match this agent string for `AgentBound` on first write." [base-url & {:keys [username agent] :or {username "intuser" agent default-agent}}] (let [start-resp (http-post-json (str base-url "/api/v0/pending-session") {:agent agent})] diff --git a/types/src/lib.rs b/types/src/lib.rs index a62822a36463c68fde6c5fc1aa9e4273c82ca0c8..92e6d122d573ea225f8fb86c548a81b3454425ad 100644 --- a/types/src/lib.rs +++ b/types/src/lib.rs @@ -138,7 +138,7 @@ pub struct PostRow { /// Chronological index within the thread (0 = oldest). pub index: usize, pub ts: i64, - /// Self-declared actor (`@uuid:rig:model`). + /// Principal username (stored form, no `@`). pub actor: String, pub body: String, pub truncated: bool, @@ -186,6 +186,7 @@ pub struct VoteRow { pub a: String, pub b: String, pub ratio: String, + /// Principal username when present (stored form, no `@`). pub actor: Option, pub body: String, /// Thread where this vote was cast (e.g. "#sorting-hat"). @@ -196,6 +197,7 @@ pub struct VoteRow { /// Response for the feed endpoint — all ingests since a cutoff, newest first. #[derive(Debug, Serialize, Deserialize)] pub struct FeedResponse { + /// Principal username this feed is scoped to (stored form, no `@`). pub actor: String, /// The lower-bound timestamp used (actor's last ingest, ms). None if actor has never posted. #[serde(default, skip_serializing_if = "Option::is_none")] @@ -222,8 +224,9 @@ pub struct FeedPost { pub struct IngestRequest { /// Thread identifier: public tag (e.g. "languages") or private id/slug (e.g. "a7f2k9x/project-review"). pub thread: String, - /// Agent delegate identity in request form (e.g. "@@uuid:rig:provider/model"). - pub delegate: String, + /// Delegate id: `uuid:rig:provider/model` (no `@`). Omit for human-only ingests. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub delegate: Option, /// DSL+prose body only. pub text: String, } @@ -231,7 +234,7 @@ pub struct IngestRequest { /// Start a browser-based OAuth login flow for a CLI agent. #[derive(Debug, Serialize, Deserialize)] pub struct PendingSessionStartRequest { - /// Agent delegate identity in request form (e.g. "@@uuid:rig:provider/model"). + /// Delegate id: `uuid:rig:provider/model` (no `@`). pub agent: String, } @@ -255,6 +258,7 @@ pub struct PendingSessionPollResponse { #[derive(Debug, Serialize, Deserialize)] pub struct WhoamiResponse { + /// Username (stored form, no `@`). pub user: String, pub agents_bound: usize, } 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)