{"messages":[{"content":"You are a constitutional council ranking individual git commits for ownership allocation.\n\nCompare these two commits. Decide which contributed more lasting value to the project.\n\nJudge substance, not spectacle:\n- Prefer correct, lasting design and real bugfixes over churn, formatting, renames, or generated noise.\n- Prefer clarity and necessity over sheer line count. A small precise change can beat a large diffuse one.\n- Do not favor a side merely because its patch is longer or noisier.\n- Weight what the change does for the project, not the contributor's name.\n\nReturn ONLY a JSON object: {\"winner\": \"A\" or \"B\", \"ratio\": \"N:M\", \"explanation\": \"...\"}\nThe explanation must cite concrete differences in the patches (1-3 sentences).\n\nSide A — contributor: tommy-mor\nSide A — commit message:\n[40b975bf] nice\n\nSide A — unified diff (full patch):\ndiff --git a/Cargo.lock b/Cargo.lock\nindex 266e876bb7ccbe788beb1d5bd53ad5b45ee5825b..2cea973082716e761ef6f5dd5886acc08ff9aac0 100644\n--- a/Cargo.lock\n+++ b/Cargo.lock\n@@ -222,6 +222,12 @@ dependencies = [\n \"syn\",\n ]\n \n+[[package]]\n+name = \"dotenvy\"\n+version = \"0.15.7\"\n+source = \"registry+https://github.com/rust-lang/crates.io-index\"\n+checksum = \"1aaf95b3e5c8f23aa320147307562d361db0ae0d51242340f558153b4eb2439b\"\n+\n [[package]]\n name = \"encoding_rs\"\n version = \"0.8.35\"\n@@ -1238,6 +1244,7 @@ version = \"0.0.1\"\n dependencies = [\n \"axum\",\n \"axum-extra\",\n+ \"dotenvy\",\n \"maud\",\n \"reqwest\",\n \"serde\",\ndiff --git a/server/Cargo.toml b/server/Cargo.toml\nindex 4677fedcb45292eebebe7e9cf6ce2f5738f18ddf..bd600138b613bd0f546bdec217a5334cdcb20aa5 100644\n--- a/server/Cargo.toml\n+++ b/server/Cargo.toml\n@@ -17,6 +17,7 @@ tower-http = { version = \"0.5\", features = [\"trace\"] }\n tracing = \"0.1\"\n tracing-subscriber = { version = \"0.3\", features = [\"env-filter\"] }\n reqwest = { version = \"0.12\", features = [\"json\"] }\n+dotenvy = \"0.15\"\n \n [dev-dependencies]\n reqwest = { version = \"0.12\", features = [\"json\"] }\ndiff --git a/server/src/api/ui_html.rs b/server/src/api/ui_html.rs\nindex d2024bd4582bcc8482b461b2ba4fedbd8bff7c66..b33a84e8bb5e817b26592868d88090e6d664d950 100644\n--- a/server/src/api/ui_html.rs\n+++ b/server/src/api/ui_html.rs\n@@ -6,7 +6,7 @@ use axum::{\n use std::collections::HashMap;\n \n use crate::{\n- html::{input_panel, js_string_literal, ranking_panel, JsBuilder},\n+ html::{entity_section, input_panel, js_string_literal, ranking_panel, JsBuilder},\n parser::parse_reddit_url,\n path_types::ItemId,\n reddit::ensure_partial_tree,\n@@ -87,6 +87,20 @@ pub async fn post_ui_html(\n .into_response()\n }\n },\n+ HtmlUiAction::FetchEntity { item } => {\n+ let id = parse_item_param(&item);\n+ if id.is_root() {\n+ return ui_js_warn(\"nothing to fetch for the root\").into_response();\n+ }\n+ state.queue_entity_fetch(id.clone());\n+ let tree = state.tree.read().await;\n+ let empty = crate::reducer::NodeState::default();\n+ let node = tree.get(&id).unwrap_or(&empty);\n+ let panel = entity_section(&id, node, true);\n+ JsBuilder::new()\n+ .morph_selector(\"#entity-section\", panel)\n+ .into_response()\n+ },\n }\n }\n \ndiff --git a/server/src/events.rs b/server/src/events.rs\nindex ed5be6b13b9d46e838831d6ce0f96f569b401730..07ce24b5e56cf72b0b442c3c3241efbf6c3b006a 100644\n--- a/server/src/events.rs\n+++ b/server/src/events.rs\n@@ -1,4 +1,5 @@\n use serde::{Deserialize, Serialize};\n+use serde_json::Value;\n \n #[derive(Debug, Clone, Serialize, Deserialize)]\n #[serde(tag = \"type\", rename_all = \"snake_case\")]\n@@ -18,4 +19,10 @@ pub enum Event {\n },\n /// Register a node path in the fractal tree (no external fetch).\n NodeEnsured { id: String },\n+ /// Full upstream API payload for a node (domain-specific view derived at replay/render time).\n+ EntityImported {\n+ id: String,\n+ ts: i64,\n+ payload: Value,\n+ },\n }\ndiff --git a/server/src/html/mod.rs b/server/src/html/mod.rs\nindex df6505021d9f446c2b453e20e3eb3cf696a111f9..caf1309c8d93b47104499c57f9cc35ee7631fbb9 100644\n--- a/server/src/html/mod.rs\n+++ b/server/src/html/mod.rs\n@@ -10,6 +10,7 @@ use crate::{\n form_template::template_json_compact,\n path_types::ItemId,\n ranking::{top_bottom, RankedItem},\n+ reddit::is_fetchable,\n reducer::{GroupState, NodeState},\n state::AppState,\n ui_action::UI_RPC_FIELD,\n@@ -151,7 +152,7 @@ pub fn breadcrumb_path(item: &ItemId) -> Markup {\n fn entity_panel(node: &NodeState) -> Markup {\n html! {\n @if let Some(data) = &node.data {\n- section id=\"entity-panel\" class=\"demo-panel entity-card\" {\n+ div id=\"entity-panel\" class=\"entity-card\" {\n h2 { (data.title) }\n @if let Some(author) = &data.author {\n p class=\"muted small\" { \"by \" (author) }\n@@ -164,6 +165,42 @@ fn entity_panel(node: &NodeState) -> Markup {\n }\n }\n \n+/// Reddit/API import control — only shown on fetchable pages; never auto-fires.\n+pub fn fetch_entity_panel(item: &ItemId, has_data: bool, fetching: bool) -> Markup {\n+ if !is_fetchable(item) {\n+ return html! {};\n+ }\n+ let label = if fetching {\n+ \"Fetching…\"\n+ } else if has_data {\n+ \"Fetch more\"\n+ } else {\n+ \"Fetch from Reddit\"\n+ };\n+ let rpc = template_json_compact(&serde_json::json!({\n+ \"action\": \"fetch_entity\",\n+ \"item\": item.as_str(),\n+ }))\n+ .expect(\"fetch_entity rpc template\");\n+ html! {\n+ form method=\"post\" action=\"/ui\" id=\"fetch-entity-form\" class=\"fetch-entity-form\" {\n+ input type=\"hidden\" name=(UI_RPC_FIELD) value=(rpc);\n+ button type=\"submit\" class=\"btn-secondary\" disabled=(fetching) { (label) }\n+ }\n+ }\n+}\n+\n+/// Entity card + explicit fetch control (morphed as `#entity-section`).\n+pub fn entity_section(item: &ItemId, node: &NodeState, fetching: bool) -> Markup {\n+ let has_data = node.data.is_some();\n+ html! {\n+ section id=\"entity-section\" class=\"demo-panel\" {\n+ (entity_panel(node))\n+ (fetch_entity_panel(item, has_data, fetching))\n+ }\n+ }\n+}\n+\n fn rank_list(label: &str, items: &[RankedItem], start_rank: usize) -> Markup {\n html! {\n @if !items.is_empty() {\n@@ -260,7 +297,7 @@ async fn item_page(state: AppState, uri: Uri, item: ItemId) -> Markup {\n h1 { \"sorter\" }\n (input_panel(\"\", None))\n (breadcrumb_path(&item))\n- (entity_panel(node))\n+ (entity_section(&item, node, false))\n (ranking_panel(&item, group))\n };\n layout(\"sorter2\", body, views)\n@@ -272,16 +309,5 @@ pub async fn home(State(state): State, uri: Uri) -> impl IntoResponse\n \n pub async fn browse(State(state): State, uri: Uri) -> impl IntoResponse {\n let item = ItemId::from_browse_uri(uri.path()).unwrap_or(ItemId::root());\n- if item.as_str().starts_with(\"reddit.com\") {\n- let needs_fetch = {\n- let tree = state.tree.read().await;\n- tree.get(&item)\n- .map(|n| n.data.is_none())\n- .unwrap_or(true)\n- };\n- if needs_fetch {\n- state.reddit.request_fetch(item.clone());\n- }\n- }\n item_page(state, uri, item).await\n }\ndiff --git a/server/src/main.rs b/server/src/main.rs\nindex c22ec6c9f5358e5ec99fb83210dc351938505a93..1f0cddc39302b35b0cd6a6219f44c9d59202facf 100644\n--- a/server/src/main.rs\n+++ b/server/src/main.rs\n@@ -2,6 +2,10 @@ use sorter2_server::state::AppConfig;\n \n #[tokio::main]\n async fn main() -> Result<(), Box> {\n+ if std::env::var(\"SORTER2_SKIP_DOTENV\").is_err() {\n+ let _ = dotenvy::dotenv();\n+ }\n+\n tracing_subscriber::fmt()\n .with_env_filter(\n tracing_subscriber::EnvFilter::try_from_default_env()\ndiff --git a/server/src/reddit.rs b/server/src/reddit.rs\nindex 90053ad03b1d7c8e94f325dd4ee64c2b4f7da900..ff0f01e57b18af878eb5be3efc47204a7673589d 100644\n--- a/server/src/reddit.rs\n+++ b/server/src/reddit.rs\n@@ -6,11 +6,15 @@ use std::time::{Duration, Instant};\n \n use reqwest::{header, Client, StatusCode};\n use serde::Deserialize;\n+use serde_json::Value;\n use tokio::sync::{mpsc, RwLock};\n \n use crate::{\n+ event_log::EventLog,\n+ events::Event,\n+ html::now_ms,\n path_types::ItemId,\n- reducer::{EntityData, GlobalTree},\n+ reducer::GlobalTree,\n };\n \n /// Bootstrap blank nodes along a URL path so breadcrumbs and voting work before fetch.\n@@ -20,6 +24,8 @@ pub fn ensure_partial_tree(tree: &mut GlobalTree, id: &ItemId) {\n \n pub struct RedditCommand {\n pub id: ItemId,\n+ /// User-initiated fetch bypasses the in-memory \"recently fetched\" cache.\n+ pub force: bool,\n }\n \n #[derive(Clone)]\n@@ -33,19 +39,31 @@ struct RedditCredentials {\n client_secret: String,\n }\n \n+#[derive(Clone)]\n+pub struct RedditApiConfig {\n+ pub api_base: String,\n+ pub oauth_base: String,\n+ pub user_agent: String,\n+ creds: Option,\n+}\n+\n struct OAuthToken {\n access_token: String,\n expires_at: Instant,\n }\n \n impl RedditBroker {\n- pub fn spawn(tree: Arc>, user_agent: &str) -> Self {\n+ pub fn spawn(\n+ tree: Arc>,\n+ event_log: Arc,\n+ config: RedditApiConfig,\n+ ) -> Self {\n let (tx, rx) = mpsc::channel(100);\n \n let mut headers = header::HeaderMap::new();\n headers.insert(\n header::USER_AGENT,\n- header::HeaderValue::from_str(user_agent).expect(\"valid user agent\"),\n+ header::HeaderValue::from_str(&config.user_agent).expect(\"valid user agent\"),\n );\n \n let client = Client::builder()\n@@ -54,22 +72,38 @@ impl RedditBroker {\n .build()\n .expect(\"reqwest client\");\n \n- let creds = RedditCredentials::from_env();\n- tokio::spawn(reddit_worker(rx, tree, client, creds));\n+ tokio::spawn(reddit_worker(rx, tree, event_log, client, config));\n \n Self { tx }\n }\n \n- /// Fire-and-forget: queue a fetch; worker updates the tree when done.\n- pub fn request_fetch(&self, id: ItemId) {\n- let _ = self.tx.try_send(RedditCommand { id });\n+ /// Queue a fetch; drops when the channel is full (backpressure).\n+ pub fn request_fetch(&self, id: ItemId, force: bool) {\n+ let _ = self.tx.try_send(RedditCommand { id, force });\n+ }\n+}\n+\n+impl RedditApiConfig {\n+ pub fn from_env() -> Self {\n+ Self {\n+ api_base: reddit_api_base(),\n+ oauth_base: reddit_oauth_base(),\n+ user_agent: default_user_agent(),\n+ creds: RedditCredentials::from_env(),\n+ }\n }\n }\n \n impl RedditCredentials {\n+ /// Reddit's OAuth docs call these \"client id\" and \"client secret\"; the app\n+ /// registration UI often labels them \"app id\" / \"app secret\" — same values.\n fn from_env() -> Option {\n- let client_id = std::env::var(\"REDDIT_CLIENT_ID\").ok()?;\n- let client_secret = std::env::var(\"REDDIT_CLIENT_SECRET\").ok()?;\n+ let client_id = std::env::var(\"REDDIT_CLIENT_ID\")\n+ .or_else(|_| std::env::var(\"REDDIT_APP_ID\"))\n+ .ok()?;\n+ let client_secret = std::env::var(\"REDDIT_CLIENT_SECRET\")\n+ .or_else(|_| std::env::var(\"REDDIT_APP_SECRET\"))\n+ .ok()?;\n if client_id.is_empty() || client_secret.is_empty() {\n return None;\n }\n@@ -80,29 +114,63 @@ impl RedditCredentials {\n }\n }\n \n+pub fn reddit_api_base() -> String {\n+ std::env::var(\"REDDIT_API_BASE\").unwrap_or_else(|_| \"https://www.reddit.com\".into())\n+}\n+\n+pub fn reddit_oauth_base() -> String {\n+ std::env::var(\"REDDIT_OAUTH_BASE\").unwrap_or_else(|_| \"https://www.reddit.com\".into())\n+}\n+\n pub fn default_user_agent() -> String {\n std::env::var(\"REDDIT_USER_AGENT\").unwrap_or_else(|_| {\n \"web:sorter2.social:v0.0.1 (by /u/sorter2)\".to_string()\n })\n }\n \n+/// True when this node can be loaded from the Reddit JSON API.\n+pub fn is_fetchable(id: &ItemId) -> bool {\n+ !map_item_to_reddit_api(id, \"https://example.com\").is_empty()\n+}\n+\n+/// Derive UI-facing fields from a stored payload (Reddit-specific when under reddit.com).\n+pub fn entity_view_from_payload(id: &ItemId, payload: &Value) -> Option {\n+ if id.as_str().starts_with(\"reddit.com\") {\n+ return parse_reddit_view(id, payload);\n+ }\n+ None\n+}\n+\n+/// Apply a full API payload to the in-memory tree (view derived for known domains).\n+pub fn apply_entity_import(tree: &mut GlobalTree, id: &ItemId, payload: Value) {\n+ let view = entity_view_from_payload(id, &payload);\n+ tree.apply_entity_raw(id, payload, view);\n+}\n+\n async fn reddit_worker(\n mut rx: mpsc::Receiver,\n tree: Arc>,\n+ event_log: Arc,\n client: Client,\n- creds: Option,\n+ config: RedditApiConfig,\n ) {\n let mut in_flight = HashSet::new();\n let mut recently_fetched: HashMap = HashMap::new();\n let mut current_delay = Duration::from_secs(1);\n let mut oauth: Option = None;\n let cache_ttl = Duration::from_secs(300);\n+ let creds = config.creds.clone();\n+ let api_base = config.api_base.clone();\n+ let oauth_base = config.oauth_base.clone();\n \n while let Some(cmd) = rx.recv().await {\n let now = Instant::now();\n recently_fetched.retain(|_, t| now.duration_since(*t) < cache_ttl);\n \n- if in_flight.contains(&cmd.id) || recently_fetched.contains_key(&cmd.id) {\n+ if in_flight.contains(&cmd.id) {\n+ continue;\n+ }\n+ if !cmd.force && recently_fetched.contains_key(&cmd.id) {\n continue;\n }\n \n@@ -112,18 +180,26 @@ async fn reddit_worker(\n tokio::time::sleep(current_delay).await;\n \n if let Some(c) = &creds {\n- oauth = ensure_oauth_token(&client, c, oauth.take()).await;\n+ oauth = ensure_oauth_token(&client, &oauth_base, c, oauth.take()).await;\n }\n \n let token = oauth.as_ref().map(|t| t.access_token.as_str());\n- let use_oauth = token.is_some();\n \n- match do_fetch(&client, &fetch_id, use_oauth, token).await {\n- Ok(FetchOutcome::Entity(data)) => {\n- let mut w = tree.write().await;\n- w.set_entity_data(&fetch_id, data);\n- recently_fetched.insert(fetch_id.clone(), Instant::now());\n- current_delay = Duration::from_millis(600);\n+ match do_fetch(&client, &api_base, &fetch_id, token).await {\n+ Ok(FetchOutcome::Payload(payload)) => {\n+ let ts = now_ms();\n+ let event = Event::EntityImported {\n+ id: fetch_id.as_str().to_string(),\n+ ts,\n+ payload: payload.clone(),\n+ };\n+ if let Err(e) = event_log.append(&event).await {\n+ tracing::warn!(\"event log append failed for {}: {}\", fetch_id, e);\n+ } else {\n+ apply_entity_import(&mut *tree.write().await, &fetch_id, payload);\n+ recently_fetched.insert(fetch_id.clone(), Instant::now());\n+ current_delay = Duration::from_millis(600);\n+ }\n }\n Ok(FetchOutcome::NotFound) => {\n recently_fetched.insert(fetch_id.clone(), Instant::now());\n@@ -149,13 +225,14 @@ async fn reddit_worker(\n }\n \n enum FetchOutcome {\n- Entity(EntityData),\n+ Payload(Value),\n NotFound,\n RateLimited { reset_secs: u64 },\n }\n \n async fn ensure_oauth_token(\n client: &Client,\n+ oauth_base: &str,\n creds: &RedditCredentials,\n existing: Option,\n ) -> Option {\n@@ -165,8 +242,13 @@ async fn ensure_oauth_token(\n }\n }\n \n+ let url = format!(\n+ \"{}/api/v1/access_token\",\n+ oauth_base.trim_end_matches('/')\n+ );\n+\n let resp = client\n- .post(\"https://www.reddit.com/api/v1/access_token\")\n+ .post(&url)\n .basic_auth(&creds.client_id, Some(&creds.client_secret))\n .form(&[(\"grant_type\", \"client_credentials\")])\n .send()\n@@ -207,11 +289,11 @@ async fn ensure_oauth_token(\n \n async fn do_fetch(\n client: &Client,\n+ api_base: &str,\n id: &ItemId,\n- use_oauth: bool,\n bearer: Option<&str>,\n ) -> Result {\n- let url = map_item_to_reddit_api(id, use_oauth);\n+ let url = map_item_to_reddit_api(id, api_base);\n if url.is_empty() {\n return Ok(FetchOutcome::NotFound);\n }\n@@ -241,10 +323,8 @@ async fn do_fetch(\n return Ok(FetchOutcome::RateLimited { reset_secs: reset });\n }\n \n- let bytes = resp.bytes().await.map_err(|e| e.to_string())?;\n- Ok(parse_reddit_json(id, &bytes)\n- .map(FetchOutcome::Entity)\n- .unwrap_or(FetchOutcome::NotFound))\n+ let payload: Value = resp.json().await.map_err(|e| e.to_string())?;\n+ Ok(FetchOutcome::Payload(payload))\n }\n \n fn rate_limit_remaining(resp: &reqwest::Response) -> Option {\n@@ -264,18 +344,14 @@ fn rate_limit_reset_secs(resp: &reqwest::Response) -> u64 {\n .unwrap_or(5)\n }\n \n-/// Map canonical item id to Reddit JSON API URL.\n-pub fn map_item_to_reddit_api(id: &ItemId, oauth: bool) -> String {\n+/// Map canonical item id to a Reddit JSON API URL under `api_base`.\n+pub fn map_item_to_reddit_api(id: &ItemId, api_base: &str) -> String {\n let path = id.as_str();\n if !path.starts_with(\"reddit.com/\") && path != \"reddit.com\" {\n return String::new();\n }\n \n- let base = if oauth {\n- \"https://oauth.reddit.com\"\n- } else {\n- \"https://www.reddit.com\"\n- };\n+ let base = api_base.trim_end_matches('/');\n \n let segments: Vec<&str> = path.split('/').collect();\n \n@@ -293,18 +369,17 @@ pub fn map_item_to_reddit_api(id: &ItemId, oauth: bool) -> String {\n String::new()\n }\n \n-fn parse_reddit_json(id: &ItemId, bytes: &[u8]) -> Option {\n- let v: serde_json::Value = serde_json::from_slice(bytes).ok()?;\n+fn parse_reddit_view(id: &ItemId, v: &Value) -> Option {\n let segments: Vec<&str> = id.as_str().split('/').collect();\n \n if segments.iter().any(|&p| p == \"comments\") {\n- parse_post_listing(&v)\n+ parse_post_listing(v)\n } else {\n- parse_subreddit_about(&v)\n+ parse_subreddit_about(v)\n }\n }\n \n-fn parse_subreddit_about(v: &serde_json::Value) -> Option {\n+fn parse_subreddit_about(v: &Value) -> Option {\n let data = v.get(\"data\")?;\n let title = data\n .get(\"title\")\n@@ -323,7 +398,7 @@ fn parse_subreddit_about(v: &serde_json::Value) -> Option {\n .filter(|s| !s.is_empty())\n .map(|s| s.to_string());\n \n- Some(EntityData {\n+ Some(crate::reducer::EntityData {\n title,\n author: None,\n body_html,\n@@ -331,10 +406,9 @@ fn parse_subreddit_about(v: &serde_json::Value) -> Option {\n })\n }\n \n-fn parse_post_listing(v: &serde_json::Value) -> Option {\n+fn parse_post_listing(v: &Value) -> Option {\n let listing = v.as_array()?.first()?;\n- let child = listing\n- .pointer(\"/data/children/0/data\")?;\n+ let child = listing.pointer(\"/data/children/0/data\")?;\n let title = child.get(\"title\")?.as_str()?.to_string();\n let author = child\n .get(\"author\")\n@@ -352,7 +426,7 @@ fn parse_post_listing(v: &serde_json::Value) -> Option {\n .filter(|s| s.starts_with(\"http\"))\n .map(|s| s.to_string());\n \n- Some(EntityData {\n+ Some(crate::reducer::EntityData {\n title,\n author,\n body_html,\n@@ -368,48 +442,51 @@ mod tests {\n fn map_subreddit_about_url() {\n let id = ItemId::parse(\"reddit.com/r/rust\").unwrap();\n assert_eq!(\n- map_item_to_reddit_api(&id, false),\n+ map_item_to_reddit_api(&id, \"https://www.reddit.com\"),\n \"https://www.reddit.com/r/rust/about.json?raw_json=1\"\n );\n assert_eq!(\n- map_item_to_reddit_api(&id, true),\n- \"https://oauth.reddit.com/r/rust/about.json?raw_json=1\"\n+ map_item_to_reddit_api(&id, \"http://127.0.0.1:9999\"),\n+ \"http://127.0.0.1:9999/r/rust/about.json?raw_json=1\"\n );\n }\n \n #[test]\n fn map_post_url() {\n- let id =\n- ItemId::parse(\"reddit.com/r/amitheasshole/comments/1trnvdl\").unwrap();\n+ let id = ItemId::parse(\"reddit.com/r/amitheasshole/comments/1trnvdl\").unwrap();\n assert_eq!(\n- map_item_to_reddit_api(&id, false),\n+ map_item_to_reddit_api(&id, \"https://www.reddit.com\"),\n \"https://www.reddit.com/r/amitheasshole/comments/1trnvdl.json?raw_json=1\"\n );\n }\n \n #[test]\n- fn map_non_reddit_empty() {\n- let id = ItemId::opaque(\"example.com/foo\");\n- assert!(map_item_to_reddit_api(&id, false).is_empty());\n+ fn is_fetchable_reddit_sub() {\n+ let id = ItemId::parse(\"reddit.com/r/rust\").unwrap();\n+ assert!(is_fetchable(&id));\n+ assert!(!is_fetchable(&ItemId::opaque(\"example.com/x\")));\n }\n \n #[test]\n fn parse_subreddit_fixture() {\n- let json = r#\"{\"kind\":\"t5\",\"data\":{\"title\":\"Rust\",\"display_name\":\"rust\",\"public_description\":\"systems\"}}\"#;\n- let entity = parse_reddit_json(\n+ let json = include_str!(\"../../test/fixtures/reddit/r_rust_about.json\");\n+ let v: Value = serde_json::from_str(json).unwrap();\n+ let entity = entity_view_from_payload(\n &ItemId::parse(\"reddit.com/r/rust\").unwrap(),\n- json.as_bytes(),\n+ &v,\n )\n .unwrap();\n- assert_eq!(entity.title, \"Rust\");\n+ assert_eq!(entity.title, \"The Rust Programming Language\");\n+ assert!(entity.body_html.as_ref().is_some_and(|b| b.contains(\"Rust\")));\n }\n \n #[test]\n fn parse_post_fixture() {\n 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\"}}]}}]\"#;\n- let entity = parse_reddit_json(\n+ let v: Value = serde_json::from_str(json).unwrap();\n+ let entity = entity_view_from_payload(\n &ItemId::parse(\"reddit.com/r/x/comments/abc\").unwrap(),\n- json.as_bytes(),\n+ &v,\n )\n .unwrap();\n assert_eq!(entity.title, \"AITA\");\ndiff --git a/server/src/reducer.rs b/server/src/reducer.rs\nindex 077f700bf00ddefe18ffd004bb5288bdc7c4adaf..60db81a562e3590775012573225a76ba82a14563 100644\n--- a/server/src/reducer.rs\n+++ b/server/src/reducer.rs\n@@ -1,6 +1,7 @@\n use std::collections::{HashMap, HashSet, VecDeque};\n \n use serde::{Deserialize, Serialize};\n+use serde_json::Value;\n \n use crate::path_types::ItemId;\n \n@@ -128,6 +129,9 @@ pub struct EntityData {\n #[derive(Debug, Clone, Default)]\n pub struct NodeState {\n pub id: ItemId,\n+ /// Full imported API JSON (persisted in the event log).\n+ pub entity_raw: Option,\n+ /// Domain-specific view derived from `entity_raw` (e.g. Reddit title/author).\n pub data: Option,\n pub children: HashSet,\n pub local_ranking: GroupState,\n@@ -197,10 +201,11 @@ impl GlobalTree {\n }\n }\n \n- pub fn set_entity_data(&mut self, id: &ItemId, data: EntityData) {\n+ pub fn apply_entity_raw(&mut self, id: &ItemId, payload: Value, view: Option) {\n self.ensure_path(id);\n if let Some(node) = self.nodes.get_mut(id) {\n- node.data = Some(data);\n+ node.entity_raw = Some(payload);\n+ node.data = view;\n }\n }\n }\ndiff --git a/server/src/state.rs b/server/src/state.rs\nindex 8c03aa60c15aee803a534439400b69935b1a3d84..d71d1079a8f486cdab15795384aef0b81b32544d 100644\n--- a/server/src/state.rs\n+++ b/server/src/state.rs\n@@ -7,7 +7,7 @@ use crate::{\n events::Event,\n journal::JournalClient,\n path_types::ItemId,\n- reddit::{default_user_agent, RedditBroker},\n+ reddit::{apply_entity_import, RedditApiConfig, RedditBroker},\n reducer::{GlobalTree, VoteData},\n views::ViewStore,\n };\n@@ -108,13 +108,18 @@ impl AppState {\n tree.ensure_path(&parsed);\n }\n }\n+ Event::EntityImported { id, payload, .. } => {\n+ if let Some(parsed) = ItemId::parse(&id).or_else(|| ItemId::from_url(&id)) {\n+ apply_entity_import(&mut tree, &parsed, payload);\n+ }\n+ }\n }\n }\n }\n \n let tree = Arc::new(RwLock::new(tree));\n let journal = JournalClient::spawn(tree.clone(), event_log.clone());\n- let reddit = RedditBroker::spawn(tree.clone(), &default_user_agent());\n+ let reddit = RedditBroker::spawn(tree.clone(), event_log.clone(), RedditApiConfig::from_env());\n \n Self {\n cfg: Arc::new(cfg),\n@@ -135,10 +140,14 @@ impl AppState {\n let mut w = self.tree.write().await;\n w.ensure_path(id);\n }\n- self.reddit.request_fetch(id.clone());\n Ok(())\n }\n \n+ /// User-initiated Reddit/API import (via \"Fetch more\" — never on paste or navigate).\n+ pub fn queue_entity_fetch(&self, id: ItemId) {\n+ self.reddit.request_fetch(id, true);\n+ }\n+\n pub async fn record_vote(\n &self,\n parent: &ItemId,\n@@ -169,6 +178,44 @@ impl AppState {\n #[cfg(test)]\n mod tests {\n use super::{normalize_scope, parse_item_param};\n+ use crate::{\n+ event_log::EventLog,\n+ events::Event,\n+ path_types::ItemId,\n+ reddit::apply_entity_import,\n+ reducer::GlobalTree,\n+ };\n+ use serde_json::json;\n+\n+ #[tokio::test]\n+ async fn replay_entity_imported_restores_view() {\n+ let tmp = tempfile::tempdir().unwrap();\n+ let log_path = tmp.path().join(\"events.jsonl\");\n+ let log = EventLog::new(log_path.to_string_lossy().into_owned());\n+ let payload = json!({\"kind\":\"t5\",\"data\":{\"title\":\"Rust\",\"display_name\":\"rust\"}});\n+ log.append(&Event::EntityImported {\n+ id: \"reddit.com/r/rust\".into(),\n+ ts: 1,\n+ payload: payload.clone(),\n+ })\n+ .await\n+ .unwrap();\n+\n+ let mut tree = GlobalTree::new();\n+ let (events, _) = log.load_all().await.unwrap();\n+ for ev in events {\n+ if let Event::EntityImported { id, payload, .. } = ev {\n+ let parsed = ItemId::parse(&id).unwrap();\n+ apply_entity_import(&mut tree, &parsed, payload);\n+ }\n+ }\n+ let node = tree.get(&ItemId::parse(\"reddit.com/r/rust\").unwrap()).unwrap();\n+ assert_eq!(node.data.as_ref().unwrap().title, \"Rust\");\n+ assert_eq!(\n+ node.entity_raw.as_ref().unwrap()[\"data\"][\"display_name\"],\n+ \"rust\"\n+ );\n+ }\n \n #[test]\n fn normalize_scope_strips_prefix_and_lowercases() {\ndiff --git a/server/src/ui_action.rs b/server/src/ui_action.rs\nindex d798874d1d94c0dfee59ec1ff703f9ee6ef432c0..3d6a49a2a3fb950752827efe1f5a049308f09baf 100644\n--- a/server/src/ui_action.rs\n+++ b/server/src/ui_action.rs\n@@ -26,6 +26,10 @@ pub enum HtmlUiAction {\n ParseQuery {\n query: String,\n },\n+ /// Fetch upstream entity data for the current page (explicit user action only).\n+ FetchEntity {\n+ item: String,\n+ },\n }\n \n #[derive(Debug, Error)]\ndiff --git a/server/static/sorter.css b/server/static/sorter.css\nindex 36f68d415cea33c8562d5d02c0c9d4c5b225a782..1f9f5158414902d43a380dd899232c70d8cf5d63 100644\n--- a/server/static/sorter.css\n+++ b/server/static/sorter.css\n@@ -52,6 +52,25 @@ body {\n filter: brightness(1.1);\n }\n \n+.btn-secondary {\n+ background: transparent;\n+ color: var(--accent);\n+ border: 1px solid var(--border);\n+ padding: 0.4rem 0.85rem;\n+ font-size: 0.9rem;\n+ cursor: pointer;\n+ margin-top: 0.75rem;\n+}\n+\n+.btn-secondary:disabled {\n+ opacity: 0.6;\n+ cursor: wait;\n+}\n+\n+.fetch-entity-form {\n+ margin-top: 0.5rem;\n+}\n+\n code {\n font-size: 0.85em;\n background: var(--bg);\ndiff --git a/test/fixtures/reddit/r_rust_about.json b/test/fixtures/reddit/r_rust_about.json\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..bb7fd74ab00d63356d9b672a775996ed155b984d\n--- /dev/null\n+++ b/test/fixtures/reddit/r_rust_about.json\n@@ -0,0 +1,21 @@\n+{\n+ \"kind\": \"t5\",\n+ \"data\": {\n+ \"user_flair_background_color\": null,\n+ \"submit_text_html\": null,\n+ \"banner_img\": \"\",\n+ \"user_is_banned\": null,\n+ \"wiki_enabled\": true,\n+ \"show_media\": true,\n+ \"id\": \"2qh0i\",\n+ \"display_name\": \"rust\",\n+ \"title\": \"The Rust Programming Language\",\n+ \"public_description\": \"A place for all things related to the Rust programming language.\",\n+ \"public_description_html\": \"<!-- SC_OFF --><div class=\\\"md\\\"><p>A place for all things related to the Rust programming language.</p>\\n</div><!-- SC_ON -->\",\n+ \"icon_img\": \"https://styles.redditmedia.com/t5_2qh0i/styles/communityIcon_kfqpknb2j6j51.png\",\n+ \"community_icon\": \"https://styles.redditmedia.com/t5_2qh0i/styles/communityIcon_kfqpknb2j6j51.png\",\n+ \"subscribers\": 350000,\n+ \"active_user_count\": 1200,\n+ \"over18\": false\n+ }\n+}\ndiff --git a/test/reddit_import.clj b/test/reddit_import.clj\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..54edaedeef08718c8f184d6aa468d8e6415068c7\n--- /dev/null\n+++ b/test/reddit_import.clj\n@@ -0,0 +1,116 @@\n+(ns test.reddit-import\n+ (:require [babashka.process :as process]\n+ [clojure.java.io :as io]\n+ [clojure.string :as str]\n+ [clojure.test :refer [deftest is testing]])\n+ (:import [com.sun.net.httpserver HttpServer HttpHandler HttpExchange]\n+ [java.net InetSocketAddress]))\n+\n+(defn- repo-root []\n+ (.getCanonicalPath (io/file (System/getProperty \"user.dir\"))))\n+\n+(defn- pick-port []\n+ (with-open [s (java.net.ServerSocket. 0)]\n+ (.getLocalPort s)))\n+\n+(defn- start-mock-reddit [port fixtures-dir]\n+ (let [fixture (io/file fixtures-dir \"r_rust_about.json\")\n+ body (.getBytes (slurp fixture) \"UTF-8\")\n+ server (HttpServer/create (InetSocketAddress. \"127.0.0.1\" port) 0)\n+ handler\n+ (proxy [HttpHandler] []\n+ (handle [^HttpExchange exchange]\n+ (.sendResponseHeaders exchange 200 (alength body))\n+ (let [out (.getResponseBody exchange)]\n+ (.write out body)\n+ (.close out))))]\n+ (.createContext server \"/\" handler)\n+ (.setExecutor server nil)\n+ (.start server)\n+ (fn stop []\n+ (.stop server 0))))\n+\n+(defn- wait-health [base-url ms]\n+ (let [deadline (+ (System/currentTimeMillis) ms)\n+ url (str base-url \"/healthz\")]\n+ (loop []\n+ (let [resp (try\n+ (process/shell {:out :string :err :string}\n+ \"curl\" \"-sf\" url)\n+ (catch Exception _ nil))]\n+ (if (and resp (zero? (:exit resp)) (= \"ok\" (str/trim (:out resp \"\"))))\n+ true\n+ (if (< (System/currentTimeMillis) deadline)\n+ (do (Thread/sleep 200) (recur))\n+ false))))))\n+\n+(defn- curl-post-ui [base rpc-json]\n+ (process/shell {:out :string :err :string}\n+ \"curl\" \"-sf\" \"-X\" \"POST\" (str base \"/ui\")\n+ \"--data-urlencode\" (str \"__rpc__=\" rpc-json)))\n+\n+(defn- wait-event-log [path ms]\n+ (let [deadline (+ (System/currentTimeMillis) ms)]\n+ (loop []\n+ (if (.exists (io/file path))\n+ true\n+ (if (< (System/currentTimeMillis) deadline)\n+ (do (Thread/sleep 200) (recur))\n+ false)))))\n+\n+(deftest reddit-fetch-via-mock-api\n+ (testing \"Fetch more queues import; event log stores full payload; page shows title\"\n+ (let [root (repo-root)\n+ fixtures (str root \"/test/fixtures/reddit\")\n+ data-dir (.getAbsolutePath\n+ (doto (io/file (System/getProperty \"java.io.tmpdir\")\n+ (str \"sorter2-reddit-\" (System/currentTimeMillis)))\n+ (.mkdirs)))\n+ reddit-port (pick-port)\n+ app-port (pick-port)\n+ reddit-base (str \"http://127.0.0.1:\" reddit-port)\n+ app-base (str \"http://127.0.0.1:\" app-port)\n+ bin (str root \"/target/release/sorter2-server\")\n+ stop-mock (start-mock-reddit reddit-port fixtures)]\n+ (try\n+ (is (zero? (:exit (process/shell {:dir root}\n+ \"cargo\" \"build\" \"--release\" \"--package\" \"sorter2-server\")))\n+ \"release build succeeds\")\n+ (let [proc (process/process {:dir root\n+ :env (into (into {} (System/getenv))\n+ {\"SORTER2_SKIP_DOTENV\" \"1\"\n+ \"SORTER2_DATA_DIR\" data-dir\n+ \"SORTER2_EVENT_LOG\" (str data-dir \"/events.jsonl\")\n+ \"PORT\" (str app-port)\n+ \"REDDIT_API_BASE\" reddit-base\n+ \"REDDIT_OAUTH_BASE\" reddit-base\n+ \"REDDIT_CLIENT_ID\" \"\"\n+ \"REDDIT_CLIENT_SECRET\" \"\"\n+ \"REDDIT_APP_ID\" \"\"\n+ \"REDDIT_APP_SECRET\" \"\"})\n+ :out :string\n+ :err :string}\n+ bin)]\n+ (try\n+ (is (wait-health app-base 20000) \"app healthz\")\n+ (let [browse-url (str app-base \"/~/https://reddit.com/r/rust\")\n+ before (:out (process/shell {:out :string :err :string}\n+ \"curl\" \"-sf\" browse-url))]\n+ (is (str/includes? before \"Fetch from Reddit\"))\n+ (is (not (str/includes? before \"The Rust Programming Language\")))\n+ (let [rpc \"{\\\"action\\\":\\\"fetch_entity\\\",\\\"item\\\":\\\"reddit.com/r/rust\\\"}\"\n+ post (curl-post-ui app-base rpc)\n+ log-path (str data-dir \"/events.jsonl\")]\n+ (is (zero? (:exit post)) \"fetch_entity POST succeeds\")\n+ (is (wait-event-log log-path 10000) \"event log written\")\n+ (let [after (:out (process/shell {:out :string :err :string}\n+ \"curl\" \"-sf\" browse-url))\n+ log (slurp (io/file log-path))]\n+ (is (str/includes? after \"The Rust Programming Language\"))\n+ (is (str/includes? log \"\\\"type\\\":\\\"entity_imported\\\"\"))\n+ (is (str/includes? log \"\\\"subscribers\\\":350000\"))\n+ (is (str/includes? log \"\\\"display_name\\\":\\\"rust\\\"\")))))\n+ (finally\n+ (process/destroy proc))))\n+ (finally\n+ (stop-mock))))))\n\n\nSide B — contributor: tommy-mor\nSide B — commit message:\n[96b6da05] rpc + reducer changes first pass\n\nSide B — unified diff (full patch):\ndiff --git a/cli/src/main.rs b/cli/src/main.rs\nindex 630c5dea1f78c0ec9bc53e6b96234a0dc75bb705..8d0442959f4332bafe499a2a8cdf364731a06871 100644\n--- a/cli/src/main.rs\n+++ b/cli/src/main.rs\n@@ -21,8 +21,9 @@ struct Cli {\n cmd: Option,\n }\n \n+/// Commands scoped to a room (`public` or `shortid/slug`).\n #[derive(Subcommand, Debug)]\n-enum Command {\n+enum ScopedCmd {\n /// Browse the garden (ontology) — light mode, ranked by votes\n Garden {\n #[command(subcommand)]\n@@ -174,6 +175,23 @@ enum Command {\n #[arg(long)]\n json: bool,\n },\n+}\n+\n+#[derive(Subcommand, Debug)]\n+enum Command {\n+ /// Public site (same as room `public`)\n+ Public {\n+ #[command(subcommand)]\n+ sub: ScopedCmd,\n+ },\n+ /// Private room id (`shortid/slug` from `room create`)\n+ Private {\n+ /// Room id, e.g. `a1b2c3d/my-project`\n+ #[arg(value_name = \"ROOM_ID\")]\n+ room: String,\n+ #[command(subcommand)]\n+ sub: ScopedCmd,\n+ },\n \n /// Show all activity since you last posted (global feed)\n ///\n@@ -628,6 +646,37 @@ fn http_client() -> Result {\n .build()?)\n }\n \n+async fn send_rpc(\n+ client: &reqwest::Client,\n+ base: &str,\n+ bearer: Option<&str>,\n+ commands: Vec,\n+) -> Result {\n+ let url = format!(\"{}/api/v0/rpc\", base.trim_end_matches('/'));\n+ let mut req = client.post(url).json(&RpcBatch(commands));\n+ if let Some(b) = bearer {\n+ req = req.header(\"Authorization\", format!(\"Bearer {}\", b));\n+ }\n+ let resp = req.send().await?;\n+ let status = resp.status();\n+ let text = resp.text().await.unwrap_or_default();\n+ if !status.is_success() {\n+ return Err(anyhow!(\"rpc HTTP {}: {}\", status, text.trim()));\n+ }\n+ serde_json::from_str(&text).map_err(|e| anyhow!(\"rpc response: {e}\"))\n+}\n+\n+fn rpc_line_ok(line: &RpcLine) -> Result<&RpcResult> {\n+ if !line.ok {\n+ let mut m = line.error.clone().unwrap_or_else(|| \"rpc error\".into());\n+ if let Some(h) = &line.hint {\n+ m.push_str(&format!(\"\\nhint: {h}\"));\n+ }\n+ return Err(anyhow!(m));\n+ }\n+ line.result.as_ref().ok_or_else(|| anyhow!(\"rpc missing result\"))\n+}\n+\n /// Normalize ontology path for API. Accepts path with or without ~/ (shell expands ~ to $HOME).\n /// Returns a bare slug path (e.g. `languages/python`) with no leading `/` or `~/`.\n /// Call `ontology_path_for_api_query` before sending `item=` / `parent=` params so the server\n@@ -729,231 +778,262 @@ fn write_secret_file(name: &str, contents: &str) -> Result<()> {\n Ok(())\n }\n \n-#[tokio::main]\n-async fn main() -> Result<()> {\n- let Cli { cmd, server } = Cli::parse();\n-\n- // If no command provided, print the guide\n- let Some(cmd) = cmd else {\n- print!(\"{}\", include_str!(\"../GUIDE.sorter\"));\n- return Ok(());\n- };\n-\n- let base = server.trim_end_matches('/');\n-\n- match cmd {\n- Command::Healthz { json } => {\n- let client = http_client()?;\n- let url = format!(\"{base}/healthz\");\n- let body = client.get(url).send().await?.text().await?;\n- if json {\n- // Wrap plain text response in a JSON object\n- println!(\"{}\", serde_json::json!({ \"ok\": true, \"body\": body.trim() }));\n- } else {\n- println!(\"{body}\");\n- }\n- }\n-\n- Command::Search { query, json } => {\n- let client = http_client()?;\n- let url = format!(\"{base}/api/v0/search?q={}\", urlencoding::encode(&query));\n- let resp: slug_types::SearchResponse = expect_json(client.get(url).send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- if !resp.items.is_empty() {\n- println!(\"items ({})\", resp.items.len());\n- for item in &resp.items {\n- print!(\" {}\", item.path);\n- if let Some(body) = &item.body {\n- let first_line = body.lines().next().unwrap_or(\"\").trim();\n- if !first_line.is_empty() {\n- print!(\" {}\", first_line);\n+async fn run_scoped(base: &str, room: &str, sub: ScopedCmd) -> Result<()> {\n+ let room = room.trim();\n+ let client = http_client()?;\n+ match sub {\n+ ScopedCmd::Garden { sub } => match sub {\n+ GardenCmd::Tree { json } => {\n+ let batch = send_rpc(&client, base, None, vec![RpcCommand::GetLeaves { room: room.to_string() }]).await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::Leaves(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ for p in &resp.paths {\n+ println!(\"~/{}\", p);\n }\n }\n- println!();\n- }\n- }\n- if !resp.threads.is_empty() {\n- if !resp.items.is_empty() { println!(); }\n- println!(\"threads ({})\", resp.threads.len());\n- let now_ms = std::time::SystemTime::now()\n- .duration_since(std::time::UNIX_EPOCH)\n- .unwrap_or_default()\n- .as_millis() as i64;\n- for t in &resp.threads {\n- println!(\" {} {}n {}\", t.tag, t.post_count, slug_types::timeago::timeago(now_ms, t.last_activity));\n- }\n- }\n- if !resp.posts.is_empty() {\n- if !resp.items.is_empty() || !resp.threads.is_empty() { println!(); }\n- println!(\"posts ({})\", resp.posts.len());\n- let now_ms = std::time::SystemTime::now()\n- .duration_since(std::time::UNIX_EPOCH)\n- .unwrap_or_default()\n- .as_millis() as i64;\n- for p in &resp.posts {\n- let first_line = p.snippet.lines().next().unwrap_or(\"\").trim();\n- println!(\" {} · {} {}\", p.thread, slug_types::timeago::timeago(now_ms, p.ts), first_line);\n- }\n- }\n- if resp.items.is_empty() && resp.threads.is_empty() && resp.posts.is_empty() {\n- println!(\"no results\");\n- }\n- }\n- }\n-\n- Command::Garden { sub } => match sub {\n- GardenCmd::Tree { json } => {\n- let client = http_client()?;\n- let url = format!(\"{base}/api/v0/leaves\");\n- let builder = client.get(url);\n- let resp: LeavesResponse = expect_json(builder.send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- for p in &resp.paths {\n- println!(\"~/{}\", p);\n }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n GardenCmd::Body { path, json, full } => {\n let path = normalize_ontology_path_input(&path).map_err(anyhow::Error::msg)?;\n let item_q = ontology_path_for_api_query(&path);\n- let client = http_client()?;\n- let mut url = format!(\"{base}/api/v0/item?item={}\", urlencoding::encode(&item_q));\n- if full {\n- url.push_str(\"&full=true\");\n- }\n- let builder = client.get(url);\n- let resp: ItemResponse = expect_json(builder.send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_item_response(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetGardenItem {\n+ room: room.to_string(),\n+ item_path: item_q,\n+ full: Some(full),\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::GardenItem(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_item_response(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n GardenCmd::Children { paths, depth, json } => {\n let paths: Vec = paths\n .iter()\n .map(|p| normalize_ontology_path_input(p).map_err(anyhow::Error::msg))\n .collect::>>()?;\n- let client = http_client()?;\n let parent_param = paths\n .iter()\n .map(|p| ontology_path_for_api_query(p))\n .collect::>()\n .join(\",\");\n- let mut url = format!(\"{base}/api/v0/rank?parent={}\", urlencoding::encode(&parent_param));\n- if let Some(d) = depth {\n- url.push_str(&format!(\"&depth={d}\"));\n- }\n- let builder = client.get(url);\n- let resp: RankResponse = expect_json(builder.send().await?).await?;\n-\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_rank_response(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetGardenRank {\n+ room: room.to_string(),\n+ parent_path: parent_param,\n+ depth,\n+ offset: None,\n+ limit: None,\n+ percent: None,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::GardenRank(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_rank_response(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n GardenCmd::Pair { path, json } => {\n let path = normalize_ontology_path_input(&path).map_err(anyhow::Error::msg)?;\n let parent_q = ontology_path_for_api_query(&path);\n- let client = http_client()?;\n- let url = format!(\"{base}/api/v0/pair?parent={}\", urlencoding::encode(&parent_q));\n- let builder = client.get(url);\n- let resp: PairResponse = expect_json(builder.send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_pair_response(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetPair {\n+ room: room.to_string(),\n+ parent_path: parent_q,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::Pair(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_pair_response(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n GardenCmd::Matchup { path, json } => {\n let path = normalize_ontology_path_input(&path).map_err(anyhow::Error::msg)?;\n let item_q = ontology_path_for_api_query(&path);\n- let client = http_client()?;\n- let url = format!(\"{base}/api/v0/matchup?item={}\", urlencoding::encode(&item_q));\n- let builder = client.get(url);\n- let resp: MatchupResponse = expect_json(builder.send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_matchup_response(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetMatchup {\n+ room: room.to_string(),\n+ item_path: item_q,\n+ limit: None,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::Matchup(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_matchup_response(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n GardenCmd::History { path, json } => {\n let path = normalize_ontology_path_input(&path).map_err(anyhow::Error::msg)?;\n let item_q = ontology_path_for_api_query(&path);\n- let client = http_client()?;\n- let url = format!(\"{base}/api/v0/rank-history?item={}\", urlencoding::encode(&item_q));\n- let builder = client.get(url);\n- let resp: slug_types::RankHistoryResponse = expect_json(builder.send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_rank_history_response(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetRankHistory {\n+ room: room.to_string(),\n+ item_path: item_q,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::RankHistory(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_rank_history_response(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n GardenCmd::Rank { limit, offset, percent, json } => {\n- let client = http_client()?;\n- let url = format!(\n- \"{base}/api/v0/global-rank?limit={limit}&offset={offset}&percent={percent}\"\n- );\n- let builder = client.get(url);\n- let resp: GlobalRankResponse = expect_json(builder.send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_global_rank_response(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetGlobalRank {\n+ room: room.to_string(),\n+ limit: Some(limit),\n+ offset: Some(offset),\n+ percent: Some(percent),\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::GlobalRank(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_global_rank_response(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n },\n-\n- Command::Forum { title, json, offset, limit, since, before, actor, post } => {\n- let client = http_client()?;\n+ ScopedCmd::Forum {\n+ title,\n+ json,\n+ offset,\n+ limit,\n+ since,\n+ before,\n+ actor,\n+ post,\n+ } => {\n match title {\n None => {\n- let url = format!(\"{base}/api/v0/threads\");\n- let mut resp: ThreadsResponse =\n- expect_json(client.get(url).send().await?).await?;\n- resp.threads.truncate(10);\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_threads(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::ListForumThreads {\n+ room: room.to_string(),\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::ForumThreads(resp) => {\n+ let limited = ThreadsResponse {\n+ threads: resp.threads.iter().take(10).cloned().collect(),\n+ };\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&limited)?);\n+ } else {\n+ print_threads(&limited);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n Some(name) => {\n let tag = normalize_thread_input(&name);\n- let mut url = format!(\"{base}/api/v0/thread?tag={}\", urlencoding::encode(&tag));\n- if let Some(o) = offset { url.push_str(&format!(\"&offset={o}\")); }\n- if let Some(l) = limit { url.push_str(&format!(\"&limit={l}\")); }\n- if let Some(s) = since { url.push_str(&format!(\"&since={}\", parse_ts(&s)?)); }\n- if let Some(b) = before { url.push_str(&format!(\"&before={}\", parse_ts(&b)?)); }\n- if let Some(a) = actor { url.push_str(&format!(\"&actor={}\", urlencoding::encode(&a))); }\n- if let Some(p) = post { url.push_str(&format!(\"&post_id={}\", urlencoding::encode(&p))); }\n- let resp: ThreadDetailResponse =\n- expect_json(client.get(url).send().await?).await?;\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- print_thread(&resp);\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetForumThread {\n+ room: room.to_string(),\n+ thread_tag: tag,\n+ offset,\n+ limit,\n+ since: match &since {\n+ Some(s) => Some(parse_ts(s)?),\n+ None => None,\n+ },\n+ before: match &before {\n+ Some(s) => Some(parse_ts(s)?),\n+ None => None,\n+ },\n+ actor,\n+ post_id: post,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::ForumThread(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ print_thread(&resp);\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n }\n }\n-\n- Command::Ingest { file, thread, delegate, json } => {\n- let client = http_client()?;\n-\n+ ScopedCmd::Ingest {\n+ file,\n+ thread,\n+ delegate,\n+ json,\n+ } => {\n let mut text = String::new();\n match file {\n Some(path) => {\n@@ -966,51 +1046,74 @@ async fn main() -> Result<()> {\n .context(\"failed to read stdin\")?;\n }\n }\n-\n if text.trim().is_empty() {\n return Err(anyhow!(\"no input provided (empty)\"));\n }\n-\n- let req = IngestRequest { thread, delegate, text };\n- let url = format!(\"{base}/api/v0/ingest\");\n let bearer = effective_bearer().ok_or_else(|| {\n anyhow!(\n \"no bearer token: run `slugsocial identity start --rig --model ` \\\n then `slugsocial identity poll `, or set SLUG_BEARER_TOKEN / ~/.config/slugsocial/token\"\n )\n })?;\n- let builder = client\n- .post(url)\n- .json(&req)\n- .header(\"Authorization\", format!(\"Bearer {bearer}\"));\n- let resp: IngestResponse = expect_json(builder.send().await?).await?;\n- if resp.ok {\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- println!(\"✓ ingested\");\n- println!(\"events: {}\", resp.events_appended);\n- if !resp.threads.is_empty() {\n- println!(\"threads:\");\n- for t in &resp.threads {\n- println!(\" {t}\");\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ Some(&bearer),\n+ vec![RpcCommand::Post {\n+ room: room.to_string(),\n+ thread_tag: thread,\n+ delegate,\n+ text,\n+ return_rank_diff: true,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::PostOk {\n+ events_appended,\n+ ranking_changes,\n+ threads,\n+ next,\n+ } => {\n+ if json {\n+ println!(\n+ \"{}\",\n+ serde_json::to_string_pretty(&serde_json::json!({\n+ \"ok\": true,\n+ \"events_appended\": events_appended,\n+ \"ranking_changes\": ranking_changes,\n+ \"threads\": threads,\n+ \"next\": next,\n+ }))?\n+ );\n+ } else {\n+ println!(\"✓ ingested\");\n+ println!(\"events: {}\", events_appended);\n+ if !threads.is_empty() {\n+ println!(\"threads:\");\n+ for t in threads {\n+ println!(\" {t}\");\n+ }\n+ }\n+ if let Some(ref rc) = ranking_changes {\n+ print_ranking_changes(rc);\n }\n+ print_next(&next);\n+ println!();\n+ println!(\"---\");\n+ println!(\"For your next comparison: remember to ask your human first. Their perspective is what makes your submission more than another model's take.\");\n+ println!(\"---\");\n }\n- print_ranking_changes(&resp.ranking_changes);\n- print_next(&resp.next);\n- println!();\n- println!(\"---\");\n- println!(\"For your next comparison: remember to ask your human first. Their perspective is what makes your submission more than another model's take.\");\n- println!(\"---\");\n }\n- } else {\n- return Err(anyhow!(\"ingest failed\"));\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n }\n }\n-\n- Command::Check { file, thread, delegate, json } => {\n- let client = http_client()?;\n-\n+ ScopedCmd::Check {\n+ file,\n+ thread: _thread,\n+ delegate: _delegate,\n+ json,\n+ } => {\n let mut text = String::new();\n match file {\n Some(path) => {\n@@ -1023,57 +1126,181 @@ async fn main() -> Result<()> {\n .context(\"failed to read stdin\")?;\n }\n }\n-\n if text.trim().is_empty() {\n return Err(anyhow!(\"no input provided (empty)\"));\n }\n-\n- let req = IngestRequest { thread, delegate, text };\n- let url = format!(\"{base}/api/v0/check\");\n- let builder = client.post(url).json(&req);\n- let resp: CheckResponse = expect_json(builder.send().await?).await?;\n- if resp.ok {\n- if json {\n- println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n- } else {\n- println!(\"✓ check ok (dry-run)\");\n- if !resp.threads.is_empty() {\n- println!(\"threads:\");\n- for t in &resp.threads {\n- println!(\" {t}\");\n- }\n- }\n- if resp.rankings.is_empty() {\n- println!();\n- println!(\"(no ranking touched by this doc yet)\");\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::Check {\n+ room: room.to_string(),\n+ text,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::CheckOk {\n+ rankings,\n+ threads,\n+ next,\n+ } => {\n+ if json {\n+ println!(\n+ \"{}\",\n+ serde_json::to_string_pretty(&serde_json::json!({\n+ \"ok\": true,\n+ \"rankings\": rankings,\n+ \"threads\": threads,\n+ \"next\": next,\n+ }))?\n+ );\n } else {\n+ println!(\"✓ check ok (dry-run)\");\n+ if !threads.is_empty() {\n+ println!(\"threads:\");\n+ for t in threads {\n+ println!(\" {t}\");\n+ }\n+ }\n+ if rankings.is_empty() {\n+ println!();\n+ println!(\"(no ranking touched by this doc yet)\");\n+ } else {\n+ println!();\n+ print_check_rankings(&rankings);\n+ }\n println!();\n- print_check_rankings(&resp.rankings);\n- }\n- println!();\n- println!(\"---\");\n- println!(\"Does this submission reflect input from your human? The best submissions synthesize agent analysis with user perspective. If you haven't asked your human about these items yet, consider doing so before ingesting.\");\n- println!(\"---\");\n- if !resp.next.is_empty() {\n- println!();\n- println!(\"next:\");\n- for n in &resp.next {\n- println!(\" {n}\");\n+ println!(\"---\");\n+ println!(\"Does this submission reflect input from your human? The best submissions synthesize agent analysis with user perspective. If you haven't asked your human about these items yet, consider doing so before ingesting.\");\n+ println!(\"---\");\n+ if !next.is_empty() {\n+ println!();\n+ println!(\"next:\");\n+ for n in next {\n+ println!(\" {n}\");\n+ }\n }\n }\n }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n+ }\n+ }\n+ }\n+ Ok(())\n+}\n+\n+#[tokio::main]\n+async fn main() -> Result<()> {\n+ let Cli { cmd, server } = Cli::parse();\n+\n+ // If no command provided, print the guide\n+ let Some(cmd) = cmd else {\n+ print!(\"{}\", include_str!(\"../GUIDE.sorter\"));\n+ return Ok(());\n+ };\n+\n+ let base = server.trim_end_matches('/');\n+\n+ match cmd {\n+ Command::Public { sub } => run_scoped(base, \"public\", sub).await?,\n+ Command::Private { room, sub } => run_scoped(base, &room, sub).await?,\n+\n+ Command::Healthz { json } => {\n+ let client = http_client()?;\n+ let url = format!(\"{base}/healthz\");\n+ let body = client.get(url).send().await?.text().await?;\n+ if json {\n+ // Wrap plain text response in a JSON object\n+ println!(\"{}\", serde_json::json!({ \"ok\": true, \"body\": body.trim() }));\n } else {\n- return Err(anyhow!(\"check failed\"));\n+ println!(\"{body}\");\n }\n }\n \n- Command::Feed { actor, since, limit, json } => {\n+ Command::Search { query, json } => {\n let client = http_client()?;\n- let mut url = format!(\"{base}/api/v0/feed?actor={}&limit={}\", urlencoding::encode(&actor), limit);\n- if let Some(s) = since {\n- url.push_str(&format!(\"&since={}\", parse_ts(&s)?));\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::Search { query }],\n+ )\n+ .await?;\n+ let resp = match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::Search(s) => s,\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n+ };\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ if !resp.items.is_empty() {\n+ println!(\"items ({})\", resp.items.len());\n+ for item in &resp.items {\n+ print!(\" {}\", item.path);\n+ if let Some(body) = &item.body {\n+ let first_line = body.lines().next().unwrap_or(\"\").trim();\n+ if !first_line.is_empty() {\n+ print!(\" {}\", first_line);\n+ }\n+ }\n+ println!();\n+ }\n+ }\n+ if !resp.threads.is_empty() {\n+ if !resp.items.is_empty() { println!(); }\n+ println!(\"threads ({})\", resp.threads.len());\n+ let now_ms = std::time::SystemTime::now()\n+ .duration_since(std::time::UNIX_EPOCH)\n+ .unwrap_or_default()\n+ .as_millis() as i64;\n+ for t in &resp.threads {\n+ println!(\n+ \" {} · {} posts · {}\",\n+ t.tag,\n+ t.post_count,\n+ slug_types::timeago::timeago(now_ms, t.last_activity)\n+ );\n+ }\n+ }\n+ if !resp.posts.is_empty() {\n+ if !resp.items.is_empty() || !resp.threads.is_empty() { println!(); }\n+ println!(\"posts ({})\", resp.posts.len());\n+ let now_ms = std::time::SystemTime::now()\n+ .duration_since(std::time::UNIX_EPOCH)\n+ .unwrap_or_default()\n+ .as_millis() as i64;\n+ for p in &resp.posts {\n+ let first_line = p.snippet.lines().next().unwrap_or(\"\").trim();\n+ println!(\" {} · {} {}\", p.thread, slug_types::timeago::timeago(now_ms, p.ts), first_line);\n+ }\n+ }\n+ if resp.items.is_empty() && resp.threads.is_empty() && resp.posts.is_empty() {\n+ println!(\"no results\");\n+ }\n }\n- let resp: slug_types::FeedResponse = expect_json(client.get(url).send().await?).await?;\n+ }\n+\n+ Command::Feed { actor, since, limit, json } => {\n+ let client = http_client()?;\n+ let batch = send_rpc(\n+ &client,\n+ base,\n+ None,\n+ vec![RpcCommand::GetFeed {\n+ actor,\n+ since: match since {\n+ Some(s) => Some(parse_ts(&s)?),\n+ None => None,\n+ },\n+ limit: Some(limit),\n+ }],\n+ )\n+ .await?;\n+ let resp = match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::Feed(f) => f,\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n+ };\n if json {\n println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n } else {\ndiff --git a/server/src/api/debug.rs b/server/src/api/debug.rs\ndeleted file mode 100644\nindex ccb1730947b5ff40d2f77a5abfa05c2cde6056d9..0000000000000000000000000000000000000000\n--- a/server/src/api/debug.rs\n+++ /dev/null\n@@ -1,25 +0,0 @@\n-use axum::{response::IntoResponse, Json};\n-use axum_extra::extract::Query;\n-use serde::{Deserialize, Serialize};\n-\n-/// Debug endpoint to prove querystring decoding + repeated params behavior.\n-///\n-/// Example:\n-/// `/debug/query?open=~%2Fmodels%2Fllms&open=~%2Fmusic%2Fjazz&q=hello+world`\n-#[derive(Debug, Clone, Default, Deserialize, Serialize)]\n-pub struct DebugQuery {\n- /// Supports repeated `open=` params: `?open=a&open=b`.\n- #[serde(default)]\n- pub open: Vec,\n- /// Optional single selection.\n- #[serde(default)]\n- pub selected: Option,\n- /// A freeform string to show percent-decoding / + decoding.\n- #[serde(default)]\n- pub q: Option,\n-}\n-\n-pub async fn get_debug_query(Query(q): Query) -> impl IntoResponse {\n- Json(q)\n-}\n-\ndiff --git a/server/src/api/feed.rs b/server/src/api/feed.rs\ndeleted file mode 100644\nindex 983b7397a846bd540b6baf3773dd54c55ade4c45..0000000000000000000000000000000000000000\n--- a/server/src/api/feed.rs\n+++ /dev/null\n@@ -1,79 +0,0 @@\n-use axum::{\n- extract::{Query, State},\n- http::StatusCode,\n- response::IntoResponse,\n- Json,\n-};\n-use serde::Deserialize;\n-\n-use crate::{api::helpers::api_error, identity::parse_username, state::AppState};\n-\n-// ============================================================================\n-// Feed -- global reverse-chronological ingest stream since a cutoff\n-// ============================================================================\n-\n-#[derive(Debug, Deserialize)]\n-pub struct FeedQuery {\n- pub actor: String,\n- /// Override the lower bound (ms). Defaults to actor's last ingest timestamp.\n- #[serde(default)]\n- pub since: Option,\n- #[serde(default)]\n- pub limit: Option,\n- #[serde(default)]\n- pub offset: Option,\n-}\n-\n-pub async fn get_feed(\n- State(state): State,\n- Query(q): Query,\n-) -> impl IntoResponse {\n- const DEFAULT_LIMIT: usize = 50;\n- const MAX_LIMIT: usize = 200;\n-\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n- let actor = match parse_username(&q.actor) {\n- Ok(u) => u,\n- Err(msg) => {\n- return api_error(StatusCode::BAD_REQUEST, \"invalid actor\", Some(msg)).into_response();\n- }\n- };\n- let since = q.since.or_else(|| reduced.actor_last_post_ts.get(&actor).copied());\n- let cutoff = since.unwrap_or(0);\n- let limit = q.limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);\n- let offset = q.offset.unwrap_or(0);\n-\n- // Scan ingests_ordered backwards (newest first), stop once ts <= cutoff.\n- let matching: Vec<&str> = reduced.ingests_ordered.iter().rev()\n- .map(|id| id.as_str())\n- .take_while(|id| reduced.ingests_by_id.get(*id).map_or(false, |ing| ing.ts > cutoff))\n- .collect();\n-\n- let total = matching.len();\n- let posts: Vec = matching.into_iter()\n- .skip(offset)\n- .take(limit)\n- .filter_map(|id| reduced.ingests_by_id.get(id))\n- .map(|ing| {\n- let thread = Some(ing.thread_id.clone());\n- let thread_post_index = reduced.ingests_by_thread.get(&ing.thread_id).and_then(|q| {\n- q.iter().rev().position(|id| id == &ing.id).map(|i| i + 1)\n- });\n- slug_types::FeedPost {\n- ts: ing.ts,\n- id: ing.id.clone(),\n- thread,\n- thread_post_index,\n- body: ing.raw.clone(),\n- }\n- })\n- .collect();\n-\n- Json(slug_types::FeedResponse {\n- actor,\n- since,\n- posts,\n- total,\n- }).into_response()\n-}\ndiff --git a/server/src/api/forum.rs b/server/src/api/forum.rs\ndeleted file mode 100644\nindex ab6f3ec4936979fed9b3a9c47bef609610c751dd..0000000000000000000000000000000000000000\n--- a/server/src/api/forum.rs\n+++ /dev/null\n@@ -1,142 +0,0 @@\n-use axum::{\n- extract::{Query, State},\n- http::StatusCode,\n- response::IntoResponse,\n- Json,\n-};\n-use serde::Deserialize;\n-use slug_types::*;\n-\n-use crate::{\n- canonical_path::canonicalize_tag,\n- identity::parse_username,\n- state::AppState,\n-};\n-\n-use super::helpers::api_error;\n-\n-pub async fn get_threads(State(state): State) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n-\n- let mut out: Vec = reduced\n- .threads\n- .iter()\n- .map(|(thread, ts)| {\n- ThreadSummary {\n- thread: format!(\"#{thread}\"),\n- last_activity_ts: ts.last_activity_ts,\n- web: format!(\"https://slug.social/t/{}\", thread),\n- }\n- })\n- .collect();\n-\n- out.sort_by(|a, b| b.last_activity_ts.cmp(&a.last_activity_ts));\n- Json(ThreadsResponse { threads: out }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct ThreadDetailQuery {\n- pub tag: String,\n- #[serde(default)]\n- /// Chronological offset (0 = oldest). Default: 0.\n- pub offset: Option,\n- /// Number of posts to return. Default: 10, max: 500.\n- pub limit: Option,\n- /// Only posts at or after this Unix ms timestamp.\n- pub since: Option,\n- /// Only posts strictly before this Unix ms timestamp.\n- pub before: Option,\n- /// Filter to posts whose principal username starts with this prefix (stored form, no `@`).\n- pub actor: Option,\n- /// Return the single post with this ingest ID.\n- pub post_id: Option,\n-}\n-\n-/// Thread (forum) detail by tag -- paginated, filterable posts, full body.\n-pub async fn get_thread(State(state): State, Query(q): Query) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let tag = canonicalize_tag(&q.tag);\n- let actor_prefix = match q.actor.as_deref().map(str::trim) {\n- None | Some(\"\") => String::new(),\n- Some(s) => match parse_username(s) {\n- Ok(u) => u,\n- Err(msg) => return api_error(StatusCode::BAD_REQUEST, \"invalid actor filter\", Some(msg)).into_response(),\n- },\n- };\n- let reduced = reduced_arc.read().await;\n-\n- // Single post lookup by ingest ID -- return full body untruncated.\n- if let Some(ref post_id) = q.post_id {\n- let thread_ids = reduced.ingests_by_thread.get(&tag);\n- let index = thread_ids.and_then(|ids| {\n- ids.iter().rev().enumerate().find(|(_, id)| *id == post_id).map(|(i, _)| i)\n- });\n- return match index.and_then(|idx| reduced.ingests_by_id.get(post_id).map(|ing| (idx, ing))) {\n- None => api_error(StatusCode::NOT_FOUND, \"post not found\", None),\n- Some((idx, ing)) => Json(ThreadDetailResponse {\n- thread: format!(\"#{}\", tag),\n- posts: vec![PostRow {\n- id: ing.id.clone(),\n- index: idx,\n- ts: ing.ts,\n- actor: ing.principal.clone(),\n- body: ing.raw.clone(),\n- truncated: false,\n- }],\n- total: 1,\n- offset: idx,\n- }).into_response(),\n- };\n- }\n-\n- let offset = q.offset.unwrap_or(0);\n- let limit = q.limit.unwrap_or(10).clamp(1, 500);\n-\n- // ingests_by_thread is newest-first; reverse to chronological order.\n- let all_ids: Vec = reduced\n- .ingests_by_thread\n- .get(&tag)\n- .map(|q| q.iter().rev().cloned().collect())\n- .unwrap_or_default();\n-\n- let filtered: Vec<(usize, _)> = all_ids\n- .into_iter()\n- .enumerate()\n- .filter_map(|(idx, id)| reduced.ingests_by_id.get(&id).map(|ing| (idx, ing.clone())))\n- .filter(|(_, ing)| q.since.map_or(true, |s| ing.ts >= s))\n- .filter(|(_, ing)| q.before.map_or(true, |b| ing.ts < b))\n- .filter(|(_, ing)| actor_prefix.is_empty() || ing.principal.to_lowercase().starts_with(&actor_prefix))\n- .collect();\n-\n- let total = filtered.len();\n-\n- const MAX_BODY: usize = 2000;\n- let posts: Vec = filtered\n- .into_iter()\n- .skip(offset)\n- .take(limit)\n- .map(|(idx, ing)| {\n- let (body, truncated) = if ing.raw.len() > MAX_BODY {\n- (ing.raw[..MAX_BODY].to_string(), true)\n- } else {\n- (ing.raw.clone(), false)\n- };\n- PostRow {\n- id: ing.id.clone(),\n- index: idx,\n- ts: ing.ts,\n- actor: ing.principal.clone(),\n- body,\n- truncated,\n- }\n- })\n- .collect();\n-\n- Json(ThreadDetailResponse {\n- thread: format!(\"#{}\", tag),\n- posts,\n- total,\n- offset,\n- }).into_response()\n-}\ndiff --git a/server/src/api/garden.rs b/server/src/api/garden.rs\ndeleted file mode 100644\nindex ca0e9e735b8e02e31075ffb77d99941129a045f0..0000000000000000000000000000000000000000\n--- a/server/src/api/garden.rs\n+++ /dev/null\n@@ -1,223 +0,0 @@\n-use axum::{\n- extract::{Query, State},\n- http::StatusCode,\n- response::IntoResponse,\n- Json,\n-};\n-use serde::Deserialize;\n-use slug_types::*;\n-use std::collections::HashSet;\n-\n-use crate::{\n- canonical_path::canonicalize_item,\n- path_types::CanonicalItemUrl,\n- state::AppState,\n-};\n-\n-use super::helpers::{api_error, item_path_for_api, vote_touches_path};\n-\n-// ============================================================================\n-// Exploration APIs (read-only)\n-// ============================================================================\n-\n-#[derive(Debug, Deserialize)]\n-pub struct PathsQuery {}\n-\n-/// List root paths (items with parent \"\").\n-pub async fn get_paths(State(state): State, Query(_q): Query) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n-\n- let out: Vec = reduced\n- .public()\n- .item_children\n- .get(&CanonicalItemUrl::ontology_root())\n- .map(|roots| {\n- let mut v: Vec = roots.iter()\n- .map(|path| {\n- let children = content.item_children.get(path.as_str()).map(|s| s.len()).unwrap_or(0);\n- PathSummary {\n- path: format!(\"~/{}\", path),\n- children,\n- web: format!(\"https://slug.social/~/{}\", path),\n- }\n- }).collect();\n- v.sort_by(|a, b| a.path.cmp(&b.path));\n- v\n- })\n- .unwrap_or_default();\n-\n- Json(PathsResponse { paths: out }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct LeavesQuery {}\n-\n-/// List every leaf item (full path). Items that have no children. Does not scale; works for now.\n-pub async fn get_leaves(State(state): State, Query(_q): Query) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- let parents: HashSet<&str> = content.item_children.keys().map(|s| s.as_str()).collect();\n- let mut paths: Vec = content\n- .items\n- .iter()\n- .filter(|p| !parents.contains(p.as_str()))\n- .map(|p| p.as_str().to_string())\n- .collect();\n- paths.sort();\n- Json(LeavesResponse { paths }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct ItemQuery {\n- pub item: String,\n- /// Return the full body without truncation (default: false, bodies >10k chars are truncated).\n- #[serde(default)]\n- pub full: Option,\n-}\n-\n-pub async fn get_item(State(state): State, Query(q): Query) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let item_str = canonicalize_item(&q.item);\n- let item = CanonicalItemUrl(item_str.clone());\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n-\n- if !content.items.contains(&item) {\n- return api_error(\n- StatusCode::NOT_FOUND,\n- \"item not found\",\n- Some(format!(\"{} does not exist\", item_path_for_api(&item_str))),\n- );\n- }\n-\n- const MAX_ITEM_BODY: usize = 10_000;\n- let want_full = q.full.unwrap_or(false);\n-\n- let (body, truncated, body_len) = match content.item_bodies.get(&item) {\n- None => (None, false, 0),\n- Some(raw) => {\n- let char_len = raw.chars().count();\n- if !want_full && char_len > MAX_ITEM_BODY {\n- let byte_end = raw.char_indices().nth(MAX_ITEM_BODY).map(|(i, _)| i).unwrap_or(raw.len());\n- (Some(raw[..byte_end].to_string()), true, char_len)\n- } else {\n- (Some(raw.clone()), false, 0)\n- }\n- }\n- };\n-\n- let threads: Vec = reduced\n- .public()\n- .item_threads\n- .get(&item)\n- .map(|s| s.iter().cloned().collect())\n- .unwrap_or_default();\n-\n- Json(ItemResponse {\n- item: item_str.clone(),\n- body,\n- truncated,\n- body_len,\n- threads,\n- })\n- .into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct RecentVotesQuery {\n- #[serde(default)]\n- pub parent: Option,\n- #[serde(default)]\n- pub limit: Option,\n-}\n-\n-pub async fn get_recent_votes(\n- State(state): State,\n- Query(q): Query,\n-) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let limit = q.limit.unwrap_or(25).clamp(1, 200);\n-\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- let group = &content.ranking_group;\n-\n- let iter = group.recent_votes.iter();\n- let iter: Box> = if let Some(parent) = &q.parent {\n- let parent_can = canonicalize_item(parent);\n- Box::new(iter.filter(move |v| vote_touches_path(v.a.as_str(), v.b.as_str(), &parent_can)))\n- } else {\n- Box::new(iter)\n- };\n-\n- let out: Vec = iter\n- .take(limit)\n- .map(|v| VoteRow {\n- ts: v.ts,\n- a: v.a.as_str().to_string(),\n- b: v.b.as_str().to_string(),\n- ratio: format!(\"{}:{}\", v.ratio_left, v.ratio_right),\n- actor: Some(v.principal.clone()),\n- body: v.body.clone(),\n- thread: Some(v.thread_id.clone()),\n- }).collect();\n-\n- Json(RecentVotesResponse { votes: out }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct MatchupQuery {\n- pub item: String,\n- #[serde(default)]\n- pub limit: Option,\n-}\n-\n-/// Vote history for one item (matchup: wins/losses with thread per vote).\n-pub async fn get_matchup(\n- State(state): State,\n- Query(q): Query,\n-) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let item_str = canonicalize_item(&q.item);\n- let item = CanonicalItemUrl(item_str.clone());\n- let limit = q.limit.unwrap_or(50).clamp(1, 200);\n-\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n-\n- if !content.items.contains(&item) {\n- return api_error(\n- StatusCode::NOT_FOUND,\n- \"item not found\",\n- Some(format!(\"{} does not exist\", item_path_for_api(&item_str))),\n- );\n- }\n-\n- let votes: Vec = content\n- .item_votes\n- .get(&item)\n- .map(|q| {\n- q.iter()\n- .take(limit)\n- .map(|v| VoteRow {\n- ts: v.ts,\n- a: v.a.as_str().to_string(),\n- b: v.b.as_str().to_string(),\n- ratio: format!(\"{}:{}\", v.ratio_left, v.ratio_right),\n- actor: Some(v.principal.clone()),\n- body: v.body.clone(),\n- thread: Some(v.thread_id.clone()),\n- })\n- .collect()\n- })\n- .unwrap_or_default();\n-\n- Json(MatchupResponse {\n- item: item_path_for_api(&item_str),\n- votes,\n- })\n- .into_response()\n-}\ndiff --git a/server/src/api/helpers.rs b/server/src/api/helpers.rs\nindex a16c3df6f9d8eaa2e587f5575b78660880062f81..81e2a55fa3abb8609b4099f91a989480336e11eb 100644\n--- a/server/src/api/helpers.rs\n+++ b/server/src/api/helpers.rs\n@@ -11,7 +11,6 @@ use crate::{\n canonical_path::canonicalize_item,\n path_types::CanonicalItemUrl,\n ranking::connected_components_from_voted_pairs,\n- reducer::ReducerState,\n };\n \n pub fn api_error(status: StatusCode, error: impl Into, hint: Option) -> axum::response::Response {\n@@ -136,8 +135,7 @@ pub fn is_pair_voted(group: &crate::reducer::GroupState, a: &str, b: &str) -> bo\n }\n \n /// Compute graph connectivity stats for a set of items within the ranking group.\n-pub fn compute_connectivity_stats(reduced: &ReducerState, pool: &[String]) -> ConnectivityStats {\n- let group = &reduced.public().ranking_group;\n+pub fn compute_connectivity_stats(group: &crate::reducer::GroupState, pool: &[String]) -> ConnectivityStats {\n let n = pool.len();\n \n // Map pool items to global indices (items not yet in the group get no index)\ndiff --git a/server/src/api/ingest.rs b/server/src/api/ingest.rs\ndeleted file mode 100644\nindex 5a58754bae8ed5f7ae9ce89e659a61de81c1bc01..0000000000000000000000000000000000000000\n--- a/server/src/api/ingest.rs\n+++ /dev/null\n@@ -1,519 +0,0 @@\n-use axum::{\n- extract::State,\n- http::{HeaderMap, StatusCode},\n- response::IntoResponse,\n- Json,\n-};\n-use serde::Deserialize;\n-use slug_types::*;\n-use std::collections::{HashMap, HashSet};\n-\n-use crate::{\n- canonical_path::canonicalize_tag,\n- dsl,\n- events::{AgentBound, Event, Ingest, ThreadCapability, ThreadVisibility},\n- identity::parse_agent,\n- path_types::CanonicalItemUrl,\n- reducer::{ReducerState, ScopeId},\n- state::AppState,\n-};\n-use super::auth::verify_bearer_principal;\n-use super::helpers::{api_error, item_path_for_api, now_ms, resolve_item};\n-\n-/// Result of validating an ingest document. Shared by post_ingest and post_check.\n-#[derive(Debug)]\n-pub struct ValidatedIngest {\n- pub doc: dsl::Document,\n- pub ts: i64,\n- pub raw_text: String,\n-}\n-\n-/// Parse and validate an ingest document against current reduced state.\n-/// Returns a validated struct for commit (ingest) or dry-run (check).\n-/// Error is (StatusCode, error message, optional hint).\n-///\n-/// Item existence is checked in `scope` first, then falls back to public.\n-pub fn validate_ingest_document(\n- reduced: &ReducerState,\n- text: &str,\n- scope: &ScopeId,\n-) -> Result)> {\n- let public_content = reduced.public();\n- // For a private scope, also check items defined there.\n- let scoped_content = match scope {\n- ScopeId::Public => None,\n- _ => reduced.content_for_scope(scope),\n- };\n- let item_exists = |key: &CanonicalItemUrl| {\n- scoped_content.map(|c| c.items.contains(key)).unwrap_or(false)\n- || public_content.items.contains(key)\n- };\n- let body_exists = |key: &CanonicalItemUrl| {\n- scoped_content.map(|c| c.item_bodies.contains_key(key)).unwrap_or(false)\n- || public_content.item_bodies.contains_key(key)\n- };\n- let doc = match dsl::parse_full(text) {\n- Ok(d) => d,\n- Err(e) => {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- \"parse error\".to_string(),\n- Some(format!(\"{}\", e)),\n- ));\n- }\n- };\n-\n- let ts = now_ms();\n- let mut defined_in_doc: HashSet = HashSet::new();\n-\n- for s in &doc.statements {\n- match s {\n- dsl::Stmt::Item { title, body } => {\n- let item = match resolve_item(title) {\n- Ok(v) => v,\n- Err(msg) => {\n- return Err((StatusCode::BAD_REQUEST, \"invalid item path\".to_string(), Some(msg)));\n- }\n- };\n- let Some(body_text) = body else {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- format!(\"item missing body: {}\", item_path_for_api(&item)),\n- Some(\"items must be declared with bodies, e.g. `~/path/item { ... }`\".to_string()),\n- ));\n- };\n- if body_text.trim().is_empty() {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- format!(\"item body is empty: {}\", item_path_for_api(&item)),\n- Some(\"write at least one sentence inside `{ ... }`\".to_string()),\n- ));\n- }\n- defined_in_doc.insert(item);\n- }\n- dsl::Stmt::Vote { item1, item2, .. } => {\n- let a = match resolve_item(item1) {\n- Ok(v) => v,\n- Err(msg) => {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- \"invalid vote item path\".to_string(),\n- Some(msg),\n- ));\n- }\n- };\n- let b = match resolve_item(item2) {\n- Ok(v) => v,\n- Err(msg) => {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- \"invalid vote item path\".to_string(),\n- Some(msg),\n- ));\n- }\n- };\n- let missing: Vec = [&a, &b]\n- .into_iter()\n- .filter(|it| {\n- let key = CanonicalItemUrl((*it).clone());\n- !defined_in_doc.contains(*it) && !item_exists(&key)\n- })\n- .map(|it| item_path_for_api(it))\n- .collect();\n- if !missing.is_empty() {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- \"vote references undefined item(s)\".to_string(),\n- Some(format!(\n- \"define items with bodies before voting. missing: {}\",\n- missing.join(\", \")\n- )),\n- ));\n- }\n- let missing_body: Vec = [&a, &b]\n- .into_iter()\n- .filter(|it| {\n- let key = CanonicalItemUrl((*it).clone());\n- !defined_in_doc.contains(*it) && !body_exists(&key)\n- })\n- .map(|it| item_path_for_api(it))\n- .collect();\n- if !missing_body.is_empty() {\n- return Err((\n- StatusCode::BAD_REQUEST,\n- \"vote references item(s) without bodies\".to_string(),\n- Some(format!(\n- \"missing bodies: {}. Declare each item as `/item {{ ... }}`\",\n- missing_body.join(\", \")\n- )),\n- ));\n- }\n- }\n- _ => {}\n- }\n- }\n-\n- Ok(ValidatedIngest {\n- doc,\n- ts,\n- raw_text: text.to_string(),\n- })\n-}\n-\n-/// Compute ranking changes between two snapshots of a parent scope's rankings.\n-/// Returns None if nothing changed.\n-fn compute_scope_rank_changes(\n- parent: &str,\n- before: &crate::scope_rank::ChildrenRankings,\n- after: &crate::scope_rank::ChildrenRankings,\n-) -> Option {\n- fn build_positions(rankings: &crate::scope_rank::ChildrenRankings) -> HashMap> {\n- let mut map = HashMap::new();\n- for comp in &rankings.component_rankings {\n- let total = comp.ranked.len();\n- for (i, item) in comp.ranked.iter().enumerate() {\n- map.insert(item.item.as_str().to_string(), Some(RankPosition { rank: i + 1, of: total }));\n- }\n- }\n- for item in &rankings.unranked_items {\n- map.insert(item.as_str().to_string(), None);\n- }\n- map\n- }\n-\n- let before_pos = build_positions(before);\n- let after_pos = build_positions(after);\n-\n- let all_items: std::collections::BTreeSet = before_pos.keys().cloned()\n- .chain(after_pos.keys().cloned())\n- .collect();\n-\n- let mut changes: Vec = Vec::new();\n- for item in all_items {\n- let b = before_pos.get(&item).cloned().flatten();\n- let a = after_pos.get(&item).cloned().flatten();\n- let changed = match (&b, &a) {\n- (Some(bp), Some(ap)) => bp.rank != ap.rank || bp.of != ap.of,\n- (None, Some(_)) => true,\n- (Some(_), None) => true,\n- (None, None) => false,\n- };\n- if changed {\n- changes.push(RankChange {\n- item: item_path_for_api(&item),\n- before: b,\n- after: a,\n- });\n- }\n- }\n-\n- if changes.is_empty() {\n- return None;\n- }\n-\n- changes.sort_by(|a, b| match (&a.after, &b.after) {\n- (Some(ap), Some(bp)) => ap.rank.cmp(&bp.rank),\n- (Some(_), None) => std::cmp::Ordering::Less,\n- (None, Some(_)) => std::cmp::Ordering::Greater,\n- (None, None) => a.item.cmp(&b.item),\n- });\n-\n- Some(ScopeRankChanges {\n- parent: if parent.is_empty() {\n- \"/\".to_string()\n- } else {\n- item_path_for_api(parent)\n- },\n- changes,\n- })\n-}\n-\n-pub async fn post_ingest(\n- State(state): State,\n- headers: HeaderMap,\n- Json(req): Json,\n-) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let event_log = state.event_log.clone();\n- let reduced = reduced_arc.read().await;\n-\n- let delegate: Option = match &req.delegate {\n- None => None,\n- Some(s) if s.trim().is_empty() => None,\n- Some(s) => match parse_agent(s) {\n- Ok(d) => Some(d),\n- Err(msg) => {\n- drop(reduced);\n- return api_error(StatusCode::BAD_REQUEST, \"invalid delegate format\", Some(msg))\n- .into_response();\n- }\n- },\n- };\n- let thread_id = canonicalize_tag(&req.thread);\n-\n- let principal = match verify_bearer_principal(&headers, &reduced) {\n- Ok(u) => u,\n- Err((st, msg)) => {\n- drop(reduced);\n- return api_error(st, msg, None).into_response();\n- }\n- };\n-\n- // Determine thread visibility and scope for item lookups and cap enforcement.\n- let visibility = reduced.threads.get(&thread_id)\n- .map(|t| t.visibility)\n- .unwrap_or(ThreadVisibility::Public);\n- let scope = match visibility {\n- ThreadVisibility::Public => ScopeId::Public,\n- ThreadVisibility::Private => ScopeId::Thread(thread_id.clone()),\n- };\n-\n- let v = match validate_ingest_document(&reduced, &req.text, &scope) {\n- Ok(x) => x,\n- Err((status, msg, hint)) => {\n- drop(reduced);\n- return api_error(status, msg, hint).into_response();\n- }\n- };\n-\n- // Enforce capabilities for private threads.\n- if matches!(visibility, ThreadVisibility::Private) {\n- let mut required: HashSet = HashSet::new();\n- required.insert(ThreadCapability::View);\n- for stmt in &v.doc.statements {\n- match stmt {\n- dsl::Stmt::Vote { .. } => { required.insert(ThreadCapability::Vote); }\n- dsl::Stmt::Item { .. } => { required.insert(ThreadCapability::AddItem); }\n- dsl::Stmt::Prose { .. } => { required.insert(ThreadCapability::Post); }\n- }\n- }\n- let missing: Vec<_> = required.iter()\n- .filter(|cap| !reduced.user_has_cap(&thread_id, &principal, **cap))\n- .collect();\n- if !missing.is_empty() {\n- drop(reduced);\n- return api_error(StatusCode::FORBIDDEN, \"insufficient capabilities for this thread\", None).into_response();\n- }\n- }\n-\n- if let Some(ref d) = delegate {\n- match reduced.agent_bindings.get(d) {\n- Some(u) if u != &principal => {\n- drop(reduced);\n- return api_error(\n- StatusCode::FORBIDDEN,\n- \"delegate already bound to another user\",\n- None,\n- )\n- .into_response();\n- }\n- _ => {}\n- }\n- }\n- let need_agent_bind = delegate\n- .as_ref()\n- .map(|d| reduced.agent_bindings.get(d).is_none())\n- .unwrap_or(false);\n- drop(reduced);\n-\n- let mut events_appended: usize = 0;\n-\n- if need_agent_bind {\n- if let Some(agent_id) = delegate.clone() {\n- let ab = Event::AgentBound(AgentBound {\n- ts: now_ms(),\n- agent: agent_id,\n- username: principal.clone(),\n- });\n- if let Err(err) = event_log.append(&ab).await {\n- return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!(\"{err}\"), None);\n- }\n- events_appended += 1;\n- {\n- let mut reduced = reduced_arc.write().await;\n- reduced.apply_event(ab);\n- }\n- }\n- }\n-\n- // Collect parent scopes for all voted items so we can compute ranking deltas.\n- let voted_parent_scopes: Vec = {\n- let mut parents: HashSet = HashSet::new();\n- for s in &v.doc.statements {\n- if let dsl::Stmt::Vote { item1, item2, .. } = s {\n- if let (Ok(a), Ok(b)) = (resolve_item(item1), resolve_item(item2)) {\n- if let Some(p) = CanonicalItemUrl::parse(&a).and_then(|c| c.parent()) { parents.insert(p); }\n- if let Some(p) = CanonicalItemUrl::parse(&b).and_then(|c| c.parent()) { parents.insert(p); }\n- }\n- }\n- }\n- let mut out: Vec = parents.into_iter().collect();\n- out.sort();\n- out\n- };\n-\n- // Snapshot rankings before the ingest event is applied.\n- let pre_rankings: HashMap =\n- if !voted_parent_scopes.is_empty() {\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- voted_parent_scopes\n- .iter()\n- .map(|p| (p.clone(), crate::scope_rank::build_children_rankings(content, p)))\n- .collect()\n- } else {\n- HashMap::new()\n- };\n-\n- let ingest_event = Event::Ingest(Ingest {\n- ts: v.ts,\n- id: uuid::Uuid::new_v4().to_string(),\n- raw: v.raw_text.clone(),\n- principal: principal.clone(),\n- delegate: delegate.clone(),\n- thread_id: thread_id.clone(),\n- });\n-\n- if let Err(err) = event_log.append(&ingest_event).await {\n- return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!(\"{err}\"), None);\n- }\n- events_appended += 1;\n-\n- {\n- let mut reduced = reduced_arc.write().await;\n- reduced.apply_event(ingest_event);\n- }\n-\n- // Snapshot rankings after the event and compute per-scope deltas.\n- let ranking_changes: Vec = if !voted_parent_scopes.is_empty() {\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- voted_parent_scopes\n- .iter()\n- .filter_map(|p| {\n- let before = pre_rankings.get(p)?;\n- let after = crate::scope_rank::build_children_rankings(content, p);\n- compute_scope_rank_changes(p.as_str(), before, &after)\n- })\n- .collect()\n- } else {\n- vec![]\n- };\n-\n- Json(IngestResponse {\n- ok: true,\n- threads: vec![format!(\"#{}\", thread_id)],\n- events_appended,\n- next: NextMoves {\n- pair: \"npx slugsocial pair\".to_string(),\n- rank: \"npx slugsocial rank\".to_string(),\n- web: format!(\"https://slug.social/t/{}\", thread_id),\n- },\n- ranking_changes,\n- }).into_response()\n-}\n-\n-pub async fn post_check(\n- State(state): State,\n- Json(req): Json,\n-) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n- // check is auth-free; always validates against public scope.\n- let v = match validate_ingest_document(&reduced, &req.text, &ScopeId::Public) {\n- Ok(x) => x,\n- Err((status, msg, hint)) => return api_error(status, msg, hint).into_response(),\n- };\n- drop(reduced);\n-\n- let delegate: Option = match &req.delegate {\n- None => None,\n- Some(s) if s.trim().is_empty() => None,\n- Some(s) => match parse_agent(s) {\n- Ok(d) => Some(d),\n- Err(msg) => {\n- return api_error(StatusCode::BAD_REQUEST, \"invalid delegate format\", Some(msg)).into_response();\n- }\n- },\n- };\n- let thread_id = canonicalize_tag(&req.thread);\n- let principal = \"placeholder\".to_string();\n-\n- let event = Event::Ingest(Ingest {\n- ts: v.ts,\n- id: uuid::Uuid::new_v4().to_string(),\n- raw: v.raw_text.clone(),\n- principal,\n- delegate,\n- thread_id: thread_id.clone(),\n- });\n-\n- let mut simulated = { reduced_arc.read().await.clone() };\n- simulated.apply_event(event);\n-\n- // Collect parent scopes touched by votes in this document.\n- let voted_parents: Vec = {\n- let mut parents: std::collections::HashSet = std::collections::HashSet::new();\n- for s in &v.doc.statements {\n- if let dsl::Stmt::Vote { item1, item2, .. } = s {\n- if let (Ok(a), Ok(b)) = (resolve_item(item1), resolve_item(item2)) {\n- if let Some(p) = CanonicalItemUrl::parse(&a).and_then(|c| c.parent()) { parents.insert(p); }\n- if let Some(p) = CanonicalItemUrl::parse(&b).and_then(|c| c.parent()) { parents.insert(p); }\n- }\n- }\n- }\n- let mut out: Vec = parents.into_iter().collect();\n- out.sort();\n- out\n- };\n-\n- // Show the scoped rankings for affected parent paths; fall back to empty if no votes.\n- let rankings: Vec = voted_parents\n- .iter()\n- .map(|parent| {\n- let scoped = crate::scope_rank::build_children_rankings(simulated.public(), parent);\n- let components: Vec = scoped\n- .component_rankings\n- .into_iter()\n- .map(|comp| RankComponent {\n- pairs: comp.pairs,\n- ranking: comp\n- .ranked\n- .into_iter()\n- .map(|r| RankRow {\n- item: item_path_for_api(r.item.as_str()),\n- score: r.score,\n- percent: None,\n- })\n- .collect(),\n- })\n- .collect();\n- slug_types::CheckScopeRanking {\n- parent: item_path_for_api(parent.as_str()),\n- components,\n- unranked_items: scoped\n- .unranked_items\n- .into_iter()\n- .map(|it| item_path_for_api(it.as_str()))\n- .collect(),\n- }\n- })\n- .collect();\n-\n- Json(CheckResponse {\n- ok: true,\n- threads: vec![format!(\"#{}\", thread_id)],\n- rankings,\n- next: vec![\n- \"npx slugsocial ingest \".to_string(),\n- \"npx slugsocial threads\".to_string(),\n- format!(\"https://slug.social/t/{}\", thread_id),\n- ],\n- }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct IngestQuery {\n- #[serde(default)]\n- pub text: Option,\n-}\ndiff --git a/server/src/api/mod.rs b/server/src/api/mod.rs\nindex e3299ee353f8c62935fa82708762490a6d5c54e3..a1ea432932cdca0e42dc7377f0a25075fd5644c0 100644\n--- a/server/src/api/mod.rs\n+++ b/server/src/api/mod.rs\n@@ -1,18 +1,7 @@\n-mod debug;\n mod auth;\n-mod feed;\n-mod forum;\n-mod garden;\n mod helpers;\n-mod ingest;\n-mod rank;\n-mod search;\n-mod thread;\n-// mod stream; // disabled while HTML/SSE views are offline\n-\n-// Re-export all public items so `api::get_rank`, `api::post_ingest`, etc. still work.\n-\n-pub use debug::{get_debug_query, DebugQuery};\n+mod rpc;\n+mod validate;\n \n pub use auth::{\n get_pending_session,\n@@ -25,32 +14,15 @@ pub use auth::{\n get_choose_username,\n };\n \n-pub use feed::{get_feed, FeedQuery};\n-\n-pub use forum::{get_thread, get_threads, ThreadDetailQuery};\n-\n-pub use garden::{\n- get_item, get_leaves, get_matchup, get_paths, get_recent_votes,\n- ItemQuery, LeavesQuery, MatchupQuery, PathsQuery, RecentVotesQuery,\n-};\n-\n pub use helpers::{\n api_error, compute_connectivity_stats, is_pair_voted, now_ms, paginate_rankings,\n parse_parent_specs, pick_random_distinct, sha256_hex, resolve_item, vote_touches_path,\n+ item_path_for_api,\n };\n \n-pub use ingest::{post_check, post_ingest, validate_ingest_document, IngestQuery, ValidatedIngest};\n-\n-pub use rank::{\n- get_global_rank, get_pair, get_rank, get_rank_history,\n- GlobalRankQuery, PairQuery, RankHistoryQuery, RankQuery,\n-};\n-\n-pub use search::{get_search, SearchApiQuery};\n-\n-pub use thread::{post_create_thread, post_thread_grants};\n+pub use rpc::handle_rpc_batch;\n \n-// pub use stream::{get_html_stream, get_stream};\n+pub use validate::{normalize_room_and_thread, validate_ingest_document, ValidatedIngest};\n \n #[cfg(test)]\n mod tests {\n@@ -66,7 +38,8 @@ mod tests {\n raw: raw.to_string(),\n principal: \"test\".to_string(),\n delegate: Some(\"00000000-0000-0000-0000-000000000000:test:local/test\".to_string()),\n- thread_id: \"t\".to_string(),\n+ room_id: \"public\".to_string(),\n+ thread_tag: \"t\".to_string(),\n }));\n }\n \n@@ -94,7 +67,7 @@ mod tests {\n #[test]\n fn validate_ingest_document_rejects_vote_on_undefined_item() {\n let reduced = ReducerState::default();\n- let text = \"~/t/a {a}\\n~/t/b 1:1 ~/t/missing {why}\\n\";\n+ let text = \"~/t/a {x}\\n~/t/b 1:1 ~/t/missing {why}\\n\";\n let err = validate_ingest_document(&reduced, text, &crate::reducer::ScopeId::Public).unwrap_err();\n assert_eq!(err.0, StatusCode::BAD_REQUEST);\n assert!(err.1.contains(\"undefined item\"));\n@@ -103,7 +76,6 @@ mod tests {\n #[test]\n fn validate_ingest_document_accepts_quoted_thread_title() {\n let reduced = ReducerState::default();\n- // quoted thread title is now prose; should still parse.\n let text = \"\\\"This is a title\\\" { This is the body of the post }\\n~/t/a {a}\\n\";\n validate_ingest_document(&reduced, text, &crate::reducer::ScopeId::Public).unwrap();\n }\n@@ -111,7 +83,6 @@ mod tests {\n #[test]\n fn validate_ingest_document_rejects_multiple_threads() {\n let reduced = ReducerState::default();\n- // multiple thread declarations are now prose; should still parse.\n let text = \"#one\\n#two\\n~/t/a {a}\\n\";\n validate_ingest_document(&reduced, text, &crate::reducer::ScopeId::Public).unwrap();\n }\ndiff --git a/server/src/api/rank.rs b/server/src/api/rank.rs\ndeleted file mode 100644\nindex f2d348280ec4d70a32044601964cb5dc9382d9d7..0000000000000000000000000000000000000000\n--- a/server/src/api/rank.rs\n+++ /dev/null\n@@ -1,453 +0,0 @@\n-use axum::{\n- extract::{Query, State},\n- http::StatusCode,\n- response::IntoResponse,\n- Json,\n-};\n-use rand::seq::SliceRandom;\n-use serde::Deserialize;\n-use slug_types::*;\n-use std::collections::HashSet;\n-\n-use crate::{\n- path_types::CanonicalItemUrl,\n- ranking::{connected_components_from_voted_pairs, ranked_items_subset},\n- state::AppState,\n-};\n-\n-use super::helpers::{\n- api_error, compute_connectivity_stats, is_pair_voted, item_path_for_api, paginate_rankings,\n- parse_parent_specs, pick_random_distinct,\n-};\n-\n-// ============================================================================\n-// Rank / Pair -- one-ranking (parent-path-filtered)\n-// ============================================================================\n-\n-#[derive(Debug, Deserialize)]\n-pub struct RankQuery {\n- /// Parent path(s). Comma-separated to merge scopes: ~/a,~/b (e.g. rank ~/models ~/ai-models).\n- /// Use `~` alone for global ranking across all items.\n- #[serde(default)]\n- pub parent: Option,\n- #[serde(default)]\n- pub limit: Option,\n- /// Skip the first N items (for pagination).\n- #[serde(default)]\n- pub offset: Option,\n- /// When true, each RankRow includes a `percent` field: score as % of max score in its component.\n- #[serde(default)]\n- pub percent: Option,\n- /// How many levels deep to resolve children (default 1 = direct children only).\n- #[serde(default)]\n- pub depth: Option,\n-}\n-\n-pub async fn get_rank(State(state): State, Query(q): Query) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n-\n- let specs = parse_parent_specs(q.parent.as_ref());\n- let is_global = q.parent.as_deref().map(|p| p.trim() == \"~\").unwrap_or(false);\n- let depth = q.depth.unwrap_or(1).max(1);\n-\n- // 404 when a specific named path is requested but doesn't exist as a known item or parent.\n- if !is_global && !specs.is_empty() {\n- let none_exist = specs.iter().all(|spec| {\n- let Some(canon) = CanonicalItemUrl::parse(spec) else { return true };\n- !content.items.contains(&canon) && !content.item_children.contains_key(&canon)\n- });\n- if none_exist {\n- return api_error(\n- StatusCode::NOT_FOUND,\n- \"path not found\",\n- Some(format!(\"{} does not exist\", specs.join(\", \"))),\n- );\n- }\n- }\n-\n- let rankings = if is_global {\n- let all_items: Vec = content.items.iter().cloned().collect();\n- crate::scope_rank::build_rankings_for_item_set(content, &all_items)\n- } else if specs.is_empty() {\n- crate::scope_rank::build_children_rankings(content, &CanonicalItemUrl::ontology_root())\n- } else if depth > 1 {\n- let items = crate::scope_rank::resolve_scope_recursive(content, &specs, depth);\n- crate::scope_rank::build_rankings_for_item_set(content, &items)\n- } else {\n- let items = crate::scope_rank::resolve_scope(content, &specs);\n- crate::scope_rank::build_rankings_for_item_set(content, &items)\n- };\n-\n- let want_percent = q.percent.unwrap_or(false);\n-\n- let components: Vec = rankings\n- .component_rankings\n- .into_iter()\n- .map(|c| {\n- let max_score = c.ranked.first().map(|r| r.score).unwrap_or(1.0).max(1e-12);\n- RankComponent {\n- pairs: c.pairs,\n- ranking: c\n- .ranked\n- .into_iter()\n- .map(|r| RankRow {\n- item: item_path_for_api(r.item.as_str()),\n- percent: if want_percent {\n- Some((r.score / max_score) * 100.0)\n- } else {\n- None\n- },\n- score: r.score,\n- })\n- .collect(),\n- }\n- })\n- .collect();\n-\n- let prefixed_unranked: Vec = rankings\n- .unranked_items\n- .into_iter()\n- .map(|s| item_path_for_api(s.as_str()))\n- .collect();\n-\n- let offset = q.offset.unwrap_or(0);\n- let (components, unranked_items) = if offset > 0 || q.limit.is_some() {\n- paginate_rankings(components, prefixed_unranked, offset, q.limit)\n- } else {\n- (components, prefixed_unranked)\n- };\n-\n- Json(RankResponse {\n- components,\n- unranked_items,\n- })\n- .into_response()\n-}\n-\n-// ============================================================================\n-// Global rank -- flat, paginated ranking across all items\n-// ============================================================================\n-\n-#[derive(Debug, Deserialize)]\n-pub struct GlobalRankQuery {\n- #[serde(default)]\n- pub limit: Option,\n- #[serde(default)]\n- pub offset: Option,\n- /// When true, add a `percent` field to each row (top item = 100, bottom = 0).\n- #[serde(default)]\n- pub percent: Option,\n-}\n-\n-pub async fn get_global_rank(\n- State(state): State,\n- Query(q): Query,\n-) -> impl IntoResponse {\n- const DEFAULT_LIMIT: usize = 50;\n- const MAX_LIMIT: usize = 500;\n-\n- let limit = q.limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);\n- let offset = q.offset.unwrap_or(0);\n- let want_percent = q.percent.unwrap_or(false);\n-\n- let reduced_arc = state.reduced.clone();\n- // Write lock needed to update cached_scores when dirty.\n- let mut reduced = reduced_arc.write().await;\n- let content = reduced.content.get_mut(&crate::reducer::ScopeId::Public).expect(\"public scope missing\");\n-\n- // Find connected components in global graph; rank within each; largest component first.\n- let n = content.ranking_group.idx_to_item.len();\n- let (mut comps, _) = connected_components_from_voted_pairs(\n- n, content.ranking_group.voted_pairs.iter().copied(),\n- );\n- comps.sort_by(|a, b| b.len().cmp(&a.len()));\n-\n- let mut ranked: Vec = Vec::new();\n- for comp in &comps {\n- let items = ranked_items_subset(&content.ranking_group, comp, 10000, 1e-8);\n- let top = items.first().map(|r| r.score).unwrap_or(1.0);\n- let bot = items.last().map(|r| r.score).unwrap_or(0.0);\n- let range = (top - bot).max(1e-12);\n- for r in items {\n- let pct = want_percent.then(|| ((r.score - bot) / range * 100.0).clamp(0.0, 100.0));\n- ranked.push(RankRow { item: item_path_for_api(r.item.as_str()), score: r.score, percent: pct });\n- }\n- }\n-\n- let ranked_total = ranked.len();\n-\n- // Items that exist but have never been voted on.\n- let mut unranked: Vec = content\n- .items\n- .iter()\n- .filter(|it| !content.ranking_group.item_to_idx.contains_key(*it))\n- .map(|it| it.as_str().to_string())\n- .collect();\n- unranked.sort();\n- let unranked_total = unranked.len();\n-\n- let page: Vec = ranked\n- .into_iter()\n- .chain(unranked.into_iter().map(|it| RankRow {\n- item: item_path_for_api(&it),\n- score: 0.0,\n- percent: want_percent.then_some(0.0),\n- }))\n- .skip(offset)\n- .take(limit)\n- .collect();\n-\n- Json(GlobalRankResponse {\n- ranked_total,\n- unranked_total,\n- offset,\n- limit,\n- items: page,\n- })\n- .into_response()\n-}\n-\n-// ============================================================================\n-// Rank history -- per-item rank ledger\n-// ============================================================================\n-\n-#[derive(Debug, Deserialize)]\n-pub struct RankHistoryQuery {\n- pub item: String,\n-}\n-\n-pub async fn get_rank_history(\n- State(state): State,\n- Query(q): Query,\n-) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n-\n- let item_str = crate::canonical_path::canonicalize_item(&q.item);\n- let item = CanonicalItemUrl(item_str.clone());\n-\n- let entries = content.rank_history.get(&item).cloned().unwrap_or_default();\n-\n- // Resolve caused_by at query time: look up the ingest and filter relevant votes.\n- let history: Vec = entries.iter().map(|e| {\n- let caused_by: Vec = reduced.ingests_by_id.get(&e.post_id)\n- .and_then(|ing| crate::dsl::parse_full(&ing.raw).ok())\n- .map(|doc| {\n- doc.statements.into_iter().filter_map(|s| {\n- if let crate::dsl::Stmt::Vote { item1, item2, ratio_left, ratio_right, explanation } = s {\n- let a = crate::canonical_path::canonicalize_item(&item1);\n- let b = crate::canonical_path::canonicalize_item(&item2);\n- if a == item_str || b == item_str {\n- Some(VoteRow {\n- ts: e.ts,\n- a: item_path_for_api(&a),\n- b: item_path_for_api(&b),\n- ratio: format!(\"{}:{}\", ratio_left, ratio_right),\n- actor: reduced.ingests_by_id.get(&e.post_id).map(|ing| ing.principal.clone()),\n- body: explanation,\n- thread: Some(format!(\"#{}\", e.thread)),\n- })\n- } else {\n- None\n- }\n- } else {\n- None\n- }\n- }).collect()\n- })\n- .unwrap_or_default();\n-\n- // Compute 1-indexed chronological position within the thread.\n- // ingests_by_thread is most-recent-first, so index from back + 1.\n- let thread_post_index = reduced.ingests_by_thread\n- .get(&e.thread)\n- .and_then(|q| q.iter().rev().position(|id| id == &e.post_id))\n- .map(|i| i + 1)\n- .unwrap_or(0);\n-\n- slug_types::RankHistoryRow {\n- ts: e.ts,\n- scope_rank: e.scope_rank,\n- scope_rank_delta: e.scope_rank_delta,\n- scope_total: e.scope_total,\n- global_rank: e.global_rank,\n- global_rank_delta: e.global_rank_delta,\n- global_total: e.global_total,\n- score: e.score,\n- thread: format!(\"#{}\", e.thread),\n- thread_post_index,\n- caused_by,\n- }\n- }).collect();\n-\n- Json(slug_types::RankHistoryResponse {\n- item: item_path_for_api(&item_str),\n- history,\n- }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct PairQuery {\n- /// Parent path(s). Comma-separated to merge scopes: ~/a,~/b (e.g. rank ~/models ~/ai-models).\n- #[serde(default)]\n- pub parent: Option,\n- #[serde(default)]\n- pub random: Option,\n-}\n-\n-pub async fn get_pair(State(state): State, Query(q): Query) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let force_random = q.random.unwrap_or(false);\n-\n- let pool: Vec = {\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- let specs = parse_parent_specs(q.parent.as_ref());\n- let raw_pool: Vec = if specs.is_empty() {\n- content.ranking_group.idx_to_item.clone()\n- } else {\n- crate::scope_rank::resolve_scope(content, &specs)\n- };\n- raw_pool.into_iter().map(|it| it.0).collect()\n- };\n-\n- if pool.len() < 2 {\n- let scope = q.parent.as_deref().unwrap_or(\"(all)\");\n- return api_error(\n- axum::http::StatusCode::BAD_REQUEST,\n- format!(\"need at least 2 items under parent /{}\", scope),\n- Some(\"add items via ingest\".to_string()),\n- );\n- }\n-\n- if force_random {\n- let Some((left, right)) = pick_random_distinct(&pool) else {\n- return api_error(axum::http::StatusCode::BAD_REQUEST, \"need at least 2 items\", None);\n- };\n- let (left_body, right_body, threads, connectivity) = {\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- let left_key = CanonicalItemUrl(left.clone());\n- let right_key = CanonicalItemUrl(right.clone());\n- let lb = content.item_bodies.get(&left_key).cloned();\n- let rb = content.item_bodies.get(&right_key).cloned();\n- let th: Vec = content\n- .item_threads\n- .get(&left_key)\n- .into_iter()\n- .chain(content.item_threads.get(&right_key))\n- .flat_map(|s| s.iter().cloned())\n- .collect::>()\n- .into_iter()\n- .collect();\n- let cs = compute_connectivity_stats(&reduced, &pool);\n- (lb, rb, th, cs)\n- };\n- return Json(PairResponse {\n- left: item_path_for_api(&left),\n- right: item_path_for_api(&right),\n- left_body,\n- right_body,\n- threads,\n- connectivity: Some(connectivity),\n- })\n- .into_response();\n- }\n-\n- let selected: Option<(String, String)> = {\n- let mut reduced = reduced_arc.write().await;\n- let group = &mut reduced\n- .content\n- .get_mut(&crate::reducer::ScopeId::Public)\n- .expect(\"public scope missing\")\n- .ranking_group;\n- if group.idx_to_item.is_empty() {\n- pick_random_distinct(&pool)\n- } else {\n- let mut rng = rand::thread_rng();\n- let idxs: Vec = pool.iter()\n- .filter_map(|it| {\n- let key = CanonicalItemUrl(it.clone());\n- group.item_to_idx.get(&key).copied()\n- })\n- .collect();\n- let ranked = ranked_items_subset(group, &idxs, 10000, 1e-8);\n- let ranked_set: std::collections::HashSet =\n- ranked.iter().map(|r| r.item.as_str().to_string()).collect();\n- let unsorted: Vec = pool.iter()\n- .filter(|it| !ranked_set.contains(*it))\n- .cloned()\n- .collect();\n-\n- let mut pick: Option<(String, String)> = None;\n- if !unsorted.is_empty() {\n- if let Some(left) = unsorted.choose(&mut rng).cloned() {\n- let mut candidates: Vec = if !ranked.is_empty() {\n- ranked.iter().map(|r| r.item.as_str().to_string()).collect()\n- } else {\n- pool.clone()\n- };\n- candidates.retain(|c| c != &left);\n- if let Some(right) = candidates.choose(&mut rng).cloned() {\n- pick = Some((left, right));\n- }\n- }\n- } else if ranked.len() >= 2 {\n- for i in 0..(ranked.len().saturating_sub(1)) {\n- let a = ranked[i].item.as_str();\n- let b = ranked[i + 1].item.as_str();\n- if a != b && !is_pair_voted(group, a, b) {\n- pick = Some((a.to_string(), b.to_string()));\n- break;\n- }\n- }\n- if pick.is_none() {\n- for _ in 0..64 {\n- let (Some(a), Some(b)) = (pool.choose(&mut rng).cloned(), pool.choose(&mut rng).cloned()) else { break; };\n- if a != b && !is_pair_voted(group, &a, &b) {\n- pick = Some((a, b));\n- break;\n- }\n- }\n- }\n- }\n- pick.or_else(|| pick_random_distinct(&pool))\n- }\n- };\n-\n- let Some((left, right)) = selected else {\n- return api_error(axum::http::StatusCode::BAD_REQUEST, \"need at least 2 items\", None);\n- };\n-\n- let (left_body, right_body, threads, connectivity) = {\n- let reduced = reduced_arc.read().await;\n- let content = reduced.public();\n- let left_key = CanonicalItemUrl(left.clone());\n- let right_key = CanonicalItemUrl(right.clone());\n- let lb = content.item_bodies.get(&left_key).cloned();\n- let rb = content.item_bodies.get(&right_key).cloned();\n- let th: Vec = content\n- .item_threads\n- .get(&left_key)\n- .into_iter()\n- .chain(content.item_threads.get(&right_key))\n- .flat_map(|s| s.iter().cloned())\n- .collect::>()\n- .into_iter()\n- .collect();\n- let cs = compute_connectivity_stats(&reduced, &pool);\n- (lb, rb, th, cs)\n- };\n- Json(PairResponse {\n- left: item_path_for_api(&left),\n- right: item_path_for_api(&right),\n- left_body,\n- right_body,\n- threads,\n- connectivity: Some(connectivity),\n- })\n- .into_response()\n-}\ndiff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..5675dcd2e28062acbe3c037b94b77c33adf32474\n--- /dev/null\n+++ b/server/src/api/rpc.rs\n@@ -0,0 +1,1330 @@\n+use std::collections::{HashMap, HashSet};\n+use std::sync::OnceLock;\n+\n+use axum::{\n+ extract::State,\n+ http::HeaderMap,\n+ response::IntoResponse,\n+ Json,\n+};\n+use rand::seq::SliceRandom;\n+use slug_types::*;\n+\n+use crate::{\n+ canonical_path::{canonicalize_item, canonicalize_tag},\n+ dsl,\n+ events::{AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability, ThreadVisibility},\n+ identity::{parse_agent, parse_username},\n+ path_types::CanonicalItemUrl,\n+ ranking::{connected_components_from_voted_pairs, ranked_items_subset},\n+ reducer::{scope_from_room_wire, ReducerState, ScopeId},\n+ state::AppState,\n+};\n+\n+use super::auth::verify_bearer_principal;\n+use super::helpers::{\n+ compute_connectivity_stats, is_pair_voted, item_path_for_api, now_ms, paginate_rankings,\n+ parse_parent_specs, pick_random_distinct, resolve_item, vote_touches_path,\n+};\n+use super::validate::{normalize_room_and_thread, validate_ingest_document};\n+\n+fn empty_content() -> &'static crate::reducer::ContentState {\n+ static E: OnceLock = OnceLock::new();\n+ E.get_or_init(Default::default)\n+}\n+\n+fn content_for_room<'a>(reduced: &'a ReducerState, room: &str) -> &'a crate::reducer::ContentState {\n+ let scope = scope_from_room_wire(room);\n+ reduced.content.get(&scope).unwrap_or(empty_content())\n+}\n+\n+type RpcErr = (String, Option);\n+\n+fn line_ok(r: RpcResult) -> RpcLine {\n+ RpcLine {\n+ ok: true,\n+ result: Some(r),\n+ error: None,\n+ hint: None,\n+ }\n+}\n+\n+fn line_err(msg: impl Into, hint: Option) -> RpcLine {\n+ RpcLine {\n+ ok: false,\n+ result: None,\n+ error: Some(msg.into()),\n+ hint,\n+ }\n+}\n+\n+fn compute_scope_rank_changes(\n+ parent: &str,\n+ before: &crate::scope_rank::ChildrenRankings,\n+ after: &crate::scope_rank::ChildrenRankings,\n+) -> Option {\n+ fn build_positions(rankings: &crate::scope_rank::ChildrenRankings) -> HashMap> {\n+ let mut map = HashMap::new();\n+ for comp in &rankings.component_rankings {\n+ let total = comp.ranked.len();\n+ for (i, item) in comp.ranked.iter().enumerate() {\n+ map.insert(item.item.as_str().to_string(), Some(RankPosition { rank: i + 1, of: total }));\n+ }\n+ }\n+ for item in &rankings.unranked_items {\n+ map.insert(item.as_str().to_string(), None);\n+ }\n+ map\n+ }\n+\n+ let before_pos = build_positions(before);\n+ let after_pos = build_positions(after);\n+\n+ let all_items: std::collections::BTreeSet = before_pos.keys().cloned()\n+ .chain(after_pos.keys().cloned())\n+ .collect();\n+\n+ let mut changes: Vec = Vec::new();\n+ for item in all_items {\n+ let b = before_pos.get(&item).cloned().flatten();\n+ let a = after_pos.get(&item).cloned().flatten();\n+ let changed = match (&b, &a) {\n+ (Some(bp), Some(ap)) => bp.rank != ap.rank || bp.of != ap.of,\n+ (None, Some(_)) => true,\n+ (Some(_), None) => true,\n+ (None, None) => false,\n+ };\n+ if changed {\n+ changes.push(RankChange {\n+ item: item_path_for_api(&item),\n+ before: b,\n+ after: a,\n+ });\n+ }\n+ }\n+\n+ if changes.is_empty() {\n+ return None;\n+ }\n+\n+ changes.sort_by(|a, b| match (&a.after, &b.after) {\n+ (Some(ap), Some(bp)) => ap.rank.cmp(&bp.rank),\n+ (Some(_), None) => std::cmp::Ordering::Less,\n+ (None, Some(_)) => std::cmp::Ordering::Greater,\n+ (None, None) => a.item.cmp(&b.item),\n+ });\n+\n+ Some(ScopeRankChanges {\n+ parent: if parent.is_empty() {\n+ \"/\".to_string()\n+ } else {\n+ item_path_for_api(parent)\n+ },\n+ changes,\n+ })\n+}\n+\n+fn parse_capability(s: &str) -> Result {\n+ match s {\n+ \"view\" => Ok(ThreadCapability::View),\n+ \"post\" => Ok(ThreadCapability::Post),\n+ \"vote\" => Ok(ThreadCapability::Vote),\n+ \"add_item\" => Ok(ThreadCapability::AddItem),\n+ \"manage\" => Ok(ThreadCapability::Manage),\n+ other => Err(format!(\"unknown capability: {other}\")),\n+ }\n+}\n+\n+fn gen_short_id() -> String {\n+ use rand::Rng;\n+ const ALPHABET: &[u8] = b\"0123456789abcdefghijklmnopqrstuvwxyz\";\n+ let mut rng = rand::thread_rng();\n+ (0..7).map(|_| ALPHABET[rng.gen_range(0..ALPHABET.len())] as char).collect()\n+}\n+\n+fn build_rank_response_for_content(\n+ content: &crate::reducer::ContentState,\n+ parent: Option<&str>,\n+ depth: usize,\n+ offset: usize,\n+ limit: Option,\n+ want_percent: bool,\n+) -> Result {\n+ let parent_owned = parent.map(|s| s.to_string());\n+ let specs = parse_parent_specs(parent_owned.as_ref());\n+ let is_global = parent.map(|p| p.trim() == \"~\").unwrap_or(false);\n+\n+ if !is_global && !specs.is_empty() {\n+ let none_exist = specs.iter().all(|spec| {\n+ let Some(canon) = CanonicalItemUrl::parse(spec) else { return true };\n+ !content.items.contains(&canon) && !content.item_children.contains_key(&canon)\n+ });\n+ if none_exist {\n+ return Err((\n+ \"path not found\".into(),\n+ Some(format!(\"{} does not exist\", specs.join(\", \"))),\n+ ));\n+ }\n+ }\n+\n+ let depth = depth.max(1);\n+ let rankings = if is_global {\n+ let all_items: Vec = content.items.iter().cloned().collect();\n+ crate::scope_rank::build_rankings_for_item_set(content, &all_items)\n+ } else if specs.is_empty() {\n+ crate::scope_rank::build_children_rankings(content, &CanonicalItemUrl::ontology_root())\n+ } else if depth > 1 {\n+ let items = crate::scope_rank::resolve_scope_recursive(content, &specs, depth);\n+ crate::scope_rank::build_rankings_for_item_set(content, &items)\n+ } else {\n+ let items = crate::scope_rank::resolve_scope(content, &specs);\n+ crate::scope_rank::build_rankings_for_item_set(content, &items)\n+ };\n+\n+ let components: Vec = rankings\n+ .component_rankings\n+ .into_iter()\n+ .map(|c| {\n+ let max_score = c.ranked.first().map(|r| r.score).unwrap_or(1.0).max(1e-12);\n+ RankComponent {\n+ pairs: c.pairs,\n+ ranking: c\n+ .ranked\n+ .into_iter()\n+ .map(|r| RankRow {\n+ item: item_path_for_api(r.item.as_str()),\n+ percent: if want_percent {\n+ Some((r.score / max_score) * 100.0)\n+ } else {\n+ None\n+ },\n+ score: r.score,\n+ })\n+ .collect(),\n+ }\n+ })\n+ .collect();\n+\n+ let prefixed_unranked: Vec = rankings\n+ .unranked_items\n+ .into_iter()\n+ .map(|s| item_path_for_api(s.as_str()))\n+ .collect();\n+\n+ let (components, unranked_items) = if offset > 0 || limit.is_some() {\n+ paginate_rankings(components, prefixed_unranked, offset, limit)\n+ } else {\n+ (components, prefixed_unranked)\n+ };\n+\n+ Ok(RankResponse {\n+ components,\n+ unranked_items,\n+ })\n+}\n+\n+async fn rpc_post(\n+ state: &AppState,\n+ headers: &HeaderMap,\n+ room: String,\n+ thread_tag: String,\n+ delegate_opt: Option,\n+ text: String,\n+ return_rank_diff: bool,\n+) -> Result {\n+ let reduced_arc = state.reduced.clone();\n+ let reduced = reduced_arc.read().await;\n+\n+ let principal = verify_bearer_principal(headers, &reduced).map_err(|(_, m)| (m, None))?;\n+\n+ let delegate: Option = match delegate_opt {\n+ None => None,\n+ Some(ref s) if s.trim().is_empty() => None,\n+ Some(s) => Some(parse_agent(&s).map_err(|msg| (format!(\"invalid delegate format\"), Some(msg)))?),\n+ };\n+\n+ let (room_key, thread_id) = normalize_room_and_thread(&room, &thread_tag);\n+ let scope = scope_from_room_wire(&room_key);\n+\n+ let is_private = !matches!(scope, ScopeId::Public);\n+ if is_private && !reduced.rooms.contains_key(&room_key) {\n+ drop(reduced);\n+ return Err((\"unknown room\".into(), Some(format!(\"room `{}` does not exist\", room_key))));\n+ }\n+\n+ let v = validate_ingest_document(&reduced, &text, &scope).map_err(|(st, m, h)| {\n+ let _ = st;\n+ (m, h)\n+ })?;\n+\n+ if is_private {\n+ let mut required: HashSet = HashSet::new();\n+ required.insert(ThreadCapability::View);\n+ for stmt in &v.doc.statements {\n+ match stmt {\n+ dsl::Stmt::Vote { .. } => { required.insert(ThreadCapability::Vote); }\n+ dsl::Stmt::Item { .. } => { required.insert(ThreadCapability::AddItem); }\n+ dsl::Stmt::Prose { .. } => { required.insert(ThreadCapability::Post); }\n+ }\n+ }\n+ let missing: Vec<_> = required.iter()\n+ .filter(|cap| !reduced.user_has_cap(&room_key, &principal, **cap))\n+ .collect();\n+ if !missing.is_empty() {\n+ drop(reduced);\n+ return Err((\"insufficient capabilities for this room\".into(), None));\n+ }\n+ }\n+\n+ if let Some(ref d) = delegate {\n+ match reduced.agent_bindings.get(d) {\n+ Some(u) if u != &principal => {\n+ drop(reduced);\n+ return Err((\"delegate already bound to another user\".into(), None));\n+ }\n+ _ => {}\n+ }\n+ }\n+ let need_agent_bind = delegate\n+ .as_ref()\n+ .map(|d| reduced.agent_bindings.get(d).is_none())\n+ .unwrap_or(false);\n+ drop(reduced);\n+\n+ let mut events_appended: usize = 0;\n+\n+ if need_agent_bind {\n+ if let Some(agent_id) = delegate.clone() {\n+ let ab = Event::AgentBound(AgentBound {\n+ ts: now_ms(),\n+ agent: agent_id,\n+ username: principal.clone(),\n+ });\n+ if let Err(err) = state.event_log.append(&ab).await {\n+ return Err((format!(\"{err}\"), None));\n+ }\n+ events_appended += 1;\n+ let mut reduced = reduced_arc.write().await;\n+ reduced.apply_event(ab);\n+ }\n+ }\n+\n+ let voted_parent_scopes: Vec = {\n+ let mut parents: HashSet = HashSet::new();\n+ for s in &v.doc.statements {\n+ if let dsl::Stmt::Vote { item1, item2, .. } = s {\n+ if let (Ok(a), Ok(b)) = (resolve_item(item1), resolve_item(item2)) {\n+ if let Some(p) = CanonicalItemUrl::parse(&a).and_then(|c| c.parent()) { parents.insert(p); }\n+ if let Some(p) = CanonicalItemUrl::parse(&b).and_then(|c| c.parent()) { parents.insert(p); }\n+ }\n+ }\n+ }\n+ let mut out: Vec = parents.into_iter().collect();\n+ out.sort();\n+ out\n+ };\n+\n+ let pre_rankings: HashMap =\n+ if !voted_parent_scopes.is_empty() {\n+ let reduced = reduced_arc.read().await;\n+ let content = content_for_room(&reduced, &room_key);\n+ voted_parent_scopes\n+ .iter()\n+ .map(|p| (p.clone(), crate::scope_rank::build_children_rankings(content, p)))\n+ .collect()\n+ } else {\n+ HashMap::new()\n+ };\n+\n+ let ingest_event = Event::Ingest(Ingest {\n+ ts: v.ts,\n+ id: uuid::Uuid::new_v4().to_string(),\n+ raw: v.raw_text.clone(),\n+ principal: principal.clone(),\n+ delegate: delegate.clone(),\n+ room_id: room_key.clone(),\n+ thread_tag: thread_id.clone(),\n+ });\n+\n+ if let Err(err) = state.event_log.append(&ingest_event).await {\n+ return Err((format!(\"{err}\"), None));\n+ }\n+ events_appended += 1;\n+\n+ {\n+ let mut reduced = reduced_arc.write().await;\n+ reduced.apply_event(ingest_event);\n+ }\n+\n+ let ranking_changes: Option> = if return_rank_diff && !voted_parent_scopes.is_empty() {\n+ let reduced = reduced_arc.read().await;\n+ let content = content_for_room(&reduced, &room_key);\n+ let v: Vec = voted_parent_scopes\n+ .iter()\n+ .filter_map(|p| {\n+ let before = pre_rankings.get(p)?;\n+ let after = crate::scope_rank::build_children_rankings(content, p);\n+ compute_scope_rank_changes(p.as_str(), before, &after)\n+ })\n+ .collect();\n+ if v.is_empty() { None } else { Some(v) }\n+ } else {\n+ None\n+ };\n+\n+ Ok(RpcResult::PostOk {\n+ events_appended,\n+ ranking_changes,\n+ threads: vec![format!(\"#{}\", thread_id)],\n+ next: NextMoves {\n+ pair: \"npx slugsocial pair\".to_string(),\n+ rank: \"npx slugsocial rank\".to_string(),\n+ web: format!(\"https://slug.social/t/{}\", thread_id),\n+ },\n+ })\n+}\n+\n+async fn rpc_check(\n+ state: &AppState,\n+ _room: String,\n+ text: String,\n+) -> Result {\n+ let reduced_arc = state.reduced.clone();\n+ let reduced = reduced_arc.read().await;\n+ let thread_id = \"check\".to_string();\n+ let v = validate_ingest_document(&reduced, &text, &ScopeId::Public).map_err(|(_, m, h)| (m, h))?;\n+ drop(reduced);\n+\n+ let delegate: Option = None;\n+ let principal = \"placeholder\".to_string();\n+ let event = Event::Ingest(Ingest {\n+ ts: v.ts,\n+ id: uuid::Uuid::new_v4().to_string(),\n+ raw: v.raw_text.clone(),\n+ principal,\n+ delegate,\n+ room_id: \"public\".into(),\n+ thread_tag: thread_id.clone(),\n+ });\n+\n+ let mut simulated = { reduced_arc.read().await.clone() };\n+ simulated.apply_event(event);\n+\n+ let voted_parents: Vec = {\n+ let mut parents: HashSet = HashSet::new();\n+ for s in &v.doc.statements {\n+ if let dsl::Stmt::Vote { item1, item2, .. } = s {\n+ if let (Ok(a), Ok(b)) = (resolve_item(item1), resolve_item(item2)) {\n+ if let Some(p) = CanonicalItemUrl::parse(&a).and_then(|c| c.parent()) { parents.insert(p); }\n+ if let Some(p) = CanonicalItemUrl::parse(&b).and_then(|c| c.parent()) { parents.insert(p); }\n+ }\n+ }\n+ }\n+ let mut out: Vec = parents.into_iter().collect();\n+ out.sort();\n+ out\n+ };\n+\n+ let rankings: Vec = voted_parents\n+ .iter()\n+ .map(|parent| {\n+ let scoped = crate::scope_rank::build_children_rankings(simulated.public(), parent);\n+ let components: Vec = scoped\n+ .component_rankings\n+ .into_iter()\n+ .map(|comp| RankComponent {\n+ pairs: comp.pairs,\n+ ranking: comp\n+ .ranked\n+ .into_iter()\n+ .map(|r| RankRow {\n+ item: item_path_for_api(r.item.as_str()),\n+ score: r.score,\n+ percent: None,\n+ })\n+ .collect(),\n+ })\n+ .collect();\n+ CheckScopeRanking {\n+ parent: item_path_for_api(parent.as_str()),\n+ components,\n+ unranked_items: scoped\n+ .unranked_items\n+ .into_iter()\n+ .map(|it| item_path_for_api(it.as_str()))\n+ .collect(),\n+ }\n+ })\n+ .collect();\n+\n+ Ok(RpcResult::CheckOk {\n+ rankings,\n+ threads: vec![format!(\"#{}\", thread_id)],\n+ next: vec![\n+ \"npx slugsocial ingest\".to_string(),\n+ \"npx slugsocial forum\".to_string(),\n+ format!(\"https://slug.social/t/{}\", thread_id),\n+ ],\n+ })\n+}\n+\n+fn forum_scope_from_room(room: &str) -> ScopeId {\n+ scope_from_room_wire(room)\n+}\n+\n+fn rpc_list_forum_threads(reduced: &ReducerState, room: &str) -> ThreadsResponse {\n+ let scope = forum_scope_from_room(room);\n+ let mut out: Vec = reduced\n+ .forum_threads\n+ .iter()\n+ .filter(|((s, _), _)| s == &scope)\n+ .map(|((_, tag), ts)| ThreadSummary {\n+ thread: format!(\"#{tag}\"),\n+ last_activity_ts: ts.last_activity_ts,\n+ web: format!(\"https://slug.social/t/{}\", tag),\n+ })\n+ .collect();\n+ out.sort_by(|a, b| b.last_activity_ts.cmp(&a.last_activity_ts));\n+ ThreadsResponse { threads: out }\n+}\n+\n+fn rpc_forum_thread_detail(\n+ reduced: &ReducerState,\n+ room: &str,\n+ thread_tag: &str,\n+ offset: usize,\n+ limit: usize,\n+ since: Option,\n+ before: Option,\n+ actor_prefix: &str,\n+ post_id: Option<&str>,\n+) -> Result {\n+ let scope = forum_scope_from_room(room);\n+ let tag = canonicalize_tag(thread_tag);\n+ let key = (scope.clone(), tag.clone());\n+\n+ if let Some(pid) = post_id {\n+ let thread_ids = reduced.ingests_by_scope_thread.get(&key);\n+ let index = thread_ids.and_then(|ids| {\n+ ids.iter().rev().enumerate().find(|(_, id)| *id == pid).map(|(i, _)| i)\n+ });\n+ return match index.and_then(|idx| reduced.ingests_by_id.get(pid).map(|ing| (idx, ing))) {\n+ None => Err((\"post not found\".into(), None)),\n+ Some((idx, ing)) => Ok(ThreadDetailResponse {\n+ thread: format!(\"#{}\", tag),\n+ posts: vec![PostRow {\n+ id: ing.id.clone(),\n+ index: idx,\n+ ts: ing.ts,\n+ actor: ing.principal.clone(),\n+ body: ing.raw.clone(),\n+ truncated: false,\n+ }],\n+ total: 1,\n+ offset: idx,\n+ }),\n+ };\n+ }\n+\n+ let all_ids: Vec = reduced\n+ .ingests_by_scope_thread\n+ .get(&key)\n+ .map(|q| q.iter().rev().cloned().collect())\n+ .unwrap_or_default();\n+\n+ let filtered: Vec<(usize, _)> = all_ids\n+ .into_iter()\n+ .enumerate()\n+ .filter_map(|(idx, id)| reduced.ingests_by_id.get(&id).map(|ing| (idx, ing.clone())))\n+ .filter(|(_, ing)| since.map_or(true, |s| ing.ts >= s))\n+ .filter(|(_, ing)| before.map_or(true, |b| ing.ts < b))\n+ .filter(|(_, ing)| actor_prefix.is_empty() || ing.principal.to_lowercase().starts_with(actor_prefix))\n+ .collect();\n+\n+ let total = filtered.len();\n+ const MAX_BODY: usize = 2000;\n+ let posts: Vec = filtered\n+ .into_iter()\n+ .skip(offset)\n+ .take(limit)\n+ .map(|(idx, ing)| {\n+ let (body, truncated) = if ing.raw.len() > MAX_BODY {\n+ (ing.raw[..MAX_BODY].to_string(), true)\n+ } else {\n+ (ing.raw.clone(), false)\n+ };\n+ PostRow {\n+ id: ing.id.clone(),\n+ index: idx,\n+ ts: ing.ts,\n+ actor: ing.principal.clone(),\n+ body,\n+ truncated,\n+ }\n+ })\n+ .collect();\n+\n+ Ok(ThreadDetailResponse {\n+ thread: format!(\"#{}\", tag),\n+ posts,\n+ total,\n+ offset,\n+ })\n+}\n+\n+fn tokenize_query(q: &str) -> Vec {\n+ q.split_whitespace()\n+ .map(|w| w.to_lowercase())\n+ .filter(|w| !w.is_empty())\n+ .collect()\n+}\n+\n+fn text_contains_all(text: &str, words: &[String]) -> bool {\n+ let lower = text.to_lowercase();\n+ words.iter().all(|w| lower.contains(w.as_str()))\n+}\n+\n+fn text_contains_any(text: &str, words: &[String]) -> usize {\n+ let lower = text.to_lowercase();\n+ words.iter().filter(|w| lower.contains(w.as_str())).count()\n+}\n+\n+fn snippet_around(text: &str, words: &[String], max_len: usize) -> String {\n+ let lower = text.to_lowercase();\n+ let first_pos = words.iter().filter_map(|w| lower.find(w.as_str())).min().unwrap_or(0);\n+ let start = first_pos.saturating_sub(max_len / 3);\n+ let start = if start > 0 {\n+ let mut i = start;\n+ while i < text.len() && !text.is_char_boundary(i) { i += 1; }\n+ text[i..].find(' ').map(|j| i + j + 1).unwrap_or(i)\n+ } else {\n+ 0\n+ };\n+ let mut end = (start + max_len).min(text.len());\n+ while end < text.len() && !text.is_char_boundary(end) { end += 1; }\n+ text[start..end].to_string()\n+}\n+\n+fn rpc_search(reduced: &ReducerState, q: &str, limit: usize) -> SearchResponse {\n+ let words = tokenize_query(q);\n+ if words.is_empty() {\n+ return SearchResponse {\n+ items: vec![],\n+ threads: vec![],\n+ posts: vec![],\n+ };\n+ }\n+ let content = reduced.public();\n+ let mut scored_items: Vec<(u32, SearchItemHit)> = Vec::new();\n+ for item in &content.items {\n+ let mut score: u32 = 0;\n+ if text_contains_all(item.as_str(), &words) { score += 10; }\n+ else if text_contains_any(item.as_str(), &words) > 0 { score += 5; }\n+ if let Some(body) = content.item_bodies.get(item) {\n+ if text_contains_all(body, &words) { score += 6; }\n+ else {\n+ let any = text_contains_any(body, &words);\n+ if any > 0 { score += any as u32; }\n+ }\n+ }\n+ if score > 0 {\n+ scored_items.push((score, SearchItemHit {\n+ path: item_path_for_api(item.as_str()),\n+ body: content.item_bodies.get(item).map(|b| snippet_around(b, &words, 120)),\n+ }));\n+ }\n+ }\n+ let mut scored_threads: Vec<(u32, i64, SearchThreadHit)> = Vec::new();\n+ for ((scope, tag), ts) in &reduced.forum_threads {\n+ if scope != &ScopeId::Public { continue; }\n+ let mut score: u32 = 0;\n+ if text_contains_all(tag, &words) { score += 8; }\n+ else if text_contains_any(tag, &words) > 0 { score += 4; }\n+ if score > 0 {\n+ let post_count = reduced\n+ .ingests_by_scope_thread\n+ .get(&(ScopeId::Public, tag.clone()))\n+ .map(|q| q.len())\n+ .unwrap_or(0);\n+ scored_threads.push((score, ts.last_activity_ts, SearchThreadHit {\n+ tag: format!(\"#{tag}\"),\n+ post_count,\n+ last_activity: ts.last_activity_ts,\n+ }));\n+ }\n+ }\n+ let mut scored_posts: Vec<(u32, i64, SearchPostHit)> = Vec::new();\n+ for (id, ingest) in &reduced.ingests_by_id {\n+ let mut score: u32 = 0;\n+ if text_contains_all(&ingest.raw, &words) { score += 4; }\n+ else {\n+ let any = text_contains_any(&ingest.raw, &words);\n+ if any > 0 { score += any as u32; }\n+ }\n+ if score > 0 {\n+ let thread = reduced\n+ .ingests_by_scope_thread\n+ .iter()\n+ .find(|(_, ids)| ids.contains(id))\n+ .map(|((scope, tag), _)| match scope {\n+ ScopeId::Public => format!(\"#{tag}\"),\n+ ScopeId::Room(rid) => format!(\"{rid}/#{tag}\"),\n+ })\n+ .unwrap_or_else(|| \"#unknown\".to_string());\n+ scored_posts.push((score, ingest.ts, SearchPostHit {\n+ thread,\n+ actor: ingest.principal.clone(),\n+ snippet: snippet_around(&ingest.raw, &words, 160),\n+ ts: ingest.ts,\n+ }));\n+ }\n+ }\n+ scored_items.sort_by(|a, b| b.0.cmp(&a.0));\n+ scored_threads.sort_by(|a, b| b.0.cmp(&a.0).then_with(|| b.1.cmp(&a.1)));\n+ scored_posts.sort_by(|a, b| b.0.cmp(&a.0).then_with(|| b.1.cmp(&a.1)));\n+ scored_items.truncate(limit);\n+ scored_threads.truncate(limit);\n+ scored_posts.truncate(limit);\n+ SearchResponse {\n+ items: scored_items.into_iter().map(|(_, h)| h).collect(),\n+ threads: scored_threads.into_iter().map(|(_, _, h)| h).collect(),\n+ posts: scored_posts.into_iter().map(|(_, _, h)| h).collect(),\n+ }\n+}\n+\n+async fn rpc_get_pair(state: &AppState, room: String, parent_path: String) -> Result {\n+ let scope = scope_from_room_wire(&room);\n+ let reduced_arc = state.reduced.clone();\n+ let pool: Vec = {\n+ let reduced = reduced_arc.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let tmp = if parent_path.trim().is_empty() {\n+ None\n+ } else {\n+ Some(parent_path.clone())\n+ };\n+ let specs = parse_parent_specs(tmp.as_ref());\n+ let raw_pool: Vec = if specs.is_empty() {\n+ content.ranking_group.idx_to_item.clone()\n+ } else {\n+ crate::scope_rank::resolve_scope(content, &specs)\n+ };\n+ raw_pool.into_iter().map(|it| it.0).collect()\n+ };\n+ if pool.len() < 2 {\n+ return Err((\n+ format!(\"need at least 2 items under parent /{}\", parent_path.trim()),\n+ Some(\"add items via ingest\".into()),\n+ ));\n+ }\n+ let selected: Option<(String, String)> = {\n+ let mut reduced = reduced_arc.write().await;\n+ let content = reduced.content.entry(scope.clone()).or_default();\n+ let group = &mut content.ranking_group;\n+ if group.idx_to_item.is_empty() {\n+ pick_random_distinct(&pool)\n+ } else {\n+ let mut rng = rand::thread_rng();\n+ let idxs: Vec = pool.iter()\n+ .filter_map(|it| {\n+ let key = CanonicalItemUrl(it.clone());\n+ group.item_to_idx.get(&key).copied()\n+ })\n+ .collect();\n+ let ranked = ranked_items_subset(group, &idxs, 10000, 1e-8);\n+ let ranked_set: HashSet = ranked.iter().map(|r| r.item.as_str().to_string()).collect();\n+ let unsorted: Vec = pool.iter()\n+ .filter(|it| !ranked_set.contains(*it))\n+ .cloned()\n+ .collect();\n+ let mut pick: Option<(String, String)> = None;\n+ if !unsorted.is_empty() {\n+ if let Some(left) = unsorted.choose(&mut rng).cloned() {\n+ let mut candidates: Vec = if !ranked.is_empty() {\n+ ranked.iter().map(|r| r.item.as_str().to_string()).collect()\n+ } else {\n+ pool.clone()\n+ };\n+ candidates.retain(|c| c != &left);\n+ if let Some(right) = candidates.choose(&mut rng).cloned() {\n+ pick = Some((left, right));\n+ }\n+ }\n+ } else if ranked.len() >= 2 {\n+ for i in 0..(ranked.len().saturating_sub(1)) {\n+ let a = ranked[i].item.as_str();\n+ let b = ranked[i + 1].item.as_str();\n+ if a != b && !is_pair_voted(group, a, b) {\n+ pick = Some((a.to_string(), b.to_string()));\n+ break;\n+ }\n+ }\n+ if pick.is_none() {\n+ for _ in 0..64 {\n+ let (Some(a), Some(b)) = (pool.choose(&mut rng).cloned(), pool.choose(&mut rng).cloned()) else { break; };\n+ if a != b && !is_pair_voted(group, &a, &b) {\n+ pick = Some((a, b));\n+ break;\n+ }\n+ }\n+ }\n+ }\n+ pick.or_else(|| pick_random_distinct(&pool))\n+ }\n+ };\n+ let Some((left, right)) = selected else {\n+ return Err((\"need at least 2 items\".into(), None));\n+ };\n+ let reduced = reduced_arc.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let left_key = CanonicalItemUrl(left.clone());\n+ let right_key = CanonicalItemUrl(right.clone());\n+ let lb = content.item_bodies.get(&left_key).cloned();\n+ let rb = content.item_bodies.get(&right_key).cloned();\n+ let th: Vec = content\n+ .item_threads\n+ .get(&left_key)\n+ .into_iter()\n+ .chain(content.item_threads.get(&right_key))\n+ .flat_map(|s| s.iter().cloned())\n+ .collect::>()\n+ .into_iter()\n+ .collect();\n+ let cs = compute_connectivity_stats(&content.ranking_group, &pool);\n+ Ok(RpcResult::Pair(PairResponse {\n+ left: item_path_for_api(&left),\n+ right: item_path_for_api(&right),\n+ left_body: lb,\n+ right_body: rb,\n+ threads: th,\n+ connectivity: Some(cs),\n+ }))\n+}\n+\n+pub async fn handle_rpc_batch(\n+ State(state): State,\n+ headers: HeaderMap,\n+ Json(RpcBatch(commands)): Json,\n+) -> impl IntoResponse {\n+ let mut results = Vec::with_capacity(commands.len());\n+\n+ for cmd in commands {\n+ let line = match cmd {\n+ RpcCommand::Post {\n+ room,\n+ thread_tag,\n+ delegate,\n+ text,\n+ return_rank_diff,\n+ } => match rpc_post(&state, &headers, room, thread_tag, delegate, text, return_rank_diff).await {\n+ Ok(r) => line_ok(r),\n+ Err((e, h)) => line_err(e, h),\n+ },\n+ RpcCommand::Check { room, text } => match rpc_check(&state, room, text).await {\n+ Ok(r) => line_ok(r),\n+ Err((e, h)) => line_err(e, h),\n+ },\n+ RpcCommand::GetGardenRank {\n+ room,\n+ parent_path,\n+ depth,\n+ offset,\n+ limit,\n+ percent,\n+ } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let popt = if parent_path.trim().is_empty() {\n+ None\n+ } else {\n+ Some(parent_path.as_str())\n+ };\n+ match build_rank_response_for_content(\n+ content,\n+ popt,\n+ depth.unwrap_or(1),\n+ offset.unwrap_or(0),\n+ limit,\n+ percent.unwrap_or(false),\n+ ) {\n+ Ok(r) => line_ok(RpcResult::GardenRank(r)),\n+ Err((e, h)) => line_err(e, h),\n+ }\n+ }\n+ RpcCommand::GetGardenItem {\n+ room,\n+ item_path,\n+ full,\n+ } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let item_str = canonicalize_item(&item_path);\n+ let item = CanonicalItemUrl(item_str.clone());\n+ if !content.items.contains(&item) {\n+ line_err(\n+ \"item not found\",\n+ Some(format!(\"{} does not exist\", item_path_for_api(&item_str))),\n+ )\n+ } else {\n+ const MAX_ITEM_BODY: usize = 10_000;\n+ let want_full = full.unwrap_or(false);\n+ let (body, truncated, body_len) = match content.item_bodies.get(&item) {\n+ None => (None, false, 0),\n+ Some(raw) => {\n+ let char_len = raw.chars().count();\n+ if !want_full && char_len > MAX_ITEM_BODY {\n+ let byte_end = raw.char_indices().nth(MAX_ITEM_BODY).map(|(i, _)| i).unwrap_or(raw.len());\n+ (Some(raw[..byte_end].to_string()), true, char_len)\n+ } else {\n+ (Some(raw.clone()), false, 0)\n+ }\n+ }\n+ };\n+ let threads: Vec = content\n+ .item_threads\n+ .get(&item)\n+ .map(|s| s.iter().cloned().collect())\n+ .unwrap_or_default();\n+ line_ok(RpcResult::GardenItem(ItemResponse {\n+ item: item_str,\n+ body,\n+ truncated,\n+ body_len,\n+ threads,\n+ }))\n+ }\n+ }\n+ RpcCommand::GetForumThread {\n+ room,\n+ thread_tag,\n+ offset,\n+ limit,\n+ since,\n+ before,\n+ actor,\n+ post_id,\n+ } => {\n+ let reduced = state.reduced.read().await;\n+ let actor_prefix = match actor.as_deref().map(str::trim) {\n+ None | Some(\"\") => Ok(String::new()),\n+ Some(s) => parse_username(s).map_err(|msg| (\"invalid actor filter\".to_string(), Some(msg))),\n+ };\n+ match actor_prefix {\n+ Err((e, h)) => line_err(e, h),\n+ Ok(actor_prefix) => {\n+ let offset = offset.unwrap_or(0);\n+ let limit = limit.unwrap_or(10).clamp(1, 500);\n+ match rpc_forum_thread_detail(\n+ &reduced,\n+ &room,\n+ &thread_tag,\n+ offset,\n+ limit,\n+ since,\n+ before,\n+ &actor_prefix,\n+ post_id.as_deref(),\n+ ) {\n+ Ok(r) => line_ok(RpcResult::ForumThread(r)),\n+ Err((e, h)) => line_err(e, h),\n+ }\n+ }\n+ }\n+ }\n+ RpcCommand::ListForumThreads { room } => {\n+ let reduced = state.reduced.read().await;\n+ line_ok(RpcResult::ForumThreads(rpc_list_forum_threads(&reduced, &room)))\n+ }\n+ RpcCommand::RoomCreate { slug, visibility } => {\n+ match verify_bearer_principal(&headers, &*state.reduced.read().await) {\n+ Ok(principal) => {\n+ let slug = slug.trim().to_lowercase();\n+ if slug.is_empty() || slug.len() > 64 {\n+ line_err(\"slug must be 1-64 characters\", None)\n+ } else if !slug.chars().all(|c| c.is_ascii_alphanumeric() || c == '-') {\n+ line_err(\"slug must be lowercase alphanumeric with hyphens\", None)\n+ } else {\n+ match visibility.as_deref().unwrap_or(\"private\") {\n+ \"private\" | \"public\" => {\n+ let vis = if visibility.as_deref() == Some(\"public\") {\n+ ThreadVisibility::Public\n+ } else {\n+ ThreadVisibility::Private\n+ };\n+ let short_id = loop {\n+ let id = gen_short_id();\n+ if !state.reduced.read().await.rooms.contains_key(&format!(\"{id}/{slug}\")) {\n+ break id;\n+ }\n+ };\n+ let room_id = format!(\"{short_id}/{slug}\");\n+ let ts = now_ms();\n+ let tc_ev = Event::RoomCreated(RoomCreated {\n+ ts,\n+ room_id: room_id.clone(),\n+ slug: slug.clone(),\n+ owner: principal.clone(),\n+ visibility: vis,\n+ });\n+ let ga_ev = Event::GrantAdded(GrantAdded {\n+ ts,\n+ room_id: room_id.clone(),\n+ username: principal.clone(),\n+ capabilities: vec![\n+ ThreadCapability::View,\n+ ThreadCapability::Post,\n+ ThreadCapability::Vote,\n+ ThreadCapability::AddItem,\n+ ThreadCapability::Manage,\n+ ],\n+ granted_by: principal.clone(),\n+ });\n+ if let Err(e) = state.event_log.append(&tc_ev).await {\n+ line_err(format!(\"{e}\"), None)\n+ } else if let Err(e) = state.event_log.append(&ga_ev).await {\n+ line_err(format!(\"{e}\"), None)\n+ } else {\n+ let mut r = state.reduced.write().await;\n+ r.apply_event(tc_ev);\n+ r.apply_event(ga_ev);\n+ line_ok(RpcResult::RoomCreated { room_id })\n+ }\n+ }\n+ other => line_err(format!(\"unknown visibility: {other}\"), None),\n+ }\n+ }\n+ }\n+ Err((_, m)) => line_err(m, None),\n+ }\n+ }\n+ RpcCommand::RoomGrant {\n+ room,\n+ username,\n+ capability,\n+ } => {\n+ match verify_bearer_principal(&headers, &*state.reduced.read().await) {\n+ Err((_, m)) => line_err(m, None),\n+ Ok(principal) => {\n+ let reduced = state.reduced.read().await;\n+ if !reduced.user_has_cap(&room, &principal, ThreadCapability::Manage) {\n+ line_err(\"requires Manage capability\", None)\n+ } else {\n+ match parse_username(&username) {\n+ Err(msg) => line_err(\"invalid username\", Some(msg)),\n+ Ok(target) => {\n+ if !reduced.users_by_provider.values().any(|u| u == &target) {\n+ line_err(format!(\"user @{target} not found\"), None)\n+ } else {\n+ match parse_capability(&capability) {\n+ Err(msg) => line_err(msg, None),\n+ Ok(cap) => {\n+ let ga_ev = Event::GrantAdded(GrantAdded {\n+ ts: now_ms(),\n+ room_id: room,\n+ username: target,\n+ capabilities: vec![cap],\n+ granted_by: principal,\n+ });\n+ if let Err(e) = state.event_log.append(&ga_ev).await {\n+ line_err(format!(\"{e}\"), None)\n+ } else {\n+ let mut r = state.reduced.write().await;\n+ r.apply_event(ga_ev);\n+ line_ok(RpcResult::GrantOk {})\n+ }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ RpcCommand::GetGlobalRank {\n+ room,\n+ limit,\n+ offset,\n+ percent,\n+ } => {\n+ const DEFAULT_LIMIT: usize = 50;\n+ const MAX_LIMIT: usize = 500;\n+ let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);\n+ let offset = offset.unwrap_or(0);\n+ let want_percent = percent.unwrap_or(false);\n+ let mut reduced = state.reduced.write().await;\n+ let content = reduced.content.entry(scope_from_room_wire(&room)).or_default();\n+ let group = &mut content.ranking_group;\n+ let n = group.idx_to_item.len();\n+ let (mut comps, _) = connected_components_from_voted_pairs(\n+ n, group.voted_pairs.iter().copied(),\n+ );\n+ comps.sort_by(|a, b| b.len().cmp(&a.len()));\n+\n+ let mut ranked: Vec = Vec::new();\n+ for comp in &comps {\n+ let items = ranked_items_subset(group, comp, 10000, 1e-8);\n+ let top = items.first().map(|r| r.score).unwrap_or(1.0);\n+ let bot = items.last().map(|r| r.score).unwrap_or(0.0);\n+ let range = (top - bot).max(1e-12);\n+ for r in items {\n+ let pct = want_percent.then(|| ((r.score - bot) / range * 100.0).clamp(0.0, 100.0));\n+ ranked.push(RankRow { item: item_path_for_api(r.item.as_str()), score: r.score, percent: pct });\n+ }\n+ }\n+\n+ let ranked_total = ranked.len();\n+ let mut unranked: Vec = content\n+ .items\n+ .iter()\n+ .filter(|it| !group.item_to_idx.contains_key(*it))\n+ .map(|it| it.as_str().to_string())\n+ .collect();\n+ unranked.sort();\n+ let unranked_total = unranked.len();\n+\n+ let page: Vec = ranked\n+ .into_iter()\n+ .chain(unranked.into_iter().map(|it| RankRow {\n+ item: item_path_for_api(&it),\n+ score: 0.0,\n+ percent: want_percent.then_some(0.0),\n+ }))\n+ .skip(offset)\n+ .take(limit)\n+ .collect();\n+\n+ line_ok(RpcResult::GlobalRank(GlobalRankResponse {\n+ ranked_total,\n+ unranked_total,\n+ offset,\n+ limit,\n+ items: page,\n+ }))\n+ }\n+ RpcCommand::GetPair { room, parent_path } => match rpc_get_pair(&state, room, parent_path).await {\n+ Ok(r) => line_ok(r),\n+ Err((e, h)) => line_err(e, h),\n+ },\n+ RpcCommand::GetMatchup {\n+ room,\n+ item_path,\n+ limit,\n+ } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let item_str = canonicalize_item(&item_path);\n+ let item = CanonicalItemUrl(item_str.clone());\n+ let limit = limit.unwrap_or(50).clamp(1, 200);\n+ if !content.items.contains(&item) {\n+ line_err(\n+ \"item not found\",\n+ Some(format!(\"{} does not exist\", item_path_for_api(&item_str))),\n+ )\n+ } else {\n+ let votes: Vec = content\n+ .item_votes\n+ .get(&item)\n+ .map(|q| {\n+ q.iter()\n+ .take(limit)\n+ .map(|v| VoteRow {\n+ ts: v.ts,\n+ a: v.a.as_str().to_string(),\n+ b: v.b.as_str().to_string(),\n+ ratio: format!(\"{}:{}\", v.ratio_left, v.ratio_right),\n+ actor: Some(v.principal.clone()),\n+ body: v.body.clone(),\n+ thread: Some(v.thread_tag.clone()),\n+ })\n+ .collect()\n+ })\n+ .unwrap_or_default();\n+ line_ok(RpcResult::Matchup(MatchupResponse {\n+ item: item_path_for_api(&item_str),\n+ votes,\n+ }))\n+ }\n+ },\n+ RpcCommand::GetRankHistory { room, item_path } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let scope = scope_from_room_wire(&room);\n+ let item_str = canonicalize_item(&item_path);\n+ let item = CanonicalItemUrl(item_str.clone());\n+ let entries = content.rank_history.get(&item).cloned().unwrap_or_default();\n+ let history: Vec = entries.iter().map(|e| {\n+ let caused_by: Vec = reduced.ingests_by_id.get(&e.post_id)\n+ .and_then(|ing| crate::dsl::parse_full(&ing.raw).ok())\n+ .map(|doc| {\n+ doc.statements.into_iter().filter_map(|s| {\n+ if let crate::dsl::Stmt::Vote { item1, item2, ratio_left, ratio_right, explanation } = s {\n+ let a = canonicalize_item(&item1);\n+ let b = canonicalize_item(&item2);\n+ if a == item_str || b == item_str {\n+ Some(VoteRow {\n+ ts: e.ts,\n+ a: item_path_for_api(&a),\n+ b: item_path_for_api(&b),\n+ ratio: format!(\"{}:{}\", ratio_left, ratio_right),\n+ actor: reduced.ingests_by_id.get(&e.post_id).map(|ing| ing.principal.clone()),\n+ body: explanation,\n+ thread: Some(format!(\"#{}\", e.thread)),\n+ })\n+ } else {\n+ None\n+ }\n+ } else {\n+ None\n+ }\n+ }).collect()\n+ })\n+ .unwrap_or_default();\n+ let thread_post_index = reduced\n+ .ingests_by_scope_thread\n+ .get(&(scope.clone(), e.thread.clone()))\n+ .and_then(|q| q.iter().rev().position(|id| id == &e.post_id))\n+ .map(|i| i + 1)\n+ .unwrap_or(0);\n+ RankHistoryRow {\n+ ts: e.ts,\n+ scope_rank: e.scope_rank,\n+ scope_rank_delta: e.scope_rank_delta,\n+ scope_total: e.scope_total,\n+ global_rank: e.global_rank,\n+ global_rank_delta: e.global_rank_delta,\n+ global_total: e.global_total,\n+ score: e.score,\n+ thread: format!(\"#{}\", e.thread),\n+ thread_post_index,\n+ caused_by,\n+ }\n+ }).collect();\n+ line_ok(RpcResult::RankHistory(RankHistoryResponse {\n+ item: item_path_for_api(&item_str),\n+ history,\n+ }))\n+ },\n+ RpcCommand::GetLeaves { room } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let parents: HashSet<&str> = content.item_children.keys().map(|s| s.as_str()).collect();\n+ let mut paths: Vec = content\n+ .items\n+ .iter()\n+ .filter(|p| !parents.contains(p.as_str()))\n+ .map(|p| p.as_str().to_string())\n+ .collect();\n+ paths.sort();\n+ line_ok(RpcResult::Leaves(LeavesResponse { paths }))\n+ },\n+ RpcCommand::GetPaths { room } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let out: Vec = content\n+ .item_children\n+ .get(&CanonicalItemUrl::ontology_root())\n+ .map(|roots| {\n+ let mut v: Vec = roots.iter()\n+ .map(|path| {\n+ let children = content.item_children.get(path.as_str()).map(|s| s.len()).unwrap_or(0);\n+ PathSummary {\n+ path: format!(\"~/{}\", path),\n+ children,\n+ web: format!(\"https://slug.social/~/{}\", path),\n+ }\n+ }).collect();\n+ v.sort_by(|a, b| a.path.cmp(&b.path));\n+ v\n+ })\n+ .unwrap_or_default();\n+ line_ok(RpcResult::Paths(PathsResponse { paths: out }))\n+ },\n+ RpcCommand::GetRecentVotes {\n+ room,\n+ parent,\n+ limit,\n+ } => {\n+ let reduced = state.reduced.read().await;\n+ let content = content_for_room(&reduced, &room);\n+ let group = &content.ranking_group;\n+ let limit = limit.unwrap_or(25).clamp(1, 200);\n+ let iter = group.recent_votes.iter();\n+ let iter: Box> = if let Some(p) = &parent {\n+ let parent_can = canonicalize_item(p);\n+ Box::new(iter.filter(move |v| vote_touches_path(v.a.as_str(), v.b.as_str(), &parent_can)))\n+ } else {\n+ Box::new(iter)\n+ };\n+ let out: Vec = iter\n+ .take(limit)\n+ .map(|v| VoteRow {\n+ ts: v.ts,\n+ a: v.a.as_str().to_string(),\n+ b: v.b.as_str().to_string(),\n+ ratio: format!(\"{}:{}\", v.ratio_left, v.ratio_right),\n+ actor: Some(v.principal.clone()),\n+ body: v.body.clone(),\n+ thread: Some(v.thread_tag.clone()),\n+ }).collect();\n+ line_ok(RpcResult::RecentVotes(RecentVotesResponse { votes: out }))\n+ },\n+ RpcCommand::Search { query } => {\n+ let reduced = state.reduced.read().await;\n+ let limit = 50usize.min(200);\n+ line_ok(RpcResult::Search(rpc_search(&reduced, &query, limit)))\n+ },\n+ RpcCommand::GetFeed {\n+ actor,\n+ since,\n+ limit,\n+ } => {\n+ const DEFAULT_LIMIT: usize = 50;\n+ const MAX_LIMIT: usize = 200;\n+ match parse_username(&actor) {\n+ Err(msg) => line_err(\"invalid actor\", Some(msg)),\n+ Ok(actor_stored) => {\n+ let reduced = state.reduced.read().await;\n+ let since = since.or_else(|| reduced.actor_last_post_ts.get(&actor_stored).copied());\n+ let cutoff = since.unwrap_or(0);\n+ let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);\n+ let matching: Vec<&str> = reduced.ingests_ordered.iter().rev()\n+ .map(|id| id.as_str())\n+ .take_while(|id| reduced.ingests_by_id.get(*id).map_or(false, |ing| ing.ts > cutoff))\n+ .collect();\n+ let total = matching.len();\n+ let posts: Vec = matching.into_iter()\n+ .take(limit)\n+ .filter_map(|id| reduced.ingests_by_id.get(id))\n+ .map(|ing| {\n+ let scope = scope_from_room_wire(&ing.room_id);\n+ let thread_post_index = reduced\n+ .ingests_by_scope_thread\n+ .get(&(scope, ing.thread_tag.clone()))\n+ .and_then(|q| {\n+ q.iter().rev().position(|pid| pid == &ing.id).map(|i| i + 1)\n+ });\n+ FeedPost {\n+ ts: ing.ts,\n+ id: ing.id.clone(),\n+ thread: Some(ing.thread_tag.clone()),\n+ thread_post_index,\n+ body: ing.raw.clone(),\n+ }\n+ })\n+ .collect();\n+ line_ok(RpcResult::Feed(FeedResponse {\n+ actor: actor_stored,\n+ since,\n+ posts,\n+ total,\n+ }))\n+ }\n+ }\n+ },\n+ };\n+ results.push(line);\n+ }\n+\n+ Json(RpcBatchResponse { results }).into_response()\n+}\ndiff --git a/server/src/api/search.rs b/server/src/api/search.rs\ndeleted file mode 100644\nindex 87621e6e4210a28c7802bb941200788f28518940..0000000000000000000000000000000000000000\n--- a/server/src/api/search.rs\n+++ /dev/null\n@@ -1,142 +0,0 @@\n-use axum::{\n- extract::{Query, State},\n- response::IntoResponse,\n- Json,\n-};\n-use serde::Deserialize;\n-\n-use crate::state::AppState;\n-\n-use super::helpers::item_path_for_api;\n-\n-// ============================================================================\n-// Search\n-// ============================================================================\n-\n-#[derive(Debug, Deserialize)]\n-pub struct SearchApiQuery {\n- pub q: Option,\n- pub limit: Option,\n-}\n-\n-fn tokenize_query(q: &str) -> Vec {\n- q.split_whitespace()\n- .map(|w| w.to_lowercase())\n- .filter(|w| !w.is_empty())\n- .collect()\n-}\n-\n-fn text_contains_all(text: &str, words: &[String]) -> bool {\n- let lower = text.to_lowercase();\n- words.iter().all(|w| lower.contains(w.as_str()))\n-}\n-\n-fn text_contains_any(text: &str, words: &[String]) -> usize {\n- let lower = text.to_lowercase();\n- words.iter().filter(|w| lower.contains(w.as_str())).count()\n-}\n-\n-fn snippet_around(text: &str, words: &[String], max_len: usize) -> String {\n- let lower = text.to_lowercase();\n- let first_pos = words.iter().filter_map(|w| lower.find(w.as_str())).min().unwrap_or(0);\n- let start = first_pos.saturating_sub(max_len / 3);\n- let start = if start > 0 {\n- let mut i = start;\n- while i < text.len() && !text.is_char_boundary(i) { i += 1; }\n- text[i..].find(' ').map(|j| i + j + 1).unwrap_or(i)\n- } else {\n- 0\n- };\n- let mut end = (start + max_len).min(text.len());\n- while end < text.len() && !text.is_char_boundary(end) { end += 1; }\n- text[start..end].to_string()\n-}\n-\n-pub async fn get_search(\n- State(state): State,\n- Query(params): Query,\n-) -> impl IntoResponse {\n- let q = params.q.unwrap_or_default();\n- let limit = params.limit.unwrap_or(50).min(200);\n-\n- let words = tokenize_query(&q);\n- if words.is_empty() {\n- return Json(slug_types::SearchResponse {\n- items: vec![], threads: vec![], posts: vec![],\n- }).into_response();\n- }\n-\n- let reduced = state.reduced.read().await;\n- let content = reduced.public();\n-\n- let mut scored_items: Vec<(u32, slug_types::SearchItemHit)> = Vec::new();\n- for item in &content.items {\n- let mut score: u32 = 0;\n- if text_contains_all(item.as_str(), &words) { score += 10; }\n- else if text_contains_any(item.as_str(), &words) > 0 { score += 5; }\n- if let Some(body) = content.item_bodies.get(item) {\n- if text_contains_all(body, &words) { score += 6; }\n- else {\n- let any = text_contains_any(body, &words);\n- if any > 0 { score += any as u32; }\n- }\n- }\n- if score > 0 {\n- scored_items.push((score, slug_types::SearchItemHit {\n- path: item_path_for_api(item.as_str()),\n- body: content.item_bodies.get(item).map(|b| snippet_around(b, &words, 120)),\n- }));\n- }\n- }\n-\n- let mut scored_threads: Vec<(u32, i64, slug_types::SearchThreadHit)> = Vec::new();\n- for (tag, ts) in &reduced.threads {\n- let mut score: u32 = 0;\n- if text_contains_all(tag, &words) { score += 8; }\n- else if text_contains_any(tag, &words) > 0 { score += 4; }\n- if score > 0 {\n- let post_count = reduced.ingests_by_thread.get(tag).map(|q| q.len()).unwrap_or(0);\n- scored_threads.push((score, ts.last_activity_ts, slug_types::SearchThreadHit {\n- tag: format!(\"#{tag}\"),\n- post_count,\n- last_activity: ts.last_activity_ts,\n- }));\n- }\n- }\n-\n- let mut scored_posts: Vec<(u32, i64, slug_types::SearchPostHit)> = Vec::new();\n- for (id, ingest) in &reduced.ingests_by_id {\n- let mut score: u32 = 0;\n- if text_contains_all(&ingest.raw, &words) { score += 4; }\n- else {\n- let any = text_contains_any(&ingest.raw, &words);\n- if any > 0 { score += any as u32; }\n- }\n- if score > 0 {\n- let thread = reduced.ingests_by_thread.iter()\n- .find(|(_, ids)| ids.contains(id))\n- .map(|(tag, _)| format!(\"#{tag}\"))\n- .unwrap_or_else(|| \"#unknown\".to_string());\n- scored_posts.push((score, ingest.ts, slug_types::SearchPostHit {\n- thread,\n- actor: ingest.principal.clone(),\n- snippet: snippet_around(&ingest.raw, &words, 160),\n- ts: ingest.ts,\n- }));\n- }\n- }\n-\n- scored_items.sort_by(|a, b| b.0.cmp(&a.0));\n- scored_threads.sort_by(|a, b| b.0.cmp(&a.0).then_with(|| b.1.cmp(&a.1)));\n- scored_posts.sort_by(|a, b| b.0.cmp(&a.0).then_with(|| b.1.cmp(&a.1)));\n-\n- scored_items.truncate(limit);\n- scored_threads.truncate(limit);\n- scored_posts.truncate(limit);\n-\n- Json(slug_types::SearchResponse {\n- items: scored_items.into_iter().map(|(_, h)| h).collect(),\n- threads: scored_threads.into_iter().map(|(_, _, h)| h).collect(),\n- posts: scored_posts.into_iter().map(|(_, _, h)| h).collect(),\n- }).into_response()\n-}\ndiff --git a/server/src/api/thread.rs b/server/src/api/thread.rs\ndeleted file mode 100644\nindex be6ef446e316fef352a7eef5b35b40583dca4442..0000000000000000000000000000000000000000\n--- a/server/src/api/thread.rs\n+++ /dev/null\n@@ -1,189 +0,0 @@\n-use axum::{\n- extract::{Path, State},\n- http::{HeaderMap, StatusCode},\n- response::IntoResponse,\n- Json,\n-};\n-use serde::{Deserialize, Serialize};\n-\n-use crate::{\n- api::helpers::{api_error, now_ms},\n- events::{Event, GrantAdded, ThreadCapability, ThreadCreated, ThreadVisibility},\n- identity::parse_username,\n- state::AppState,\n-};\n-use super::auth::verify_bearer_principal;\n-\n-fn gen_short_id() -> String {\n- use rand::Rng;\n- const ALPHABET: &[u8] = b\"0123456789abcdefghijklmnopqrstuvwxyz\";\n- let mut rng = rand::thread_rng();\n- (0..7).map(|_| ALPHABET[rng.gen_range(0..ALPHABET.len())] as char).collect()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct CreateThreadRequest {\n- pub slug: String,\n- pub visibility: Option,\n-}\n-\n-#[derive(Debug, Serialize)]\n-pub struct CreateThreadResponse {\n- pub ok: bool,\n- pub thread_id: String,\n-}\n-\n-pub async fn post_create_thread(\n- State(state): State,\n- headers: HeaderMap,\n- Json(req): Json,\n-) -> impl IntoResponse {\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n-\n- let principal = match verify_bearer_principal(&headers, &reduced) {\n- Ok(u) => u,\n- Err((st, msg)) => return api_error(st, msg, None).into_response(),\n- };\n-\n- let slug = req.slug.trim().to_lowercase();\n- if slug.is_empty() || slug.len() > 64 {\n- return api_error(StatusCode::BAD_REQUEST, \"slug must be 1-64 characters\", None).into_response();\n- }\n- if !slug.chars().all(|c| c.is_ascii_alphanumeric() || c == '-') {\n- return api_error(StatusCode::BAD_REQUEST, \"slug must be lowercase alphanumeric with hyphens\", None).into_response();\n- }\n-\n- let visibility = match req.visibility.as_deref().unwrap_or(\"private\") {\n- \"private\" => ThreadVisibility::Private,\n- \"public\" => ThreadVisibility::Public,\n- other => return api_error(StatusCode::BAD_REQUEST, format!(\"unknown visibility: {other}\"), None).into_response(),\n- };\n-\n- // Generate a collision-free short id.\n- let short_id = loop {\n- let id = gen_short_id();\n- // Threads are keyed by full \"shortid/slug\" in the map.\n- if !reduced.threads.contains_key(&format!(\"{id}/{slug}\")) {\n- break id;\n- }\n- };\n- let thread_id = format!(\"{short_id}/{slug}\");\n-\n- drop(reduced);\n-\n- let ts = now_ms();\n-\n- let tc_ev = Event::ThreadCreated(ThreadCreated {\n- ts,\n- thread_id: thread_id.clone(),\n- slug: slug.clone(),\n- owner: principal.clone(),\n- visibility,\n- });\n-\n- // Owner gets all capabilities.\n- let ga_ev = Event::GrantAdded(GrantAdded {\n- ts,\n- thread_id: thread_id.clone(),\n- username: principal.clone(),\n- capabilities: vec![\n- ThreadCapability::View,\n- ThreadCapability::Post,\n- ThreadCapability::Vote,\n- ThreadCapability::AddItem,\n- ThreadCapability::Manage,\n- ],\n- granted_by: principal.clone(),\n- });\n-\n- if let Err(err) = state.event_log.append(&tc_ev).await {\n- return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!(\"{err}\"), None).into_response();\n- }\n- if let Err(err) = state.event_log.append(&ga_ev).await {\n- return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!(\"{err}\"), None).into_response();\n- }\n-\n- {\n- let mut reduced = reduced_arc.write().await;\n- reduced.apply_event(tc_ev);\n- reduced.apply_event(ga_ev);\n- }\n-\n- Json(CreateThreadResponse { ok: true, thread_id }).into_response()\n-}\n-\n-#[derive(Debug, Deserialize)]\n-pub struct AddGrantRequest {\n- pub username: String,\n- pub capabilities: Vec,\n-}\n-\n-fn parse_capability(s: &str) -> Result {\n- match s {\n- \"view\" => Ok(ThreadCapability::View),\n- \"post\" => Ok(ThreadCapability::Post),\n- \"vote\" => Ok(ThreadCapability::Vote),\n- \"add_item\" => Ok(ThreadCapability::AddItem),\n- \"manage\" => Ok(ThreadCapability::Manage),\n- other => Err(format!(\"unknown capability: {other}\")),\n- }\n-}\n-\n-pub async fn post_thread_grants(\n- Path((short_id, slug)): Path<(String, String)>,\n- State(state): State,\n- headers: HeaderMap,\n- Json(req): Json,\n-) -> impl IntoResponse {\n- let thread_id = format!(\"{short_id}/{slug}\");\n- let reduced_arc = state.reduced.clone();\n- let reduced = reduced_arc.read().await;\n-\n- let principal = match verify_bearer_principal(&headers, &reduced) {\n- Ok(u) => u,\n- Err((st, msg)) => return api_error(st, msg, None).into_response(),\n- };\n-\n- if !reduced.user_has_cap(&thread_id, &principal, ThreadCapability::Manage) {\n- return api_error(StatusCode::FORBIDDEN, \"requires Manage capability\", None).into_response();\n- }\n-\n- let target = match parse_username(&req.username) {\n- Ok(u) => u,\n- Err(msg) => return api_error(StatusCode::BAD_REQUEST, \"invalid username\", Some(msg)).into_response(),\n- };\n- if !reduced.users_by_provider.values().any(|u| u == &target) {\n- return api_error(StatusCode::NOT_FOUND, format!(\"user @{target} not found\"), None).into_response();\n- }\n-\n- if req.capabilities.is_empty() {\n- return api_error(StatusCode::BAD_REQUEST, \"capabilities must not be empty\", None).into_response();\n- }\n- let capabilities: Result, _> = req.capabilities.iter().map(|s| parse_capability(s)).collect();\n- let capabilities = match capabilities {\n- Ok(c) => c,\n- Err(msg) => return api_error(StatusCode::BAD_REQUEST, msg, None).into_response(),\n- };\n-\n- drop(reduced);\n-\n- let ga_ev = Event::GrantAdded(GrantAdded {\n- ts: now_ms(),\n- thread_id,\n- username: target,\n- capabilities,\n- granted_by: principal,\n- });\n-\n- if let Err(err) = state.event_log.append(&ga_ev).await {\n- return api_error(StatusCode::INTERNAL_SERVER_ERROR, format!(\"{err}\"), None).into_response();\n- }\n-\n- {\n- let mut reduced = reduced_arc.write().await;\n- reduced.apply_event(ga_ev);\n- }\n-\n- Json(serde_json::json!({\"ok\": true})).into_response()\n-}\ndiff --git a/server/src/api/validate.rs b/server/src/api/validate.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..27150577982f52cbd6d11bb93654dc3d4cf75cc5\n--- /dev/null\n+++ b/server/src/api/validate.rs\n@@ -0,0 +1,151 @@\n+use axum::http::StatusCode;\n+use std::collections::HashSet;\n+\n+use crate::{\n+ canonical_path::canonicalize_tag,\n+ dsl,\n+ path_types::CanonicalItemUrl,\n+ reducer::{ReducerState, ScopeId},\n+};\n+\n+use super::helpers::{item_path_for_api, resolve_item};\n+\n+#[derive(Debug)]\n+pub struct ValidatedIngest {\n+ pub doc: dsl::Document,\n+ pub ts: i64,\n+ pub raw_text: String,\n+}\n+\n+pub fn validate_ingest_document(\n+ reduced: &ReducerState,\n+ text: &str,\n+ scope: &ScopeId,\n+) -> Result)> {\n+ let public_content = reduced.public();\n+ let scoped_content = match scope {\n+ ScopeId::Public => None,\n+ _ => reduced.content_for_scope(scope),\n+ };\n+ let item_exists = |key: &CanonicalItemUrl| {\n+ scoped_content.map(|c| c.items.contains(key)).unwrap_or(false)\n+ || public_content.items.contains(key)\n+ };\n+ let body_exists = |key: &CanonicalItemUrl| {\n+ scoped_content.map(|c| c.item_bodies.contains_key(key)).unwrap_or(false)\n+ || public_content.item_bodies.contains_key(key)\n+ };\n+ let doc = match dsl::parse_full(text) {\n+ Ok(d) => d,\n+ Err(e) => {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ \"parse error\".to_string(),\n+ Some(format!(\"{}\", e)),\n+ ));\n+ }\n+ };\n+\n+ let ts = super::helpers::now_ms();\n+ let mut defined_in_doc: HashSet = HashSet::new();\n+\n+ for s in &doc.statements {\n+ match s {\n+ dsl::Stmt::Item { title, body } => {\n+ let item = match resolve_item(title) {\n+ Ok(v) => v,\n+ Err(msg) => {\n+ return Err((StatusCode::BAD_REQUEST, \"invalid item path\".to_string(), Some(msg)));\n+ }\n+ };\n+ let Some(body_text) = body else {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ format!(\"item missing body: {}\", item_path_for_api(&item)),\n+ Some(\"items must be declared with bodies, e.g. `~/path/item { ... }`\".to_string()),\n+ ));\n+ };\n+ if body_text.trim().is_empty() {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ format!(\"item body is empty: {}\", item_path_for_api(&item)),\n+ Some(\"write at least one sentence inside `{ ... }`\".to_string()),\n+ ));\n+ }\n+ defined_in_doc.insert(item);\n+ }\n+ dsl::Stmt::Vote { item1, item2, .. } => {\n+ let a = match resolve_item(item1) {\n+ Ok(v) => v,\n+ Err(msg) => {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ \"invalid vote item path\".to_string(),\n+ Some(msg),\n+ ));\n+ }\n+ };\n+ let b = match resolve_item(item2) {\n+ Ok(v) => v,\n+ Err(msg) => {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ \"invalid vote item path\".to_string(),\n+ Some(msg),\n+ ));\n+ }\n+ };\n+ let missing: Vec = [&a, &b]\n+ .into_iter()\n+ .filter(|it| {\n+ let key = CanonicalItemUrl((*it).clone());\n+ !defined_in_doc.contains(*it) && !item_exists(&key)\n+ })\n+ .map(|it| item_path_for_api(it))\n+ .collect();\n+ if !missing.is_empty() {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ \"vote references undefined item(s)\".to_string(),\n+ Some(format!(\n+ \"define items with bodies before voting. missing: {}\",\n+ missing.join(\", \")\n+ )),\n+ ));\n+ }\n+ let missing_body: Vec = [&a, &b]\n+ .into_iter()\n+ .filter(|it| {\n+ let key = CanonicalItemUrl((*it).clone());\n+ !defined_in_doc.contains(*it) && !body_exists(&key)\n+ })\n+ .map(|it| item_path_for_api(it))\n+ .collect();\n+ if !missing_body.is_empty() {\n+ return Err((\n+ StatusCode::BAD_REQUEST,\n+ \"vote references item(s) without bodies\".to_string(),\n+ Some(format!(\n+ \"missing bodies: {}. Declare each item as `/item {{ ... }}`\",\n+ missing_body.join(\", \")\n+ )),\n+ ));\n+ }\n+ }\n+ _ => {}\n+ }\n+ }\n+\n+ Ok(ValidatedIngest {\n+ doc,\n+ ts,\n+ raw_text: text.to_string(),\n+ })\n+}\n+\n+pub fn normalize_room_and_thread(room: &str, thread_tag: &str) -> (String, String) {\n+ (\n+ room.trim().to_string(),\n+ canonicalize_tag(thread_tag),\n+ )\n+}\ndiff --git a/server/src/events.rs b/server/src/events.rs\nindex 26edac7505ce4e67a042a5511113efba5fea4950..125005d89ade7b22c9b6d997f0ed81781df5e7ec 100644\n--- a/server/src/events.rs\n+++ b/server/src/events.rs\n@@ -23,7 +23,7 @@ pub enum Event {\n UserRegistered(UserRegistered),\n TokenIssued(TokenIssued),\n AgentBound(AgentBound),\n- ThreadCreated(ThreadCreated),\n+ RoomCreated(RoomCreated),\n GrantAdded(GrantAdded),\n GrantRevoked(GrantRevoked),\n /// Ingest of a DSL+prose body. Identity and routing live in event metadata.\n@@ -56,9 +56,9 @@ pub struct AgentBound {\n }\n \n #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n-pub struct ThreadCreated {\n+pub struct RoomCreated {\n pub ts: i64,\n- pub thread_id: String,\n+ pub room_id: String,\n pub slug: String,\n pub owner: String,\n pub visibility: ThreadVisibility,\n@@ -67,7 +67,7 @@ pub struct ThreadCreated {\n #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n pub struct GrantAdded {\n pub ts: i64,\n- pub thread_id: String,\n+ pub room_id: String,\n pub username: String,\n pub capabilities: Vec,\n pub granted_by: String,\n@@ -76,7 +76,7 @@ pub struct GrantAdded {\n #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n pub struct GrantRevoked {\n pub ts: i64,\n- pub thread_id: String,\n+ pub room_id: String,\n pub username: String,\n pub capabilities: Vec,\n pub revoked_by: String,\n@@ -96,8 +96,10 @@ pub struct Ingest {\n /// AI delegate id `uuid:rig:model` (wire and storage: no `@`). Omitted when absent.\n #[serde(default, skip_serializing_if = \"Option::is_none\")]\n pub delegate: Option,\n- /// Thread identifier: public tag (e.g. \"languages\") or private id/slug (e.g. \"a7f2k9x/project-review\").\n- pub thread_id: String,\n+ /// Permission boundary: `\"public\"` or private room id (`shortid/slug`).\n+ pub room_id: String,\n+ /// Forum channel within the room (e.g. `languages`).\n+ pub thread_tag: String,\n }\n \n fn generate_id() -> String {\ndiff --git a/server/src/html/editor.rs b/server/src/html/editor.rs\nindex 4f206a0a895b47592df7669cd8d2fb4d83286fe7..fae8f31adbf4f4aaf5920f08e2764892b3719eea 100644\n--- a/server/src/html/editor.rs\n+++ b/server/src/html/editor.rs\n@@ -111,7 +111,8 @@ pub async fn editor_check(\n raw: form.text.clone(),\n principal: String::new(),\n delegate: None,\n- thread_id: String::new(),\n+ room_id: \"public\".to_string(),\n+ thread_tag: String::new(),\n });\n let mut simulated = { reduced_arc.read().await.clone() };\n simulated.apply_event(event);\ndiff --git a/server/src/html/forum.rs b/server/src/html/forum.rs\nindex a4e17ddf0064b08dfa221c964ccaa0eb99948b5b..228a120065dbacd8a308ba6e8bfc6ab111877c4b 100644\n--- a/server/src/html/forum.rs\n+++ b/server/src/html/forum.rs\n@@ -8,7 +8,7 @@ use maud::{html, Markup};\n \n use crate::{\n canonical_path::canonicalize_tag,\n- reducer::ReducerState,\n+ reducer::{ReducerState, ScopeId},\n state::AppState,\n timeago,\n };\n@@ -30,10 +30,15 @@ struct ThreadRow {\n fn collect_thread_rows(reduced: &ReducerState, now: i64) -> Vec {\n let _ = now;\n reduced\n- .threads\n+ .forum_threads\n .iter()\n- .map(|(tag, thread)| {\n- let ingests = reduced.ingests_by_thread.get(tag).map(|q| q.len()).unwrap_or(0);\n+ .filter(|((scope, _), _)| scope == &ScopeId::Public)\n+ .map(|((_, tag), thread)| {\n+ let ingests = reduced\n+ .ingests_by_scope_thread\n+ .get(&(ScopeId::Public, tag.clone()))\n+ .map(|q| q.len())\n+ .unwrap_or(0);\n ThreadRow {\n tag: tag.clone(),\n subtitle: None,\n@@ -163,12 +168,12 @@ pub async fn thread_view(\n ) -> impl IntoResponse {\n let tag = canonicalize_tag(&tag);\n \n- // ingests_by_thread is newest-first; collect all IDs in chronological order.\n+ // Newest-first queue → chronological for the page.\n let all_ids: Vec = {\n let reduced = state.reduced.read().await;\n reduced\n- .ingests_by_thread\n- .get(&tag)\n+ .ingests_by_scope_thread\n+ .get(&(ScopeId::Public, tag.clone()))\n .map(|q| q.iter().rev().cloned().collect())\n .unwrap_or_default()\n };\n@@ -243,7 +248,7 @@ pub async fn thread_post_view(\n let index: usize = index_str.parse().unwrap_or(0);\n let (ing, subtitle) = {\n let reduced = state.reduced.read().await;\n- let ing = reduced.ingests_by_thread.get(&tag)\n+ let ing = reduced.ingests_by_scope_thread.get(&(ScopeId::Public, tag.clone()))\n .and_then(|q| q.iter().rev().nth(index))\n .and_then(|id| reduced.ingests_by_id.get(id).cloned());\n let subtitle: Option = None;\n@@ -288,8 +293,8 @@ pub async fn thread_post_expand(\n let ing = {\n let reduced = state.reduced.read().await;\n reduced\n- .ingests_by_thread\n- .get(&tag)\n+ .ingests_by_scope_thread\n+ .get(&(ScopeId::Public, tag.clone()))\n .and_then(|q| q.iter().rev().nth(index))\n .and_then(|id| reduced.ingests_by_id.get(id).cloned())\n };\ndiff --git a/server/src/html/garden.rs b/server/src/html/garden.rs\nindex dcd7f0d4a2d7b602026b643c564d212cb084ff9f..af9cc5be89a3a259c719247c3c83aa33804cdf65 100644\n--- a/server/src/html/garden.rs\n+++ b/server/src/html/garden.rs\n@@ -217,7 +217,7 @@ fn build_rank_history(\n .map(|ing| ing.principal.clone())\n .unwrap_or_default(),\n delegate: reduced.ingests_by_id.get(&e.post_id).and_then(|ing| ing.delegate.clone()),\n- thread_id: e.thread.clone(),\n+ thread_tag: e.thread.clone(),\n })\n } else { None }\n } else { None }\n@@ -225,8 +225,8 @@ fn build_rank_history(\n })\n .unwrap_or_default();\n \n- let thread_post_index = reduced.ingests_by_thread\n- .get(&e.thread)\n+ let thread_post_index = reduced.ingests_by_scope_thread\n+ .get(&(crate::reducer::ScopeId::Public, e.thread.clone()))\n .and_then(|q| q.iter().rev().position(|id| id == &e.post_id))\n .map(|i| i + 1)\n .unwrap_or(0);\n@@ -481,7 +481,8 @@ mod tests {\n raw: raw.to_string(),\n principal: \"testuser\".to_string(),\n delegate: Some(\"00000000-0000-0000-0000-000000000000:test:local/test\".to_string()),\n- thread_id: String::new(),\n+ room_id: \"public\".to_string(),\n+ thread_tag: String::new(),\n }));\n }\n \ndiff --git a/server/src/html/search.rs b/server/src/html/search.rs\nindex df5cb0a59ff72956f4bc94e53b86ab8337452817..9353cf229fbbd2cfbab5b89ba63ebe192042c91f 100644\n--- a/server/src/html/search.rs\n+++ b/server/src/html/search.rs\n@@ -6,7 +6,7 @@ use maud::{html, Markup, PreEscaped};\n use serde::Deserialize;\n \n use crate::{\n- reducer::ReducerState,\n+ reducer::{ReducerState, ScopeId},\n state::AppState,\n timeago,\n };\n@@ -113,8 +113,11 @@ fn search(state: &ReducerState, q: &str, limit: usize) -> SearchResults {\n }\n }\n \n- // Search threads\n- for (tag, thread_state) in &state.threads {\n+ // Search threads (public room only in HTML)\n+ for ((scope, tag), thread_state) in &state.forum_threads {\n+ if scope != &ScopeId::Public {\n+ continue;\n+ }\n let mut score: u32 = 0;\n if contains_all(tag, &words) {\n score += 8;\n@@ -122,7 +125,11 @@ fn search(state: &ReducerState, q: &str, limit: usize) -> SearchResults {\n score += 4;\n }\n if score > 0 {\n- let post_count = state.ingests_by_thread.get(tag).map(|q| q.len()).unwrap_or(0);\n+ let post_count = state\n+ .ingests_by_scope_thread\n+ .get(&(ScopeId::Public, tag.clone()))\n+ .map(|q| q.len())\n+ .unwrap_or(0);\n scored_threads.push((score, ThreadRow {\n tag: tag.clone(),\n subtitle: None,\n@@ -145,10 +152,13 @@ fn search(state: &ReducerState, q: &str, limit: usize) -> SearchResults {\n }\n if score > 0 {\n let thread = state\n- .ingests_by_thread\n+ .ingests_by_scope_thread\n .iter()\n .find(|(_, ids)| ids.contains(id))\n- .map(|(tag, _)| tag.clone())\n+ .map(|((scope, tag), _)| match scope {\n+ ScopeId::Public => format!(\"#{tag}\"),\n+ ScopeId::Room(rid) => format!(\"{rid}/#{tag}\"),\n+ })\n .unwrap_or_else(|| \"unknown\".to_string());\n scored_posts.push((score, PostRow {\n thread,\ndiff --git a/server/src/html/tree.rs b/server/src/html/tree.rs\nindex 339f69af3b3279918a97b1dced62e4e3e0b528e3..6aed7db08682f088df3d52d0358b8dc954e9e1df 100644\n--- a/server/src/html/tree.rs\n+++ b/server/src/html/tree.rs\n@@ -713,7 +713,8 @@ mod tests {\n raw: raw.to_string(),\n principal: \"tester\".to_string(),\n delegate: None,\n- thread_id: String::new(),\n+ room_id: \"public\".to_string(),\n+ thread_tag: String::new(),\n })\n }\n \ndiff --git a/server/src/lib.rs b/server/src/lib.rs\nindex 173fbdd6b3f23c062a200b268d04491d5e030b3f..d00b56d909e1fabd1d122d054dd02785e13e0bd8 100644\n--- a/server/src/lib.rs\n+++ b/server/src/lib.rs\n@@ -16,6 +16,7 @@ pub mod state;\n pub mod timeago;\n \n use axum::Router;\n+use axum::routing::post;\n use tower_http::trace::TraceLayer;\n \n use crate::state::AppState;\n@@ -34,29 +35,7 @@ pub fn create_app(state: AppState) -> Router {\n .route(\"/api/v0/pending-session\", axum::routing::post(api::post_pending_session))\n .route(\"/api/v0/pending-session/:id\", axum::routing::get(api::get_pending_session))\n .route(\"/api/v0/whoami\", axum::routing::get(api::get_whoami))\n- .route(\"/debug/query\", axum::routing::get(api::get_debug_query))\n- .route(\"/api/v0/paths\", axum::routing::get(api::get_paths))\n- .route(\"/api/v0/leaves\", axum::routing::get(api::get_leaves))\n- .route(\"/api/v0/threads\", axum::routing::get(api::get_threads))\n- .route(\"/api/v0/thread\", axum::routing::get(api::get_thread))\n- .route(\"/api/v0/item\", axum::routing::get(api::get_item))\n- .route(\"/api/v0/matchup\", axum::routing::get(api::get_matchup))\n- .route(\"/api/v0/recent_votes\", axum::routing::get(api::get_recent_votes))\n- // REMOVED: /api/v0/vote - use /api/v0/ingest instead\n- .route(\"/api/v0/thread\", axum::routing::post(api::post_create_thread))\n- .route(\"/api/v0/thread/:short_id/:slug/grants\", axum::routing::post(api::post_thread_grants))\n- .route(\"/api/v0/ingest\", axum::routing::post(api::post_ingest))\n- .route(\"/api/v0/check\", axum::routing::post(api::post_check))\n- .route(\"/api/v0/pair\", axum::routing::get(api::get_pair))\n- .route(\"/api/v0/rank\", axum::routing::get(api::get_rank))\n- .route(\"/api/v0/global-rank\", axum::routing::get(api::get_global_rank))\n- .route(\"/api/v0/rank-history\", axum::routing::get(api::get_rank_history))\n- .route(\"/api/v0/feed\", axum::routing::get(api::get_feed))\n- .route(\"/api/v0/search\", axum::routing::get(api::get_search))\n- // .route(\"/api/v0/stream\", axum::routing::get(api::get_stream))\n+ .route(\"/api/v0/rpc\", post(api::handle_rpc_batch))\n .with_state(state)\n- // view_count_middleware currently disabled while HTML views are offline\n .layer(TraceLayer::new_for_http())\n }\n-\n-\ndiff --git a/server/src/paths.rs b/server/src/paths.rs\nindex 23fb67cab17658dab655c30f1e601207a8263c24..85dd839db2d5d3f1924fbff5ae9d83d61c0690bd 100644\n--- a/server/src/paths.rs\n+++ b/server/src/paths.rs\n@@ -266,26 +266,21 @@ mod tests {\n }\n \n // -----------------------------------------------------------------------\n- // Mutation 7: self.ingests_by_thread.entry(thread).or_default().push_front(ing.id.clone())\n- //\n- // Original (reducer.rs line 419-420):\n- // let q = self.ingests_by_thread.entry(thread).or_default();\n- // q.push_front(ing.id.clone());\n- //\n- // Path version:\n- // nav!(self.ingests_by_thread, keypath(thread), push_front(ing_id));\n+ // Mutation 7: ingests_by_scope_thread (scope, thread_tag) push_front\n // -----------------------------------------------------------------------\n #[test]\n- fn mutation_7_ingests_by_thread() {\n- let mut ingests_by_thread: HashMap> = HashMap::new();\n+ fn mutation_7_ingests_by_scope_thread() {\n+ use crate::reducer::ScopeId;\n+ let mut ingests_by_scope_thread: HashMap<(ScopeId, String), VecDeque> = HashMap::new();\n let thread = \"debate-1\".to_string();\n let id1 = \"ing-1\".to_string();\n let id2 = \"ing-2\".to_string();\n+ let k = (ScopeId::Public, thread.clone());\n \n- nav!(ingests_by_thread, keypath(thread.clone()), push_front(id1));\n- nav!(ingests_by_thread, keypath(thread), push_front(id2));\n+ nav!(ingests_by_scope_thread, keypath(k.clone()), push_front(id1));\n+ nav!(ingests_by_scope_thread, keypath(k), push_front(id2));\n \n- let q = ingests_by_thread.get(\"debate-1\").unwrap();\n+ let q = ingests_by_scope_thread.get(&(ScopeId::Public, \"debate-1\".to_string())).unwrap();\n assert_eq!(q[0], \"ing-2\"); // most recent first\n assert_eq!(q[1], \"ing-1\");\n }\ndiff --git a/server/src/ranking.rs b/server/src/ranking.rs\nindex f70da0d076da2ba8f0abb5a4415babae1c2305f9..b24e80ddd5f34b4e454bb81b541e8f54bd97a541 100644\n--- a/server/src/ranking.rs\n+++ b/server/src/ranking.rs\n@@ -266,7 +266,7 @@ mod tests {\n body: \"because\".to_string(),\n principal: \"test\".to_string(),\n delegate: Some(\"00000000-0000-0000-0000-000000000000:test:local/test\".to_string()),\n- thread_id: \"untagged\".to_string(),\n+ thread_tag: \"untagged\".to_string(),\n }\n }\n \ndiff --git a/server/src/reducer.rs b/server/src/reducer.rs\nindex 550bcda1a1068df13908d19484d86ad46d65ddf1..118cb19f4df4a083df2878ca7a1424f3d63de851 100644\n--- a/server/src/reducer.rs\n+++ b/server/src/reducer.rs\n@@ -9,11 +9,20 @@ use crate::path_types::CanonicalItemUrl;\n #[derive(Debug, Clone, Hash, PartialEq, Eq, PartialOrd, Ord)]\n pub enum ScopeId {\n Public,\n- Thread(String),\n+ Room(String),\n+}\n+\n+/// Wire `room` field → content scope (`\"public\"` → [`ScopeId::Public`]).\n+pub fn scope_from_room_wire(room: &str) -> ScopeId {\n+ let r = room.trim();\n+ if r.is_empty() || r == \"public\" {\n+ ScopeId::Public\n+ } else {\n+ ScopeId::Room(r.to_string())\n+ }\n }\n \n /// Parsed vote data (internal representation).\n-// Note: NOT a wire-format struct — this is internal only. `a` and `b` are CanonicalItemUrl.\n #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n pub struct VoteData {\n pub ts: i64,\n@@ -22,15 +31,12 @@ pub struct VoteData {\n pub ratio_left: i32,\n pub ratio_right: i32,\n pub body: String,\n- /// Human principal username (no `@` in stored events).\n pub principal: String,\n- /// AI delegate id, if any (no `@` in stored events).\n pub delegate: Option,\n- /// Thread id where this vote was cast (public tag or private id/slug).\n- pub thread_id: String,\n+ /// Forum channel where this vote was cast (tag only, not room id).\n+ pub thread_tag: String,\n }\n \n-/// Single ranking group (one-ranking model). All votes contribute to this group.\n #[derive(Debug, Clone)]\n pub struct GroupState {\n pub item_to_idx: HashMap,\n@@ -41,10 +47,8 @@ pub struct GroupState {\n \n /// Unordered pairs that have at least one vote recorded between them (i,\n-\n pub dirty: bool,\n pub cached_scores: Vec,\n-\n pub recent_votes: VecDeque,\n }\n \n@@ -100,7 +104,7 @@ impl GroupState {\n pub fn apply_vote(&mut self, mut vote: VoteData) {\n vote.a = CanonicalItemUrl::parse(vote.a.as_str()).unwrap_or(vote.a);\n vote.b = CanonicalItemUrl::parse(vote.b.as_str()).unwrap_or(vote.b);\n- vote.thread_id = canonicalize_tag(&vote.thread_id);\n+ vote.thread_tag = canonicalize_tag(&vote.thread_tag);\n if vote.ratio_left < 0 {\n vote.ratio_left = 0;\n }\n@@ -147,18 +151,20 @@ pub struct RankHistoryEntry {\n pub post_id: String,\n }\n \n-/// First-class thread metadata.\n #[derive(Debug, Clone)]\n-pub struct ThreadState {\n- pub last_activity_ts: i64,\n+pub struct RoomState {\n pub visibility: crate::events::ThreadVisibility,\n }\n \n-impl Default for ThreadState {\n+#[derive(Debug, Clone)]\n+pub struct ForumThreadState {\n+ pub last_activity_ts: i64,\n+}\n+\n+impl Default for ForumThreadState {\n fn default() -> Self {\n Self {\n last_activity_ts: 0,\n- visibility: crate::events::ThreadVisibility::Public,\n }\n }\n }\n@@ -188,23 +194,21 @@ pub struct ReducerState {\n pub users_by_provider: HashMap<(String, String), String>,\n /// token_id -> (username, salt, token_hash)\n pub tokens_by_id: HashMap,\n- /// agent id (naked `uuid:rig:model`) -> username\n pub agent_bindings: HashMap,\n \n- pub ingests_by_id: HashMap,\n- /// Thread -> recent ingest ids (most recent first).\n- pub ingests_by_thread: HashMap>,\n-\n- /// First-class thread state: bump time, subscriber count.\n- pub threads: HashMap,\n \n- /// Last ingest timestamp (ms) per actor. Used by feed to default `since`.\n+ pub ingests_by_id: HashMap,\n+ /// (scope, thread_tag) → ingest ids, newest first.\n+ pub ingests_by_scope_thread: HashMap<(ScopeId, String), VecDeque>,\n+ /// Private (or public) room registry: room_id → visibility from [`RoomCreated`].\n+ pub rooms: HashMap,\n+ /// (scope, thread_tag) → last activity.\n+ pub forum_threads: HashMap<(ScopeId, String), ForumThreadState>,\n pub actor_last_post_ts: HashMap,\n \n /// All ingest IDs in chronological order (oldest first). Used by the feed endpoint.\n pub ingests_ordered: Vec,\n-\n- /// thread_id -> username -> set of capabilities\n+ /// room_id → username → capabilities\n pub grants: HashMap>>,\n }\n \n@@ -213,10 +217,9 @@ impl ReducerState {\n self.content.get(scope)\n }\n \n- /// Check whether `username` has `cap` in `thread_id`.\n- pub fn user_has_cap(&self, thread_id: &str, username: &str, cap: ThreadCapability) -> bool {\n+ pub fn user_has_cap(&self, room_id: &str, username: &str, cap: ThreadCapability) -> bool {\n self.grants\n- .get(thread_id)\n+ .get(room_id)\n .and_then(|t| t.get(username))\n .map(|caps| caps.contains(&cap))\n .unwrap_or(false)\n@@ -346,24 +349,23 @@ impl ReducerState {\n }\n self.agent_bindings.insert(ab.agent, ab.username);\n }\n- Event::ThreadCreated(tc) => {\n- self.threads.entry(tc.thread_id.clone()).or_default().visibility = tc.visibility;\n+ Event::RoomCreated(rc) => {\n+ self.rooms.insert(\n+ rc.room_id.clone(),\n+ RoomState {\n+ visibility: rc.visibility,\n+ },\n+ );\n }\n Event::Ingest(mut ing) => {\n- ing.thread_id = canonicalize_tag(&ing.thread_id);\n+ ing.thread_tag = canonicalize_tag(&ing.thread_tag);\n+ let room_key = ing.room_id.trim().to_string();\n+ let scope = scope_from_room_wire(&room_key);\n+ let canonical_thread = ing.thread_tag.clone();\n+ let scope_thread_key = (scope.clone(), canonical_thread.clone());\n \n self.ingests_by_id.insert(ing.id.clone(), ing.clone());\n \n- let visibility = self\n- .threads\n- .entry(ing.thread_id.clone())\n- .or_default()\n- .visibility;\n- let scope = match visibility {\n- crate::events::ThreadVisibility::Public => ScopeId::Public,\n- crate::events::ThreadVisibility::Private => ScopeId::Thread(ing.thread_id.clone()),\n- };\n-\n let doc = match crate::dsl::parse_full(&ing.raw) {\n Ok(d) => d,\n Err(e) => {\n@@ -386,10 +388,9 @@ impl ReducerState {\n \n let principal = ing.principal.clone();\n let delegate = ing.delegate.clone();\n- let canonical_thread: String = ing.thread_id.clone();\n \n {\n- let content = self.content_for_scope_mut(scope);\n+ let content = self.content_for_scope_mut(scope.clone());\n \n // Snapshot ranks before votes are applied (force-compute if stale).\n let before: HashMap = if !voted_items.is_empty() {\n@@ -446,7 +447,7 @@ impl ReducerState {\n body: explanation,\n principal: principal.clone(),\n delegate: delegate.clone(),\n- thread_id: canonical_thread.clone(),\n+ thread_tag: canonical_thread.clone(),\n };\n \n ingest_items.insert(item_a.clone());\n@@ -513,12 +514,11 @@ impl ReducerState {\n }\n }\n \n- // Bump thread state for the canonical thread.\n- let ts = self.threads.entry(canonical_thread.clone()).or_default();\n- nav!(ts.last_activity_ts, selected(ing.ts > ts.last_activity_ts, setval(ing.ts)));\n+ let ft = self.forum_threads.entry(scope_thread_key.clone()).or_default();\n+ let prev_ts = ft.last_activity_ts;\n+ nav!(ft.last_activity_ts, selected(ing.ts > prev_ts, setval(ing.ts)));\n \n- // Index ingest by canonical thread.\n- nav!(self.ingests_by_thread, keypath(canonical_thread.clone()), push_front(ing.id.clone()));\n+ nav!(self.ingests_by_scope_thread, keypath(scope_thread_key), push_front(ing.id.clone()));\n \n nav!(self.ingests_ordered, push_back(ing.id.clone()));\n \n@@ -526,7 +526,7 @@ impl ReducerState {\n }\n Event::GrantAdded(ga) => {\n let caps = self.grants\n- .entry(ga.thread_id)\n+ .entry(ga.room_id)\n .or_default()\n .entry(ga.username)\n .or_default();\n@@ -535,18 +535,18 @@ impl ReducerState {\n }\n }\n Event::GrantRevoked(gr) => {\n- if let Some(thread_grants) = self.grants.get_mut(&gr.thread_id) {\n+ if let Some(room_grants) = self.grants.get_mut(&gr.room_id) {\n let username = gr.username;\n- if let Some(caps) = thread_grants.get_mut(&username) {\n+ if let Some(caps) = room_grants.get_mut(&username) {\n for cap in &gr.capabilities {\n caps.remove(cap);\n }\n if caps.is_empty() {\n- thread_grants.remove(&username);\n+ room_grants.remove(&username);\n }\n }\n- if thread_grants.is_empty() {\n- self.grants.remove(&gr.thread_id);\n+ if room_grants.is_empty() {\n+ self.grants.remove(&gr.room_id);\n }\n }\n }\n@@ -564,8 +564,9 @@ impl Default for ReducerState {\n tokens_by_id: HashMap::new(),\n agent_bindings: HashMap::new(),\n ingests_by_id: HashMap::new(),\n- ingests_by_thread: HashMap::new(),\n- threads: HashMap::new(),\n+ ingests_by_scope_thread: HashMap::new(),\n+ rooms: HashMap::new(),\n+ forum_threads: HashMap::new(),\n actor_last_post_ts: HashMap::new(),\n ingests_ordered: Vec::new(),\n grants: HashMap::new(),\ndiff --git a/server/tests/basic.rs b/server/tests/basic.rs\nindex 21cacd3049c9f5a307e75eefd5ebfc7bf0097b5d..b39bef181451088ca948b044b5ee2f7468b707f9 100644\n--- a/server/tests/basic.rs\n+++ b/server/tests/basic.rs\n@@ -3,7 +3,7 @@ use slugsocial_server::{\n canonical_path::{canonicalize_item, canonicalize_tag},\n events::{Event, Ingest},\n ranking::ranked_items,\n- reducer::{GroupState, ReducerState},\n+ reducer::{GroupState, ReducerState, ScopeId},\n };\n \n \n@@ -17,7 +17,8 @@ fn ingest_event(ts: i64, raw: &str) -> Event {\n raw: raw.to_string(),\n principal: \"test\".to_string(),\n delegate: Some(\"00000000-0000-0000-0000-000000000000:test:local/test\".to_string()),\n- thread_id: \"t\".to_string(),\n+ room_id: \"public\".to_string(),\n+ thread_tag: \"t\".to_string(),\n })\n }\n \n@@ -94,12 +95,12 @@ fn reducer_handles_item_and_body_from_ingest() {\n #[test]\n fn reducer_indexes_item_threads_and_vote_thread() {\n let mut state = ReducerState::default();\n- // Thread routing is metadata now (ingest.thread_id), not parsed from raw.\n+ // Thread routing is metadata (ingest.thread_tag), not parsed from raw.\n let mut ev = match ingest_event(1, \"~/sorts/insertion { O(n^2) }\\n~/sorts/mergesort { O(n log n) }\\n~/sorts/insertion 3:1 ~/sorts/mergesort { simpler for small n }\\n\") {\n Event::Ingest(i) => i,\n _ => unreachable!(),\n };\n- ev.thread_id = \"sorting-hat\".to_string();\n+ ev.thread_tag = \"sorting-hat\".to_string();\n state.apply_event(Event::Ingest(ev));\n \n let content = state.public();\n@@ -109,7 +110,7 @@ fn reducer_indexes_item_threads_and_vote_thread() {\n assert!(threads_for_mergesort.contains(\"sorting-hat\"));\n \n let vote = content.item_votes.get(\"https://slug.social/~/sorts/insertion\").unwrap().front().unwrap();\n- assert_eq!(vote.thread_id, \"sorting-hat\");\n+ assert_eq!(vote.thread_tag, \"sorting-hat\");\n }\n \n #[test]\n@@ -493,7 +494,7 @@ fn reducer_negative_ratio_clamped_to_zero() {\n body: \"negative\".to_string(),\n principal: \"test\".to_string(),\n delegate: Some(\"00000000-0000-0000-0000-000000000000:test:local/test\".to_string()),\n- thread_id: \"t\".to_string(),\n+ thread_tag: \"t\".to_string(),\n });\n assert_eq!(group.idx_to_item.len(), 2);\n // Both edges should exist (negatives clamped to 0, then 0:0 -> 1:1)\n@@ -664,22 +665,23 @@ fn test_actor_last_post_ts() {\n #[test]\n fn test_thread_timestamp_bump() {\n let mut state = ReducerState::default();\n+ let key = (ScopeId::Public, \"my-thread\".to_string());\n let mut ev1 = match ingest_event(100, \"~/t/a {a}\\n\") { Event::Ingest(i) => i, _ => unreachable!() };\n- ev1.thread_id = \"my-thread\".to_string();\n+ ev1.thread_tag = \"my-thread\".to_string();\n state.apply_event(Event::Ingest(ev1));\n- assert_eq!(state.threads.get(\"my-thread\").unwrap().last_activity_ts, 100);\n+ assert_eq!(state.forum_threads.get(&key).unwrap().last_activity_ts, 100);\n \n let mut ev2 = match ingest_event(200, \"~/t/b {b}\\n\") { Event::Ingest(i) => i, _ => unreachable!() };\n- ev2.thread_id = \"my-thread\".to_string();\n+ ev2.thread_tag = \"my-thread\".to_string();\n state.apply_event(Event::Ingest(ev2));\n- assert_eq!(state.threads.get(\"my-thread\").unwrap().last_activity_ts, 200);\n+ assert_eq!(state.forum_threads.get(&key).unwrap().last_activity_ts, 200);\n \n // Ingest with earlier timestamp should NOT regress\n let mut ev3 = match ingest_event(150, \"~/t/c {c}\\n\") { Event::Ingest(i) => i, _ => unreachable!() };\n- ev3.thread_id = \"my-thread\".to_string();\n+ ev3.thread_tag = \"my-thread\".to_string();\n state.apply_event(Event::Ingest(ev3));\n assert_eq!(\n- state.threads.get(\"my-thread\").unwrap().last_activity_ts,\n+ state.forum_threads.get(&key).unwrap().last_activity_ts,\n 200,\n \"thread timestamp should not regress to an earlier value\"\n );\n@@ -692,7 +694,7 @@ fn test_thread_id_is_used_for_votes_and_indexes() {\n Event::Ingest(i) => i,\n _ => unreachable!(),\n };\n- ev.thread_id = \"first\".to_string();\n+ ev.thread_tag = \"first\".to_string();\n state.apply_event(Event::Ingest(ev));\n \n let content = state.public();\n@@ -702,8 +704,8 @@ fn test_thread_id_is_used_for_votes_and_indexes() {\n .unwrap()\n .front()\n .unwrap();\n- assert_eq!(vote.thread_id, \"first\");\n- assert!(state.ingests_by_thread.contains_key(\"first\"));\n+ assert_eq!(vote.thread_tag, \"first\");\n+ assert!(state.ingests_by_scope_thread.contains_key(&(ScopeId::Public, \"first\".to_string())));\n assert!(state\n .public()\n .item_threads\n@@ -772,10 +774,13 @@ fn test_ingests_by_thread_ordering() {\n let mut state = ReducerState::default();\n for ts in [100, 200, 300] {\n let mut ev = match ingest_event(ts, \"~/t/a {a}\\n\") { Event::Ingest(i) => i, _ => unreachable!() };\n- ev.thread_id = \"order-thread\".to_string();\n+ ev.thread_tag = \"order-thread\".to_string();\n state.apply_event(Event::Ingest(ev));\n }\n- let thread_ingests = state.ingests_by_thread.get(\"order-thread\").unwrap();\n+ let thread_ingests = state\n+ .ingests_by_scope_thread\n+ .get(&(ScopeId::Public, \"order-thread\".to_string()))\n+ .unwrap();\n assert_eq!(thread_ingests.len(), 3);\n // Most recent first (push_front ordering)\n assert_eq!(thread_ingests[0], \"test-300\");\ndiff --git a/server/tests/debug_query_params.rs b/server/tests/debug_query_params.rs\ndeleted file mode 100644\nindex 42ec000a57716060a3d14ffb5adc121405f37714..0000000000000000000000000000000000000000\n--- a/server/tests/debug_query_params.rs\n+++ /dev/null\n@@ -1,49 +0,0 @@\n-use axum::{\n- body::Body,\n- http::{Request, StatusCode},\n-};\n-use tower::ServiceExt as _;\n-\n-use slugsocial_server::state::{AppConfig, AppState};\n-\n-#[tokio::test]\n-async fn debug_query_parses_repeated_open_array_and_decodes() {\n- let tmp = tempfile::tempdir().expect(\"tempdir\");\n- let data_dir = tmp.path().to_string_lossy().to_string();\n- let event_log_path = format!(\"{}/events.jsonl\", data_dir);\n-\n- let state = AppState::new(AppConfig {\n- data_dir,\n- event_log_path,\n- });\n-\n- let app = slugsocial_server::create_app(state);\n-\n- // ~%2Fmodels%2Fllms == \"~/models/llms\"\n- let uri = \"/debug/query?open=~%2Fmodels%2Fllms&open=~%2Fmusic%2Fjazz&selected=~%2Fmodels%2Fllms&q=hello+world\";\n-\n- let resp = app\n- .oneshot(\n- Request::builder()\n- .method(\"GET\")\n- .uri(uri)\n- .body(Body::empty())\n- .unwrap(),\n- )\n- .await\n- .unwrap();\n-\n- assert_eq!(resp.status(), StatusCode::OK);\n- let body = axum::body::to_bytes(resp.into_body(), usize::MAX)\n- .await\n- .unwrap();\n- let v: serde_json::Value = serde_json::from_slice(&body).unwrap();\n-\n- assert_eq!(\n- v[\"open\"],\n- serde_json::json!([\"~/models/llms\", \"~/music/jazz\"])\n- );\n- assert_eq!(v[\"selected\"], serde_json::json!(\"~/models/llms\"));\n- assert_eq!(v[\"q\"], serde_json::json!(\"hello world\"));\n-}\n-\ndiff --git a/server/tests/integration.rs b/server/tests/integration.rs\nindex eaa25f2c5cc21df3096d9b8f54c060c83207bfd6..9162bd5bee84035e3908bfb9c8e201b7878ca339 100644\n--- a/server/tests/integration.rs\n+++ b/server/tests/integration.rs\n@@ -21,6 +21,27 @@ fn test_bearer() -> String {\n format!(\"slug_{token_id}_{secret}\")\n }\n \n+/// `commands` is a JSON array of RPC commands (`RpcBatch` is a transparent `Vec`).\n+async fn rpc_batch(\n+ client: &reqwest::Client,\n+ addr: SocketAddr,\n+ bearer: Option<&str>,\n+ commands: serde_json::Value,\n+) -> serde_json::Value {\n+ let url = format!(\"http://{}/api/v0/rpc\", addr);\n+ let mut req = client.post(url).json(&commands);\n+ if let Some(b) = bearer {\n+ req = req.header(\"Authorization\", format!(\"Bearer {b}\"));\n+ }\n+ let response = req.send().await.unwrap();\n+ assert!(\n+ response.status().is_success(),\n+ \"rpc http {}\",\n+ response.status()\n+ );\n+ response.json().await.unwrap()\n+}\n+\n async fn seed_test_token(state: &AppState) {\n let token_id = \"testtok\";\n let secret = \"secret\";\n@@ -83,47 +104,41 @@ async fn test_index_page() {\n async fn test_ingest_actor_with_colons_is_detected_and_validated() {\n let (addr, _tmp, _log, _handle) = create_test_server().await;\n let client = reqwest::Client::new();\n+ let b = test_bearer();\n \n // Old archive style: agent includes colons but UUID is only a prefix (invalid).\n- // We should detect the agent line, then fail with \"invalid agent format\".\n- let ingest_payload = serde_json::json!({\n- \"delegate\": \"aec1e31c:claudecode:anthropic/claude-sonnet-4.5\",\n- \"thread\": \"t\",\n- \"text\": \"~/x {x}\\n\",\n- });\n-\n- let response = client\n- .post(&format!(\"http://{}/api/v0/ingest\", addr))\n- .json(&ingest_payload)\n- .send()\n- .await\n- .unwrap();\n-\n- assert_eq!(response.status(), reqwest::StatusCode::BAD_REQUEST);\n- let body: serde_json::Value = response.json().await.unwrap();\n- assert_eq!(body[\"ok\"], false);\n- assert_eq!(body[\"error\"], \"invalid delegate format\");\n- let hint = body[\"hint\"].as_str().unwrap_or_default();\n+ let bad_delegate = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"t\",\n+ \"delegate\": \"aec1e31c:claudecode:anthropic/claude-sonnet-4.5\",\n+ \"text\": \"~/x {x}\\n\",\n+ \"return_rank_diff\": false\n+ }\n+ }]);\n+ let body = rpc_batch(&client, addr, Some(&b), bad_delegate).await;\n+ let line = &body[\"results\"][0];\n+ assert_eq!(line[\"ok\"], false);\n+ assert_eq!(line[\"error\"], \"invalid delegate format\");\n+ let hint = line[\"hint\"].as_str().unwrap_or_default();\n assert!(\n hint.to_lowercase().contains(\"uuid\"),\n \"hint should mention uuid, got: {hint}\"\n );\n \n- let at_payload = serde_json::json!({\n- \"delegate\": \"@00000000-0000-0000-0000-000000000000:test:local/test\",\n- \"thread\": \"t\",\n- \"text\": \"~/x {x}\\n\",\n- });\n- let at_resp = client\n- .post(&format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {}\", test_bearer()))\n- .json(&at_payload)\n- .send()\n- .await\n- .unwrap();\n- assert_eq!(at_resp.status(), reqwest::StatusCode::BAD_REQUEST);\n- let at_body: serde_json::Value = at_resp.json().await.unwrap();\n- let at_hint = at_body[\"hint\"].as_str().unwrap_or_default();\n+ let at_batch = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"t\",\n+ \"delegate\": \"@00000000-0000-0000-0000-000000000000:test:local/test\",\n+ \"text\": \"~/x {x}\\n\",\n+ \"return_rank_diff\": false\n+ }\n+ }]);\n+ let at_body = rpc_batch(&client, addr, Some(&b), at_batch).await;\n+ let at_line = &at_body[\"results\"][0];\n+ assert_eq!(at_line[\"ok\"], false);\n+ let at_hint = at_line[\"hint\"].as_str().unwrap_or_default();\n assert!(\n at_hint.contains('@'),\n \"hint should reject '@' in delegate, got: {at_hint}\"\n@@ -135,25 +150,20 @@ async fn test_vote_endpoint() {\n let (addr, _tmp, _log, _handle) = create_test_server().await;\n let client = reqwest::Client::new();\n \n- // /api/v0/vote was removed; all votes are submitted via ingest.\n- let ingest_payload = serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n- \"thread\": \"cli\",\n- \"text\": \"~/clap {cli parser}\\n~/argh {cli parser}\\n~/clap 3:1 ~/argh {because clap is more full-featured}\\n\",\n- });\n-\n- let response = client\n- .post(format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {}\", test_bearer()))\n- .json(&ingest_payload)\n- .send()\n- .await\n- .unwrap();\n-\n- assert!(response.status().is_success());\n- let body: serde_json::Value = response.json().await.unwrap();\n- assert_eq!(body[\"ok\"], true);\n- assert!(body[\"next\"].is_object());\n+ let batch = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"cli\",\n+ \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n+ \"text\": \"~/clap {cli parser}\\n~/argh {cli parser}\\n~/clap 3:1 ~/argh {because clap is more full-featured}\\n\",\n+ \"return_rank_diff\": true\n+ }\n+ }]);\n+ let body = rpc_batch(&client, addr, Some(&test_bearer()), batch).await;\n+ let line = &body[\"results\"][0];\n+ assert_eq!(line[\"ok\"], true);\n+ let post = &line[\"result\"][\"PostOk\"];\n+ assert!(post[\"next\"].is_object());\n }\n \n #[tokio::test]\n@@ -161,37 +171,35 @@ async fn test_rank_endpoint() {\n let (addr, _tmp, _log, _handle) = create_test_server().await;\n let client = reqwest::Client::new();\n \n- // Ingest items + vote (vote endpoint removed).\n- let ingest_payload = serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n- \"thread\": \"langs\",\n- \"text\": \"~/rust {systems}\\n~/go {concurrency}\\n~/rust 3:1 ~/go {because i prefer rust for systems work}\\n\",\n- });\n- client\n- .post(format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {}\", test_bearer()))\n- .json(&ingest_payload)\n- .send()\n- .await\n- .unwrap();\n-\n- // Then query ranking (global: parent=~ — default empty parent only ranks direct children of \"\")\n- let response = client\n- .get(&format!(\"http://{}/api/v0/rank\", addr))\n- .query(&[(\"parent\", \"~\")])\n- .send()\n- .await\n- .unwrap();\n-\n- assert!(response.status().is_success());\n- let body: serde_json::Value = response.json().await.unwrap();\n- assert!(body[\"components\"].is_array());\n- let components = body[\"components\"].as_array().unwrap();\n+ let ingest = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"langs\",\n+ \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n+ \"text\": \"~/rust {systems}\\n~/go {concurrency}\\n~/rust 3:1 ~/go {because i prefer rust for systems work}\\n\",\n+ \"return_rank_diff\": false\n+ }\n+ }]);\n+ rpc_batch(&client, addr, Some(&test_bearer()), ingest).await;\n+\n+ let rank_batch = serde_json::json!([{\n+ \"GetGardenRank\": {\n+ \"room\": \"public\",\n+ \"parent_path\": \"~\",\n+ \"depth\": 1\n+ }\n+ }]);\n+ let body = rpc_batch(&client, addr, None, rank_batch).await;\n+ let line = &body[\"results\"][0];\n+ assert_eq!(line[\"ok\"], true);\n+ let rank = &line[\"result\"][\"GardenRank\"];\n+ assert!(rank[\"components\"].is_array());\n+ let components = rank[\"components\"].as_array().unwrap();\n assert_eq!(components.len(), 1);\n let ranking = components[0][\"ranking\"].as_array().unwrap();\n assert_eq!(ranking.len(), 2);\n assert_eq!(ranking[0][\"item\"], \"https://slug.social/~/rust\");\n- assert!(body[\"unranked_items\"].is_array());\n+ assert!(rank[\"unranked_items\"].is_array());\n }\n \n #[tokio::test]\n@@ -199,32 +207,30 @@ async fn test_check_endpoint_does_not_commit() {\n let (addr, tmp, _log, _handle) = create_test_server().await;\n let client = reqwest::Client::new();\n \n- let check_payload = serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n- \"thread\": \"t\",\n- \"text\": \"~/a {x}\\n~/b {y}\\n~/a 2:1 ~/b {because}\\n\",\n- });\n- let resp = client\n- .post(&format!(\"http://{}/api/v0/check\", addr))\n- .json(&check_payload)\n- .send()\n- .await\n- .unwrap();\n- assert!(resp.status().is_success());\n+ let check_batch = serde_json::json!([{\n+ \"Check\": {\n+ \"room\": \"public\",\n+ \"text\": \"~/a {x}\\n~/b {y}\\n~/a 2:1 ~/b {because}\\n\",\n+ }\n+ }]);\n+ let resp_body = rpc_batch(&client, addr, None, check_batch).await;\n+ let line = &resp_body[\"results\"][0];\n+ assert_eq!(line[\"ok\"], true);\n \n // It should not write events.jsonl.\n let log_path = tmp.path().join(\"events.jsonl\");\n let content = std::fs::read_to_string(&log_path).unwrap_or_default();\n assert!(content.trim().is_empty(), \"check must not append events\");\n \n- // And the live state should still have no threads (check is dry-run).\n- let threads_resp = client\n- .get(&format!(\"http://{}/api/v0/threads\", addr))\n- .send()\n- .await\n- .unwrap();\n- let body: serde_json::Value = threads_resp.json().await.unwrap();\n- assert!(body[\"threads\"].as_array().unwrap().is_empty());\n+ // And the live state should still have no forum threads (check is dry-run).\n+ let list_batch = serde_json::json!([{\n+ \"ListForumThreads\": { \"room\": \"public\" }\n+ }]);\n+ let threads_body = rpc_batch(&client, addr, None, list_batch).await;\n+ let tline = &threads_body[\"results\"][0];\n+ assert_eq!(tline[\"ok\"], true);\n+ let threads = tline[\"result\"][\"ForumThreads\"][\"threads\"].as_array().unwrap();\n+ assert!(threads.is_empty());\n }\n \n #[tokio::test]\n@@ -232,36 +238,38 @@ async fn test_ontology_item_page_shows_body_children_and_collapsible_votes() {\n // HTML routes are offline during the auth-v3 refactor.\n }\n \n-/// Thread connective tissue: item_threads and VoteData.thread are exposed by item, pair, and matchup APIs.\n+/// Thread connective tissue: item_threads and VoteData.thread_tag are exposed by item, pair, and matchup RPC.\n #[tokio::test]\n async fn test_garden_item_pair_matchup_include_threads() {\n let (addr, _tmp, _log, _handle) = create_test_server().await;\n let client = reqwest::Client::new();\n \n- // Ingest with thread_id metadata so item_threads and vote.thread_id are populated.\n- let ingest_payload = serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n- \"thread\": \"sorting-hat\",\n- \"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\",\n- });\n- let ingest_resp = client\n- .post(format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {}\", test_bearer()))\n- .json(&ingest_payload)\n- .send()\n- .await\n- .unwrap();\n- assert!(ingest_resp.status().is_success(), \"ingest should succeed\");\n-\n- // GET item: body + threads\n- let item_resp = client\n- .get(&format!(\"http://{}/api/v0/item\", addr))\n- .query(&[(\"item\", \"~/sorts/insertion\")])\n- .send()\n- .await\n- .unwrap();\n- assert!(item_resp.status().is_success());\n- let item: serde_json::Value = item_resp.json().await.unwrap();\n+ let ingest_batch = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"sorting-hat\",\n+ \"delegate\": \"00000000-0000-0000-0000-000000000000:test:local/test\",\n+ \"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\",\n+ \"return_rank_diff\": false\n+ }\n+ }]);\n+ let ing = rpc_batch(&client, addr, Some(&test_bearer()), ingest_batch).await;\n+ assert_eq!(ing[\"results\"][0][\"ok\"], true, \"ingest should succeed\");\n+\n+ let item_body = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetGardenItem\": {\n+ \"room\": \"public\",\n+ \"item_path\": \"~/sorts/insertion\",\n+ \"full\": true\n+ }\n+ }]),\n+ )\n+ .await;\n+ let item = &item_body[\"results\"][0][\"result\"][\"GardenItem\"];\n assert_eq!(item[\"item\"], \"https://slug.social/~/sorts/insertion\");\n assert!(item[\"body\"].as_str().unwrap().contains(\"O(n^2)\"));\n let threads: Vec<&str> = item[\"threads\"]\n@@ -272,15 +280,19 @@ async fn test_garden_item_pair_matchup_include_threads() {\n .collect();\n assert!(threads.contains(&\"sorting-hat\"), \"item threads should contain sorting-hat: {:?}\", threads);\n \n- // GET pair (under sorts): left, right, threads\n- let pair_resp = client\n- .get(&format!(\"http://{}/api/v0/pair\", addr))\n- .query(&[(\"parent\", \"~/sorts\")])\n- .send()\n- .await\n- .unwrap();\n- assert!(pair_resp.status().is_success());\n- let pair: serde_json::Value = pair_resp.json().await.unwrap();\n+ let pair_body = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetPair\": {\n+ \"room\": \"public\",\n+ \"parent_path\": \"~/sorts\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let pair = &pair_body[\"results\"][0][\"result\"][\"Pair\"];\n let pair_threads: Vec<&str> = pair[\"threads\"]\n .as_array()\n .unwrap()\n@@ -289,15 +301,19 @@ async fn test_garden_item_pair_matchup_include_threads() {\n .collect();\n assert!(pair_threads.contains(&\"sorting-hat\"), \"pair threads should contain sorting-hat: {:?}\", pair_threads);\n \n- // GET matchup: vote history with thread per vote\n- let matchup_resp = client\n- .get(&format!(\"http://{}/api/v0/matchup\", addr))\n- .query(&[(\"item\", \"~/sorts/insertion\")])\n- .send()\n- .await\n- .unwrap();\n- assert!(matchup_resp.status().is_success());\n- let matchup: serde_json::Value = matchup_resp.json().await.unwrap();\n+ let matchup_body = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetMatchup\": {\n+ \"room\": \"public\",\n+ \"item_path\": \"~/sorts/insertion\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let matchup = &matchup_body[\"results\"][0][\"result\"][\"Matchup\"];\n assert_eq!(matchup[\"item\"], \"https://slug.social/~/sorts/insertion\");\n let votes = matchup[\"votes\"].as_array().unwrap();\n assert!(!votes.is_empty(), \"matchup should have at least one vote\");\n@@ -326,42 +342,50 @@ async fn test_rank_history() {\n let client = reqwest::Client::new();\n \n let bearer = test_bearer();\n- let ingest = |payload: serde_json::Value| {\n+ let ingest = |delegate: &str, text: &str| {\n let client = client.clone();\n let addr = addr;\n let bearer = bearer.clone();\n+ let text = text.to_string();\n+ let delegate = delegate.to_string();\n async move {\n- client\n- .post(format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {bearer}\"))\n- .json(&payload)\n- .send()\n- .await\n- .unwrap()\n+ rpc_batch(\n+ &client,\n+ addr,\n+ Some(&bearer),\n+ serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"hist-test\",\n+ \"delegate\": delegate,\n+ \"text\": text,\n+ \"return_rank_diff\": false\n+ }\n+ }]),\n+ )\n+ .await\n }\n };\n \n- // First ingest: rust vs python — two votes on rust in one doc (the multi-vote case).\n ingest(\n- serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000001:rig:test/model\",\n- \"thread\": \"hist-test\",\n- \"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\",\n- }),\n+ \"00000000-0000-0000-0000-000000000001:rig:test/model\",\n+ \"~/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\",\n )\n .await;\n \n- // History for rust should have one entry with two caused_by votes.\n- let resp: serde_json::Value = client\n- .get(&format!(\"http://{}/api/v0/rank-history\", addr))\n- .query(&[(\"item\", \"~/hist/rust\")])\n- .send()\n- .await\n- .unwrap()\n- .json()\n- .await\n- .unwrap();\n-\n+ let hist_rust = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetRankHistory\": {\n+ \"room\": \"public\",\n+ \"item_path\": \"~/hist/rust\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let resp = &hist_rust[\"results\"][0][\"result\"][\"RankHistory\"];\n assert_eq!(resp[\"item\"], \"https://slug.social/~/hist/rust\");\n let history = resp[\"history\"].as_array().unwrap();\n assert_eq!(history.len(), 1, \"one ingest → one history entry\");\n@@ -371,45 +395,46 @@ async fn test_rank_history() {\n let caused_by = entry[\"caused_by\"].as_array().unwrap();\n assert_eq!(caused_by.len(), 2, \"both votes in the ingest touched rust\");\n \n- // Second ingest: python beats go — rust not directly touched, so python gets a new entry.\n ingest(\n- serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000002:rig:test/model\",\n- \"thread\": \"hist-test\",\n- \"text\": \"~/hist/python 3:1 ~/hist/go { dynamic typing is worth it }\\n\",\n- }),\n+ \"00000000-0000-0000-0000-000000000002:rig:test/model\",\n+ \"~/hist/python 3:1 ~/hist/go { dynamic typing is worth it }\\n\",\n )\n .await;\n \n- // Python history: two entries (appeared in first ingest, then voted again here).\n- let resp2: serde_json::Value = client\n- .get(&format!(\"http://{}/api/v0/rank-history\", addr))\n- .query(&[(\"item\", \"~/hist/python\")])\n- .send()\n- .await\n- .unwrap()\n- .json()\n- .await\n- .unwrap();\n+ let hist_py = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetRankHistory\": {\n+ \"room\": \"public\",\n+ \"item_path\": \"~/hist/python\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let resp2 = &hist_py[\"results\"][0][\"result\"][\"RankHistory\"];\n let hist2 = resp2[\"history\"].as_array().unwrap();\n assert_eq!(hist2.len(), 2, \"python touched in both ingests\");\n \n- // Second entry for python: caused_by has one vote (python 3:1 go).\n let caused_by2 = hist2[1][\"caused_by\"].as_array().unwrap();\n assert_eq!(caused_by2.len(), 1);\n assert!(caused_by2[0][\"a\"].as_str().unwrap().ends_with(\"python\") ||\n caused_by2[0][\"b\"].as_str().unwrap().ends_with(\"python\"));\n \n- // Rust was NOT directly voted on in the second ingest — no new history entry.\n- let resp3: serde_json::Value = client\n- .get(&format!(\"http://{}/api/v0/rank-history\", addr))\n- .query(&[(\"item\", \"~/hist/rust\")])\n- .send()\n- .await\n- .unwrap()\n- .json()\n- .await\n- .unwrap();\n+ let hist_rust2 = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetRankHistory\": {\n+ \"room\": \"public\",\n+ \"item_path\": \"~/hist/rust\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let resp3 = &hist_rust2[\"results\"][0][\"result\"][\"RankHistory\"];\n assert_eq!(\n resp3[\"history\"].as_array().unwrap().len(),\n 1,\n@@ -422,66 +447,71 @@ async fn pair_returns_connectivity_stats() {\n let (addr, _tmp, _log, _handle) = create_test_server().await;\n let client = reqwest::Client::new();\n \n- // Ingest 4 items with 1 vote (a vs b), leaving c and d as isolates.\n- let doc = serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000001:testrig:test/model\",\n- \"thread\": \"connectivity-test\",\n- \"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\",\n- });\n- let resp = client\n- .post(format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {}\", test_bearer()))\n- .json(&doc)\n- .send()\n- .await\n- .unwrap();\n- assert!(resp.status().is_success(), \"ingest should succeed\");\n-\n- // Request pair under ~/conn — should include connectivity stats.\n- let pair_resp = client\n- .get(&format!(\"http://{}/api/v0/pair\", addr))\n- .query(&[(\"parent\", \"~/conn\")])\n- .send()\n- .await\n- .unwrap();\n- assert!(pair_resp.status().is_success());\n- let pair: serde_json::Value = pair_resp.json().await.unwrap();\n+ let doc = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"connectivity-test\",\n+ \"delegate\": \"00000000-0000-0000-0000-000000000001:testrig:test/model\",\n+ \"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\",\n+ \"return_rank_diff\": false\n+ }\n+ }]);\n+ let resp = rpc_batch(&client, addr, Some(&test_bearer()), doc).await;\n+ assert_eq!(resp[\"results\"][0][\"ok\"], true, \"ingest should succeed\");\n+\n+ let pair_body = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetPair\": {\n+ \"room\": \"public\",\n+ \"parent_path\": \"~/conn\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let pair = &pair_body[\"results\"][0][\"result\"][\"Pair\"];\n \n let conn = &pair[\"connectivity\"];\n assert!(!conn.is_null(), \"pair response should include connectivity stats\");\n assert_eq!(conn[\"items\"].as_u64().unwrap(), 4, \"4 items in scope\");\n- // 1 component (a,b) + 2 isolates (c,d) = 3 components\n assert_eq!(conn[\"components\"].as_u64().unwrap(), 3, \"3 components (1 connected + 2 isolates)\");\n assert_eq!(conn[\"comparisons_until_connected\"].as_u64().unwrap(), 2, \"need 2 more comparisons\");\n assert_eq!(conn[\"pairs_voted\"].as_u64().unwrap(), 1, \"1 pair voted\");\n assert_eq!(conn[\"pairs_possible\"].as_u64().unwrap(), 6, \"4*3/2 = 6 possible pairs\");\n \n- // Add a vote connecting c to a — should reduce components.\n- let doc2 = serde_json::json!({\n- \"delegate\": \"00000000-0000-0000-0000-000000000001:testrig:test/model\",\n- \"thread\": \"connectivity-test\",\n- \"text\": \"~/conn/c 2:1 ~/conn/a { c beats a }\\n\",\n- });\n- let resp2 = client\n- .post(format!(\"http://{}/api/v0/ingest\", addr))\n- .header(\"Authorization\", format!(\"Bearer {}\", test_bearer()))\n- .json(&doc2)\n- .send()\n- .await\n- .unwrap();\n- let resp2_status = resp2.status();\n- let resp2_body: serde_json::Value = resp2.json().await.unwrap();\n- assert!(resp2_status.is_success(), \"second ingest failed: {}\", resp2_body);\n-\n- let pair_resp2 = client\n- .get(&format!(\"http://{}/api/v0/pair\", addr))\n- .query(&[(\"parent\", \"~/conn\")])\n- .send()\n- .await\n- .unwrap();\n- let pair2: serde_json::Value = pair_resp2.json().await.unwrap();\n+ let doc2 = serde_json::json!([{\n+ \"Post\": {\n+ \"room\": \"public\",\n+ \"thread_tag\": \"connectivity-test\",\n+ \"delegate\": \"00000000-0000-0000-0000-000000000001:testrig:test/model\",\n+ \"text\": \"~/conn/c 2:1 ~/conn/a { c beats a }\\n\",\n+ \"return_rank_diff\": false\n+ }\n+ }]);\n+ let resp2 = rpc_batch(&client, addr, Some(&test_bearer()), doc2).await;\n+ assert_eq!(\n+ resp2[\"results\"][0][\"ok\"],\n+ true,\n+ \"second ingest failed: {:?}\",\n+ resp2\n+ );\n+\n+ let pair_body2 = rpc_batch(\n+ &client,\n+ addr,\n+ None,\n+ serde_json::json!([{\n+ \"GetPair\": {\n+ \"room\": \"public\",\n+ \"parent_path\": \"~/conn\"\n+ }\n+ }]),\n+ )\n+ .await;\n+ let pair2 = &pair_body2[\"results\"][0][\"result\"][\"Pair\"];\n let conn2 = &pair2[\"connectivity\"];\n- // Now: component (a,b,c) + isolate (d) = 2 components\n assert_eq!(conn2[\"components\"].as_u64().unwrap(), 2, \"2 components after connecting c\");\n assert_eq!(conn2[\"comparisons_until_connected\"].as_u64().unwrap(), 1, \"1 more comparison to connect\");\n assert_eq!(conn2[\"pairs_voted\"].as_u64().unwrap(), 2, \"2 pairs voted now\");\ndiff --git a/test/grants.bb b/test/grants.bb\nindex f72799ae6fe099f48bdf316deee07f15882e9d45..7066c1f2a57f39a362052f8524206dccf37bb7b1 100644\n--- a/test/grants.bb\n+++ b/test/grants.bb\n@@ -1,12 +1,12 @@\n (ns test.grants\n- \"Grant enforcement integration test: private thread access control on POST /api/v0/ingest.\n+ \"Grant enforcement integration test: private room capabilities via POST /api/v0/rpc.\n \n Covers:\n- - user without any grant -> 403 on private thread\n- - user with View but no Post -> 403 for prose\n- - user with View + Post -> 200 for prose\n- - user with View + Post but no Vote -> 403 for vote (requires items in scope)\n- - user with View + Post + Vote -> 200 for vote\"\n+ - user without any grant -> RPC line ok=false on private room post\n+ - user with View but no Post -> ok=false for prose\n+ - user with View + Post -> ok=true for prose\n+ - user with View + Post but no Vote -> ok=false for vote\n+ - user with View + Post + Vote -> ok=true for vote\"\n (:require [babashka.process :as p]\n [babashka.fs :as fs]\n [cheshire.core :as json]\n@@ -141,10 +141,21 @@\n \n (defn- bearer [token] {\"Authorization\" (str \"Bearer \" token)})\n \n-(defn- ingest! [base-url token thread delegate text]\n- (http-post-json (str base-url \"/api/v0/ingest\")\n- {:thread thread :delegate delegate :text text}\n- :headers (bearer token)))\n+(defn- rpc-batch! [base-url token cmds]\n+ (let [resp (http-post-json (str base-url \"/api/v0/rpc\") cmds :headers (bearer token))]\n+ {:status (:status resp)\n+ :parsed (json/parse-string (:body resp) false)}))\n+\n+(defn- rpc-line-ok? [parsed]\n+ (true? (get-in parsed [\"results\" 0 \"ok\"])))\n+\n+(defn- ingest! [base-url token room thread delegate text]\n+ (rpc-batch! base-url token\n+ [{\"Post\" {\"room\" room\n+ \"thread_tag\" thread\n+ \"delegate\" delegate\n+ \"text\" text\n+ \"return_rank_diff\" false}}]))\n \n (defn grants-test [& _args]\n (println \"\\n━━━ grants enforcement integration check ━━━\\n\")\n@@ -194,80 +205,77 @@\n \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n \"bob\")\n \n- ;; Alice creates a private thread.\n- _ (println \"\\nalice creates private thread…\")\n- create-resp (http-post-json (str base-url \"/api/v0/thread\")\n- {:slug \"secret-project\" :visibility \"private\"}\n- :headers (bearer alice-token))\n- _ (assert! (= 200 (:status create-resp)) \"thread create returns 200\")\n- thread-id (-> create-resp :body (json/parse-string true) :thread_id)\n- _ (assert! (some? thread-id) \"thread_id present in response\")]\n-\n- ;; Alice (owner) can post prose to her own private thread.\n- (println \"\\nalice posts prose to her private thread…\")\n- (assert! (= 200 (:status (ingest! base-url alice-token thread-id\n- \"00000000-0000-0000-0000-000000000001:test:local/dev\"\n- \"Hello from alice.\")))\n+ ;; Alice creates a private room.\n+ _ (println \"\\nalice creates private room…\")\n+ create (rpc-batch! base-url alice-token\n+ [{\"RoomCreate\" {\"slug\" \"secret-project\" \"visibility\" \"private\"}}])\n+ _ (assert! (= 200 (:status create)) \"room create HTTP 200\")\n+ _ (assert! (rpc-line-ok? (:parsed create)) \"room create RPC ok\")\n+ room-id (get-in (:parsed create) [\"results\" 0 \"result\" \"RoomCreated\" \"room_id\"])\n+ _ (assert! (some? room-id) \"room_id present in RPC result\")]\n+\n+ ;; Alice (owner) can post prose to her own private room.\n+ (println \"\\nalice posts prose to her private room…\")\n+ (assert! (rpc-line-ok? (:parsed (ingest! base-url alice-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000001:test:local/dev\"\n+ \"Hello from alice.\")))\n \"alice prose post succeeds\")\n \n- ;; Bob has no grants at all — should get 403.\n+ ;; Bob has no grants — RPC line fails.\n (println \"\\nbob (no grants) tries to post prose…\")\n- (assert! (= 403 (:status (ingest! base-url bob-token thread-id\n- \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n- \"Hello from bob, unauthorized.\")))\n- \"bob without grants gets 403\")\n+ (assert! (not (rpc-line-ok? (:parsed (ingest! base-url bob-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n+ \"Hello from bob, unauthorized.\"))))\n+ \"bob without grants gets RPC failure\")\n \n- ;; Alice grants bob View only — still not enough to post prose.\n+ ;; Alice grants bob View only.\n (println \"\\nalice grants bob View only…\")\n- (assert! (= 200 (:status (http-post-json (str base-url \"/api/v0/thread/\" thread-id \"/grants\")\n- {:username \"bob\" :capabilities [\"view\"]}\n- :headers (bearer alice-token))))\n- \"grant View returns 200\")\n+ (assert! (rpc-line-ok? (:parsed (rpc-batch! base-url alice-token\n+ [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capability\" \"view\"}}])))\n+ \"grant View RPC ok\")\n \n (println \"\\nbob (View only) tries to post prose…\")\n- (assert! (= 403 (:status (ingest! base-url bob-token thread-id\n- \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n- \"Hello from bob, view only.\")))\n- \"bob with View but no Post gets 403\")\n+ (assert! (not (rpc-line-ok? (:parsed (ingest! base-url bob-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n+ \"Hello from bob, view only.\"))))\n+ \"bob with View but no Post gets RPC failure\")\n \n- ;; Alice grants bob Post — now prose should work.\n+ ;; Alice grants bob Post.\n (println \"\\nalice grants bob Post…\")\n- (assert! (= 200 (:status (http-post-json (str base-url \"/api/v0/thread/\" thread-id \"/grants\")\n- {:username \"bob\" :capabilities [\"post\"]}\n- :headers (bearer alice-token))))\n- \"grant Post returns 200\")\n+ (assert! (rpc-line-ok? (:parsed (rpc-batch! base-url alice-token\n+ [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capability\" \"post\"}}])))\n+ \"grant Post RPC ok\")\n \n (println \"\\nbob (View + Post) posts prose…\")\n- (assert! (= 200 (:status (ingest! base-url bob-token thread-id\n- \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n- \"Hello from bob, now authorised.\")))\n+ (assert! (rpc-line-ok? (:parsed (ingest! base-url bob-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n+ \"Hello from bob, now authorised.\")))\n \"bob with View + Post succeeds for prose\")\n \n- ;; Alice defines two items and votes on them in the private thread.\n- (println \"\\nalice posts items + vote to private thread…\")\n- (assert! (= 200 (:status (ingest! base-url alice-token thread-id\n- \"00000000-0000-0000-0000-000000000001:test:local/dev\"\n- \"~/fruits/apple { A crisp red apple. }\\n~/fruits/banana { A yellow banana. }\\n~/fruits/apple > ~/fruits/banana { apples are better }\")))\n- \"alice vote in private thread succeeds\")\n+ ;; Alice defines two items and votes in the private room.\n+ (println \"\\nalice posts items + vote to private room…\")\n+ (assert! (rpc-line-ok? (:parsed (ingest! base-url alice-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000001:test:local/dev\"\n+ \"~/fruits/apple { A crisp red apple. }\\n~/fruits/banana { A yellow banana. }\\n~/fruits/apple > ~/fruits/banana { apples are better }\")))\n+ \"alice vote in private room succeeds\")\n \n- ;; Bob (View + Post, no Vote) tries to vote — should be 403.\n+ ;; Bob (View + Post, no Vote) tries to vote.\n (println \"\\nbob (no Vote) tries to vote…\")\n- (assert! (= 403 (:status (ingest! base-url bob-token thread-id\n- \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n- \"~/fruits/apple > ~/fruits/banana { bob's take }\")))\n- \"bob without Vote gets 403\")\n+ (assert! (not (rpc-line-ok? (:parsed (ingest! base-url bob-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n+ \"~/fruits/apple > ~/fruits/banana { bob's take }\"))))\n+ \"bob without Vote gets RPC failure\")\n \n- ;; Alice grants bob Vote — now voting should work.\n+ ;; Alice grants bob Vote.\n (println \"\\nalice grants bob Vote…\")\n- (assert! (= 200 (:status (http-post-json (str base-url \"/api/v0/thread/\" thread-id \"/grants\")\n- {:username \"bob\" :capabilities [\"vote\"]}\n- :headers (bearer alice-token))))\n- \"grant Vote returns 200\")\n+ (assert! (rpc-line-ok? (:parsed (rpc-batch! base-url alice-token\n+ [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capability\" \"vote\"}}])))\n+ \"grant Vote RPC ok\")\n \n (println \"\\nbob (View + Post + Vote) votes…\")\n- (assert! (= 200 (:status (ingest! base-url bob-token thread-id\n- \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n- \"~/fruits/apple > ~/fruits/banana { bob's take }\")))\n+ (assert! (rpc-line-ok? (:parsed (ingest! base-url bob-token room-id \"main\"\n+ \"00000000-0000-0000-0000-000000000002:test:local/dev\"\n+ \"~/fruits/apple > ~/fruits/banana { bob's take }\")))\n \"bob with Vote succeeds\"))\n \n (finally\ndiff --git a/test/integration.bb b/test/integration.bb\nindex 4ff5d625840be3bf5e4162cef107a3cbd45aec4d..c71dcf20021056421ab27b1313e4dec94570ddaa 100644\n--- a/test/integration.bb\n+++ b/test/integration.bb\n@@ -146,7 +146,7 @@\n \n ;; 2.5 check endpoint — disconnected components not flattened (before OAuth; no events.jsonl yet)\n (println \"\\nchecking (dry-run) returns discrete ranking groups…\")\n- (bind check-result (common/run-cli cli-bin base-url [\"check\" \"--json\" \"--thread\" \"integration-test\"] :input check-doc-disconnected))\n+ (bind check-result (common/run-cli cli-bin base-url [\"public\" \"check\" \"--json\" \"--thread\" \"integration-test\"] :input check-doc-disconnected))\n (assert! (zero? (:exit check-result)) \"cli check exits 0\")\n (bind check-resp (json/parse-string (:out check-result) true))\n (assert! (:ok check-resp) \"check response ok=true\")\n@@ -166,7 +166,7 @@\n \n ;; 3. ingest via CLI (bearer required)\n (println \"\\ningesting .sorter document via CLI…\")\n- (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))\n+ (bind ingest1-result (common/run-cli cli-bin base-url [\"public\" \"ingest\" \"--json\" \"--thread\" \"integration-test\" \"--delegate\" \"00000000-0000-0000-0000-000000000000:cli:local/dev\"] :input sorter-doc :extra-env token-env))\n (assert! (zero? (:exit ingest1-result)) \"cli ingest exits 0\")\n (bind ingest1-resp (json/parse-string (:out ingest1-result) true))\n (assert! (:ok ingest1-resp) \"ingest response ok=true\")\n@@ -176,7 +176,7 @@\n \n ;; 4. query rankings via CLI\n (println \"\\nquerying garden children via CLI…\")\n- (bind children-result (common/run-cli cli-bin base-url [\"garden\" \"children\" \"languages\" \"--json\"]))\n+ (bind children-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"children\" \"languages\" \"--json\"]))\n (assert! (zero? (:exit children-result)) \"cli garden children exits 0\")\n (bind children-resp (json/parse-string (:out children-result) true))\n (bind ranked (mapv :item (mapcat :ranking (:components children-resp))))\n@@ -195,14 +195,14 @@\n \n ;; 6. query forum via CLI\n (println \"\\nquerying forum via CLI…\")\n- (bind forum-result (common/run-cli cli-bin base-url [\"forum\" \"--json\"]))\n+ (bind forum-result (common/run-cli cli-bin base-url [\"public\" \"forum\" \"--json\"]))\n (assert! (zero? (:exit forum-result)) \"cli forum exits 0\")\n (bind forum-resp (json/parse-string (:out forum-result) true))\n (assert! (some (fn [t] (= \"#integration-test\" (:thread t))) (:threads forum-resp))\n \"thread '#integration-test' visible in forum\")\n \n ;; 7. query item body via CLI\n- (bind body-result (common/run-cli cli-bin base-url [\"garden\" \"body\" \"languages/rust\" \"--json\"]))\n+ (bind body-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"body\" \"languages/rust\" \"--json\"]))\n (assert! (zero? (:exit body-result)) \"cli garden body exits 0\")\n (bind body-resp (json/parse-string (:out body-result) true))\n (assert! (str/includes? (or (:body body-resp) \"\") \"ownership\")\n@@ -210,7 +210,7 @@\n \n ;; 9. global rank endpoint\n (println \"\\ntesting global rank endpoint…\")\n- (bind grank-result (common/run-cli cli-bin base-url [\"garden\" \"rank\" \"--json\"]))\n+ (bind grank-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"rank\" \"--json\"]))\n (assert! (zero? (:exit grank-result)) \"cli garden rank exits 0\")\n (bind grank-resp (json/parse-string (:out grank-result) true))\n (assert! (pos? (:ranked_total grank-resp)) \"global rank has ranked items\")\n@@ -218,13 +218,13 @@\n (assert! (str/ends-with? (:item (first (:items grank-resp))) \"languages/rust\")\n (str \"rust is #1 globally (got \" (:item (first (:items grank-resp))) \")\"))\n \n- (bind grank-pct-result (common/run-cli cli-bin base-url [\"garden\" \"rank\" \"--percent\" \"--limit\" \"2\" \"--json\"]))\n+ (bind grank-pct-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"rank\" \"--percent\" \"--limit\" \"2\" \"--json\"]))\n (assert! (zero? (:exit grank-pct-result)) \"cli garden rank --percent --limit 2 exits 0\")\n (bind grank-pct-resp (json/parse-string (:out grank-pct-result) true))\n (assert! (= 2 (count (:items grank-pct-resp))) \"limit=2 returns 2 items\")\n (assert! (= 100.0 (:percent (first (:items grank-pct-resp)))) \"top item has 100.0% score\")\n \n- (bind grank-off-result (common/run-cli cli-bin base-url [\"garden\" \"rank\" \"--limit\" \"1\" \"--offset\" \"1\" \"--json\"]))\n+ (bind grank-off-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"rank\" \"--limit\" \"1\" \"--offset\" \"1\" \"--json\"]))\n (assert! (zero? (:exit grank-off-result)) \"cli garden rank --offset 1 exits 0\")\n (bind grank-off-resp (json/parse-string (:out grank-off-result) true))\n (assert! (= (:item (second (:items grank-pct-resp))) (:item (first (:items grank-off-resp))))\n@@ -237,7 +237,7 @@\n \"#integration-test\"\n \"~/languages/rust 4:1 ~/languages/python { type safety }\"\n \"~/languages/rust 3:1 ~/languages/go { zero-cost abstractions }\"]))\n- (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))\n+ (bind hist-ingest (common/run-cli cli-bin base-url [\"public\" \"ingest\" \"--json\" \"--thread\" \"integration-test\" \"--delegate\" \"00000000-0000-0000-0000-000000000000:cli:local/dev\"] :input two-vote-doc :extra-env token-env))\n (assert! (zero? (:exit hist-ingest))\n (str \"two-vote ingest exits 0 (err: \" (:err hist-ingest) \")\"))\n \n@@ -245,7 +245,7 @@\n (assert! (= 5 (count event-lines-after-hist))\n (str \"5 events in JSONL after second ingest (got \" (count event-lines-after-hist) \")\"))\n \n- (bind hist-result (common/run-cli cli-bin base-url [\"garden\" \"history\" \"languages/rust\" \"--json\"]))\n+ (bind hist-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"history\" \"languages/rust\" \"--json\"]))\n (assert! (zero? (:exit hist-result))\n (str \"cli garden history exits 0 (err: \" (:err hist-result) \")\"))\n (bind hist-resp (json/parse-string (:out hist-result) true))\n@@ -270,7 +270,7 @@\n \n ;; 12. same query, same result\n (println \"\\nquerying rankings after replay…\")\n- (bind replay-result (common/run-cli cli-bin base-url [\"garden\" \"children\" \"languages\" \"--json\"]))\n+ (bind replay-result (common/run-cli cli-bin base-url [\"public\" \"garden\" \"children\" \"languages\" \"--json\"]))\n (assert! (zero? (:exit replay-result)) \"cli garden children exits 0 after restart\")\n (bind replay-resp (json/parse-string (:out replay-result) true))\n (bind replay-ranked (mapv :item (mapcat :ranking (:components replay-resp))))\ndiff --git a/types/src/lib.rs b/types/src/lib.rs\nindex 92e6d122d573ea225f8fb86c548a81b3454425ad..a715fe5639f4a3a7bfedcc706b5fb312b98c9ac5 100644\n--- a/types/src/lib.rs\n+++ b/types/src/lib.rs\n@@ -100,12 +100,12 @@ pub struct PathSummary {\n pub web: String,\n }\n \n-#[derive(Debug, Serialize, Deserialize)]\n+#[derive(Debug, Clone, Serialize, Deserialize)]\n pub struct ThreadsResponse {\n pub threads: Vec,\n }\n \n-#[derive(Debug, Serialize, Deserialize)]\n+#[derive(Debug, Clone, Serialize, Deserialize)]\n pub struct ThreadSummary {\n pub thread: String,\n pub last_activity_ts: i64,\n@@ -220,15 +220,173 @@ pub struct FeedPost {\n pub body: String,\n }\n \n+// ---------------------------------------------------------------------------\n+// RPC batch API (`POST /api/v0/rpc`)\n+// ---------------------------------------------------------------------------\n+\n+#[derive(Debug, Serialize, Deserialize)]\n+#[serde(transparent)]\n+pub struct RpcBatch(pub Vec);\n+\n+#[derive(Debug, Serialize, Deserialize)]\n+#[serde(rename_all = \"PascalCase\")]\n+pub enum RpcCommand {\n+ Post {\n+ room: String,\n+ thread_tag: String,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ delegate: Option,\n+ text: String,\n+ #[serde(default)]\n+ return_rank_diff: bool,\n+ },\n+ Check {\n+ room: String,\n+ text: String,\n+ },\n+ GetGardenRank {\n+ room: String,\n+ parent_path: String,\n+ #[serde(default)]\n+ depth: Option,\n+ #[serde(default)]\n+ offset: Option,\n+ #[serde(default)]\n+ limit: Option,\n+ #[serde(default)]\n+ percent: Option,\n+ },\n+ GetGardenItem {\n+ room: String,\n+ item_path: String,\n+ #[serde(default)]\n+ full: Option,\n+ },\n+ GetForumThread {\n+ room: String,\n+ thread_tag: String,\n+ #[serde(default)]\n+ offset: Option,\n+ #[serde(default)]\n+ limit: Option,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ since: Option,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ before: Option,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ actor: Option,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ post_id: Option,\n+ },\n+ ListForumThreads {\n+ room: String,\n+ },\n+ RoomCreate {\n+ slug: String,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ visibility: Option,\n+ },\n+ RoomGrant {\n+ room: String,\n+ username: String,\n+ capability: String,\n+ },\n+ GetGlobalRank {\n+ room: String,\n+ #[serde(default)]\n+ limit: Option,\n+ #[serde(default)]\n+ offset: Option,\n+ #[serde(default)]\n+ percent: Option,\n+ },\n+ GetPair {\n+ room: String,\n+ parent_path: String,\n+ },\n+ GetMatchup {\n+ room: String,\n+ item_path: String,\n+ #[serde(default)]\n+ limit: Option,\n+ },\n+ GetRankHistory {\n+ room: String,\n+ item_path: String,\n+ },\n+ GetLeaves {\n+ room: String,\n+ },\n+ GetPaths {\n+ room: String,\n+ },\n+ GetRecentVotes {\n+ room: String,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ parent: Option,\n+ #[serde(default)]\n+ limit: Option,\n+ },\n+ Search {\n+ query: String,\n+ },\n+ GetFeed {\n+ actor: String,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ since: Option,\n+ #[serde(default)]\n+ limit: Option,\n+ },\n+}\n+\n+#[derive(Debug, Serialize, Deserialize)]\n+#[serde(rename_all = \"PascalCase\")]\n+pub enum RpcResult {\n+ PostOk {\n+ events_appended: usize,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ ranking_changes: Option>,\n+ threads: Vec,\n+ next: NextMoves,\n+ },\n+ CheckOk {\n+ rankings: Vec,\n+ threads: Vec,\n+ next: Vec,\n+ },\n+ GardenRank(RankResponse),\n+ GardenItem(ItemResponse),\n+ ForumThread(ThreadDetailResponse),\n+ ForumThreads(ThreadsResponse),\n+ RoomCreated {\n+ room_id: String,\n+ },\n+ GrantOk {},\n+ GlobalRank(GlobalRankResponse),\n+ Pair(PairResponse),\n+ Matchup(MatchupResponse),\n+ RankHistory(RankHistoryResponse),\n+ Leaves(LeavesResponse),\n+ Paths(PathsResponse),\n+ RecentVotes(RecentVotesResponse),\n+ Search(SearchResponse),\n+ Feed(FeedResponse),\n+}\n+\n+#[derive(Debug, Serialize, Deserialize)]\n+pub struct RpcLine {\n+ pub ok: bool,\n+ #[serde(skip_serializing_if = \"Option::is_none\")]\n+ pub result: Option,\n+ #[serde(skip_serializing_if = \"Option::is_none\")]\n+ pub error: Option,\n+ #[serde(skip_serializing_if = \"Option::is_none\")]\n+ pub hint: Option,\n+}\n+\n #[derive(Debug, Serialize, Deserialize)]\n-pub struct IngestRequest {\n- /// Thread identifier: public tag (e.g. \"languages\") or private id/slug (e.g. \"a7f2k9x/project-review\").\n- pub thread: String,\n- /// Delegate id: `uuid:rig:provider/model` (no `@`). Omit for human-only ingests.\n- #[serde(default, skip_serializing_if = \"Option::is_none\")]\n- pub delegate: Option,\n- /// DSL+prose body only.\n- pub text: String,\n+pub struct RpcBatchResponse {\n+ pub results: Vec,\n }\n \n /// Start a browser-based OAuth login flow for a CLI agent.\n","role":"user"}],"model":"openai/gpt-chat-latest"}