Side A performs a genuine architectural cleanup: separating identity parsing/validation into a dedicated module, removing ad-hoc identity rewriting from the reducer, making delegate optional, and consistently updating wire/HTML boundaries with matching test and CLI updates across the whole codebase. Side B is a solid, scoped bugfix (forcing OAuth and retry-on-401/403 for Reddit fetches) with good error truncation and a test, but it's narrower in scope and impact compared to A's broader identity-model correctness fix that touches API contracts, storage, and display logic throughout the system.
constitution · epochs · watch · epoch 3
c_e2ee16c7ada5 (tommy-mor) vs c_ca9169f732b8 (tommy-mor)
download prompt · raw event · cmp_40a4942a7bc10a
council reasoning
A permanently separates path vs identity concerns, makes wire/storage identities strict and naked (rejecting `@`), optional delegates, and stops the reducer from rewriting usernames/agents—core event-model correctness across auth, ingest, APIs, and HTML display. B is a real, tightly scoped production fix (mandatory OAuth when creds exist, 401/403 refresh, no public www fallback, fly base URL), but it only touches the Reddit worker path versus A’s lasting domain contract.
Side A makes a broad architectural improvement by separating path normalization from identity parsing (`canonical_path.rs` and `identity.rs`), removing identity rewriting from the reducer, enforcing strict stored-form usernames/agent IDs without `@`, and updating APIs and event structures (including optional delegates) to use a consistent wire format. Side B is a solid operational bug fix—requiring OAuth when configured, refreshing tokens after 401/403, and improving Reddit error handling—but it is confined to one subsystem, whereas Side A establishes cleaner long-term boundaries and data invariants across the project.
sides
A — c_e2ee16c7ada5 (tommy-mor)
message
[80ad7753] refactor: split canonical_path and identity; strict wire identity without @ - Add canonical_path.rs (tag + item URL normalization) and identity.rs (parse_username/parse_agent; reject @ in API input). - Slim events.rs to event types only; reducer applies no identity rewriting. - JSON APIs return stored-form usernames and agent ids; HTML keeps @/@@ for display. - Optional delegate on ingest; CLI and tests use naked uuid:rig:model. Made-with: Cursor
diff preview
diff --git a/cli/src/main.rs b/cli/src/main.rs
index 5ac8e289f2b7a366b5959d9338d02f9b546f408a..630c5dea1f78c0ec9bc53e6b96234a0dc75bb705 100644
--- a/cli/src/main.rs
+++ b/cli/src/main.rs
@@ -61,8 +61,8 @@ enum Command {
/// Example: --before 2026-06-01
#[arg(long, value_name = "DATE_OR_MS")]
before: Option<String>,
- /// Filter to posts from this actor (UUID prefix match).
- /// Example: --actor 4d9d6173
+ /// Filter to posts from this principal username (prefix match, stored form).
+ /// Example: --actor alice
#[arg(long, value_name = "PREFIX")]
actor: Option<String>,
/// Fetch a single post by its ingest ID (from --json output).
@@ -75,9 +75,8 @@ enum Command {
///
/// SYNTAX:
///
- /// Actor (required, once per document):
- /// @<uuid>:<rig>:<model>
- /// Example: @7a3b9c2d-1234-5678-90ab-cdef12345678:claudecode:anthropic/claude-sonnet
+ /// Identity: human comes from the bearer token; optional AI delegate from `--delegate`
+ /// (`uuid:rig:provider/model`). The document body is DSL only (items, votes, prose) — no `@` lines.
///
/// Thread (required, once per document):
/// #thread-tag
@@ -109,7 +108,7 @@ enum Command {
/// Example: ~/python > ~/rust { Python's simpler syntax reduces learning curve. }
///
/// Prose (optional, anywhere):
- /// Any line that doesn't start with @, #, or ~ is prose.
+ /// Any line that doesn't start with # or ~ (or `http`) is prose.
/// Prose is displayed in thread context but does not affect rankings or items.
/// Use prose to write blog posts, reasoning, or notes within your ingest.
///
@@ -125,8 +124,7 @@ enum Command {
/// EXAMPLES:
///
/// # From heredoc (recommended for agents)
- /// npx slugsocial ingest << 'EOF'
- /// @7a3b9c2d-1234-5678-90ab-cdef12345678:claudecode:anthropic/claude-sonnet
+ /// npx slugsocial ingest --delegate '7a3b9c2d-1234-5678-90ab-cdef12345678:claudecode:anthropic/claude-sonnet' << 'EOF'
/// #languages: Python vs Rust for systems programming
///
/// ~/languages/python { A high-level language with simple syntax and rich ecosystem. }
@@ -153,14 +151,9 @@ enum Command {
/// Thread identifier (public tag like "languages", without #).
#[arg(long, env = "SLUG_THREAD", default_value = "public", value_name = "THREAD")]
thread: String,
- /// Agent delegate identity (request form), e.g. @@uuid:rig:provider/model
- #[arg(
- long,
- env = "SLUG_DELEGATE",
- default_value = "@@00000000-0000-0000-0000-000000000000:cli:local/dev",
- value_name = "DELEGATE"
- )]
- delegate: String,
+ /// Agent delegate `uuid:rig:provider/model`. Omit for human-only ingests.
+ #[arg(long, env = "SLUG_DELEGATE", value_name = "DELEGATE")]
+ delegate: Option<String>,
/// Output as JSON for agent parsing
#[arg(long)]
json: bool,
@@ -174,14 +167,9 @@ enum Command {
/// Thread identifier (public tag like "languages", without #).
#[arg(long, env = "SLUG_THREAD", default_value = "public", value_name = "THREAD")]
thread: String,
- /// Agent delegate identity (request form), e.g. @@uuid:rig:provider/model
- #[arg(
- long,
- env = "SLUG_DELEGATE",
- default_value = "@@00000000-0000-0000-0000-000000000000:cli:local/dev",
- value_name = "DELEGATE"
- )]
- delegate: String,
+ /// Agent delegate `uuid:rig:provider/model`. Omit for human-only ingests.
+ #[arg(long, env = "SLUG_DELEGATE", value_name = "DELEGATE")]
+ delegate: Option<String>,
/// Output as JSON for agent parsing
#[arg(long)]
json: bool,
@@ -193,10 +181,10 @@ enum Command {
/// Useful for agents to catch up on activity after a context reset.
///
/// Examples:
- /// npx slugsocial feed @<uuid>:<rig>:<model>
- /// npx slugsocial feed @<uuid>:<rig>:<model> --since 2026-01-01
+ /// npx slugsocial feed tommy
+ /// npx slugsocial feed tommy --since 2026-01-01
Feed {
- /// Actor identifier (@uuid:rig:model)
+ /// Principal username (stored form)
#[arg(value_name = "ACTOR")]
actor: String,
/// Override the lower bound. Accepts Unix ms or YYYY-MM-DD.
@@ -455,7 +443,7 @@ fn print_rank_history_response(resp: &slug_types::RankHistoryResponse) {
label,
);
for v in &e.caused_by {
- println!(" {} {} {} {}", v.a, v.ratio, v.b, v.actor.as_deref().map(|a| format!(" (@{})", a)).unwrap_or_default());
+ println!(" {} {} {} {}", v.a, v.ratio, v.b, v.actor.as_deref().map(|a| format!(" ({})", a)).unwrap_or_default());
if !v.body.is_empty() {
println!(" {}", v.body.lines().next().unwrap_or(&v.body).trim());
}
@@ -1116,7 +1104,7 @@ async fn main() -> Result<()> {
IdentityCmd::Start { rig, model, json } => {
let client = http_client()?;
let uuid = uuid::Uuid::new_v4().to_string();
- let delegate = format!("@@{}:{}:{}", uuid, rig, model);
+ let delegate = format!("{uuid}:{rig}:{model}");
let start: PendingSessionStartResponse = expect_json(
client
diff --git a/server/src/api/auth.rs b/server/src/api/auth.rs
index cb0faa29b834931e2c2b2f5c174c875e2e2e9346..995ce4a61d29b024c399c656134f541ecfd880cf 100644
--- a/server/src/api/auth.rs
+++ b/server/src/api/auth.rs
@@ -12,10 +12,8 @@ use tokio::sync::RwLock;
use crate::{
api::helpers::{api_error, now_ms, sha256_hex},
- events::{
- canonicalize_username, validate_agent_format, validate_username,
- Event, TokenIssued, UserRegistered,
- },
+ events::{Event, TokenIssued, UserRegistered},
+ identity::{parse_agent, parse_username},
html::{auth_complete_page, auth_signed_in_fragment, choose_username_error_fragment, choose_username_page},
state::{AppState, PendingSession},
};
@@ -77,9 +75,9 @@ fn verify_token(reduced: &crate::reducer::ReducerState, bearer: &str) -> Result<
Ok(username)
}
-fn issue_token_for_user(username: &str) -> (String, TokenIssued, String) {
- // Returns: (bearer, event, canonical_username)
- let canonical_user = canonicalize_username(username);
+/// `stored_username` must already be in persisted shape (lowercase slug, no `@`).
+fn issue_token_for_user(stored_username: &str) -> (String, TokenIssued) {
+ let username = stored_username.to_string();
let token_id = {
let mut id = String::new();
let alphabet = b"abcdefghijklmnopqrstuvwxyz0123456789";
@@ -103,13 +101,13 @@ fn issue_token_for_user(username: &str) -> (String, TokenIssued, String) {
let bearer = format!("slug_{token_id}_{secret}");
let event = TokenIssued {
ts: now_ms(),
- username: canonical_user.clone(),
+ username: username.clone(),
token_id,
token_hash,
salt,
issued_via: "oauth".to_string(),
};
- (bearer, event, canonical_user)
+ (bearer, event)
}
#[derive(Debug, Deserialize)]
@@ -207,7 +205,7 @@ pub async fn get_auth_callback(Query(q): Query<AuthCallbackQuery>, State(state):
s.provider = Some("google".to_string());
s.provider_id = Some(sub.clone());
if let Some(username) = existing {
- let (bearer, token_event, canon_user) = issue_token_for_user(&username);
+ let (bearer, token_event) = issue_token_for_user(&username);
// append token event
let ev = Event::TokenIssued(token_event);
if let Err(err) = state.event_log.append(&ev).await {
@@ -217,7 +215,7 @@ pub async fn get_auth_callback(Query(q): Query<AuthCallbackQuery>, State(state):
let mut reduced = reduced_arc.write().await;
reduced.apply_event(ev);
}
- s.complete = Some((canon_user, bearer));
+ s.complete = Some((username, bearer));
return Redirect::temporary(&format!("{public_url}/auth/complete")).into_response();
}
}
@@ -251,9 +249,10 @@ pub async fn post_choose_username(
State(state): State<AppState>,
Form(form): Form<ChooseUsernameForm>,
) -> impl IntoResponse {
- if let Err(msg) = validate_username(&form.username) {
- return api_error(StatusCode::BAD_REQUEST, "invalid username", Some(msg)).into_response();
- }
+ let canon_user = match parse_username(&form.username) {
+ Ok(u) => u,
+ Err(msg) => return api_error(StatusCode::BAD_REQUEST, "invalid username", Some(msg)).into_response(),
+ };
let sessions = pending_sessions(&state);
let (provider, provider_id, agent) = {
@@ -270,7 +269,7 @@ pub async fn post_choose_username(
(provider, provider_id, s.agent.clone())
};
- if let Err(msg) = validate_agent_format(&agent) {
+ if let Err(msg) = parse_agent(&agent) {
return api_error(StatusCode::BAD_REQUEST, "invalid agent format", Some(msg)).into_response();
}
@@ -280,7 +279,7 @@ pub async fn post_choose_username(
if reduced.users_by_provider.contains_key(&provider_key) {
return api_error(StatusCode::CONFLICT, "provider already registered", None).into_response();
}
- if reduced.users_by_provider.values().any(|u| u == &canonicalize_username(&form.username)) {
+ if reduced.users_by_provider.values().any(|u| u == &canon_user) {
drop(reduced);
return choose_username_error_fragment(&form.session, "that username is taken — try another").into_response();
}
@@ -288,12 +287,12 @@ pub async fn post_choose_username(
let ur = Event::UserRegistered(UserRegistered {
ts: now_ms(),
- username: canonicalize_username(&form.username),
+ username: canon_user.clone(),
provider: provider.to_lowercase(),
provider_id: provider_id.clone(),
});
- let (bearer, ti, canon_user) = issue_token_for_user(&form.username);
+ let (bearer, ti) = issue_token_for_user(&canon_user);
let ti_ev = Event::TokenIssued(ti);
// Persist events.
@@ -325,15 +324,18 @@ pub async fn post_pending_session(
State(state): State<AppState>,
Json(req): Json<PendingSessionStartRequest>,
) -> impl IntoResponse {
- if let Err(msg) = validate_agent_format(&req.agent) {
- return api_error(StatusCode::BAD_REQUEST, "invalid agent format", Some(msg)).into_response();
- }
+ let agent_naked = match parse_agent(&req.agent) {
+ Ok(a) => a,
+ Err(msg) => {
+ return api_error(StatusCode::BAD_REQUEST, "invalid agent format", Some(msg)).into_response();
+ }
+ };
let session = format!("p_{}", uuid::Uuid::new_v4().simple());
let public_url = std::env::var("SLUG_PUBLIC_URL").unwrap_or_else(|_| "http://127.0.0.1:8080".to_string());
let login_url = format!("{public_url}/auth/login?session={}", urlencoding::encode(&session));
let poll_url = format!("/api/v0/pending-session/{}", session);
let s = PendingSession {
- agent: req.agent.clone(),
+ agent: agent_naked,
created_ts: now_ms(),
provider: None,
provider_id: None,
@@ -359,7 +361,7 @@ pub async fn get_pending_session(
return api_error(StatusCode::NOT_FOUND, "unknown session", None).into_response();
};
let (complete, user, token) = match &s.complete {
- Some((u, t)) => (true, Some(format!("@{}", u)), Some(t.clone())),
+ Some((u, t)) => (true, Some(u.clone()), Some(t.clone())),
None => (false, None, None),
};
Json(PendingSessionPollResponse {
@@ -388,7 +390,7 @@ pub async fn get_whoami(State(state): State<AppState>, headers: HeaderMap)
… preview truncated; 62,772 characters omittedB — c_ca9169f732b8 (tommy-mor)
message
[8f69c309] Require Reddit OAuth when credentials are set and refresh on 401/403. Avoid falling back to the public www.reddit.com API from cloud IPs, which returns Reddit's network-security block page. Also pin SORTER2_BASE_URL in fly.toml. Co-authored-by: Cursor <cursoragent@cursor.com>
diff preview
diff --git a/fly.toml b/fly.toml
index f0c6a39f643c204987debc177234d95a8ef44b65..ca7e0088a7efd7d58d29d808f8d89080b2bff233 100644
--- a/fly.toml
+++ b/fly.toml
@@ -5,6 +5,7 @@ primary_region = "iad"
dockerfile = "Dockerfile"
[env]
+ SORTER2_BASE_URL = "https://reddit.sorter.social"
SORTER2_DATA_DIR = "/data"
SORTER2_EVENT_LOG = "/data/events.jsonl"
PORT = "8080"
diff --git a/server/src/reddit.rs b/server/src/reddit.rs
index a874814f8927192ee62cab2d0db1efd27dcd57b7..f409764c1e1f36216f1b08107043c2eab905694c 100644
--- a/server/src/reddit.rs
+++ b/server/src/reddit.rs
@@ -283,26 +283,35 @@ async fn reddit_worker(
);
tokio::time::sleep(current_delay).await;
- if let Some(c) = &creds {
- oauth = ensure_oauth_token(&client, &oauth_token_base, c, oauth.take()).await;
- }
-
- let token = oauth.as_ref().map(|t| t.access_token.as_str());
- let fetch_base = if token.is_some() {
- tracing::debug!(
- item = %fetch_id,
- base = %oauth_api_base,
- "reddit fetch using OAuth bearer"
- );
- &oauth_api_base
- } else {
- &api_base
- };
- let url = match kind {
- FetchKind::SelfEntity => map_item_to_reddit_api(&fetch_id, fetch_base),
- FetchKind::Children => map_children_url(&fetch_id, fetch_base),
+ let outcome = match &creds {
+ Some(c) => {
+ // OAuth is required when credentials are configured — never fall
+ // back to the public www.reddit.com JSON endpoints (cloud IPs
+ // get blocked with a 403 HTML interstitial).
+ fetch_with_oauth(
+ &client,
+ &oauth_token_base,
+ &oauth_api_base,
+ c,
+ &mut oauth,
+ &fetch_id,
+ kind,
+ )
+ .await
+ }
+ None => {
+ let url = match kind {
+ FetchKind::SelfEntity => map_item_to_reddit_api(&fetch_id, &api_base),
+ FetchKind::Children => map_children_url(&fetch_id, &api_base),
+ };
+ match do_fetch(&client, &url, &fetch_id, None).await {
+ Ok(FetchOutcome::AuthRejected { status, detail }) => {
+ Err(format!("Reddit API {status}: {detail}"))
+ }
+ other => other,
+ }
+ }
};
- let outcome = do_fetch(&client, &url, &fetch_id, token).await;
match outcome {
Ok(FetchOutcome::Payload(payload)) => {
@@ -342,6 +351,12 @@ async fn reddit_worker(
current_delay = (current_delay * 2).min(Duration::from_secs(60));
notify(done, FetchJobResult::RateLimited { reset_secs });
}
+ Ok(FetchOutcome::AuthRejected { status, detail }) => {
+ let e = format!("Reddit API {status}: {detail}");
+ tracing::warn!(item = %fetch_id, err = %e, "reddit fetch auth rejected");
+ current_delay = (current_delay * 2).min(Duration::from_secs(60));
+ notify(done, FetchJobResult::Failed(e));
+ }
Err(e) => {
tracing::warn!(item = %fetch_id, err = %e, "reddit fetch failed");
current_delay = (current_delay * 2).min(Duration::from_secs(60));
@@ -357,6 +372,60 @@ enum FetchOutcome {
Payload(Value),
NotFound,
RateLimited { reset_secs: u64 },
+ /// Bearer rejected — caller should drop the cached token and retry once.
+ AuthRejected { status: StatusCode, detail: String },
+}
+
+async fn fetch_with_oauth(
+ client: &Client,
+ oauth_token_base: &str,
+ oauth_api_base: &str,
+ creds: &RedditCredentials,
+ oauth: &mut Option<OAuthToken>,
+ fetch_id: &ItemId,
+ kind: FetchKind,
+) -> Result<FetchOutcome, String> {
+ for attempt in 0..2 {
+ let force_refresh = attempt > 0;
+ *oauth = Some(
+ ensure_oauth_token(client, oauth_token_base, creds, oauth.take(), force_refresh)
+ .await?,
+ );
+ let token = oauth
+ .as_ref()
+ .expect("token set above")
+ .access_token
+ .clone();
+
+ tracing::debug!(
+ item = %fetch_id,
+ base = %oauth_api_base,
+ attempt,
+ "reddit fetch using OAuth bearer"
+ );
+
+ let url = match kind {
+ FetchKind::SelfEntity => map_item_to_reddit_api(fetch_id, oauth_api_base),
+ FetchKind::Children => map_children_url(fetch_id, oauth_api_base),
+ };
+ match do_fetch(client, &url, fetch_id, Some(&token)).await? {
+ FetchOutcome::AuthRejected { status, detail } if attempt == 0 => {
+ tracing::warn!(
+ item = %fetch_id,
+ %status,
+ %detail,
+ "reddit OAuth rejected; refreshing token and retrying"
+ );
+ *oauth = None;
+ continue;
+ }
+ FetchOutcome::AuthRejected { status, detail } => {
+ return Err(format!("Reddit API {status}: {detail}"));
+ }
+ other => return Ok(other),
+ }
+ }
+ unreachable!("loop always returns")
}
async fn ensure_oauth_token(
@@ -364,35 +433,35 @@ async fn ensure_oauth_token(
oauth_base: &str,
creds: &RedditCredentials,
existing: Option<OAuthToken>,
-) -> Option<OAuthToken> {
- if let Some(t) = existing {
- if Instant::now() < t.expires_at - Duration::from_secs(60) {
- tracing::debug!("reddit OAuth token still valid");
- return Some(t);
+ force_refresh: bool,
+) -> Result<OAuthToken, String> {
+ if !force_refresh {
+ if let Some(t) = existing {
+ if Instant::now() < t.expires_at - Duration::from_secs(60) {
+ tracing::debug!("reddit OAuth token still valid");
+ return Ok(t);
+ }
}
}
let url = format!("{}/api/v1/access_token", oauth_base.trim_end_matches('/'));
- tracing::debug!(%url, "reddit OAuth token request");
+ tracing::debug!(%url, force_refresh, "reddit OAuth token request");
let resp = client
.post(&url)
.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;
- }
- };
+ .await
+ .map_err(|e| format!("Reddit OAuth token request failed: {e}"))?;
if !resp.status().is_success() {
- tracing::warn!("reddit OAuth token HTTP {}", resp.status());
- return None;
+ let status = resp.status();
+ let body = resp.text().await.unwrap_or_default();
+ return Err(format!(
+ "Reddit OAuth token HTTP {status}: {}",
+ truncate_for_error(&body)
+ ));
}
#[derive(Deserialize)]
@@ -401,21 +470,40 @@ async fn ensure_oauth_token(
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;
- }
- };
+ let body: TokenResponse = resp
+ .json()
+ .await
+ .map_err(|e| format!("Reddit OAuth token parse failed: {e}"))?;
- tracing::debug!(expires_in = body.expires_in, "reddit OAuth token acquired");
- Some(OAuthToken {
+ tracing::info!(expires_in = body.expires_in, "reddit OAuth token acquired");
+ Ok(OAuthToken {
access_token: body.access_token,
expires_at: Instant::now() + Duration::from_secs(body.expires_in),
})
}
+fn truncate_for_error(body: &str) -> String {
+ let compact: String = body.split_whitespace().collect::<Vec<_>>().join(" ");
+ if compact.is_empty() {
+ return "(empty body)".into();
+ }
+ // Prefer the human-readable block message over dumping Reddit's CSS.
+ if let Some(idx) = compact.find("You've been blocked") {
+ let slice: String = compact.chars().skip(idx).take(160).collect();
+ return if compact.chars().count() > idx + 160 {
+ format!("{slice}…")
+ } else {
+ slice
+ };
+ }
+ let chars: String = compact.chars().take(200).collect();
+ if compact.chars().count() > 200 {
+ format!("{chars}…")
+ } else {
+ chars
+ }
+}
+
async fn do_fetch(
client: &Client,
url: &str,
@@ -460,15 +548,16 @@ async fn do_fetch(
if !status.is_success() {
let body = resp.text().await.unwrap_or_default();
+ let detail = truncate_for_error(&body);
tracing::debug!(
item = %id,
%status,
body_len = body.len(),
- body_prefix = %body.chars().take(240).collect::<String>(),
+ %detail,
"reddit non-success body"
);
if status == StatusCode::FORBIDDEN || status == StatusCode::UNAUTHORIZED {
- return Err(format!("Reddit API {status}: {body}"));
+ return Ok(FetchOutcome::AuthRejected { status, detail });
}
return Ok(FetchOutcome::NotFound);
}
@@ -716,6 +805,15 @@ fn reddit_direct_image_url(url: &str) -> bool {
mod tests {
use super::*;
+ #[test]
+ fn truncate_error_prefers_block_message() {
+ let html = r#"<style>.x{color:red}</style><div>You've been blocked by network security. To continue, log in</div>"#;
+ let msg = truncate_for_error(html);
+ assert!(msg.starts_with("You've been blocked"));
+ assert!(msg.len() < 200);
+ assert!(!msg.contains(".x{color"));
+ }
+
#[test]
fn map_subreddit_about_url() {
let id = ItemId::from_url("https://reddit.com/r/rust").unwrap();
Hardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.