Side B replaces a no-op stub with a real, tested Reddit fetch worker including OAuth token handling, rate-limit/backoff logic, JSON parsing for subreddits and posts, and wiring into the item-fetch path plus unit tests for URL mapping and parsing. Side A mainly wires a new CLI room-create subcommand end-to-end and removes an unused room 'visibility' concept, which is useful cleanup and plumbing but smaller in scope and impact than the working external-integration feature added in B.
constitution · epochs · watch · epoch 3
c_1c1c8e7a2de8 (tommy-mor) vs c_c124c217f89c (tommy-mor)
download prompt · raw event · cmp_99501c916d60fe
council reasoning
Commit A lands the end-to-end private room create path (CLI RoomCmd + RPC usage) and a lasting domain cleanup (drop ThreadVisibility, rooms as HashSet) with docs and tests aligned. Commit B is a large, well-structured Reddit fetch worker, but it is a peripheral import path still marked incomplete (rendering etc. in todo), so it adds less durable product core than A.
Side A delivers a complete user-facing capability by adding a `room create` CLI command wired through the existing RPC, improving documentation, and simplifying the room model by removing the unused visibility field and replacing the room registry with a `HashSet` of room IDs. Side B introduces substantial Reddit import infrastructure (background worker, OAuth, rate limiting, parsing, and fetch triggers), but it is a larger, partially integrated feature with placeholder/TODO elements and more implementation complexity, whereas Side A provides a cohesive, immediately usable feature and durable model cleanup.
sides
A — c_1c1c8e7a2de8 (tommy-mor)
message
[62d18183] room create path
diff preview
diff --git a/cli/GUIDE.sorter b/cli/GUIDE.sorter
index dcb06a46045564f8f6f6acffbda6f88644d453cc..9828cba4d9c17b7cce3de597d8724609b2b2adbe 100644
--- a/cli/GUIDE.sorter
+++ b/cli/GUIDE.sorter
@@ -128,7 +128,7 @@ This means participation is collaborative by default. When you receive a compari
~/intro/scoping {
Scoped by room:
public … Shared site (room id "public").
- private <ROOM_ID> … Private room (e.g. abc12xy/my-project from RoomCreate over RPC).
+ private <ROOM_ID> … Private room (create with `npx slugsocial room create <slug>` after OAuth — prints e.g. abc12xy/my-project).
Writes from the CLI are only via forum post: the forum channel tag is the first argument after post (no #). Humans post through the website; CLI requires --delegate (agent identity).
@@ -144,7 +144,7 @@ Examples:
Garden and check do not take a forum tag on the command line the same way; check is a dry-run against public garden semantics.
-Global (no room prefix): identity, whoami, feed, search, healthz.
+Global (no room prefix): room, identity, whoami, feed, search, healthz.
}
~/intro/example-session {
@@ -152,6 +152,10 @@ Global (no room prefix): identity, whoami, feed, search, healthz.
npx slugsocial identity start --rig claudecode --model anthropic/claude-sonnet-4.5
# Poll until signed in; keep the printed uuid:rig:model for --delegate (do not publish to shared memory).
+# Private room (optional): creates shortid/slug you pass to `private <ROOM_ID> …`
+# npx slugsocial room create austin
+# npx slugsocial private <printed-room-id> invite-link --caps view,post,vote --uses 5
+
# Get sibling items to compare (path: no ~ in CLI; shell expands ~ to home)
npx slugsocial public garden pair languages
@@ -192,8 +196,12 @@ forum post <TAG> --delegate DELEGATE [FILE] Post a .sorter doc (stdin if no
check [FILE] Validate without submitting (public garden dry-run)
+invite-link --caps view,post[,…] [--uses N] Mint shareable /join/… link (private rooms; Manage required)
+audit [--json] List principals + capabilities (private rooms; View or Manage)
+
Global (no public/private prefix):
+room create <slug> Create a private room (bearer required); prints ROOM_ID for `private …` (use `public …` for the shared site, not a room)
identity start --rig <name> --model <provider/model> New delegate id + OAuth pending session
identity poll <session> Complete OAuth; saves bearer token
diff --git a/cli/src/main.rs b/cli/src/main.rs
index 8eda9f485bd1f7392f1e34be27176c21e20354eb..5b1a5845e90bfea9bd9a1e1af5744e97e566b5ae 100644
--- a/cli/src/main.rs
+++ b/cli/src/main.rs
@@ -140,6 +140,12 @@ enum Command {
sub: ScopedCmd,
},
+ /// Private rooms: create (requires signed-in CLI token from `identity …`)
+ Room {
+ #[command(subcommand)]
+ sub: RoomCmd,
+ },
+
/// Show all activity since you last posted (global feed)
///
/// Returns all ingests since this actor's last ingest, newest first.
@@ -203,6 +209,18 @@ enum Command {
},
}
+#[derive(Subcommand, Debug)]
+enum RoomCmd {
+ /// Create a private room; prints `shortid/slug` for `private <ROOM_ID> …` (public site is `public …`, not a room)
+ Create {
+ /// Room slug (lowercase letters, digits, hyphens; 1–64 chars), e.g. `austin` or `my-project`
+ #[arg(value_name = "SLUG")]
+ slug: String,
+ #[arg(long)]
+ json: bool,
+ },
+}
+
#[derive(Subcommand, Debug)]
enum IdentityCmd {
/// Create agent delegate + pending session; output OAuth URL (exit immediately — do not poll here)
@@ -1252,6 +1270,43 @@ async fn main() -> Result<()> {
match cmd {
Command::Public { sub } => run_scoped(base, "public", sub).await?,
Command::Private { room, sub } => run_scoped(base, &room, sub).await?,
+ Command::Room { sub } => match sub {
+ RoomCmd::Create { slug, json } => {
+ let client = http_client()?;
+ let bearer = effective_bearer().ok_or_else(|| {
+ anyhow!(
+ "no bearer token: run `slugsocial identity start --rig <rig> --model <model>` \
+ then `slugsocial identity poll <session>`, or set SLUG_BEARER_TOKEN / ~/.config/slugsocial/token"
+ )
+ })?;
+ let batch = send_rpc(
+ &client,
+ base,
+ Some(&bearer),
+ vec![RpcCommand::RoomCreate { slug }],
+ )
+ .await?;
+ match rpc_line_ok(&batch.results[0])? {
+ RpcResult::RoomCreated { room_id } => {
+ if json {
+ println!(
+ "{}",
+ serde_json::to_string_pretty(&serde_json::json!({
+ "ok": true,
+ "room_id": room_id,
+ }))?
+ );
+ } else {
+ println!("{room_id}");
+ println!();
+ println!("Next: npx slugsocial private {room_id} forum post <TAG> --delegate '…' …");
+ println!(" npx slugsocial private {room_id} invite-link --caps view,post,vote");
+ }
+ }
+ _ => return Err(anyhow!("unexpected RPC result")),
+ }
+ }
+ },
Command::Healthz { json } => {
let client = http_client()?;
diff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs
index f6bbc3df71909a2da7403cd46fe4ea6ca130c692..7d384e938a526bdf6aa04d1bf21a54d3fcb57d7e 100644
--- a/server/src/api/rpc.rs
+++ b/server/src/api/rpc.rs
@@ -14,7 +14,7 @@ use crate::{
canonical_path::{canonicalize_item, canonicalize_tag},
dsl,
events::{
- AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability, ThreadVisibility,
+ AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability,
},
identity::{parse_agent, parse_username},
path_types::CanonicalItemUrl,
@@ -270,7 +270,7 @@ async fn rpc_post(
let scope = scope_from_room_wire(&room_key);
let is_private = !matches!(scope, ScopeId::Public);
- if is_private && !reduced.rooms.contains_key(&room_key) {
+ if is_private && !reduced.rooms.contains(&room_key) {
drop(reduced);
return Err(("unknown room".into(), Some(format!("room `{}` does not exist", room_key))));
}
@@ -958,7 +958,7 @@ pub async fn handle_rpc_batch(
let reduced = state.reduced.read().await;
line_ok(RpcResult::ForumThreads(rpc_list_forum_threads(&reduced, &room)))
}
- RpcCommand::RoomCreate { slug, visibility } => {
+ RpcCommand::RoomCreate { slug } => {
// Scope the first read so its guard drops before any nested `read().await` / `write().await`.
// A guard from `match verify(..., &*state.reduced.read().await)` would otherwise live for the
// whole `match` and deadlock here (tokio::sync::RwLock is not reentrant).
@@ -975,53 +975,42 @@ pub async fn handle_rpc_batch(
} else if !slug.chars().all(|c| c.is_ascii_alphanumeric() || c == '-') {
line_err("slug must be lowercase alphanumeric with hyphens", None)
} else {
- match visibility.as_deref().unwrap_or("private") {
- "private" | "public" => {
- let vis = if visibility.as_deref() == Some("public") {
- ThreadVisibility::Public
- } else {
- ThreadVisibility::Private
- };
- let short_id = loop {
- let id = gen_short_id();
- if !state.reduced.read().await.rooms.contains_key(&format!("{id}/{slug}")) {
- break id;
- }
- };
- let room_id = format!("{short_id}/{slug}");
- let ts = now_ms();
- let tc_ev = Event::RoomCreated(RoomCreated {
- ts,
- room_id: room_id.clone(),
- slug: slug.clone(),
- owner: principal.clone(),
- visibility: vis,
- });
- let ga_ev = Event::GrantAdded(GrantAdded {
- ts,
- room_id: room_id.clone(),
- username: principal.clone(),
- capabilities: vec![
- ThreadCapability::View,
- ThreadCapability::Post,
- ThreadCapability::Vote,
- ThreadCapability::AddItem,
- ThreadCapability::Manage,
- ],
- granted_by: principal.clone(),
- });
- if let Err(e) = state.event_log.append(&tc_ev).await {
- line_err(format!("{e}"), None)
- } else if let Err(e) = state.event_log.append(&ga_ev).await {
- line_err(format!("{e}"), None)
- } else {
- let mut r = state.reduced.write().await;
- r.apply_event(tc_ev);
- r.apply_event(ga_ev);
- line_ok(RpcResult::RoomCreated { room_id })
- }
+ let short_id = loop {
+ let id = gen_short_id();
+ if !state.reduced.read().await.rooms.contains(&format!("{id}/{slug}")) {
+ break id;
}
- other => line_err(format!("unknown visibility: {other}"), None),
+ };
+ let room_id = format!("{short_id}/{slug}");
+ let ts = now_ms();
+ let tc_ev = Event::RoomCreated(RoomCreated {
+ ts,
+ room_id: room_id.clone(),
+ slug: slug.clone(),
+ owner: principal.clone(),
+ });
+ let ga_ev = Event::GrantAdded(GrantAdded {
+ ts,
+ room_id: room_id.clone(),
+ username: principal.clone(),
+ capabilities: vec![
+ ThreadCapability::View,
+ ThreadCapability::Post,
+ ThreadCapability::Vot
… preview truncated; 10,232 characters omittedB — c_c124c217f89c (tommy-mor)
message
[8d8230d1] reddit
diff preview
diff --git a/.gitignore b/.gitignore
index 4c7073f9fac0c30fd2050d79a60ef447af58ebeb..ada462e900d24a3a6d08165d158c80f79c35a5a5 100644
--- a/.gitignore
+++ b/.gitignore
@@ -7,3 +7,4 @@
data/
repomix-output.xml
dev-data/
+.env
diff --git a/server/Cargo.toml b/server/Cargo.toml
index 7906a8547d56b8e6a48ef59c37aa82a8510fdee9..4677fedcb45292eebebe7e9cf6ce2f5738f18ddf 100644
--- a/server/Cargo.toml
+++ b/server/Cargo.toml
@@ -16,6 +16,7 @@ tower = "0.5"
tower-http = { version = "0.5", features = ["trace"] }
tracing = "0.1"
tracing-subscriber = { version = "0.3", features = ["env-filter"] }
+reqwest = { version = "0.12", features = ["json"] }
[dev-dependencies]
reqwest = { version = "0.12", features = ["json"] }
diff --git a/server/src/html/mod.rs b/server/src/html/mod.rs
index 2864407ed6e8ec284a1dc663acf1805534566b62..df6505021d9f446c2b453e20e3eb3cf696a111f9 100644
--- a/server/src/html/mod.rs
+++ b/server/src/html/mod.rs
@@ -272,5 +272,16 @@ pub async fn home(State(state): State<AppState>, uri: Uri) -> impl IntoResponse
pub async fn browse(State(state): State<AppState>, uri: Uri) -> impl IntoResponse {
let item = ItemId::from_browse_uri(uri.path()).unwrap_or(ItemId::root());
+ if item.as_str().starts_with("reddit.com") {
+ let needs_fetch = {
+ let tree = state.tree.read().await;
+ tree.get(&item)
+ .map(|n| n.data.is_none())
+ .unwrap_or(true)
+ };
+ if needs_fetch {
+ state.reddit.request_fetch(item.clone());
+ }
+ }
item_page(state, uri, item).await
}
diff --git a/server/src/reddit.rs b/server/src/reddit.rs
index d203dca09245daf869b3aa942898447700ae69fb..90053ad03b1d7c8e94f325dd4ee64c2b4f7da900 100644
--- a/server/src/reddit.rs
+++ b/server/src/reddit.rs
@@ -1,4 +1,12 @@
-//! Reddit API import (async, decoupled from UI request path).
+//! Reddit API import via a single background worker (rate limits, dedup, backoff).
+
+use std::collections::{HashMap, HashSet};
+use std::sync::Arc;
+use std::time::{Duration, Instant};
+
+use reqwest::{header, Client, StatusCode};
+use serde::Deserialize;
+use tokio::sync::{mpsc, RwLock};
use crate::{
path_types::ItemId,
@@ -10,12 +18,401 @@ pub fn ensure_partial_tree(tree: &mut GlobalTree, id: &ItemId) {
tree.ensure_path(id);
}
-/// Placeholder for Reddit JSON import. Returns entity data when implemented.
-pub async fn fetch_reddit_entity(_id: &ItemId) -> Option<EntityData> {
- None
+pub struct RedditCommand {
+ pub id: ItemId,
+}
+
+#[derive(Clone)]
+pub struct RedditBroker {
+ tx: mpsc::Sender<RedditCommand>,
+}
+
+#[derive(Clone)]
+struct RedditCredentials {
+ client_id: String,
+ client_secret: String,
+}
+
+struct OAuthToken {
+ access_token: String,
+ expires_at: Instant,
+}
+
+impl RedditBroker {
+ pub fn spawn(tree: Arc<RwLock<GlobalTree>>, user_agent: &str) -> Self {
+ let (tx, rx) = mpsc::channel(100);
+
+ let mut headers = header::HeaderMap::new();
+ headers.insert(
+ header::USER_AGENT,
+ header::HeaderValue::from_str(user_agent).expect("valid user agent"),
+ );
+
+ let client = Client::builder()
+ .default_headers(headers)
+ .timeout(Duration::from_secs(15))
+ .build()
+ .expect("reqwest client");
+
+ let creds = RedditCredentials::from_env();
+ tokio::spawn(reddit_worker(rx, tree, client, creds));
+
+ Self { tx }
+ }
+
+ /// Fire-and-forget: queue a fetch; worker updates the tree when done.
+ pub fn request_fetch(&self, id: ItemId) {
+ let _ = self.tx.try_send(RedditCommand { id });
+ }
+}
+
+impl RedditCredentials {
+ fn from_env() -> Option<Self> {
+ let client_id = std::env::var("REDDIT_CLIENT_ID").ok()?;
+ let client_secret = std::env::var("REDDIT_CLIENT_SECRET").ok()?;
+ if client_id.is_empty() || client_secret.is_empty() {
+ return None;
+ }
+ Some(Self {
+ client_id,
+ client_secret,
+ })
+ }
+}
+
+pub fn default_user_agent() -> String {
+ std::env::var("REDDIT_USER_AGENT").unwrap_or_else(|_| {
+ "web:sorter2.social:v0.0.1 (by /u/sorter2)".to_string()
+ })
}
-/// Apply fetched entity data to a node (called from async worker).
-pub fn apply_entity(tree: &mut GlobalTree, id: &ItemId, data: EntityData) {
- tree.set_entity_data(id, data);
+async fn reddit_worker(
+ mut rx: mpsc::Receiver<RedditCommand>,
+ tree: Arc<RwLock<GlobalTree>>,
+ client: Client,
+ creds: Option<RedditCredentials>,
+) {
+ let mut in_flight = HashSet::new();
+ let mut recently_fetched: HashMap<ItemId, Instant> = HashMap::new();
+ let mut current_delay = Duration::from_secs(1);
+ let mut oauth: Option<OAuthToken> = None;
+ let cache_ttl = Duration::from_secs(300);
+
+ while let Some(cmd) = rx.recv().await {
+ let now = Instant::now();
+ recently_fetched.retain(|_, t| now.duration_since(*t) < cache_ttl);
+
+ if in_flight.contains(&cmd.id) || recently_fetched.contains_key(&cmd.id) {
+ continue;
+ }
+
+ in_flight.insert(cmd.id.clone());
+ let fetch_id = cmd.id.clone();
+
+ tokio::time::sleep(current_delay).await;
+
+ if let Some(c) = &creds {
+ oauth = ensure_oauth_token(&client, c, oauth.take()).await;
+ }
+
+ let token = oauth.as_ref().map(|t| t.access_token.as_str());
+ let use_oauth = token.is_some();
+
+ match do_fetch(&client, &fetch_id, use_oauth, token).await {
+ Ok(FetchOutcome::Entity(data)) => {
+ let mut w = tree.write().await;
+ w.set_entity_data(&fetch_id, data);
+ recently_fetched.insert(fetch_id.clone(), Instant::now());
+ current_delay = Duration::from_millis(600);
+ }
+ Ok(FetchOutcome::NotFound) => {
+ recently_fetched.insert(fetch_id.clone(), Instant::now());
+ }
+ Ok(FetchOutcome::RateLimited { reset_secs }) => {
+ let wait = Duration::from_secs(reset_secs.max(1));
+ tracing::warn!(
+ "Reddit rate limit for {}; sleeping {}s",
+ fetch_id,
+ wait.as_secs()
+ );
+ tokio::time::sleep(wait).await;
+ current_delay = (current_delay * 2).min(Duration::from_secs(60));
+ }
+ Err(e) => {
+ tracing::warn!("Reddit fetch failed for {}: {}", fetch_id, e);
+ current_delay = (current_delay * 2).min(Duration::from_secs(60));
+ }
+ }
+
+ in_flight.remove(&fetch_id);
+ }
+}
+
+enum FetchOutcome {
+ Entity(EntityData),
+ NotFound,
+ RateLimited { reset_secs: u64 },
+}
+
+async fn ensure_oauth_token(
+ client: &Client,
+ creds: &RedditCredentials,
+ existing: Option<OAuthToken>,
+) -> Option<OAuthToken> {
+ if let Some(t) = existing {
+ if Instant::now() < t.expires_at - Duration::from_secs(60) {
+ return Some(t);
+ }
+ }
+
+ let resp = client
+ .post("https://www.reddit.com/api/v1/access_token")
+ .basic_auth(&creds.client_id, Some(&creds.client_secret))
+ .form(&[("grant_type", "client_credentials")])
+ .send()
+ .await;
+
+ let resp = match resp {
+ Ok(r) => r,
+ Err(e) => {
+ tracing::warn!("Reddit OAuth token request failed: {e}");
+ return None;
+ }
+ };
+
+ if !resp.status().is_success() {
+ tracing::warn!("Reddit OAuth token HTTP {}", resp.status());
+ return None;
+ }
+
+ #[derive(Deserialize)]
+ struct TokenResponse {
+ access_token: String,
+ expires_in: u64,
+ }
+
+ let body: TokenResponse = match resp.json().await {
+ Ok(b) => b,
+ Err(e) => {
+ tracing::warn!("Reddit OAuth token parse failed: {e}");
+ return None;
+ }
+ };
+
+ Some(OAuthToken {
+ access_token: body.access_token,
+ expires_at: Instant::now() + Duration::from_secs(body.expires_in),
+ })
+}
+
+async fn do_fetch(
+ client: &Client,
+ id: &ItemId,
+ use_oauth: bool,
+ bearer: Option<&str>,
+) -> Result<FetchOutcome, String> {
+ let url = map_item_to_reddit_api(id, use_oauth);
+ if url.is_empty() {
+ return Ok(FetchOutcome::NotFound);
+ }
+
+ let mut req = client.get(&url);
+ if let Some(token) = bearer {
+ req = req.bearer_auth(token);
+ }
+
+ let resp = req.send().await.map_err(|e| e.to_string())?;
+
+ if resp.status() == StatusCode::TOO_MANY_REQUESTS {
+ let reset = rate_limit_reset_secs(&resp);
+ return Ok(FetchOutcome::RateLimited { reset_secs: reset });
+ }
+
+ if resp.status() == StatusCode::SERVICE_UNAVAILABLE {
+ return Err("Reddit unavailable (503)".to_string());
+ }
+
+ if !resp.status().is_success() {
+ return Ok(FetchOutcome::NotFound);
+ }
+
+ if rate_limit_remaining(&resp) == Some(0) {
+ let reset = rate_limit_reset_secs(&resp);
+ return Ok(FetchOutcome::RateLimited { reset_secs: reset });
+ }
+
+ let bytes = resp.bytes().await.map_err(|e| e.to_string())?;
+ Ok(parse_reddit_json(id, &bytes)
+ .map(FetchOutcome::Entity)
+ .unwrap_or(FetchOutcome::NotFound))
+}
+
+fn rate_limit_remaining(resp: &reqwest::Response) -> Option<u64> {
+ resp.headers()
+ .get("x-ratelimit-remaining")
+ .and_then(|v| v.to_str().ok())
+ .and_then(|s| s.parse::<f64>().ok())
+ .map(|f| f.floor() as u64)
+}
+
+fn rate_limit_reset_secs(resp: &reqwest::Response) -> u64 {
+ resp.headers()
+ .get("x-ratelimit-reset")
+ .and_then(|v| v.to_str().ok())
+ .and_then(|s| s.parse::<f64>().ok())
+ .map(|f| f.ceil() as u64)
+ .unwrap_or(5)
+}
+
+/// Map canonical item id to Reddit JSON API URL.
+pub fn map_item_to_reddit_api(id: &ItemId, oauth: bool) -> String {
+ let path = id.as_str();
+ if !path.starts_with("reddit.com/") && path != "reddit.com" {
+ return String::new();
+ }
+
+ let base = if oauth {
+ "https://oauth.reddit.com"
+ } else {
+ "https://www.reddit.com"
+ };
+
+ let segments: Vec<&str> = path.split('/').collect();
+
+ if let Some(i) = segments.iter().position(|&p| p == "comments") {
+ if segments.len() > i + 1 {
+ let api_path = segments[1..=i + 1].join("/");
+ return format!("{base}/{api_path}.json?raw_json=1");
+ }
+ }
+
+ if segments.len() == 3 && segments[1] == "r" {
+ return format!("{base}/r/{}/about.json?raw_json=1", segments[2]);
+ }
+
+ String::new()
+}
+
+fn parse_reddit_json(id: &ItemId, bytes: &[u8]) -> Option<EntityData> {
+ let v: serde_json::Value = serde_json::from_slice(bytes).ok()?;
+ let segments: Vec<&str> = id.as_str().split('/').collect();
+
+ if segments.iter().any(|&p| p == "comments") {
+ parse_post_listing(&v)
+ } else {
+ parse_subreddit_about(&v)
+ }
+}
+
+fn parse_subreddit_about(v: &serde_json::Value) -> Option<EntityData> {
+ let data = v.get("data")?;
+ let title = data
+ .get("title")
+ .or_else(|| data.get("display_name"))
+ .and_then(|t| t.as_str())?
+ .to_string();
+ let body_html = data
+ .get("public_description_html")
+ .or_else(|| data.get("public_description"))
+ .and_then(|t| t.as_str())
+ .map(|s| s.to_string());
+ let thumb_url = data
+ .get("icon_img")
+ .or_else(|| data.get("community_icon"))
+ .and_then(|t| t.as_str())
+ .filter(|s| !s.is_empty())
+ .map(|s| s.to_string());
+
+ Some(EntityData {
+ title,
+ author: None,
+ body_html,
+ thumb_url,
+ })
+}
+
+fn parse_post_listing(v:
… preview truncated; 4,570 characters omittedHardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.