{"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[00be3a29] invite system\n\nSide A — unified diff (full patch):\ndiff --git a/bb.edn b/bb.edn\nindex 50be8232847e672b3f273a2fb25ddd1d12adb7e2..818f850765370d12d58f239e287764fc3649d78b 100644\n--- a/bb.edn\n+++ b/bb.edn\n@@ -47,14 +47,16 @@\n \"RUST_LOG\" \"info\"})})))}\n \n test\n- {:doc \"Full test suite: integration + auth + grants\"\n+ {:doc \"Full test suite: integration + auth + grants + invites\"\n :requires ([test.integration :as integration]\n [test.auth :as auth]\n- [test.grants :as grants])\n+ [test.grants :as grants]\n+ [test.invites :as invites])\n :task (do\n (integration/integration)\n (auth/auth-test)\n- (grants/grants-test))}\n+ (grants/grants-test)\n+ (invites/invites-test))}\n \n perf\n {:doc \"Performance test: concurrent HTTP requests to detect blocking I/O\"\ndiff --git a/cli/src/main.rs b/cli/src/main.rs\nindex e5833b0b93d8e667b94c574ba2b0f8cb758ff3df..8eda9f485bd1f7392f1e34be27176c21e20354eb 100644\n--- a/cli/src/main.rs\n+++ b/cli/src/main.rs\n@@ -105,6 +105,23 @@ enum ScopedCmd {\n #[arg(long)]\n json: bool,\n },\n+\n+ /// Mint a shareable invite link (24h TTL, in-memory until redeemed). Requires Manage on the room.\n+ InviteLink {\n+ /// Comma-separated: view, post, vote, add_item, manage\n+ #[arg(long = \"caps\", value_delimiter = ',')]\n+ caps: Vec,\n+ #[arg(long, default_value_t = 1)]\n+ uses: usize,\n+ #[arg(long)]\n+ json: bool,\n+ },\n+\n+ /// List principals granted access in this room (requires View or Manage)\n+ Audit {\n+ #[arg(long)]\n+ json: bool,\n+ },\n }\n \n #[derive(Subcommand, Debug)]\n@@ -514,17 +531,35 @@ fn print_thread(resp: &ThreadDetailResponse) {\n .duration_since(std::time::UNIX_EPOCH)\n .unwrap_or_default()\n .as_millis() as i64;\n- if resp.total > resp.posts.len() {\n- let end = resp.offset + resp.posts.len();\n- eprintln!(\"# showing {}-{} of {} posts (--offset N --limit N to paginate)\", resp.offset, end.saturating_sub(1), resp.total);\n+ if resp.total > resp.items.len() {\n+ let end = resp.offset + resp.items.len();\n+ eprintln!(\n+ \"# showing {}-{} of {} rows (--offset N --limit N to paginate)\",\n+ resp.offset,\n+ end.saturating_sub(1),\n+ resp.total\n+ );\n }\n- for (i, post) in resp.posts.iter().enumerate() {\n- let timeago = slug_types::timeago::timeago_compact(now_ms, post.ts);\n- let body = &post.body.trim();\n- println!(\"\", post.index, timeago);\n- println!(\"{}\", body);\n- println!(\"\");\n- if i + 1 < resp.posts.len() {\n+ for (i, item) in resp.items.iter().enumerate() {\n+ match item {\n+ ThreadItem::Post {\n+ index,\n+ ts,\n+ body,\n+ ..\n+ } => {\n+ let timeago = slug_types::timeago::timeago_compact(now_ms, *ts);\n+ let body = body.trim();\n+ println!(\"\", index, timeago);\n+ println!(\"{}\", body);\n+ println!(\"\");\n+ }\n+ ThreadItem::System { ts, text } => {\n+ let timeago = slug_types::timeago::timeago_compact(now_ms, *ts);\n+ println!(\"{}\", timeago, text.trim());\n+ }\n+ }\n+ if i + 1 < resp.items.len() {\n println!();\n println!();\n }\n@@ -1036,6 +1071,95 @@ async fn run_scoped(base: &str, room: &str, sub: ScopedCmd) -> Result<()> {\n }\n }\n },\n+ ScopedCmd::InviteLink { caps, uses, json } => {\n+ let caps: Vec = caps\n+ .into_iter()\n+ .flat_map(|s| {\n+ s.split(',')\n+ .map(|p| p.trim().to_lowercase())\n+ .filter(|p| !p.is_empty())\n+ .collect::>()\n+ })\n+ .collect();\n+ if caps.is_empty() {\n+ return Err(anyhow!(\"--caps is required (e.g. --caps view,post,vote)\"));\n+ }\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 batch = send_rpc(\n+ &client,\n+ base,\n+ Some(&bearer),\n+ vec![RpcCommand::RoomMintInvite {\n+ room: room.to_string(),\n+ capabilities: caps,\n+ max_uses: uses,\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::RoomInviteMinted {\n+ invite_url,\n+ expires_at_ms,\n+ max_uses,\n+ } => {\n+ if json {\n+ println!(\n+ \"{}\",\n+ serde_json::to_string_pretty(&serde_json::json!({\n+ \"invite_url\": invite_url,\n+ \"expires_at_ms\": expires_at_ms,\n+ \"max_uses\": max_uses,\n+ }))?\n+ );\n+ } else {\n+ println!(\"{invite_url}\");\n+ println!(\"(Expires in 24 hours. Max uses: {max_uses})\");\n+ }\n+ }\n+ _ => return Err(anyhow!(\"unexpected RPC result\")),\n+ }\n+ }\n+ ScopedCmd::Audit { json } => {\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 batch = send_rpc(\n+ &client,\n+ base,\n+ Some(&bearer),\n+ vec![RpcCommand::RoomAudit {\n+ room: room.to_string(),\n+ }],\n+ )\n+ .await?;\n+ match rpc_line_ok(&batch.results[0])? {\n+ RpcResult::RoomAudit(resp) => {\n+ if json {\n+ println!(\"{}\", serde_json::to_string_pretty(&resp)?);\n+ } else {\n+ println!(\"room {}\", resp.room);\n+ if resp.grants.is_empty() {\n+ println!(\"(no grants recorded)\");\n+ } else {\n+ let w_user = resp.grants.iter().map(|g| g.username.len()).max().unwrap_or(0);\n+ for g in &resp.grants {\n+ let caps = g.capabilities.join(\", \");\n+ println!(\"{: return Err(anyhow!(\"unexpected RPC result\")),\n+ }\n+ }\n ScopedCmd::Check { file, json } => {\n let mut text = String::new();\n match file {\ndiff --git a/server/src/api/auth.rs b/server/src/api/auth.rs\nindex 995ce4a61d29b024c399c656134f541ecfd880cf..b45ba39419c84af8bf2333fc9b7d47e98525c45f 100644\n--- a/server/src/api/auth.rs\n+++ b/server/src/api/auth.rs\n@@ -12,12 +12,61 @@ use tokio::sync::RwLock;\n \n use crate::{\n api::helpers::{api_error, now_ms, sha256_hex},\n- events::{Event, TokenIssued, UserRegistered},\n+ events::{Event, GrantAdded, TokenIssued, UserRegistered},\n identity::{parse_agent, parse_username},\n html::{auth_complete_page, auth_signed_in_fragment, choose_username_error_fragment, choose_username_page},\n state::{AppState, PendingSession},\n };\n \n+/// Delegate id for browser users who land via `/join/inv_…` (no CLI agent).\n+const INVITE_BROWSER_AGENT: &str = \"00000000-0000-0000-0000-000000000000:invite:web/join\";\n+\n+async fn apply_invite_redemption(state: &AppState, invite_token: &str, grantee_username: &str) -> Result<(), String> {\n+ let now = now_ms();\n+ let ga = {\n+ let mut invites = state.invites.write().await;\n+ let Some(inv) = invites.get_mut(invite_token) else {\n+ return Err(\"invite not found\".into());\n+ };\n+ if now > inv.expires_at_ms {\n+ invites.remove(invite_token);\n+ return Err(\"invite expired\".into());\n+ }\n+ if inv.current_uses >= inv.max_uses {\n+ return Err(\"invite exhausted\".into());\n+ }\n+ inv.current_uses += 1;\n+ Event::GrantAdded(GrantAdded {\n+ ts: now,\n+ room_id: inv.room_id.clone(),\n+ username: grantee_username.to_string(),\n+ capabilities: inv.capabilities.clone(),\n+ granted_by: inv.inviter.clone(),\n+ })\n+ };\n+\n+ match state.event_log.append(&ga).await {\n+ Ok(()) => {\n+ let mut reduced = state.reduced.write().await;\n+ reduced.apply_event(ga);\n+ let mut invites = state.invites.write().await;\n+ if let Some(inv) = invites.get(invite_token) {\n+ if inv.current_uses >= inv.max_uses {\n+ invites.remove(invite_token);\n+ }\n+ }\n+ Ok(())\n+ }\n+ Err(e) => {\n+ let mut invites = state.invites.write().await;\n+ if let Some(inv) = invites.get_mut(invite_token) {\n+ inv.current_uses = inv.current_uses.saturating_sub(1);\n+ }\n+ Err(format!(\"{e}\"))\n+ }\n+ }\n+}\n+\n fn pending_sessions(state: &AppState) -> Arc>> {\n state.pending_sessions.clone()\n }\n@@ -115,6 +164,42 @@ pub struct AuthLoginQuery {\n pub session: String,\n }\n \n+pub async fn get_join_invite(Path(token): Path, State(state): State) -> impl IntoResponse {\n+ let token = token.trim().to_string();\n+ if token.is_empty() {\n+ return api_error(StatusCode::NOT_FOUND, \"invite invalid or expired\", None).into_response();\n+ }\n+ let now = now_ms();\n+ let valid = {\n+ let invites = state.invites.read().await;\n+ match invites.get(&token) {\n+ None => false,\n+ Some(inv) => now <= inv.expires_at_ms && inv.current_uses < inv.max_uses,\n+ }\n+ };\n+ if !valid {\n+ return api_error(StatusCode::NOT_FOUND, \"invite invalid or expired\", None).into_response();\n+ }\n+\n+ let session = format!(\"p_{}\", uuid::Uuid::new_v4().simple());\n+ let s = PendingSession {\n+ agent: INVITE_BROWSER_AGENT.to_string(),\n+ created_ts: now_ms(),\n+ provider: None,\n+ provider_id: None,\n+ redeem_invite: Some(token),\n+ complete: None,\n+ };\n+ state.pending_sessions.write().await.insert(session.clone(), s);\n+\n+ let public_url = std::env::var(\"SLUG_PUBLIC_URL\").unwrap_or_else(|_| \"http://127.0.0.1:8080\".to_string());\n+ Redirect::temporary(&format!(\n+ \"{public_url}/auth/login?session={}\",\n+ urlencoding::encode(&session)\n+ ))\n+ .into_response()\n+}\n+\n pub async fn get_auth_login(Query(q): Query, State(state): State) -> impl IntoResponse {\n // Redirect to Google auth endpoint.\n let sessions = pending_sessions(&state);\n@@ -205,6 +290,7 @@ pub async fn get_auth_callback(Query(q): Query, State(state):\n s.provider = Some(\"google\".to_string());\n s.provider_id = Some(sub.clone());\n if let Some(username) = existing {\n+ let invite_tok = s.redeem_invite.clone();\n let (bearer, token_event) = issue_token_for_user(&username);\n // append token event\n let ev = Event::TokenIssued(token_event);\n@@ -215,6 +301,11 @@ pub async fn get_auth_callback(Query(q): Query, State(state):\n let mut reduced = reduced_arc.write().await;\n reduced.apply_event(ev);\n }\n+ if let Some(tok) = invite_tok {\n+ if let Err(e) = apply_invite_redemption(&state, &tok, &username).await {\n+ tracing::warn!(error = %e, \"invite redemption skipped after oauth\");\n+ }\n+ }\n s.complete = Some((username, bearer));\n return Redirect::temporary(&format!(\"{public_url}/auth/complete\")).into_response();\n }\n@@ -310,6 +401,16 @@ pub async fn post_choose_username(\n reduced.apply_event(ti_ev.clone());\n }\n \n+ // Redeem invite (if any) before marking the session complete.\n+ if let Some(tok) = {\n+ let sessions_read = sessions.read().await;\n+ sessions_read.get(&form.session).and_then(|s| s.redeem_invite.clone())\n+ } {\n+ if let Err(e) = apply_invite_redemption(&state, &tok, &canon_user).await {\n+ tracing::warn!(error = %e, \"invite redemption skipped after registration\");\n+ }\n+ }\n+\n // Mark complete for polling.\n {\n let mut sessions_write = sessions.write().await;\n@@ -339,6 +440,7 @@ pub async fn post_pending_session(\n created_ts: now_ms(),\n provider: None,\n provider_id: None,\n+ redeem_invite: None,\n complete: None,\n };\n let sessions = pending_sessions(&state);\ndiff --git a/server/src/api/mod.rs b/server/src/api/mod.rs\nindex a1ea432932cdca0e42dc7377f0a25075fd5644c0..cb031aecebad1107dfa2898a98fb6084b28ba6e9 100644\n--- a/server/src/api/mod.rs\n+++ b/server/src/api/mod.rs\n@@ -4,6 +4,7 @@ mod rpc;\n mod validate;\n \n pub use auth::{\n+ get_join_invite,\n get_pending_session,\n get_whoami,\n post_pending_session,\ndiff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs\nindex 11a76ed7a02aae588edadb7a87a43660f5669d9e..f6bbc3df71909a2da7403cd46fe4ea6ca130c692 100644\n--- a/server/src/api/rpc.rs\n+++ b/server/src/api/rpc.rs\n@@ -13,12 +13,14 @@ use slug_types::*;\n use crate::{\n canonical_path::{canonicalize_item, canonicalize_tag},\n dsl,\n- events::{AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability, ThreadVisibility},\n+ events::{\n+ AgentBound, Event, GrantAdded, Ingest, RoomCreated, ThreadCapability, ThreadVisibility,\n+ },\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+ state::{AppState, InviteState},\n };\n \n use super::auth::verify_bearer_principal;\n@@ -142,6 +144,27 @@ fn gen_short_id() -> String {\n (0..7).map(|_| ALPHABET[rng.gen_range(0..ALPHABET.len())] as char).collect()\n }\n \n+fn gen_invite_token() -> String {\n+ use rand::Rng;\n+ const ALPHABET: &[u8] = b\"0123456789abcdefghijklmnopqrstuvwxyz\";\n+ let mut rng = rand::thread_rng();\n+ let tail: String = (0..16).map(|_| ALPHABET[rng.gen_range(0..ALPHABET.len())] as char).collect();\n+ format!(\"inv_{tail}\")\n+}\n+\n+const INVITE_TTL_MS: i64 = 86_400_000;\n+\n+fn capability_wire(c: ThreadCapability) -> String {\n+ match c {\n+ ThreadCapability::View => \"view\",\n+ ThreadCapability::Post => \"post\",\n+ ThreadCapability::Vote => \"vote\",\n+ ThreadCapability::AddItem => \"add_item\",\n+ ThreadCapability::Manage => \"manage\",\n+ }\n+ .to_string()\n+}\n+\n fn build_rank_response_for_content(\n content: &crate::reducer::ContentState,\n parent: Option<&str>,\n@@ -512,7 +535,7 @@ fn rpc_forum_thread_detail(\n None => Err((\"post not found\".into(), None)),\n Some((idx, ing)) => Ok(ThreadDetailResponse {\n thread: format!(\"#{}\", tag),\n- posts: vec![PostRow {\n+ items: vec![ThreadItem::Post {\n id: ing.id.clone(),\n index: idx,\n ts: ing.ts,\n@@ -543,7 +566,7 @@ fn rpc_forum_thread_detail(\n \n let total = filtered.len();\n const MAX_BODY: usize = 2000;\n- let posts: Vec = filtered\n+ let items: Vec = filtered\n .into_iter()\n .skip(offset)\n .take(limit)\n@@ -553,7 +576,7 @@ fn rpc_forum_thread_detail(\n } else {\n (ing.raw.clone(), false)\n };\n- PostRow {\n+ ThreadItem::Post {\n id: ing.id.clone(),\n index: idx,\n ts: ing.ts,\n@@ -566,7 +589,7 @@ fn rpc_forum_thread_detail(\n \n Ok(ThreadDetailResponse {\n thread: format!(\"#{}\", tag),\n- posts,\n+ items,\n total,\n offset,\n })\n@@ -1007,7 +1030,7 @@ pub async fn handle_rpc_batch(\n RpcCommand::RoomGrant {\n room,\n username,\n- capability,\n+ capabilities,\n } => {\n let principal = {\n let reduced = state.reduced.read().await;\n@@ -1022,6 +1045,8 @@ pub async fn handle_rpc_batch(\n };\n if !can_manage {\n line_err(\"requires Manage capability\", None)\n+ } else if capabilities.is_empty() {\n+ line_err(\"capabilities must not be empty\", None)\n } else {\n match parse_username(&username) {\n Err(msg) => line_err(\"invalid username\", Some(msg)),\n@@ -1033,14 +1058,18 @@ pub async fn handle_rpc_batch(\n if !user_exists {\n line_err(format!(\"user @{target} not found\"), None)\n } else {\n- match parse_capability(&capability) {\n+ let caps: Result, String> = capabilities\n+ .iter()\n+ .map(|c| parse_capability(c.trim()))\n+ .collect();\n+ match caps {\n Err(msg) => line_err(msg, None),\n- Ok(cap) => {\n+ Ok(caps) => {\n let ga_ev = Event::GrantAdded(GrantAdded {\n ts: now_ms(),\n room_id: room,\n username: target,\n- capabilities: vec![cap],\n+ capabilities: caps,\n granted_by: principal,\n });\n if let Err(e) = state.event_log.append(&ga_ev).await {\n@@ -1059,6 +1088,114 @@ pub async fn handle_rpc_batch(\n }\n }\n }\n+ RpcCommand::RoomMintInvite {\n+ room,\n+ capabilities,\n+ max_uses,\n+ } => {\n+ let principal = {\n+ let reduced = state.reduced.read().await;\n+ verify_bearer_principal(&headers, &*reduced)\n+ };\n+ match principal {\n+ Err((_, m)) => line_err(m, None),\n+ Ok(principal) => {\n+ let can_manage = {\n+ let reduced = state.reduced.read().await;\n+ reduced.user_has_cap(&room, &principal, ThreadCapability::Manage)\n+ };\n+ if !can_manage {\n+ line_err(\"requires Manage capability\", None)\n+ } else if capabilities.is_empty() {\n+ line_err(\"capabilities must not be empty\", None)\n+ } else {\n+ match capabilities\n+ .iter()\n+ .map(|c| parse_capability(c.trim()))\n+ .collect::, String>>()\n+ {\n+ Err(msg) => line_err(msg, None),\n+ Ok(caps) => {\n+ let max_uses = max_uses.max(1).min(100_000);\n+ let now = now_ms();\n+ let expires_at_ms = now + INVITE_TTL_MS;\n+ let token = loop {\n+ let t = gen_invite_token();\n+ let taken = {\n+ let invites = state.invites.read().await;\n+ invites.contains_key(&t)\n+ };\n+ if !taken {\n+ break t;\n+ }\n+ };\n+ let inv = InviteState {\n+ room_id: room.clone(),\n+ capabilities: caps,\n+ expires_at_ms,\n+ max_uses,\n+ current_uses: 0,\n+ inviter: principal,\n+ };\n+ state.invites.write().await.insert(token.clone(), inv);\n+ let public_url = std::env::var(\"SLUG_PUBLIC_URL\")\n+ .unwrap_or_else(|_| \"http://127.0.0.1:8080\".to_string());\n+ let invite_url = format!(\"{public_url}/join/{token}\");\n+ line_ok(RpcResult::RoomInviteMinted {\n+ invite_url,\n+ expires_at_ms: Some(expires_at_ms),\n+ max_uses,\n+ })\n+ }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ RpcCommand::RoomAudit { room } => {\n+ let principal = {\n+ let reduced = state.reduced.read().await;\n+ verify_bearer_principal(&headers, &*reduced)\n+ };\n+ match principal {\n+ Err((_, m)) => line_err(m, None),\n+ Ok(principal) => {\n+ let reduced = state.reduced.read().await;\n+ if !reduced.rooms.contains_key(&room) {\n+ line_err(\"unknown room\", None)\n+ } else {\n+ let can_audit = reduced.user_has_cap(&room, &principal, ThreadCapability::View)\n+ || reduced.user_has_cap(&room, &principal, ThreadCapability::Manage);\n+ if !can_audit {\n+ line_err(\"requires View or Manage capability\", None)\n+ } else {\n+ let grants: Vec = reduced\n+ .grants\n+ .get(&room)\n+ .map(|m| {\n+ let mut v: Vec = m\n+ .iter()\n+ .map(|(username, caps)| {\n+ let mut c: Vec =\n+ caps.iter().copied().map(capability_wire).collect();\n+ c.sort();\n+ RoomAuditEntry {\n+ username: username.clone(),\n+ capabilities: c,\n+ }\n+ })\n+ .collect();\n+ v.sort_by(|a, b| a.username.cmp(&b.username));\n+ v\n+ })\n+ .unwrap_or_default();\n+ line_ok(RpcResult::RoomAudit(RoomAuditResponse { room, grants }))\n+ }\n+ }\n+ }\n+ }\n+ }\n+ RpcCommand::RoomRevoke { .. } => line_err(\"RoomRevoke is not implemented yet\", None),\n RpcCommand::GetGlobalRank {\n room,\n limit,\ndiff --git a/server/src/events.rs b/server/src/events.rs\nindex 125005d89ade7b22c9b6d997f0ed81781df5e7ec..9e60f2c218f5c871a93e855e415dd984e3a017c5 100644\n--- a/server/src/events.rs\n+++ b/server/src/events.rs\n@@ -26,6 +26,8 @@ pub enum Event {\n RoomCreated(RoomCreated),\n GrantAdded(GrantAdded),\n GrantRevoked(GrantRevoked),\n+ InviteMinted(InviteMinted),\n+ InviteRedeemed(InviteRedeemed),\n /// Ingest of a DSL+prose body. Identity and routing live in event metadata.\n Ingest(Ingest),\n }\n@@ -82,6 +84,25 @@ pub struct GrantRevoked {\n pub revoked_by: String,\n }\n \n+#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n+pub struct InviteMinted {\n+ pub ts: i64,\n+ pub token: String,\n+ pub room_id: String,\n+ pub capabilities: Vec,\n+ pub inviter: String,\n+ pub max_uses: u32,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ pub expires_ts_ms: Option,\n+}\n+\n+#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n+pub struct InviteRedeemed {\n+ pub ts: i64,\n+ pub token: String,\n+ pub username: String,\n+}\n+\n #[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]\n pub struct Ingest {\n /// Unix timestamp in milliseconds.\ndiff --git a/server/src/lib.rs b/server/src/lib.rs\nindex d00b56d909e1fabd1d122d054dd02785e13e0bd8..bf2f73d4d0f26a357961e88c52c3b4f72626af81 100644\n--- a/server/src/lib.rs\n+++ b/server/src/lib.rs\n@@ -27,6 +27,7 @@ pub fn create_app(state: AppState) -> Router {\n Router::new()\n .route(\"/healthz\", axum::routing::get(|| async { \"ok\" }))\n .route(\"/static/:filename\", axum::routing::get(crate::html::serve_theme_css))\n+ .route(\"/join/:token\", axum::routing::get(api::get_join_invite))\n .route(\"/auth/login\", axum::routing::get(api::get_auth_login))\n .route(\"/auth/callback\", axum::routing::get(api::get_auth_callback))\n .route(\"/auth/complete\", axum::routing::get(api::get_auth_complete))\ndiff --git a/server/src/reducer.rs b/server/src/reducer.rs\nindex 118cb19f4df4a083df2878ca7a1424f3d63de851..5190d0a7ee3d545a786a68fef483322aafdce7f9 100644\n--- a/server/src/reducer.rs\n+++ b/server/src/reducer.rs\n@@ -3,7 +3,7 @@ use std::collections::{HashMap, HashSet, VecDeque};\n use serde::{Deserialize, Serialize};\n \n use crate::canonical_path::canonicalize_tag;\n-use crate::events::{Event, Ingest, ThreadCapability};\n+use crate::events::{Event, Ingest, ThreadCapability, ThreadVisibility};\n use crate::path_types::CanonicalItemUrl;\n \n #[derive(Debug, Clone, Hash, PartialEq, Eq, PartialOrd, Ord)]\n@@ -156,6 +156,41 @@ pub struct RoomState {\n pub visibility: crate::events::ThreadVisibility,\n }\n \n+/// Durable invite link state (from [`crate::events::InviteMinted`] / [`crate::events::InviteRedeemed`]).\n+#[derive(Debug, Clone)]\n+pub struct ActiveInviteState {\n+ pub room_id: String,\n+ pub capabilities: HashSet,\n+ pub inviter: String,\n+ pub uses_remaining: u32,\n+ pub expires_ts_ms: Option,\n+}\n+\n+#[derive(Clone, Debug)]\n+pub enum RoomTimelineKind {\n+ RoomCreated {\n+ owner: String,\n+ slug: String,\n+ visibility: ThreadVisibility,\n+ },\n+ GrantAdded {\n+ username: String,\n+ granted_by: String,\n+ capabilities: Vec,\n+ },\n+ GrantRevoked {\n+ username: String,\n+ revoked_by: String,\n+ capabilities: Vec,\n+ },\n+}\n+\n+#[derive(Clone, Debug)]\n+pub struct RoomTimelineEntry {\n+ pub ts: i64,\n+ pub kind: RoomTimelineKind,\n+}\n+\n #[derive(Debug, Clone)]\n pub struct ForumThreadState {\n pub last_activity_ts: i64,\n@@ -210,6 +245,10 @@ pub struct ReducerState {\n pub ingests_ordered: Vec,\n /// room_id → username → capabilities\n pub grants: HashMap>>,\n+ /// room_id → chronological room admin lines (for thread UI).\n+ pub room_timeline: HashMap>,\n+ /// Invite token → active invite (absent when fully consumed or never minted).\n+ pub invites: HashMap,\n }\n \n impl ReducerState {\n@@ -225,6 +264,20 @@ impl ReducerState {\n .unwrap_or(false)\n }\n \n+ /// Invite link is present, not expired, and has uses left.\n+ pub fn invite_token_active(&self, token: &str, now_ms: i64) -> Option<&ActiveInviteState> {\n+ let inv = self.invites.get(token)?;\n+ if inv.uses_remaining == 0 {\n+ return None;\n+ }\n+ if let Some(exp) = inv.expires_ts_ms {\n+ if now_ms > exp {\n+ return None;\n+ }\n+ }\n+ Some(inv)\n+ }\n+\n pub fn content_for_scope_mut(&mut self, scope: ScopeId) -> &mut ContentState {\n self.content.entry(scope).or_default()\n }\n@@ -356,6 +409,17 @@ impl ReducerState {\n visibility: rc.visibility,\n },\n );\n+ self.room_timeline\n+ .entry(rc.room_id.clone())\n+ .or_default()\n+ .push(RoomTimelineEntry {\n+ ts: rc.ts,\n+ kind: RoomTimelineKind::RoomCreated {\n+ owner: rc.owner.clone(),\n+ slug: rc.slug.clone(),\n+ visibility: rc.visibility,\n+ },\n+ });\n }\n Event::Ingest(mut ing) => {\n ing.thread_tag = canonicalize_tag(&ing.thread_tag);\n@@ -525,18 +589,31 @@ impl ReducerState {\n nav!(self.actor_last_post_ts, keypath(ing.principal.clone()), setval(ing.ts));\n }\n Event::GrantAdded(ga) => {\n+ let room_id = ga.room_id.clone();\n let caps = self.grants\n .entry(ga.room_id)\n .or_default()\n- .entry(ga.username)\n+ .entry(ga.username.clone())\n .or_default();\n- for cap in ga.capabilities {\n+ for cap in ga.capabilities.iter().copied() {\n caps.insert(cap);\n }\n+ self.room_timeline\n+ .entry(room_id)\n+ .or_default()\n+ .push(RoomTimelineEntry {\n+ ts: ga.ts,\n+ kind: RoomTimelineKind::GrantAdded {\n+ username: ga.username.clone(),\n+ granted_by: ga.granted_by.clone(),\n+ capabilities: ga.capabilities.clone(),\n+ },\n+ });\n }\n Event::GrantRevoked(gr) => {\n+ let room_id = gr.room_id.clone();\n if let Some(room_grants) = self.grants.get_mut(&gr.room_id) {\n- let username = gr.username;\n+ let username = gr.username.clone();\n if let Some(caps) = room_grants.get_mut(&username) {\n for cap in &gr.capabilities {\n caps.remove(cap);\n@@ -549,6 +626,37 @@ impl ReducerState {\n self.grants.remove(&gr.room_id);\n }\n }\n+ self.room_timeline\n+ .entry(room_id)\n+ .or_default()\n+ .push(RoomTimelineEntry {\n+ ts: gr.ts,\n+ kind: RoomTimelineKind::GrantRevoked {\n+ username: gr.username.clone(),\n+ revoked_by: gr.revoked_by.clone(),\n+ capabilities: gr.capabilities.clone(),\n+ },\n+ });\n+ }\n+ Event::InviteMinted(im) => {\n+ self.invites.insert(\n+ im.token.clone(),\n+ ActiveInviteState {\n+ room_id: im.room_id.clone(),\n+ capabilities: im.capabilities.iter().copied().collect(),\n+ inviter: im.inviter.clone(),\n+ uses_remaining: im.max_uses,\n+ expires_ts_ms: im.expires_ts_ms,\n+ },\n+ );\n+ }\n+ Event::InviteRedeemed(ir) => {\n+ if let Some(inv) = self.invites.get_mut(&ir.token) {\n+ inv.uses_remaining = inv.uses_remaining.saturating_sub(1);\n+ if inv.uses_remaining == 0 {\n+ self.invites.remove(&ir.token);\n+ }\n+ }\n }\n }\n }\n@@ -570,6 +678,8 @@ impl Default for ReducerState {\n actor_last_post_ts: HashMap::new(),\n ingests_ordered: Vec::new(),\n grants: HashMap::new(),\n+ room_timeline: HashMap::new(),\n+ invites: HashMap::new(),\n }\n }\n }\ndiff --git a/server/src/state.rs b/server/src/state.rs\nindex b1ff2330903780cbe0cdf35964cb166eb2423d7d..628fd5921b34ba53ac0ca6b6fd33957ca09129b9 100644\n--- a/server/src/state.rs\n+++ b/server/src/state.rs\n@@ -1,8 +1,20 @@\n+use std::collections::HashMap;\n use std::sync::Arc;\n \n use tokio::sync::{broadcast, RwLock};\n \n-use crate::{event_log::EventLog, reducer::ReducerState};\n+use crate::{event_log::EventLog, events::ThreadCapability, reducer::ReducerState};\n+\n+/// Ephemeral invite link (24h TTL, in-memory only; not written to the event log).\n+#[derive(Debug, Clone)]\n+pub struct InviteState {\n+ pub room_id: String,\n+ pub capabilities: Vec,\n+ pub expires_at_ms: i64,\n+ pub max_uses: usize,\n+ pub current_uses: usize,\n+ pub inviter: String,\n+}\n \n #[derive(Debug, Clone)]\n pub struct PendingSession {\n@@ -10,6 +22,8 @@ pub struct PendingSession {\n pub created_ts: i64,\n pub provider: Option,\n pub provider_id: Option,\n+ /// When set, successful OAuth completion redeems this invite token and appends [`crate::events::GrantAdded`].\n+ pub redeem_invite: Option,\n pub complete: Option<(String /*username*/, String /*bearer*/ )>,\n }\n \n@@ -42,7 +56,9 @@ pub struct AppState {\n pub cfg: Arc,\n pub event_log: Arc,\n pub reduced: Arc>,\n- pub pending_sessions: Arc>>,\n+ pub pending_sessions: Arc>>,\n+ /// Ephemeral invite tokens (`inv_…`) until expiry or exhaustion.\n+ pub invites: Arc>>,\n /// Broadcast channel for SSE live-streaming. Capacity = 64 events.\n pub stream_tx: broadcast::Sender,\n /// Broadcast channel for web SSE HTML fragments (poem pattern). Capacity = 64.\n@@ -58,7 +74,8 @@ impl AppState {\n cfg: Arc::new(cfg),\n event_log: Arc::new(event_log),\n reduced: Arc::new(RwLock::new(ReducerState::default())),\n- pending_sessions: Arc::new(RwLock::new(std::collections::HashMap::new())),\n+ pending_sessions: Arc::new(RwLock::new(HashMap::new())),\n+ invites: Arc::new(RwLock::new(HashMap::new())),\n stream_tx,\n html_tx,\n }\ndiff --git a/server/src/timeline.rs b/server/src/timeline.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..251158943ab36d3015268e35f6bdd01c3f42ab3b\n--- /dev/null\n+++ b/server/src/timeline.rs\n@@ -0,0 +1,153 @@\n+//! Room admin lines merged into forum thread views.\n+\n+use crate::{\n+ canonical_path::canonicalize_tag,\n+ reducer::{ReducerState, RoomTimelineEntry, RoomTimelineKind},\n+};\n+\n+fn cap_label(c: crate::events::ThreadCapability) -> &'static str {\n+ use crate::events::ThreadCapability::*;\n+ match c {\n+ View => \"view\",\n+ Post => \"post\",\n+ Vote => \"vote\",\n+ AddItem => \"add_item\",\n+ Manage => \"manage\",\n+ }\n+}\n+\n+fn caps_list(caps: &[crate::events::ThreadCapability]) -> String {\n+ let mut v: Vec<_> = caps.iter().map(|c| cap_label(*c)).collect();\n+ v.sort();\n+ v.join(\", \")\n+}\n+\n+/// Human-readable system line for the thread feed.\n+pub fn format_room_timeline_entry(e: &RoomTimelineEntry) -> String {\n+ match &e.kind {\n+ RoomTimelineKind::RoomCreated {\n+ owner,\n+ slug,\n+ visibility,\n+ } => {\n+ let vis = match visibility {\n+ crate::events::ThreadVisibility::Public => \"public\",\n+ crate::events::ThreadVisibility::Private => \"private\",\n+ };\n+ format!(\"@{owner} created room #{slug} ({vis})\")\n+ }\n+ RoomTimelineKind::GrantAdded {\n+ username,\n+ granted_by,\n+ capabilities,\n+ } => {\n+ format!(\n+ \"@{granted_by} granted @{} {}\",\n+ username,\n+ caps_list(capabilities)\n+ )\n+ }\n+ RoomTimelineKind::GrantRevoked {\n+ username,\n+ revoked_by,\n+ capabilities,\n+ } => {\n+ format!(\n+ \"@{revoked_by} revoked @{} {}\",\n+ username,\n+ caps_list(capabilities)\n+ )\n+ }\n+ }\n+}\n+\n+#[derive(Clone, Debug)]\n+pub enum MergedThreadRow {\n+ System { ts: i64, text: String },\n+ Post {\n+ index: usize,\n+ id: String,\n+ ts: i64,\n+ principal: String,\n+ raw: String,\n+ },\n+}\n+\n+/// Merge room admin lines with thread ingests for one room + tag. Oldest first.\n+/// `actor_prefix` filters posts only (system lines always included).\n+pub fn merge_thread_rows(\n+ reduced: &ReducerState,\n+ room_wire: &str,\n+ thread_tag: &str,\n+ since: Option,\n+ before: Option,\n+ actor_prefix: &str,\n+) -> Vec {\n+ let scope = crate::reducer::scope_from_room_wire(room_wire);\n+ let tag = canonicalize_tag(thread_tag);\n+ let key = (scope.clone(), tag.clone());\n+\n+ let mut rows: Vec = Vec::new();\n+\n+ if let Some(entries) = reduced.room_timeline.get(room_wire.trim()) {\n+ for e in entries {\n+ if since.map_or(true, |s| e.ts >= s) && before.map_or(true, |b| e.ts < b) {\n+ rows.push(MergedThreadRow::System {\n+ ts: e.ts,\n+ text: format_room_timeline_entry(e),\n+ });\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+ for (idx, id) in all_ids.into_iter().enumerate() {\n+ let Some(ing) = reduced.ingests_by_id.get(&id) else {\n+ continue;\n+ };\n+ if since.map_or(true, |s| ing.ts >= s) && before.map_or(true, |b| ing.ts < b) {\n+ if !actor_prefix.is_empty()\n+ && !ing\n+ .principal\n+ .to_lowercase()\n+ .starts_with(actor_prefix)\n+ {\n+ continue;\n+ }\n+ rows.push(MergedThreadRow::Post {\n+ index: idx,\n+ id: ing.id.clone(),\n+ ts: ing.ts,\n+ principal: ing.principal.clone(),\n+ raw: ing.raw.clone(),\n+ });\n+ }\n+ }\n+\n+ rows.sort_by(|a, b| {\n+ let ta = match a {\n+ MergedThreadRow::System { ts, .. } | MergedThreadRow::Post { ts, .. } => *ts,\n+ };\n+ let tb = match b {\n+ MergedThreadRow::System { ts, .. } | MergedThreadRow::Post { ts, .. } => *ts,\n+ };\n+ ta.cmp(&tb)\n+ });\n+ rows\n+}\n+\n+/// Public forum thread (`room_wire == \"public\"`): same merge (timeline usually empty).\n+pub fn merge_public_thread_rows(\n+ reduced: &ReducerState,\n+ thread_tag: &str,\n+ since: Option,\n+ before: Option,\n+ actor_prefix: &str,\n+) -> Vec {\n+ merge_thread_rows(reduced, \"public\", thread_tag, since, before, actor_prefix)\n+}\ndiff --git a/test/grants.bb b/test/grants.bb\nindex 10ef0f08ce1c3945e5f1634314e6aebcf6810cd2..4ded5ce92576cdea653cce0ba8a751be4374731a 100644\n--- a/test/grants.bb\n+++ b/test/grants.bb\n@@ -105,7 +105,7 @@\n ;; Alice grants bob View only.\n (println \"\\nalice grants bob View only…\")\n (assert! (rpc-line-ok? (:parsed (rpc-batch! base-url alice-token\n- [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capability\" \"view\"}}])))\n+ [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capabilities\" [\"view\"]}}])))\n \"grant View RPC ok\")\n \n (println \"\\nbob (View only) tries to post prose…\")\n@@ -117,7 +117,7 @@\n ;; Alice grants bob Post.\n (println \"\\nalice grants bob Post…\")\n (assert! (rpc-line-ok? (:parsed (rpc-batch! base-url alice-token\n- [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capability\" \"post\"}}])))\n+ [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capabilities\" [\"post\"]}}])))\n \"grant Post RPC ok\")\n \n (println \"\\nbob (View + Post) posts prose…\")\n@@ -143,7 +143,7 @@\n ;; Alice grants bob Vote.\n (println \"\\nalice grants bob Vote…\")\n (assert! (rpc-line-ok? (:parsed (rpc-batch! base-url alice-token\n- [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capability\" \"vote\"}}])))\n+ [{\"RoomGrant\" {\"room\" room-id \"username\" \"bob\" \"capabilities\" [\"vote\"]}}])))\n \"grant Vote RPC ok\")\n \n (println \"\\nbob (View + Post + Vote) votes…\")\ndiff --git a/test/invites.bb b/test/invites.bb\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..f72b1e630a6856174cdef11eb977e2f81b16e551\n--- /dev/null\n+++ b/test/invites.bb\n@@ -0,0 +1,150 @@\n+(ns test.invites\n+ \"Ephemeral invite links: mint via RPC, GET /join → pending session + OAuth, redemption → GrantAdded,\n+ RoomAudit, post succeeds, second GET /join returns 404 when max_uses exhausted.\"\n+ (:require [babashka.fs :as fs]\n+ [cheshire.core :as json]\n+ [clojure.string :as str]\n+ [test.common :as common]\n+ [test.oauth :as oauth]))\n+\n+(def ^:private counts (atom {:pass 0 :fail 0}))\n+\n+(defn- assert! [pred msg]\n+ (common/test-assert! counts pred msg))\n+\n+(defn- bearer [token] {\"Authorization\" (str \"Bearer \" token)})\n+\n+(defn- rpc-batch! [base-url token cmds]\n+ (let [resp (oauth/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- session-from-login-location [loc]\n+ (when loc\n+ (let [qpart (if (str/includes? loc \"?\")\n+ (-> loc (str/split #\"\\?\" 2) second)\n+ \"\")\n+ qpart (if (str/includes? qpart \"#\")\n+ (-> qpart (str/split #\"#\" 2) first)\n+ qpart)\n+ m (oauth/parse-query qpart)]\n+ (some-> (get m :session) str))))\n+\n+(defn- invite-token-from-url [invite-url]\n+ (second (re-find #\"/join/(inv_[^/?#]+)\" (str invite-url))))\n+\n+(defn- register-user! [base-url session-agent username]\n+ (oauth/complete-registration! base-url\n+ :agent session-agent\n+ :username username\n+ :assert! (fn [pred msg] (assert! pred msg))))\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 invites-test [& _args]\n+ (println \"\\n━━━ ephemeral invite + audit integration check ━━━\\n\")\n+ (reset! counts {:pass 0 :fail 0})\n+\n+ (println \"building server binary…\")\n+ (common/letlocals\n+ (bind build (common/run-cargo-build-release! [\"slugsocial-server\"]))\n+ (assert! (zero? (:exit build)) \"cargo build succeeds\")\n+ (bind server-bin \"target/release/slugsocial-server\")\n+\n+ (bind tmp-dir (str (fs/create-temp-dir {:prefix \"slug-invites-\"})))\n+ (bind slug-port (common/pick-port))\n+ (bind google-port (common/pick-port))\n+ (bind base-url (str \"http://127.0.0.1:\" slug-port))\n+ (bind google-url (str \"http://127.0.0.1:\" google-port))\n+\n+ (bind !server (atom nil))\n+ (bind !google (atom nil))\n+\n+ (bind server-env (common/slug-server-env tmp-dir base-url google-url slug-port))\n+ (try\n+ (println (str \"starting mock google on :\" google-port))\n+ (reset! !google (oauth/start-mock-google google-port\n+ :google-users [\"google-user-alice\" \"google-user-bob\"]))\n+\n+ (println (str \"starting server on :\" slug-port))\n+ (reset! !server (common/start-server server-bin server-env))\n+ (assert! (common/wait-for-server base-url 10000) \"server responds to /healthz\")\n+\n+ (println \"\\nregistering alice…\")\n+ (let [alice-token (register-user! base-url\n+ \"00000000-0000-0000-0000-000000000001:test:local/dev\"\n+ \"alice\")\n+\n+ _ (println \"\\nalice creates private room…\")\n+ create (rpc-batch! base-url alice-token\n+ [{\"RoomCreate\" {\"slug\" \"invite-demo\" \"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\")\n+\n+ _ (println \"\\nalice mints invite (view,post,vote uses=1)…\")\n+ mint (rpc-batch! base-url alice-token\n+ [{\"RoomMintInvite\" {\"room\" room-id\n+ \"capabilities\" [\"view\" \"post\" \"vote\"]\n+ \"max_uses\" 1}}])\n+ _ (assert! (= 200 (:status mint)) \"mint HTTP 200\")\n+ _ (assert! (rpc-line-ok? (:parsed mint)) \"mint RPC ok\")\n+ invite-url (get-in (:parsed mint) [\"results\" 0 \"result\" \"RoomInviteMinted\" \"invite_url\"])\n+ inv-tok (invite-token-from-url invite-url)\n+ _ (assert! (some? inv-tok) \"invite token parsed from URL\")\n+\n+ _ (println \"\\nGET /join/:token (expect redirect + session)…\")\n+ join-resp (oauth/http-get-no-redirect (str base-url \"/join/\" inv-tok))\n+ _ (assert! (contains? #{302 307} (:status join-resp))\n+ (str \"join returns redirect (got status \" (:status join-resp) \")\"))\n+ sess (session-from-login-location (:location join-resp))\n+ _ (assert! (and (some? sess) (str/starts-with? sess \"p_\")) \"Location carries session=p_…\")\n+\n+ _ (println \"\\nbob completes OAuth via invite session…\")\n+ bob-token (oauth/complete-pending-session! base-url sess \"bob\"\n+ :assert! (fn [pred msg] (assert! pred msg)))\n+\n+ _ (println \"\\nalice runs RoomAudit…\")\n+ audit (rpc-batch! base-url alice-token [{\"RoomAudit\" {\"room\" room-id}}])\n+ _ (assert! (rpc-line-ok? (:parsed audit)) \"audit RPC ok\")\n+ grants (get-in (:parsed audit) [\"results\" 0 \"result\" \"RoomAudit\" \"grants\"])\n+ bob-entry (first (filter #(= \"bob\" (get % \"username\")) grants))\n+ _ (assert! (some? bob-entry) \"audit lists bob\")\n+ bob-caps (set (get bob-entry \"capabilities\"))\n+ _ (assert! (= bob-caps #{\"view\" \"post\" \"vote\"}) \"bob has view, post, vote\")\n+\n+ _ (println \"\\nbob posts prose to private room…\")\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 via invite link.\")))\n+ \"bob post succeeds\")\n+\n+ _ (println \"\\nsecond GET /join (invite exhausted → 404)…\")\n+ join2 (oauth/http-get-no-redirect (str base-url \"/join/\" inv-tok))\n+ _ (assert! (= 404 (:status join2)) \"exhausted invite returns 404\")]\n+\n+ (println \"\\ninvite lifecycle OK.\"))\n+\n+ (finally\n+ (when-some [s @!server] (common/kill-server s))\n+ (when-some [g @!google] ((:stop-fn g)))\n+ (fs/delete-tree tmp-dir)))\n+\n+ (bind {pass :pass fail :fail} @counts)\n+ (if (zero? fail)\n+ (println (str \"\\n\" common/ansi-green \"━━━ \" pass \" invite checks passed ━━━\" common/ansi-reset \"\\n\"))\n+ (do (println (str \"\\n\" common/ansi-red \"━━━ \" fail \" invite checks FAILED ━━━\" common/ansi-reset \"\\n\"))\n+ (System/exit 1)))))\n+\n+(when (= *file* (System/getProperty \"babashka.file\"))\n+ (invites-test))\ndiff --git a/test/oauth.bb b/test/oauth.bb\nindex b69459acd17c23c023170d8dd93fe45bb49c8d88..efb3f3e66346ffde81bf9622d75c3a6ec89f8db2 100644\n--- a/test/oauth.bb\n+++ b/test/oauth.bb\n@@ -21,6 +21,21 @@\n resp (.send (http-client) req (java.net.http.HttpResponse$BodyHandlers/ofString))]\n {:status (.statusCode resp) :body (.body resp) :headers (.map (.headers resp))})))\n \n+(defn http-get-no-redirect\n+ \"GET without following redirects; returns `:location` from the first `Location` header when present.\"\n+ [url & {:keys [headers]}]\n+ (let [client (-> (java.net.http.HttpClient/newBuilder)\n+ (.followRedirects java.net.http.HttpClient$Redirect/NEVER)\n+ (.connectTimeout connect-timeout)\n+ (.build))\n+ b (java.net.http.HttpRequest/newBuilder (java.net.URI/create url))]\n+ (doseq [[k v] (or headers {})]\n+ (.header b k v))\n+ (let [req (-> b (.timeout request-timeout) (.GET) (.build))\n+ resp (.send client req (java.net.http.HttpResponse$BodyHandlers/ofString))\n+ loc (first (get (.map (.headers resp)) \"location\"))]\n+ {:status (.statusCode resp) :body (.body resp) :location loc})))\n+\n (defn http-post-json [url data & {:keys [headers]}]\n (let [body (json/generate-string data)\n b (java.net.http.HttpRequest/newBuilder (java.net.URI/create url))]\n@@ -50,6 +65,35 @@\n resp (.send (http-client) req (java.net.http.HttpResponse$BodyHandlers/ofString))]\n {:status (.statusCode resp) :body (.body resp) :headers (.map (.headers resp))})))\n \n+(defn complete-pending-session!\n+ \"Finish OAuth for an existing pending session id (e.g. created by `GET /join/inv_…`). Returns bearer token.\"\n+ [base-url session-id username & {:keys [assert!]}]\n+ (let [check! (fn [pred msg resp]\n+ (if assert!\n+ (assert! pred msg)\n+ (when-not pred\n+ (throw (ex-info msg {:resp resp})))))]\n+ (let [enc (java.net.URLEncoder/encode session-id \"UTF-8\")\n+ login-url (str base-url \"/auth/login?session=\" enc)\n+ login-get (http-get login-url)]\n+ (check! (= 200 (:status login-get))\n+ (str \"oauth redirect chain for session \" session-id)\n+ login-get)\n+ (let [choose (http-post-form (str base-url \"/auth/choose-username\")\n+ {:session session-id :username username})]\n+ (check! (= 200 (:status choose))\n+ (str \"choose-username for \" username \" returns 200\")\n+ choose)\n+ (let [poll (http-get (str base-url \"/api/v0/pending-session/\" session-id))]\n+ (check! (= 200 (:status poll))\n+ (str \"pending-session poll returns 200\")\n+ poll)\n+ (let [poll-json (json/parse-string (:body poll) true)]\n+ (check! (:complete poll-json)\n+ (str \"pending session complete for \" username)\n+ poll)\n+ (:token poll-json)))))))\n+\n (defn parse-query [s]\n (into {}\n (for [part (str/split (or s \"\") #\"&\")\ndiff --git a/types/src/lib.rs b/types/src/lib.rs\nindex a715fe5639f4a3a7bfedcc706b5fb312b98c9ac5..98c66a00fd27803ef6f75b5ac478ff2eb762d771 100644\n--- a/types/src/lib.rs\n+++ b/types/src/lib.rs\n@@ -123,13 +123,32 @@ pub struct PathDetailResponse {\n #[derive(Debug, Serialize, Deserialize)]\n pub struct ThreadDetailResponse {\n pub thread: String,\n- pub posts: Vec,\n- /// Total posts in this thread.\n+ /// Chronological page: prose posts and room system lines, oldest first within the window.\n+ pub items: Vec,\n+ /// Total rows (posts + system lines) in this thread after filters.\n pub total: usize,\n- /// Chronological offset of the first post in this page.\n+ /// Offset into the merged chronological list.\n pub offset: usize,\n }\n \n+/// One row in a thread timeline: a normal post or a room system line.\n+#[derive(Debug, Serialize, Deserialize)]\n+#[serde(tag = \"kind\", rename_all = \"snake_case\")]\n+pub enum ThreadItem {\n+ Post {\n+ id: String,\n+ index: usize,\n+ ts: i64,\n+ actor: String,\n+ body: String,\n+ truncated: bool,\n+ },\n+ System {\n+ ts: i64,\n+ text: String,\n+ },\n+}\n+\n /// One post in a thread. Full body, no snippet.\n #[derive(Debug, Serialize, Deserialize)]\n pub struct PostRow {\n@@ -224,6 +243,23 @@ pub struct FeedPost {\n // RPC batch API (`POST /api/v0/rpc`)\n // ---------------------------------------------------------------------------\n \n+fn default_invite_max_uses() -> usize {\n+ 1\n+}\n+\n+/// One principal's capabilities in a private room (from [`RpcCommand::RoomAudit`]).\n+#[derive(Debug, Serialize, Deserialize)]\n+pub struct RoomAuditEntry {\n+ pub username: String,\n+ pub capabilities: Vec,\n+}\n+\n+#[derive(Debug, Serialize, Deserialize)]\n+pub struct RoomAuditResponse {\n+ pub room: String,\n+ pub grants: Vec,\n+}\n+\n #[derive(Debug, Serialize, Deserialize)]\n #[serde(transparent)]\n pub struct RpcBatch(pub Vec);\n@@ -287,10 +323,27 @@ pub enum RpcCommand {\n visibility: Option,\n },\n RoomGrant {\n+ room: String,\n+ username: String,\n+ /// Capability names: `view`, `post`, `vote`, `add_item`, `manage`.\n+ capabilities: Vec,\n+ },\n+ RoomRevoke {\n room: String,\n username: String,\n capability: String,\n },\n+ /// Mint a shareable invite link (24h TTL, stored in memory only until redeemed or expiry).\n+ RoomMintInvite {\n+ room: String,\n+ capabilities: Vec,\n+ #[serde(default = \"default_invite_max_uses\")]\n+ max_uses: usize,\n+ },\n+ /// List principals granted access in a room (requires View or Manage).\n+ RoomAudit {\n+ room: String,\n+ },\n GetGlobalRank {\n room: String,\n #[serde(default)]\n@@ -361,6 +414,13 @@ pub enum RpcResult {\n RoomCreated {\n room_id: String,\n },\n+ RoomInviteMinted {\n+ invite_url: String,\n+ #[serde(default, skip_serializing_if = \"Option::is_none\")]\n+ expires_at_ms: Option,\n+ max_uses: usize,\n+ },\n+ RoomAudit(RoomAuditResponse),\n GrantOk {},\n GlobalRank(GlobalRankResponse),\n Pair(PairResponse),\n\n\nSide B — contributor: tommy-mor\nSide B — commit message:\n[4e327784] deploy live constitution dashboard\n\nExpose auditable progress and event streaming, configure the production roots and runtime, and make tested main-branch commits the deployment authority.\n\nCo-authored-by: Cursor \n\nSide B — unified diff (full patch):\ndiff --git a/.dockerignore b/.dockerignore\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..c9d63a722beba0a0297fc853089332c460ab78dd\n--- /dev/null\n+++ b/.dockerignore\n@@ -0,0 +1,7 @@\n+.git\n+.venv\n+.hypothesis\n+__pycache__\n+tests\n+*.json\n+*.bsp\ndiff --git a/.github/workflows/deploy.yml b/.github/workflows/deploy.yml\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..76dcdf82d53177c1e47d86b23a54523239d232a6\n--- /dev/null\n+++ b/.github/workflows/deploy.yml\n@@ -0,0 +1,49 @@\n+name: Test and deploy\n+\n+on:\n+ push:\n+ branches: [main]\n+\n+concurrency:\n+ group: production\n+ cancel-in-progress: false\n+\n+permissions:\n+ contents: read\n+\n+jobs:\n+ test:\n+ runs-on: ubuntu-latest\n+ steps:\n+ - uses: actions/checkout@v4\n+\n+ - uses: astral-sh/setup-uv@v6\n+ with:\n+ enable-cache: true\n+\n+ - name: Run Python tests\n+ run: uv run pytest -q\n+\n+ - name: Install Babashka\n+ run: |\n+ curl -fsSL https://raw.githubusercontent.com/babashka/babashka/master/install \\\n+ | sudo bash -s -- --dir /usr/local/bin\n+\n+ - name: Run process integration tests\n+ run: bb TEST.sh\n+\n+ deploy:\n+ needs: test\n+ runs-on: ubuntu-latest\n+ environment:\n+ name: production\n+ url: https://token.slug.social\n+ steps:\n+ - uses: actions/checkout@v4\n+\n+ - uses: superfly/flyctl-actions/setup-flyctl@master\n+\n+ - name: Deploy to Fly\n+ run: flyctl deploy --remote-only\n+ env:\n+ FLY_API_TOKEN: ${{ secrets.FLY_API_TOKEN }}\ndiff --git a/Dockerfile b/Dockerfile\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..c9a5c00782371c19ad5ab5c58cf6f5a8ffec0141\n--- /dev/null\n+++ b/Dockerfile\n@@ -0,0 +1,17 @@\n+FROM ghcr.io/astral-sh/uv:python3.11-bookworm-slim\n+\n+RUN apt-get update \\\n+ && apt-get install -y --no-install-recommends git ca-certificates \\\n+ && rm -rf /var/lib/apt/lists/*\n+\n+WORKDIR /app\n+COPY pyproject.toml uv.lock ./\n+RUN uv sync --frozen --no-install-project\n+\n+COPY constitution.py ./\n+\n+ENV PATH=\"/app/.venv/bin:${PATH}\" \\\n+ PYTHONUNBUFFERED=\"1\"\n+\n+EXPOSE 8080\n+CMD [\"python\", \"constitution.py\"]\ndiff --git a/constitution.py b/constitution.py\nindex 4bee9f83663ab7fb36db95129b92b64b4ef57258..a58257e1881b21d1d6fa8e68a3faa222e4f661ef 100644\n--- a/constitution.py\n+++ b/constitution.py\n@@ -24,12 +24,12 @@ A daily GitHub Action backs up the JSONL ledger to the same repo.\n Run: uv run constitution.py\n \"\"\"\n \n-from decimal import Decimal, getcontext\n+from decimal import Decimal, getcontext, DefaultContext\n from datetime import datetime, timezone\n from fastapi import FastAPI, Request, Response\n from fastapi.responses import PlainTextResponse, HTMLResponse\n from starlette.middleware.sessions import SessionMiddleware\n-import json, time, os, asyncio, httpx, pathlib, subprocess, hashlib, re, fcntl\n+import json, time, os, asyncio, httpx, pathlib, subprocess, hashlib, re, fcntl, base64\n import sympy as sp # type: ignore[reportMissingImports]\n from tenacity import retry, retry_if_exception, stop_after_attempt, wait_exponential\n from evaleval import (\n@@ -37,6 +37,7 @@ from evaleval import (\n exec_event, One, Two, Three, Selector, MORPH, PREPEND,\n )\n \n+DefaultContext.prec = 50\n getcontext().prec = 50\n \n app = FastAPI()\n@@ -129,14 +130,47 @@ OPENROUTER_BASE_URL = os.environ.get(\"OPENROUTER_BASE_URL\", \"https://openrouter.\n # using the exact same source; their normalized values are committed to every\n # discovery event.\n DEFAULT_REPOSITORIES = [\n+ {\n+ \"id\": \"constitution\",\n+ \"url\": \"https://github.com/sortersocial/constitution.git\",\n+ \"refs\": [\"refs/heads/**\"],\n+ },\n {\n \"id\": \"slug\",\n- \"url\": \"https://github.com/tommy-mor/slug.git\",\n+ \"url\": \"https://github.com/sortersocial/slug.git\",\n+ \"refs\": [\"refs/heads/**\"],\n+ },\n+ {\n+ \"id\": \"sorter\",\n+ \"url\": \"https://github.com/sorterisntonline/sorter.git\",\n+ \"refs\": [\"refs/heads/**\"],\n+ },\n+ {\n+ \"id\": \"sorter2\",\n+ \"url\": \"https://github.com/sortersocial/sorter2.git\",\n+ \"refs\": [\"refs/heads/**\"],\n+ },\n+ {\n+ \"id\": \"sorter-oldest\",\n+ \"url\": \"https://github.com/tommy-mor/sorter.git\",\n \"refs\": [\"refs/heads/**\"],\n },\n ]\n DEFAULT_CONTRIBUTORS = {\n \"tommy-mor\": [\"thmorriss@gmail.com\"],\n+ \"christopher-whitman\": [\n+ \"chris@cwwhitman.com\",\n+ \"7566903+cwwhitman@users.noreply.github.com\",\n+ ],\n+ \"jake-chvatal\": [\n+ \"jake+github@uln.industries\",\n+ \"jakechvatal@gmail.com\",\n+ \"jake@isnt.online\",\n+ ],\n+ \"lara\": [\"me@lara.lv\"],\n+ \"nat-reid\": [\"nathanielreid@gmail.com\"],\n+ \"zod\": [\"jason.p.mcel@gmail.com\", \"me@zod.tf\"],\n+ \"jovan\": [\"jovan@slug.social\", \"jovan@getcivicai.com\"],\n }\n \n REPOSITORIES = json.loads(\n@@ -147,6 +181,7 @@ CONTRIBUTORS = json.loads(\n )\n GIT_MIRROR_DIR = pathlib.Path(os.environ.get(\"GIT_MIRROR_DIR\", \"/data/git\"))\n GIT_TIMEOUT_SECONDS = int(os.environ.get(\"GIT_TIMEOUT_SECONDS\", \"120\"))\n+GITHUB_TOKEN = os.environ.get(\"GITHUB_TOKEN\", \"\")\n \n # Council model IDs: slug.social garden rank under this parent (bodies = OpenRouter URLs), then top-up from OpenRouter list.\n SLUG_SOCIAL_BASE_URL = os.environ.get(\"SLUG_SOCIAL_BASE_URL\", \"https://slug.social\").rstrip(\"/\")\n@@ -661,20 +696,30 @@ def _git(repo: pathlib.Path | None, *args: str, input_bytes: bytes | None = None\n if repo is not None:\n command += [\"-C\", str(repo)]\n command += list(args)\n+ git_env = {\n+ **os.environ,\n+ \"GIT_CONFIG_NOSYSTEM\": \"1\",\n+ \"GIT_CONFIG_GLOBAL\": os.devnull,\n+ \"GIT_NO_REPLACE_OBJECTS\": \"1\",\n+ \"LC_ALL\": \"C\",\n+ \"TZ\": \"UTC\",\n+ }\n+ if GITHUB_TOKEN:\n+ credential = base64.b64encode(\n+ f\"x-access-token:{GITHUB_TOKEN}\".encode()\n+ ).decode()\n+ git_env.update({\n+ \"GIT_CONFIG_COUNT\": \"1\",\n+ \"GIT_CONFIG_KEY_0\": \"http.https://github.com/.extraHeader\",\n+ \"GIT_CONFIG_VALUE_0\": f\"Authorization: Basic {credential}\",\n+ })\n try:\n result = subprocess.run(\n command,\n input=input_bytes,\n stdout=subprocess.PIPE,\n stderr=subprocess.PIPE,\n- env={\n- **os.environ,\n- \"GIT_CONFIG_NOSYSTEM\": \"1\",\n- \"GIT_CONFIG_GLOBAL\": os.devnull,\n- \"GIT_NO_REPLACE_OBJECTS\": \"1\",\n- \"LC_ALL\": \"C\",\n- \"TZ\": \"UTC\",\n- },\n+ env=git_env,\n timeout=GIT_TIMEOUT_SECONDS,\n check=False,\n )\n@@ -844,9 +889,15 @@ def _build_discovery(epoch_n: int, boundary_ms: int, events: list) -> GitDiscove\n canonical_location = min(\n locations[qualified_oid], key=lambda x: (x[0], x[1])\n )\n+ # One commit may be reachable from dozens of refs in the same mirror.\n+ # Verify its object once per repository, not once per source ref.\n+ object_locations = {\n+ (str(m), raw_oid): (m, raw_oid)\n+ for _, _, m, raw_oid in locations[qualified_oid]\n+ }\n object_hashes = {\n hashlib.sha256(_git(m, \"cat-file\", \"commit\", raw_oid)).hexdigest()\n- for _, _, m, raw_oid in locations[qualified_oid]\n+ for m, raw_oid in object_locations.values()\n }\n if len(object_hashes) != 1:\n raise RuntimeError(f\"conflicting Git objects share OID {qualified_oid}\")\n@@ -994,11 +1045,55 @@ async def discover_repositories(epoch_n: int, boundary_ms: int) -> GitDiscovery:\n \n \n SSE_CLIENTS = []\n+AUDIT_HISTORY = []\n+AUDIT_SEQUENCE = 0\n+PROCESS_STATE = {\n+ \"running\": False,\n+ \"phase\": \"idle\",\n+ \"progress\": 100,\n+ \"message\": \"Waiting for the next epoch\",\n+}\n+\n+\n+def _sse_event(event_name: str, payload: dict) -> str:\n+ return (\n+ f\"event: {event_name}\\n\"\n+ f\"data: {json.dumps(payload, separators=(',', ':'))}\\n\\n\"\n+ )\n+\n+\n+async def broadcast_audit(\n+ kind: str,\n+ message: str,\n+ *,\n+ progress: int | None = None,\n+ phase: str | None = None,\n+) -> dict:\n+ global AUDIT_SEQUENCE\n+ AUDIT_SEQUENCE += 1\n+ if progress is not None:\n+ PROCESS_STATE[\"progress\"] = max(0, min(100, int(progress)))\n+ if phase is not None:\n+ PROCESS_STATE[\"phase\"] = phase\n+ PROCESS_STATE[\"message\"] = message\n+ payload = {\n+ \"id\": AUDIT_SEQUENCE,\n+ \"timestamp_ms\": int(time.time() * 1000),\n+ \"kind\": kind,\n+ \"message\": message,\n+ **PROCESS_STATE,\n+ }\n+ AUDIT_HISTORY.append(payload)\n+ del AUDIT_HISTORY[:-200]\n+ wire = _sse_event(\"audit\", payload)\n+ for queue in list(SSE_CLIENTS):\n+ await queue.put(wire)\n+ return payload\n \n \n async def broadcast_js(js: str):\n \"\"\"Send a JS snippet to all connected SSE clients.\"\"\"\n- for queue in SSE_CLIENTS:\n+ for queue in list(SSE_CLIENTS):\n await queue.put(js)\n \n \n@@ -1006,10 +1101,29 @@ async def rank_commits(commits: list[dict]):\n if not commits:\n return {}, []\n \n- models = await fetch_top_models(n=3)\n contributors = sorted(set(c[\"contributor\"] for c in commits))\n- if len(contributors) > 1 and not models:\n+ if len(contributors) == 1:\n+ await broadcast_audit(\n+ \"ranking\",\n+ f\"Only {contributors[0]} is eligible; rank is 1.0\",\n+ progress=90,\n+ phase=\"finalizing\",\n+ )\n+ return {contributors[0]: Decimal(\"1\")}, []\n+ if not (OPENROUTER_API_KEY or \"\").strip():\n+ raise RuntimeError(\n+ \"OPENROUTER_API_KEY is required when multiple contributors need ranking\"\n+ )\n+\n+ models = await fetch_top_models(n=3)\n+ if not models:\n raise RuntimeError(\"no council models available for contributor ranking\")\n+ await broadcast_audit(\n+ \"council\",\n+ f\"Council selected: {', '.join(models)}\",\n+ progress=35,\n+ phase=\"ranking\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-log\")][PREPEND][\n [\"div.log-council\", f\"Council: {', '.join(models)} — {len(commits)} commits\"]\n ]))\n@@ -1035,6 +1149,11 @@ async def rank_commits(commits: list[dict]):\n \n async def compare_fn(i, j):\n a1, a2 = authors[i], authors[j]\n+ await broadcast_audit(\n+ \"comparison\",\n+ f\"Comparing {a1} with {a2}\",\n+ phase=\"ranking\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-status\")][MORPH][\n [\"div#emission-status\", f\"Comparing {a1} vs {a2}…\"]\n ]))\n@@ -1050,6 +1169,11 @@ async def rank_commits(commits: list[dict]):\n if winner_weight <= 0 or loser_weight <= 0:\n raise ValueError(\"ratio weights must be positive\")\n results.append((w, l, winner_weight, loser_weight))\n+ await broadcast_audit(\n+ \"vote\",\n+ f\"{model}: {authors[w]} over {authors[l]} ({result['ratio']})\",\n+ phase=\"ranking\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-log\")][PREPEND][\n [\"div.log-vote\",\n [\"span.model\", model], \" — \",\n@@ -1059,6 +1183,11 @@ async def rank_commits(commits: list[dict]):\n ]\n ]))\n except Exception as e:\n+ await broadcast_audit(\n+ \"error\",\n+ f\"{model} failed: {e}\",\n+ phase=\"error\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-log\")][PREPEND][\n [\"div.log-error\", f\"⚠ {model}: {e}\"]\n ]))\n@@ -1068,8 +1197,13 @@ async def rank_commits(commits: list[dict]):\n async def progress_fn(ev):\n if ev[\"phase\"] == \"spanning_tree\":\n label = f\"Spanning tree: {ev['step']}/{ev['total']}\"\n+ percent = 35 + round(35 * ev[\"step\"] / max(ev[\"total\"], 1))\n else:\n label = f\"Zip pass {ev['pass']}: {ev['step']}/{ev['total']}\"\n+ percent = 70 + round(20 * ev[\"step\"] / max(ev[\"total\"], 1))\n+ await broadcast_audit(\n+ \"progress\", label, progress=percent, phase=\"ranking\"\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-status\")][MORPH][\n [\"div#emission-status\", label]\n ]))\n@@ -1083,6 +1217,12 @@ async def rank_commits(commits: list[dict]):\n scores = rank_centrality(pairs)\n ranking = {authors[i]: Decimal(str(scores[i])) for i in range(len(authors))}\n ranking_rows = sorted(ranking.items(), key=lambda x: x[1], reverse=True)\n+ await broadcast_audit(\n+ \"ranking\",\n+ \"Ranking: \" + \", \".join(f\"{a} {s:.4f}\" for a, s in ranking_rows),\n+ progress=90,\n+ phase=\"finalizing\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-log\")][PREPEND][\n [\"div.log-ranking\",\n [\"b\", \"Ranking: \"],\n@@ -1104,11 +1244,33 @@ def pool_remaining(events: list) -> Decimal:\n \n \n async def run_emission(epoch_n, boundary_ms):\n+ PROCESS_STATE[\"running\"] = True\n+ await broadcast_audit(\n+ \"start\",\n+ f\"Epoch {epoch_n} emission started\",\n+ progress=2,\n+ phase=\"starting\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-log\")][PREPEND][\n [\"div.log-start\", f\"⚡ Epoch {epoch_n} emission started\"]\n ]))\n \n+ await broadcast_audit(\n+ \"discovery\",\n+ \"Fetching configured repositories and snapshotting refs\",\n+ progress=8,\n+ phase=\"discovery\",\n+ )\n discovery = await discover_repositories(epoch_n, boundary_ms)\n+ await broadcast_audit(\n+ \"discovery\",\n+ (\n+ f\"Discovered {len(discovery.observations)} new commits; \"\n+ f\"{len(discovery.commits)} are eligible\"\n+ ),\n+ progress=30,\n+ phase=\"discovery\",\n+ )\n ranking, models = await rank_commits(discovery.commits)\n \n def make_emission(events):\n@@ -1146,6 +1308,13 @@ async def run_emission(epoch_n, boundary_ms):\n \n entry = await store.atomic(make_emission)\n if entry:\n+ PROCESS_STATE[\"running\"] = False\n+ await broadcast_audit(\n+ \"complete\",\n+ f\"Epoch {entry.epoch} complete; emitted {entry.total_emitted} SLG\",\n+ progress=100,\n+ phase=\"idle\",\n+ )\n await broadcast_js(exec_event(Three[Selector(\"#emission-log\")][PREPEND][\n [\"div.log-amount\",\n f\"Pool {entry.pool_before} → emit {entry.total_emitted} → {entry.pool_after}\"]\n@@ -1189,13 +1358,24 @@ async def distribute_usdc(holdings, treasury_balance):\n async def epoch_loop():\n while True:\n epoch_n, current_start, next_boundary = current_epoch()\n+ processed = {e.epoch for e in store.read() if isinstance(e, Emission)}\n+ if epoch_n >= 0 and epoch_n not in processed:\n+ try:\n+ await run_emission(epoch_n, current_start)\n+ except Exception as exc:\n+ PROCESS_STATE[\"running\"] = False\n+ await broadcast_audit(\n+ \"error\",\n+ f\"Epoch {epoch_n} failed: {exc}; retrying in 60 seconds\",\n+ phase=\"error\",\n+ )\n+ print(f\"epoch {epoch_n} emission failed: {exc}\", flush=True)\n+ await asyncio.sleep(60)\n+ continue\n+\n now = int(time.time() * 1000)\n wait_ms = next_boundary - now\n-\n if wait_ms <= 0:\n- processed = {e.epoch for e in store.read() if isinstance(e, Emission)}\n- if epoch_n not in processed and epoch_n >= 0:\n- await run_emission(epoch_n, current_start)\n await asyncio.sleep(60)\n elif wait_ms < 86_400_000:\n await broadcast_js(exec_event(Three[Selector(\"#emission-status\")][MORPH][\n@@ -1243,6 +1423,36 @@ async def get_ranking():\n return {\"ranking\": latest.ranking, \"epoch\": latest.epoch}\n \n \n+@app.get(\"/api/status\")\n+async def get_status():\n+ events = store.read()\n+ discoveries = [e for e in events if isinstance(e, GitDiscovery)]\n+ emissions = [e for e in events if isinstance(e, Emission)]\n+ return {\n+ **PROCESS_STATE,\n+ \"epoch\": current_epoch()[0],\n+ \"openrouter_configured\": bool((OPENROUTER_API_KEY or \"\").strip()),\n+ \"sse_clients\": len(SSE_CLIENTS),\n+ \"latest_discovery\": (\n+ {\n+ \"epoch\": discoveries[-1].epoch,\n+ \"snapshot_id\": discoveries[-1].snapshot_id,\n+ \"observations\": len(discoveries[-1].observations),\n+ \"eligible_commits\": len(discoveries[-1].commits),\n+ }\n+ if discoveries else None\n+ ),\n+ \"latest_emission\": (\n+ {\n+ \"epoch\": emissions[-1].epoch,\n+ \"total_emitted\": emissions[-1].total_emitted,\n+ \"ranking\": emissions[-1].ranking,\n+ }\n+ if emissions else None\n+ ),\n+ }\n+\n+\n @app.get(\"/api/contributor/{github_username}\")\n async def get_contributor(github_username: str):\n history = [\n@@ -1306,13 +1516,6 @@ async def test_emit():\n \n # ===========================================================================\n # §9. SSE — live audit stream of the pairwise voting process\n-#\n-# TODO: the /sse emission audit page needs a real SSE-driven UI. votes arrive\n-# incrementally during rank_commits(), and the client should show a live\n-# progress bar and per-vote results as they stream in. this requires a\n-# dedicated page that connects to /sse and updates the DOM on each event\n-# (council, comparing, vote, ranking, emission_complete). defer until we\n-# have playwright tests to cover it — the incremental rendering is fiddly.\n # ===========================================================================\n \n @app.get(\"/sse\")\n@@ -1322,6 +1525,13 @@ async def sse_stream(request: Request):\n \n async def generate():\n try:\n+ yield _sse_event(\"audit\", {\n+ \"id\": AUDIT_SEQUENCE,\n+ \"timestamp_ms\": int(time.time() * 1000),\n+ \"kind\": \"connection\",\n+ \"message\": f\"Connected to epoch {current_epoch()[0]}\",\n+ **PROCESS_STATE,\n+ })\n yield exec_event(Three[Selector(\"#emission-status\")][MORPH][\n [\"div#emission-status\", f\"Connected — epoch {current_epoch()[0]}\"]\n ])\n@@ -1334,10 +1544,15 @@ async def sse_stream(request: Request):\n except asyncio.TimeoutError:\n yield \": keepalive\\n\\n\"\n finally:\n- SSE_CLIENTS.remove(queue)\n+ if queue in SSE_CLIENTS:\n+ SSE_CLIENTS.remove(queue)\n \n from starlette.responses import StreamingResponse\n- return StreamingResponse(generate(), media_type=\"text/event-stream\")\n+ return StreamingResponse(\n+ generate(),\n+ media_type=\"text/event-stream\",\n+ headers={\"Cache-Control\": \"no-cache\", \"X-Accel-Buffering\": \"no\"},\n+ )\n \n \n # ===========================================================================\n@@ -1359,6 +1574,7 @@ def _page(title: str, body: list) -> HTMLResponse:\n [\"meta\", {\"charset\": \"utf-8\"}],\n [\"meta\", {\"name\": \"viewport\", \"content\": \"width=device-width, initial-scale=1\"}],\n [\"title\", title],\n+ [\"style\", RawContent(_WATCH_CSS)],\n ],\n [\"body\",\n body,\n@@ -1367,6 +1583,387 @@ def _page(title: str, body: list) -> HTMLResponse:\n ]))\n \n \n+_WATCH_CSS = \"\"\"\n+/* ================================================================\n+ ZIGGURAT — bevel-first dark theme\n+ --spread (0→1) controls bevel depth. 0 = flat. 1 = full relief.\n+ Light source: top-left. Shadow: bottom-right.\n+ Platforms nest. Each level is raised. Nothing is rounded.\n+ ================================================================ */\n+\n+:root {\n+ color-scheme: dark;\n+ --spread: 1;\n+\n+ --g0: #080808;\n+ --g1: #131313;\n+ --g2: #1c1c1c;\n+ --g3: #252525;\n+ --g4: #2e2e2e;\n+ --g5: #383838;\n+\n+ --hi: #5e5e5e;\n+ --lo: #050505;\n+ --bv: calc(var(--spread) * 4px + 1px);\n+ --bv-lg: calc(var(--spread) * 6px + 2px);\n+\n+ --signal: #f0f0f0;\n+ --prose: #c2c2c2;\n+ --ui: #888;\n+ --meta: #4a4a4a;\n+ --link: #8899ee;\n+ --code-fg: #c8dda0;\n+\n+ --font-prose: \"Iowan Old Style\", \"Palatino Linotype\", Palatino, \"Book Antiqua\", Georgia, serif;\n+ --font-ui: system-ui, -apple-system, sans-serif;\n+ --font-code: ui-monospace, \"Cascadia Code\", \"SF Mono\", Menlo, monospace;\n+}\n+\n+*, *::before, *::after { box-sizing: border-box; }\n+html, body { margin: 0; padding: 0; }\n+\n+body {\n+ background: var(--g0);\n+ color: var(--prose);\n+ font-family: var(--font-ui);\n+ font-size: 14px;\n+ line-height: 1.6;\n+ margin: 0 auto;\n+ max-width: 560px;\n+ min-height: 100vh;\n+ padding: 0 16px 48px;\n+}\n+main { width: 100%; padding: 18px 0 48px; }\n+\n+h1, h2, h3 {\n+ color: var(--signal);\n+ font-size: 11px;\n+ font-weight: bold;\n+ letter-spacing: 0.12em;\n+ margin: 14px 0 6px;\n+ text-transform: uppercase;\n+}\n+a { color: var(--link); text-decoration: none; }\n+a:hover { color: var(--signal); }\n+.eyebrow {\n+ background: var(--g2);\n+ border: var(--bv) solid;\n+ border-color: var(--hi) var(--lo) var(--lo) var(--hi);\n+ color: var(--ui);\n+ font-size: 11px;\n+ letter-spacing: 0.12em;\n+ padding: 4px 10px;\n+ text-transform: uppercase;\n+ width: fit-content;\n+}\n+\n+/* Every dashboard section is a raised platform. */\n+.panel {\n+ background: var(--g2);\n+ border: var(--bv-lg) solid;\n+ border-color: var(--hi) var(--lo) var(--lo) var(--hi);\n+ margin: 8px 0;\n+ padding: 10px;\n+ width: 100%;\n+}\n+.status-row {\n+ align-items: center;\n+ display: flex;\n+ flex-wrap: wrap;\n+ gap: 8px;\n+ justify-content: space-between;\n+}\n+#process-status { color: var(--signal); font-family: var(--font-code); font-weight: bold; }\n+.badge {\n+ align-items: center;\n+ background: var(--g3);\n+ border: var(--bv) solid;\n+ border-color: var(--hi) var(--lo) var(--lo) var(--hi);\n+ color: var(--ui);\n+ display: inline-flex;\n+ font-size: 11px;\n+ gap: 7px;\n+ padding: 3px 8px;\n+}\n+.dot { background: var(--meta); height: 8px; width: 8px; }\n+.live .dot { background: #7acc7a; }\n+.warn .dot { background: #cc9955; }\n+\n+/* The progress track is inset; its signal is raised inside it. */\n+.progress-shell {\n+ background: var(--g1);\n+ border: var(--bv-lg) solid;\n+ border-color: var(--lo) var(--hi) var(--hi) var(--lo);\n+ height: 58px;\n+ margin: 14px 0 10px;\n+ overflow: hidden;\n+ position: relative;\n+}\n+#progress-fill {\n+ background: var(--link);\n+ border: var(--bv) solid;\n+ border-color: var(--hi) var(--lo) var(--lo) var(--hi);\n+ height: 100%;\n+ transition: width .35s steps(8, end);\n+ width: 0;\n+}\n+#progress-label {\n+ color: var(--signal);\n+ display: grid;\n+ font-family: var(--font-code);\n+ font-size: 18px;\n+ font-weight: bold;\n+ inset: 0;\n+ place-items: center;\n+ position: absolute;\n+ text-shadow: 1px 1px var(--lo);\n+}\n+\n+.controls { align-items: center; display: flex; flex-wrap: wrap; gap: 8px; }\n+button {\n+ background: var(--g5);\n+ border: var(--bv) solid;\n+ border-color: var(--hi) var(--lo) var(--lo) var(--hi);\n+ color: var(--signal);\n+ cursor: pointer;\n+ font: inherit;\n+ font-size: 12px;\n+ padding: 4px 10px;\n+}\n+button:hover { background: #404040; }\n+button:active {\n+ background: var(--g4);\n+ border-color: var(--lo) var(--hi) var(--hi) var(--lo);\n+ transform: translate(1px, 1px);\n+}\n+button:disabled { cursor: default; opacity: .4; }\n+.note { color: var(--meta); font-size: 11px; margin: 4px 0; }\n+\n+.feed-head { align-items: baseline; display: flex; justify-content: space-between; }\n+#audit-feed {\n+ background: var(--g1);\n+ border: var(--bv) solid;\n+ border-color: var(--lo) var(--hi) var(--hi) var(--lo);\n+ display: flex;\n+ flex-direction: column;\n+ gap: 5px;\n+ margin-top: 8px;\n+ padding: 6px;\n+}\n+.event {\n+ background: var(--g3);\n+ border: var(--bv) solid;\n+ border-color: var(--hi) var(--lo) var(--lo) var(--hi);\n+ display: grid;\n+ gap: 6px;\n+ grid-template-columns: 82px 88px 1fr;\n+ padding: 5px 8px;\n+}\n+.event[data-kind=\"error\"] { border-left-color: #cc5555; }\n+.event[data-kind=\"complete\"], .event[data-kind=\"ranking\"] { border-left-color: #7acc7a; }\n+.event[data-kind=\"vote\"] { border-left-color: var(--link); }\n+.event time, .event-kind { color: var(--meta); font-family: var(--font-code); font-size: 10px; }\n+.event-kind { text-transform: uppercase; }\n+.event-message { color: var(--prose); font-family: var(--font-prose); }\n+\n+code {\n+ background: var(--g1);\n+ border: 2px solid;\n+ border-color: var(--lo) var(--hi) var(--hi) var(--lo);\n+ color: var(--code-fg);\n+ font-family: var(--font-code);\n+ font-size: 12px;\n+ padding: 1px 4px;\n+}\n+\n+@media (max-width: 520px) {\n+ .event { grid-template-columns: 72px 1fr; }\n+ .event-message { grid-column: 1 / -1; }\n+}\n+\"\"\"\n+\n+\n+def _watch_initial_state() -> dict:\n+ events = store.read()\n+ feed = []\n+ for event_ in events[-40:]:\n+ if isinstance(event_, GitDiscovery):\n+ feed.append({\n+ \"id\": f\"discovery-{event_.snapshot_id}\",\n+ \"timestamp_ms\": event_.timestamp_ms,\n+ \"kind\": \"discovery\",\n+ \"message\": (\n+ f\"Epoch {event_.epoch}: observed {len(event_.observations)} commits; \"\n+ f\"{len(event_.commits)} eligible\"\n+ ),\n+ })\n+ elif isinstance(event_, Emission):\n+ feed.append({\n+ \"id\": f\"emission-{event_.epoch}\",\n+ \"timestamp_ms\": event_.timestamp_ms,\n+ \"kind\": \"complete\",\n+ \"message\": (\n+ f\"Epoch {event_.epoch}: emitted {event_.total_emitted} SLG; \"\n+ f\"ranking {event_.ranking}\"\n+ ),\n+ })\n+ feed.extend(AUDIT_HISTORY)\n+ return {\n+ \"process\": dict(PROCESS_STATE),\n+ \"openrouter_configured\": bool((OPENROUTER_API_KEY or \"\").strip()),\n+ \"epoch\": current_epoch()[0],\n+ \"feed\": feed[-200:],\n+ }\n+\n+\n+_WATCH_JS = \"\"\"\n+const initial = __INITIAL__;\n+const feed = document.querySelector('#audit-feed');\n+const processStatus = document.querySelector('#process-status');\n+const connection = document.querySelector('#connection-status');\n+const fill = document.querySelector('#progress-fill');\n+const progressLabel = document.querySelector('#progress-label');\n+const play = document.querySelector('#play');\n+const pause = document.querySelector('#pause');\n+const seen = new Set();\n+let source = null;\n+\n+function setProgress(value) {\n+ const n = Math.max(0, Math.min(100, Number(value ?? 0)));\n+ fill.style.width = `${n}%`;\n+ progressLabel.textContent = `${Math.round(n)}%`;\n+ document.querySelector('.progress-shell').setAttribute('aria-valuenow', String(n));\n+}\n+\n+function addEvent(event) {\n+ const id = String(event.id);\n+ if (seen.has(id)) return;\n+ seen.add(id);\n+ const row = document.createElement('div');\n+ row.className = 'event';\n+ row.dataset.kind = event.kind || 'event';\n+ const when = document.createElement('time');\n+ when.dateTime = new Date(event.timestamp_ms).toISOString();\n+ when.textContent = new Date(event.timestamp_ms).toLocaleTimeString();\n+ const kind = document.createElement('span');\n+ kind.className = 'event-kind';\n+ kind.textContent = event.kind || 'event';\n+ const message = document.createElement('span');\n+ message.className = 'event-message';\n+ message.textContent = event.message;\n+ row.append(when, kind, message);\n+ feed.prepend(row);\n+ while (feed.children.length > 200) feed.lastElementChild.remove();\n+}\n+\n+function applyState(event) {\n+ processStatus.textContent = event.message || 'Waiting for the next epoch';\n+ setProgress(event.progress);\n+ if (event.kind !== 'connection') addEvent(event);\n+}\n+\n+function connect() {\n+ if (source) return;\n+ source = new EventSource('/sse');\n+ connection.classList.remove('warn');\n+ connection.classList.add('live');\n+ connection.querySelector('span:last-child').textContent = 'connecting';\n+ play.disabled = true;\n+ pause.disabled = false;\n+ source.onopen = () => {\n+ connection.querySelector('span:last-child').textContent = 'live';\n+ };\n+ source.addEventListener('audit', event => applyState(JSON.parse(event.data)));\n+ source.onerror = () => {\n+ connection.classList.remove('live');\n+ connection.classList.add('warn');\n+ connection.querySelector('span:last-child').textContent = 'reconnecting';\n+ };\n+}\n+\n+function disconnect() {\n+ if (source) source.close();\n+ source = null;\n+ connection.classList.remove('live');\n+ connection.classList.add('warn');\n+ connection.querySelector('span:last-child').textContent = 'paused locally';\n+ play.disabled = false;\n+ pause.disabled = true;\n+}\n+\n+play.addEventListener('click', connect);\n+pause.addEventListener('click', disconnect);\n+initial.feed.forEach(addEvent);\n+processStatus.textContent = initial.process.message;\n+setProgress(initial.process.progress);\n+connect();\n+\"\"\"\n+\n+\n+@app.get(\"/watch\")\n+async def watch():\n+ initial = json.dumps(\n+ _watch_initial_state(), separators=(\",\", \":\")\n+ ).replace(\" 0\")\n \n- ;; 9. SSE connects and sends initial event\n+ ;; 9. watch UI exposes progress, controls, readiness, and live SSE\n+ (println \"\\nchecking /watch UI…\")\n+ (bind watch-html (slurp (str base-url \"/watch\")))\n+ (assert! (str/includes? watch-html \"role=\\\"progressbar\\\"\")\n+ \"watch page has progress bar\")\n+ (assert! (str/includes? watch-html \"id=\\\"play\\\"\")\n+ \"watch page has play control\")\n+ (assert! (str/includes? watch-html \"id=\\\"pause\\\"\")\n+ \"watch page has pause control\")\n+ (assert! (str/includes? watch-html \"OpenRouter configured\")\n+ \"watch page reports council readiness\")\n+ (bind status-resp (get-json base-url \"/api/status\"))\n+ (assert! (true? (:openrouter_configured status-resp))\n+ \"status API reports OpenRouter configuration\")\n+\n+ ;; 10. SSE connects and sends initial event\n (println \"\\nchecking /sse initial event…\")\n (bind sse-events (read-sse-events (str base-url \"/sse\") 1 5000))\n (assert! (= 1 (count sse-events)) \"received 1 SSE event\")\n (assert! (not (str/blank? (first sse-events)))\n \"initial SSE event contains executable audit data\")\n \n- ;; 10. POST /test/emit — full ranking pipeline hits mocks\n+ ;; 11. POST /test/emit — full ranking pipeline hits mocks\n (println \"\\ntriggering /test/emit (epoch 1)…\")\n (bind emit-resp (post-json! base-url \"/test/emit\"))\n (assert! (= \"emission\" (:type emit-resp)) \"emit response type is emission\")\n@@ -503,7 +518,7 @@\n (bind rank-after (get-json base-url \"/api/ranking\"))\n (assert! (= 1 (:epoch rank-after)) \"latest ranking is epoch 1\")\n \n- ;; 11. kill and restart — prove replay determinism\n+ ;; 12. kill and restart — prove replay determinism\n (println \"\\nkilling server for replay test…\")\n (.destroyForcibly (:proc server))\n (deref server)\ndiff --git a/tests/test_git_discovery.py b/tests/test_git_discovery.py\nindex 0dd31bc42a19bc8c59842dc61f193c595c474659..5a9e167e16ad2ffb25988af85820f6f19e8910cd 100644\n--- a/tests/test_git_discovery.py\n+++ b/tests/test_git_discovery.py\n@@ -393,7 +393,9 @@ def test_empty_epoch_records_zero_emission_without_burning_pool(\n monkeypatch.setattr(c, \"store\", c.JsonlStore(discovery_config / \"ledger.jsonl\"))\n \n async def discover(_epoch, _boundary):\n- return SimpleNamespace(commits=[], snapshot_id=\"empty-snapshot\")\n+ return SimpleNamespace(\n+ observations=[], commits=[], snapshot_id=\"empty-snapshot\"\n+ )\n \n async def rank(_commits):\n return {}, []\n@@ -412,7 +414,11 @@ def test_emission_distribution_sums_exactly_to_total(\n monkeypatch.setattr(c, \"store\", c.JsonlStore(discovery_config / \"ledger.jsonl\"))\n \n async def discover(_epoch, _boundary):\n- return SimpleNamespace(commits=[{\"x\": 1}], snapshot_id=\"ranked-snapshot\")\n+ return SimpleNamespace(\n+ observations=[{\"x\": 1}],\n+ commits=[{\"x\": 1}],\n+ snapshot_id=\"ranked-snapshot\",\n+ )\n \n async def rank(_commits):\n return {\n@@ -454,6 +460,7 @@ def test_any_council_failure_aborts_ranking(monkeypatch):\n \n monkeypatch.setattr(c, \"fetch_top_models\", models)\n monkeypatch.setattr(c, \"llm_pairwise_compare\", compare)\n+ monkeypatch.setattr(c, \"OPENROUTER_API_KEY\", \"test-key\")\n commits = [\n {\n \"contributor\": contributor,\n@@ -465,3 +472,54 @@ def test_any_council_failure_aborts_ranking(monkeypatch):\n ]\n with pytest.raises(RuntimeError, match=\"council model failed\"):\n asyncio.run(c.rank_commits(commits))\n+\n+\n+def test_contested_ranking_requires_openrouter_key(monkeypatch):\n+ monkeypatch.setattr(c, \"OPENROUTER_API_KEY\", \"\")\n+ commits = [\n+ {\n+ \"contributor\": contributor,\n+ \"oid\": \"sha1:\" + char * 40,\n+ \"message\": contributor,\n+ \"patch\": \"patch\",\n+ }\n+ for contributor, char in [(\"alice\", \"a\"), (\"bob\", \"b\")]\n+ ]\n+ with pytest.raises(RuntimeError, match=\"OPENROUTER_API_KEY\"):\n+ asyncio.run(c.rank_commits(commits))\n+\n+\n+def test_watch_page_has_live_controls_progress_and_key_warning(\n+ discovery_config, monkeypatch\n+):\n+ monkeypatch.setattr(c, \"store\", c.JsonlStore(discovery_config / \"ledger.jsonl\"))\n+ monkeypatch.setattr(c, \"OPENROUTER_API_KEY\", \"\")\n+ monkeypatch.setattr(c, \"current_epoch\", lambda: (3, 0, 1))\n+ response = asyncio.run(c.watch())\n+ html = response.body.decode()\n+ assert 'role=\"progressbar\"' in html\n+ assert 'id=\"play\"' in html\n+ assert 'id=\"pause\"' in html\n+ assert \"new EventSource('/sse')\" in html\n+ assert \"OpenRouter key missing\" in html\n+\n+\n+def test_audit_events_are_json_sse_and_update_process_state(monkeypatch):\n+ clients = []\n+ history = []\n+ monkeypatch.setattr(c, \"SSE_CLIENTS\", clients)\n+ monkeypatch.setattr(c, \"AUDIT_HISTORY\", history)\n+ queue = asyncio.Queue()\n+ clients.append(queue)\n+\n+ async def emit():\n+ event = await c.broadcast_audit(\n+ \"progress\", \"halfway\", progress=50, phase=\"ranking\"\n+ )\n+ return event, await queue.get()\n+\n+ event, wire = asyncio.run(emit())\n+ assert event[\"progress\"] == 50\n+ assert event[\"phase\"] == \"ranking\"\n+ assert wire.startswith(\"event: audit\\ndata: {\")\n+ assert '\"message\":\"halfway\"' in wire\n","role":"user"}],"model":"~anthropic/claude-sonnet-latest"}