{"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[9be54f7d] Add GitHub OAuth auth and move durable out of tree.\n\nGate votes behind session + pseudonym claim, project identity events into durable maps, and depend on tommy-mor/durable from git instead of the in-repo crates.\n\nCo-authored-by: Cursor \n\nSide B — unified diff (full patch):\ndiff --git a/Cargo.lock b/Cargo.lock\nindex 49a908ef935c430dbe63c6a28d8a24e38b489486..aa02997ad85777195f135bfd9456bcee0fc9a590 100644\n--- a/Cargo.lock\n+++ b/Cargo.lock\n@@ -69,12 +69,6 @@ version = \"1.1.2\"\n source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"1505bd5d3d116872e7271a6d4e16d81d0c8570876c8de68093a09ac269d8aac0\"\n \n-[[package]]\n-name = \"autocfg\"\n-version = \"1.5.1\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"f2032f911046de80f0a198e0901378627c33f59ea0ac00e363d481118bd70a53\"\n-\n [[package]]\n name = \"axum\"\n version = \"0.7.9\"\n@@ -211,21 +205,6 @@ dependencies = [\n \"syn\",\n ]\n \n-[[package]]\n-name = \"bit-set\"\n-version = \"0.8.0\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"08807e080ed7f9d5433fa9b275196cfc35414f66a0c79d864dc51a0d825231a3\"\n-dependencies = [\n- \"bit-vec\",\n-]\n-\n-[[package]]\n-name = \"bit-vec\"\n-version = \"0.8.0\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"5e764a1d40d510daf35e07be9eb06e75770908c27d411ee6c92109c9840eaaf7\"\n-\n [[package]]\n name = \"bitflags\"\n version = \"1.3.2\"\n@@ -435,19 +414,19 @@ dependencies = [\n [[package]]\n name = \"durable\"\n version = \"0.2.0\"\n+source = \"git+https://github.com/tommy-mor/durable.git?branch=main#a6c14eaa809693140eea0c22b07ef24d8e74adaf\"\n dependencies = [\n \"ciborium\",\n \"durable-derive\",\n- \"proptest\",\n \"rocksdb\",\n \"serde\",\n- \"tempfile\",\n \"thiserror\",\n ]\n \n [[package]]\n name = \"durable-derive\"\n version = \"0.2.0\"\n+source = \"git+https://github.com/tommy-mor/durable.git?branch=main#a6c14eaa809693140eea0c22b07ef24d8e74adaf\"\n dependencies = [\n \"proc-macro2\",\n \"quote\",\n@@ -1238,15 +1217,6 @@ version = \"0.2.2\"\n source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"521739c6d2bac4aa25192232afe6841231376b2b26d4d9fae5ecf8ca5772e441\"\n \n-[[package]]\n-name = \"num-traits\"\n-version = \"0.2.19\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"071dfc062690e90b734c0b2273ce72ad0ffa95f0c74596bc250dcfd960262841\"\n-dependencies = [\n- \"autocfg\",\n-]\n-\n [[package]]\n name = \"once_cell\"\n version = \"1.21.4\"\n@@ -1358,7 +1328,7 @@ source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"3c80231409c20246a13fddb31776fb942c38553c51e871f8cbd687a4cfb5843d\"\n dependencies = [\n \"phf_shared\",\n- \"rand 0.8.6\",\n+ \"rand\",\n ]\n \n [[package]]\n@@ -1467,31 +1437,6 @@ dependencies = [\n \"unicode-ident\",\n ]\n \n-[[package]]\n-name = \"proptest\"\n-version = \"1.11.0\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"4b45fcc2344c680f5025fe57779faef368840d0bd1f42f216291f0dc4ace4744\"\n-dependencies = [\n- \"bit-set\",\n- \"bit-vec\",\n- \"bitflags 2.11.1\",\n- \"num-traits\",\n- \"rand 0.9.4\",\n- \"rand_chacha 0.9.0\",\n- \"rand_xorshift\",\n- \"regex-syntax\",\n- \"rusty-fork\",\n- \"tempfile\",\n- \"unarray\",\n-]\n-\n-[[package]]\n-name = \"quick-error\"\n-version = \"1.2.3\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"a1d01941d82fa2ab50be1e79e6714289dd7cde78eba4c074bc5a4374f650dfe0\"\n-\n [[package]]\n name = \"quote\"\n version = \"1.0.45\"\n@@ -1520,18 +1465,8 @@ source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"5ca0ecfa931c29007047d1bc58e623ab12e5590e8c7cc53200d5202b69266d8a\"\n dependencies = [\n \"libc\",\n- \"rand_chacha 0.3.1\",\n- \"rand_core 0.6.4\",\n-]\n-\n-[[package]]\n-name = \"rand\"\n-version = \"0.9.4\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"44c5af06bb1b7d3216d91932aed5265164bf384dc89cd6ba05cf59a35f5f76ea\"\n-dependencies = [\n- \"rand_chacha 0.9.0\",\n- \"rand_core 0.9.5\",\n+ \"rand_chacha\",\n+ \"rand_core\",\n ]\n \n [[package]]\n@@ -1541,17 +1476,7 @@ source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"e6c10a63a0fa32252be49d21e7709d4d4baf8d231c2dbce1eaa8141b9b127d88\"\n dependencies = [\n \"ppv-lite86\",\n- \"rand_core 0.6.4\",\n-]\n-\n-[[package]]\n-name = \"rand_chacha\"\n-version = \"0.9.0\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"d3022b5f1df60f26e1ffddd6c66e8aa15de382ae63b3a0c1bfc0e4d3e3f325cb\"\n-dependencies = [\n- \"ppv-lite86\",\n- \"rand_core 0.9.5\",\n+ \"rand_core\",\n ]\n \n [[package]]\n@@ -1563,24 +1488,6 @@ dependencies = [\n \"getrandom 0.2.17\",\n ]\n \n-[[package]]\n-name = \"rand_core\"\n-version = \"0.9.5\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"76afc826de14238e6e8c374ddcc1fa19e374fd8dd986b0d2af0d02377261d83c\"\n-dependencies = [\n- \"getrandom 0.3.4\",\n-]\n-\n-[[package]]\n-name = \"rand_xorshift\"\n-version = \"0.4.0\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"513962919efc330f829edb2535844d1b912b0fbe2ca165d613e4e8788bb05a5a\"\n-dependencies = [\n- \"rand_core 0.9.5\",\n-]\n-\n [[package]]\n name = \"redox_syscall\"\n version = \"0.5.18\"\n@@ -1747,18 +1654,6 @@ version = \"1.0.22\"\n source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"b39cdef0fa800fc44525c84ccb54a029961a8215f9619753635a9c0d2538d46d\"\n \n-[[package]]\n-name = \"rusty-fork\"\n-version = \"0.3.1\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"cc6bf79ff24e648f6da1f8d1f011e9cac26491b619e6b9280f2b47f1774e6ee2\"\n-dependencies = [\n- \"fnv\",\n- \"quick-error\",\n- \"tempfile\",\n- \"wait-timeout\",\n-]\n-\n [[package]]\n name = \"ryu\"\n version = \"1.0.23\"\n@@ -1940,7 +1835,7 @@ dependencies = [\n \"durable\",\n \"futures-util\",\n \"maud\",\n- \"rand 0.8.6\",\n+ \"rand\",\n \"reqwest\",\n \"serde\",\n \"serde_json\",\n@@ -2335,12 +2230,6 @@ version = \"0.2.5\"\n source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b\"\n \n-[[package]]\n-name = \"unarray\"\n-version = \"0.1.4\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"eaea85b334db583fe3274d12b4cd1880032beab409c0d774be044d4480ab9a94\"\n-\n [[package]]\n name = \"unicode-ident\"\n version = \"1.0.24\"\n@@ -2407,15 +2296,6 @@ version = \"0.9.5\"\n source = \"registry+https://github.com/rust-lang/crates.io-index\"\n checksum = \"0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a\"\n \n-[[package]]\n-name = \"wait-timeout\"\n-version = \"0.2.1\"\n-source = \"registry+https://github.com/rust-lang/crates.io-index\"\n-checksum = \"09ac3b126d3914f9849036f826e054cbabdc8519970b8998ddaf3b5bd3c65f11\"\n-dependencies = [\n- \"libc\",\n-]\n-\n [[package]]\n name = \"want\"\n version = \"0.3.1\"\ndiff --git a/Cargo.toml b/Cargo.toml\nindex 9c387a8e106861dae210eaf07857eee3b8dad92a..156820ec5151b709a254dafc177da10488fee566 100644\n--- a/Cargo.toml\n+++ b/Cargo.toml\n@@ -1,3 +1,3 @@\n [workspace]\n-members = [\"server\", \"durable\", \"durable-derive\"]\n+members = [\"server\"]\n resolver = \"2\"\ndiff --git a/durable-derive/Cargo.toml b/durable-derive/Cargo.toml\ndeleted file mode 100644\nindex 1e9979ef32073a9fdaff80269342637ba641f798..0000000000000000000000000000000000000000\n--- a/durable-derive/Cargo.toml\n+++ /dev/null\n@@ -1,15 +0,0 @@\n-[package]\n-name = \"durable-derive\"\n-version = \"0.2.0\"\n-edition = \"2021\"\n-authors = [\"Durable Contributors\"]\n-description = \"#[derive(Durable)] macro for the durable crate\"\n-license = \"MIT OR Apache-2.0\"\n-\n-[lib]\n-proc-macro = true\n-\n-[dependencies]\n-syn = { version = \"2.0\", features = [\"full\"] }\n-quote = \"1.0\"\n-proc-macro2 = \"1.0\"\ndiff --git a/durable-derive/src/lib.rs b/durable-derive/src/lib.rs\ndeleted file mode 100644\nindex e646ca36ae9813bbe0b3d20af4ba81e4b0e306c2..0000000000000000000000000000000000000000\n--- a/durable-derive/src/lib.rs\n+++ /dev/null\n@@ -1,90 +0,0 @@\n-//! `#[derive(Durable)]` for the `durable` crate.\n-//!\n-//! Turns a struct whose fields are durable schema types into a navigable schema:\n-//!\n-//! - implements `durable::Schema` for the struct,\n-//! - generates a `{Name}Fields` extension trait (implemented for\n-//! `durable::Path`) with one navigator method per field, and\n-//! - adds `Name::root()` / `Name::namespaced(name)` constructors.\n-//!\n-//! Each field is assigned a stable numeric id from its declaration order, which\n-//! is encoded into the on-disk key. Reordering fields changes the layout; add new\n-//! fields at the end.\n-\n-use proc_macro::TokenStream;\n-use quote::quote;\n-use syn::{parse_macro_input, Data, DeriveInput, Fields, Ident};\n-\n-#[proc_macro_derive(Durable)]\n-pub fn derive_durable(input: TokenStream) -> TokenStream {\n- let input = parse_macro_input!(input as DeriveInput);\n- let name = &input.ident;\n- let vis = &input.vis;\n- let (impl_generics, ty_generics, where_clause) = input.generics.split_for_impl();\n-\n- let fields = match &input.data {\n- Data::Struct(data) => match &data.fields {\n- Fields::Named(named) => &named.named,\n- _ => {\n- return syn::Error::new_spanned(\n- name,\n- \"#[derive(Durable)] requires a struct with named fields\",\n- )\n- .to_compile_error()\n- .into();\n- }\n- },\n- _ => {\n- return syn::Error::new_spanned(name, \"#[derive(Durable)] is only supported on structs\")\n- .to_compile_error()\n- .into();\n- }\n- };\n-\n- let mut trait_methods = Vec::new();\n- let mut impl_methods = Vec::new();\n-\n- for (index, field) in fields.iter().enumerate() {\n- let field_ident = field.ident.as_ref().expect(\"named field\");\n- let field_ty = &field.ty;\n- let field_id = index as u32;\n- trait_methods.push(quote! {\n- fn #field_ident(&self) -> ::durable::Path<#field_ty>;\n- });\n- impl_methods.push(quote! {\n- fn #field_ident(&self) -> ::durable::Path<#field_ty> {\n- self.child_field(#field_id)\n- }\n- });\n- }\n-\n- let trait_name = Ident::new(&format!(\"{name}Fields\"), name.span());\n- let trait_doc = format!(\"Field navigators for [`{name}`], implemented for `durable::Path<{name}>`.\");\n-\n- let expanded = quote! {\n- impl #impl_generics ::durable::Schema for #name #ty_generics #where_clause {}\n-\n- #[doc = #trait_doc]\n- #vis trait #trait_name {\n- #(#trait_methods)*\n- }\n-\n- impl #impl_generics #trait_name for ::durable::Path<#name #ty_generics> #where_clause {\n- #(#impl_methods)*\n- }\n-\n- impl #impl_generics #name #ty_generics #where_clause {\n- /// The root path of this schema (empty prefix; one root per database).\n- #vis fn root() -> ::durable::Path<#name #ty_generics> {\n- ::durable::Path::root()\n- }\n-\n- /// A root path namespaced under `name`, to share a database between schemas.\n- #vis fn namespaced(name: &str) -> ::durable::Path<#name #ty_generics> {\n- ::durable::Path::namespaced(name)\n- }\n- }\n- };\n-\n- expanded.into()\n-}\ndiff --git a/durable/.gitignore b/durable/.gitignore\ndeleted file mode 100644\nindex 90c273ef12d6313594dde541c3465f2a47739729..0000000000000000000000000000000000000000\n--- a/durable/.gitignore\n+++ /dev/null\n@@ -1,35 +0,0 @@\n-all.txt\n-\n-# Generated by Cargo\n-# will have compiled files and executables\n-debug/\n-target/\n-\n-# These are backup files generated by rustfmt\n-**/*.rs.bk\n-\n-# MSVC Windows builds of rustc generate these, which store debugging information\n-*.pdb\n-\n-# Generated by cargo mutants\n-# Contains mutation testing data\n-**/mutants.out*/\n-\n-# RustRover\n-# JetBrains specific template is maintained in a separate JetBrains.gitignore that can\n-# be found at https://github.com/github/gitignore/blob/main/Global/JetBrains.gitignore\n-# and can be added to the global gitignore or merged into this file. For a more nuclear\n-# option (not recommended) you can uncomment the following to ignore the entire idea folder.\n-#.idea/\n-\n-\n-# Added by cargo\n-\n-/target\n-\n-\n-# Added by cargo\n-#\n-# already existing elements were commented out\n-\n-#/target\ndiff --git a/durable/Cargo.toml b/durable/Cargo.toml\ndeleted file mode 100644\nindex 95be0c71d64862c7cdb78e64fca768dcf83be961..0000000000000000000000000000000000000000\n--- a/durable/Cargo.toml\n+++ /dev/null\n@@ -1,19 +0,0 @@\n-[package]\n-name = \"durable\"\n-version = \"0.2.0\"\n-edition = \"2021\"\n-authors = [\"Durable Contributors\"]\n-description = \"Deeply nested, precisely updatable RocksDB-backed data structures with paths-as-data\"\n-license = \"MIT OR Apache-2.0\"\n-\n-[dependencies]\n-rocksdb = \"0.21\"\n-serde = { version = \"1.0\", features = [\"derive\"] }\n-thiserror = \"1.0\"\n-ciborium = \"0.2.2\"\n-durable-derive = { path = \"../durable-derive\", version = \"0.2.0\" }\n-\n-[dev-dependencies]\n-tempfile = \"3.8\"\n-proptest = \"1.4\"\n-serde = { version = \"1.0\", features = [\"derive\"] }\ndiff --git a/durable/LICENSE b/durable/LICENSE\ndeleted file mode 100644\nindex 261eeb9e9f8b2b4b0d119366dda99c6fd7d35c64..0000000000000000000000000000000000000000\n--- a/durable/LICENSE\n+++ /dev/null\n@@ -1,201 +0,0 @@\n- Apache License\n- Version 2.0, January 2004\n- http://www.apache.org/licenses/\n-\n- TERMS AND CONDITIONS FOR USE, REPRODUCTION, AND DISTRIBUTION\n-\n- 1. Definitions.\n-\n- \"License\" shall mean the terms and conditions for use, reproduction,\n- and distribution as defined by Sections 1 through 9 of this document.\n-\n- \"Licensor\" shall mean the copyright owner or entity authorized by\n- the copyright owner that is granting the License.\n-\n- \"Legal Entity\" shall mean the union of the acting entity and all\n- other entities that control, are controlled by, or are under common\n- control with that entity. For the purposes of this definition,\n- \"control\" means (i) the power, direct or indirect, to cause the\n- direction or management of such entity, whether by contract or\n- otherwise, or (ii) ownership of fifty percent (50%) or more of the\n- outstanding shares, or (iii) beneficial ownership of such entity.\n-\n- \"You\" (or \"Your\") shall mean an individual or Legal Entity\n- exercising permissions granted by this License.\n-\n- \"Source\" form shall mean the preferred form for making modifications,\n- including but not limited to software source code, documentation\n- source, and configuration files.\n-\n- \"Object\" form shall mean any form resulting from mechanical\n- transformation or translation of a Source form, including but\n- not limited to compiled object code, generated documentation,\n- and conversions to other media types.\n-\n- \"Work\" shall mean the work of authorship, whether in Source or\n- Object form, made available under the License, as indicated by a\n- copyright notice that is included in or attached to the work\n- (an example is provided in the Appendix below).\n-\n- \"Derivative Works\" shall mean any work, whether in Source or Object\n- form, that is based on (or derived from) the Work and for which the\n- editorial revisions, annotations, elaborations, or other modifications\n- represent, as a whole, an original work of authorship. For the purposes\n- of this License, Derivative Works shall not include works that remain\n- separable from, or merely link (or bind by name) to the interfaces of,\n- the Work and Derivative Works thereof.\n-\n- \"Contribution\" shall mean any work of authorship, including\n- the original version of the Work and any modifications or additions\n- to that Work or Derivative Works thereof, that is intentionally\n- submitted to Licensor for inclusion in the Work by the copyright owner\n- or by an individual or Legal Entity authorized to submit on behalf of\n- the copyright owner. For the purposes of this definition, \"submitted\"\n- means any form of electronic, verbal, or written communication sent\n- to the Licensor or its representatives, including but not limited to\n- communication on electronic mailing lists, source code control systems,\n- and issue tracking systems that are managed by, or on behalf of, the\n- Licensor for the purpose of discussing and improving the Work, but\n- excluding communication that is conspicuously marked or otherwise\n- designated in writing by the copyright owner as \"Not a Contribution.\"\n-\n- \"Contributor\" shall mean Licensor and any individual or Legal Entity\n- on behalf of whom a Contribution has been received by Licensor and\n- subsequently incorporated within the Work.\n-\n- 2. Grant of Copyright License. Subject to the terms and conditions of\n- this License, each Contributor hereby grants to You a perpetual,\n- worldwide, non-exclusive, no-charge, royalty-free, irrevocable\n- copyright license to reproduce, prepare Derivative Works of,\n- publicly display, publicly perform, sublicense, and distribute the\n- Work and such Derivative Works in Source or Object form.\n-\n- 3. Grant of Patent License. Subject to the terms and conditions of\n- this License, each Contributor hereby grants to You a perpetual,\n- worldwide, non-exclusive, no-charge, royalty-free, irrevocable\n- (except as stated in this section) patent license to make, have made,\n- use, offer to sell, sell, import, and otherwise transfer the Work,\n- where such license applies only to those patent claims licensable\n- by such Contributor that are necessarily infringed by their\n- Contribution(s) alone or by combination of their Contribution(s)\n- with the Work to which such Contribution(s) was submitted. If You\n- institute patent litigation against any entity (including a\n- cross-claim or counterclaim in a lawsuit) alleging that the Work\n- or a Contribution incorporated within the Work constitutes direct\n- or contributory patent infringement, then any patent licenses\n- granted to You under this License for that Work shall terminate\n- as of the date such litigation is filed.\n-\n- 4. Redistribution. You may reproduce and distribute copies of the\n- Work or Derivative Works thereof in any medium, with or without\n- modifications, and in Source or Object form, provided that You\n- meet the following conditions:\n-\n- (a) You must give any other recipients of the Work or\n- Derivative Works a copy of this License; and\n-\n- (b) You must cause any modified files to carry prominent notices\n- stating that You changed the files; and\n-\n- (c) You must retain, in the Source form of any Derivative Works\n- that You distribute, all copyright, patent, trademark, and\n- attribution notices from the Source form of the Work,\n- excluding those notices that do not pertain to any part of\n- the Derivative Works; and\n-\n- (d) If the Work includes a \"NOTICE\" text file as part of its\n- distribution, then any Derivative Works that You distribute must\n- include a readable copy of the attribution notices contained\n- within such NOTICE file, excluding those notices that do not\n- pertain to any part of the Derivative Works, in at least one\n- of the following places: within a NOTICE text file distributed\n- as part of the Derivative Works; within the Source form or\n- documentation, if provided along with the Derivative Works; or,\n- within a display generated by the Derivative Works, if and\n- wherever such third-party notices normally appear. The contents\n- of the NOTICE file are for informational purposes only and\n- do not modify the License. You may add Your own attribution\n- notices within Derivative Works that You distribute, alongside\n- or as an addendum to the NOTICE text from the Work, provided\n- that such additional attribution notices cannot be construed\n- as modifying the License.\n-\n- You may add Your own copyright statement to Your modifications and\n- may provide additional or different license terms and conditions\n- for use, reproduction, or distribution of Your modifications, or\n- for any such Derivative Works as a whole, provided Your use,\n- reproduction, and distribution of the Work otherwise complies with\n- the conditions stated in this License.\n-\n- 5. Submission of Contributions. Unless You explicitly state otherwise,\n- any Contribution intentionally submitted for inclusion in the Work\n- by You to the Licensor shall be under the terms and conditions of\n- this License, without any additional terms or conditions.\n- Notwithstanding the above, nothing herein shall supersede or modify\n- the terms of any separate license agreement you may have executed\n- with Licensor regarding such Contributions.\n-\n- 6. Trademarks. This License does not grant permission to use the trade\n- names, trademarks, service marks, or product names of the Licensor,\n- except as required for reasonable and customary use in describing the\n- origin of the Work and reproducing the content of the NOTICE file.\n-\n- 7. Disclaimer of Warranty. Unless required by applicable law or\n- agreed to in writing, Licensor provides the Work (and each\n- Contributor provides its Contributions) on an \"AS IS\" BASIS,\n- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or\n- implied, including, without limitation, any warranties or conditions\n- of TITLE, NON-INFRINGEMENT, MERCHANTABILITY, or FITNESS FOR A\n- PARTICULAR PURPOSE. You are solely responsible for determining the\n- appropriateness of using or redistributing the Work and assume any\n- risks associated with Your exercise of permissions under this License.\n-\n- 8. Limitation of Liability. In no event and under no legal theory,\n- whether in tort (including negligence), contract, or otherwise,\n- unless required by applicable law (such as deliberate and grossly\n- negligent acts) or agreed to in writing, shall any Contributor be\n- liable to You for damages, including any direct, indirect, special,\n- incidental, or consequential damages of any character arising as a\n- result of this License or out of the use or inability to use the\n- Work (including but not limited to damages for loss of goodwill,\n- work stoppage, computer failure or malfunction, or any and all\n- other commercial damages or losses), even if such Contributor\n- has been advised of the possibility of such damages.\n-\n- 9. Accepting Warranty or Additional Liability. While redistributing\n- the Work or Derivative Works thereof, You may choose to offer,\n- and charge a fee for, acceptance of support, warranty, indemnity,\n- or other liability obligations and/or rights consistent with this\n- License. However, in accepting such obligations, You may act only\n- on Your own behalf and on Your sole responsibility, not on behalf\n- of any other Contributor, and only if You agree to indemnify,\n- defend, and hold each Contributor harmless for any liability\n- incurred by, or claims asserted against, such Contributor by reason\n- of your accepting any such warranty or additional liability.\n-\n- END OF TERMS AND CONDITIONS\n-\n- APPENDIX: How to apply the Apache License to your work.\n-\n- To apply the Apache License to your work, attach the following\n- boilerplate notice, with the fields enclosed by brackets \"[]\"\n- replaced with your own identifying information. (Don't include\n- the brackets!) The text should be enclosed in the appropriate\n- comment syntax for the file format. We also recommend that a\n- file or class name and description of purpose be included on the\n- same \"printed page\" as the copyright notice for easier\n- identification within third-party archives.\n-\n- Copyright [yyyy] [name of copyright owner]\n-\n- Licensed under the Apache License, Version 2.0 (the \"License\");\n- you may not use this file except in compliance with the License.\n- You may obtain a copy of the License at\n-\n- http://www.apache.org/licenses/LICENSE-2.0\n-\n- Unless required by applicable law or agreed to in writing, software\n- distributed under the License is distributed on an \"AS IS\" BASIS,\n- WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\n- See the License for the specific language governing permissions and\n- limitations under the License.\ndiff --git a/durable/README.md b/durable/README.md\ndeleted file mode 100644\nindex cd3c84ad801be756b643b8f479c2ff39e9904233..0000000000000000000000000000000000000000\n--- a/durable/README.md\n+++ /dev/null\n@@ -1,168 +0,0 @@\n-# durable\n-\n-Deeply nested, precisely updatable RocksDB-backed data structures for Rust,\n-built around **paths as data**.\n-\n-Most embedded-storage wrappers make you serialize a whole struct into one blob.\n-Updating one field means reading, deserializing, mutating, re-serializing, and\n-rewriting the entire value. `durable` takes the opposite approach: you describe\n-your data as a *schema* of composable types, and address any location with a\n-typed **path**. A path lowers to a deterministic RocksDB key with no I/O, so a\n-mutation touches exactly the keys it names — nothing else.\n-\n-```rust\n-use durable::{Db, Durable, Durability, Leaf, Map, Sum};\n-\n-#[derive(Durable)]\n-struct Store {\n- scores: Map>,\n- title: Leaf,\n-}\n-\n-fn main() -> durable::Result<()> {\n- let db = Db::open(\"scores.db\")?;\n- let root = Store::root();\n- let alice = \"alice\".to_string();\n-\n- // Three precise writes, one atomic batch, one WAL flush.\n- db.apply(\n- &[\n- root.scores().key(&alice).add(10), // blind merge — no read\n- root.scores().key(&alice).add(5),\n- root.title().set(&\"leaderboard\".to_string()),\n- ],\n- Durability::SyncWal,\n- )?;\n-\n- assert_eq!(root.scores().key(&alice).get(&db)?, 15);\n- Ok(())\n-}\n-```\n-\n-## The model\n-\n-### Schema types\n-\n-A *schema* is a type-level description of a location's shape. Compose them\n-freely:\n-\n-| Type | Meaning | Key terminal ops |\n-|------|---------|------------------|\n-| `Leaf` | one CBOR-encoded value | `get`, `set`, `delete` |\n-| `Map` | keys `K` → sub-schema `V` | `key`, `keys`, `entries`, `len`, `contains`, `clear` |\n-| `List` | index-addressed sequence | `at`, `push`, `pop`, `iter`, `len`, `clear` |\n-| `Deque` | double-ended queue (O(1) ends) | `push_back`, `push_front`, `pop_front`, `pop_back`, `front`, `back`, `iter` |\n-| `Sum` | numeric accumulator | `add` (blind merge), `get`, `set` |\n-| `#[derive(Durable)] struct` | fixed named fields | one navigator method per field |\n-\n-Leaf- and `Sum`-valued maps additionally get `get`, `iter`, and\n-`transform_values` (a one-scan bulk rewrite that yields reified writes — e.g.\n-\"decay every edge weight\").\n-\n-Nest them arbitrarily:\n-\n-```rust\n-use durable::{Deque, Durable, Leaf, Map, Sum};\n-# use serde::{Serialize, Deserialize};\n-# #[derive(Serialize, Deserialize)] struct Vote;\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct GroupState {\n- edges: Map<(u32, u32), Sum>,\n- recent_votes: Deque>,\n- item_count: Sum,\n-}\n-\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct Store {\n- scopes: Map,\n-}\n-```\n-\n-Now `Store::root().scopes().key(&scope).edges().key(&(i, j)).add(1.0)` updates a\n-single edge weight without reading or rewriting anything else in the scope.\n-\n-### Paths are data\n-\n-`Path` is just a byte prefix plus a phantom schema type. Navigation is pure\n-and allocation-light; nothing hits the database until you read or apply. Because\n-paths are values you can build them once and reuse them, pass them around, and\n-compose them.\n-\n-### Mutations are reified\n-\n-Terminal mutating operations don't perform side effects — they return a\n-[`Write`], a typed wrapper around a plain-data [`Op`] (`Put` / `Delete` /\n-`DeletePrefix` / `Merge`). Collect several and apply them atomically:\n-\n-```rust,ignore\n-let writes = vec![\n- edges.key(&(0, 1)).add(2.0),\n- edges.key(&(1, 0)).add(1.0),\n- voted_pairs.key(&(0, 1)).set(&true),\n-];\n-db.apply(&writes, Durability::DisableWal)?;\n-```\n-\n-Reified writes are inspectable and testable — you can assert on the `Op` a path\n-produces, log it, or serialize it.\n-\n-### Blind vs. read-modify-write\n-\n-The cost model is explicit, not hidden:\n-\n-- **Blind** (no read): `Leaf::set`/`delete`, `Sum::add`/`set`, `Map::clear`.\n- These are pure `Op` data and compose freely in a batch.\n-- **Read-modify-write**: `List::push`/`pop`, `Deque` pushes/pops (they read a\n- length/cursor). In a batch, appends are deferred and resolved at commit so\n- several land at contiguous indices in one atomic write.\n-- **Scan**: `Map::keys`/`iter`/`len`, `transform_values`. Prefix range scans.\n-\n-`Sum` deserves a special mention: it's backed by a RocksDB associative merge\n-operator, so `add` is a blind O(1) write whose folding happens lazily during\n-compaction — ideal for counters and edge weights.\n-\n-## Durability\n-\n-Every batch commits with an explicit policy:\n-\n-- `Durability::SyncWal` — write the WAL and fsync before returning (survives\n- power loss).\n-- `Durability::WalOnly` — write the WAL without forcing an fsync.\n-- `Durability::DisableWal` — skip the WAL. Use only for projections rebuildable\n- from another durable source of truth.\n-\n-## Key layout\n-\n-Every location lowers to a key built from length-prefixed segments\n-(`uvarint(len) ++ bytes`), which makes segment sequences self-delimiting: a\n-parent prefix only ever prefixes its own descendants, so sibling subtrees never\n-collide. Within a location prefix `P`:\n-\n-- `P` (exact) holds a `Leaf`/`Sum` value;\n-- `P ++ [0x01] ++ seg` holds child data (map entries, struct fields, elements);\n-- `P ++ [0x00] ++ seg` holds collection metadata (lengths, deque cursors).\n-\n-Deleting a subtree is a single RocksDB range delete over `[P, upper_bound(P))`.\n-\n-## What this is not\n-\n-- Not multi-process safe. One writer process; serialize writes at the app layer.\n-- Not distributed, not SQL.\n-- Map iteration order is encoded-byte order, not logical key order.\n-- On-disk struct field ids come from declaration order — add new fields at the\n- end; reordering changes the layout.\n-- Schema evolution is your responsibility. Because durable shines as a\n- *rebuildable projection*, the simplest migration is often to drop the data and\n- replay from your canonical log.\n-\n-## Testing\n-\n-```bash\n-cargo test -p durable\n-```\n-\n-Covers the codec, the merge operator, every collection kind end-to-end, atomic\n-batches, durability modes, persistence across reopen, and property tests against\n-`BTreeMap`/`VecDeque`/sum-of-deltas models.\ndiff --git a/durable/docs/design.md b/durable/docs/design.md\ndeleted file mode 100644\nindex b2f63587c4ff395d49ab6b6bd73fed2e320002e6..0000000000000000000000000000000000000000\n--- a/durable/docs/design.md\n+++ /dev/null\n@@ -1,93 +0,0 @@\n-# durable — design notes\n-\n-This document describes how `durable` actually works, so the layout and cost\n-model are auditable rather than mysterious.\n-\n-## Goals\n-\n-1. **Precise updates.** A mutation touches only the keys it names. No\n- read-deserialize-mutate-reserialize-write of a whole struct.\n-2. **Deep nesting.** Maps, lists, deques, and structs compose to arbitrary\n- depth, all in one RocksDB column family.\n-3. **Type safety.** Illegal navigation and illegal operations fail to compile.\n-4. **Paths and mutations as data.** Addresses and edits are values you can\n- build, reuse, inspect, and apply in atomic batches.\n-\n-Non-goals: multi-process concurrency, distribution, SQL, ad-hoc range queries\n-over logical key order.\n-\n-## Key encoding\n-\n-A location is a sequence of **segments**. Each segment is length-prefixed:\n-`uvarint(len) ++ bytes`. The full RocksDB key is the concatenation of a parent\n-prefix and a one-byte discriminator plus a segment per step.\n-\n-Length-prefixing makes segment sequences *self-delimiting*: no segment can be a\n-byte-prefix of a different segment, so a parent prefix only ever prefixes its own\n-descendants. Sibling subtrees never overlap.\n-\n-Within a location prefix `P`:\n-\n-| Key | Holds |\n-|-----|-------|\n-| `P` (exact) | a `Leaf` / `Sum` scalar value |\n-| `P ++ [0x01] ++ seg` | child data: map entry, struct field, list/deque element |\n-| `P ++ [0x00] ++ seg` | collection metadata: list `len`, deque `head`/`tail` |\n-\n-- **Map** entry under key `k`: segment is `cbor(k)`. Iteration is a range scan\n- over `P ++ [0x01]`; logical keys are deduplicated by their first segment\n- (nested values contribute several physical keys sharing that segment).\n-- **List** element `i`: segment is `i` as 8 big-endian bytes; `len` lives in\n- metadata.\n-- **Deque** element `i` (an `i64`, possibly negative): segment is an\n- order-preserving encoding (`(i as u64) ^ (1<<63)` big-endian) so the byte order\n- matches signed numeric order. `head`/`tail` cursors live in metadata; both ends\n- are O(1) and never renumber.\n-- **Struct** field: segment is the field's declaration-order id as a uvarint.\n-\n-Deleting a subtree is one RocksDB range delete over `[P, prefix_upper_bound(P))`\n-(falling back to a scan only when the prefix is empty or all `0xff`).\n-\n-## Types and navigation\n-\n-`Path` carries the lowered prefix bytes and a phantom schema `S`. Navigation\n-methods are implemented per concrete schema, so `Path>` has `key`,\n-`Path>` has `at`, a derived struct's `Path` has its field navigators, and\n-so on. Each step appends a segment and returns a `Path` of the child schema.\n-\n-`#[derive(Durable)]` generates, for a struct, the `Schema` impl, a `{Name}Fields`\n-extension trait of navigators implemented for `Path`, and `Name::root()` /\n-`Name::namespaced(name)` constructors.\n-\n-## Mutations and the cost model\n-\n-Terminal mutating operations return reified `Write`s wrapping a plain-data `Op`\n-(`Put` / `Delete` / `DeletePrefix` / `Merge`). They are applied via `Db::apply`\n-or pushed onto a `Batch`, which commits as a single RocksDB write with an\n-explicit `Durability`.\n-\n-- **Blind** ops carry fully-determined keys and never read: `Leaf::set`/`delete`,\n- `Sum::add`/`set`/`delete`, collection `clear`.\n-- **Read-modify-write** ops read a length or cursor: list/deque pushes and pops.\n- In a `Batch`, appends are deferred and resolved at commit so contiguous appends\n- get contiguous indices and the whole batch is one atomic write.\n-- **Scans**: `keys`/`iter`/`len`/`contains`/`transform_values`.\n-\n-### Sum and the merge operator\n-\n-`Sum` is backed by a RocksDB associative merge operator registered at\n-`Db::open`. Accumulator values are stored tagged (`[type_tag, 8 LE bytes]`) so a\n-single operator folds `f64`, `i64`, and `u64` correctly. `add(delta)` is a blind\n-`Merge` write: O(1), no read, folded lazily during compaction. This is the right\n-primitive for counters and graph edge weights.\n-\n-## Durability and recovery\n-\n-`SyncWal` fsyncs the WAL before returning; `WalOnly` writes the WAL without an\n-fsync; `DisableWal` skips it. `DisableWal` is intended for projections that can\n-be rebuilt from another durable source of truth — its writes may be lost on an\n-unclean crash.\n-\n-`durable` is deliberately not a recovery plan on its own. Pair it with a\n-canonical log if you need crash recovery, and prefer \"drop and replay\" over\n-in-place migration when a schema changes.\ndiff --git a/durable/examples/ranking.rs b/durable/examples/ranking.rs\ndeleted file mode 100644\nindex 59425a8ece3556746c11fd3fda0339ec641f002a..0000000000000000000000000000000000000000\n--- a/durable/examples/ranking.rs\n+++ /dev/null\n@@ -1,95 +0,0 @@\n-//! A pairwise-ranking scope stored with precise, point-addressable updates.\n-//!\n-//! Run with: `cargo run -p durable --example ranking`\n-//!\n-//! This mirrors the motivating use case: a \"scope\" holds an edge-weight graph, a\n-//! capped window of recent votes, and a counter. A vote updates a handful of\n-//! keys in one atomic batch — it never reads or rewrites the whole scope.\n-\n-use durable::{Db, Deque, Durability, Durable, Leaf, Map, Sum};\n-use serde::{Deserialize, Serialize};\n-\n-#[derive(Serialize, Deserialize, Clone, Debug)]\n-struct Vote {\n- winner: u32,\n- loser: u32,\n- weight: f64,\n-}\n-\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct Scope {\n- /// Directed edge weights: (from, to) -> accumulated weight.\n- edges: Map<(u32, u32), Sum>,\n- /// Most recent votes, newest at the back.\n- recent: Deque>,\n- /// Total votes recorded in this scope.\n- votes: Sum,\n-}\n-\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct Store {\n- scopes: Map,\n-}\n-\n-const RECENT_CAP: u64 = 5;\n-\n-fn record_vote(db: &Db, scope: &str, vote: Vote) -> durable::Result<()> {\n- let s = Store::root().scopes().key(&scope.to_string());\n-\n- // One atomic batch: bump the winning edge, flag the count, push the vote.\n- let mut batch = db.batch();\n- batch.write(s.edges().key(&(vote.winner, vote.loser)).add(vote.weight));\n- batch.write(s.votes().add(1));\n- batch.push_back(&s.recent(), &vote)?;\n- batch.commit_with(Durability::SyncWal)?;\n-\n- // Keep only the most recent N votes (O(1) per eviction).\n- while s.recent().len(db)? > RECENT_CAP {\n- s.recent().pop_front(db)?;\n- }\n- Ok(())\n-}\n-\n-fn main() -> durable::Result<()> {\n- let dir = tempfile::tempdir().unwrap();\n- let db = Db::open(dir.path())?;\n-\n- for i in 0..8 {\n- let (winner, loser) = (i % 3, (i + 1) % 3);\n- record_vote(\n- &db,\n- \"rust\",\n- Vote {\n- winner,\n- loser,\n- weight: 1.0 + (i as f64) * 0.1,\n- },\n- )?;\n- }\n-\n- let s = Store::root().scopes().key(&\"rust\".to_string());\n-\n- println!(\"total votes: {}\", s.votes().get(&db)?);\n- println!(\"recent window (cap {RECENT_CAP}): {}\", s.recent().len(&db)?);\n-\n- let mut edges = s.edges().iter(&db)?;\n- edges.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap());\n- println!(\"edges by weight:\");\n- for ((from, to), weight) in edges {\n- println!(\" {from} -> {to}: {weight:.1}\");\n- }\n-\n- // Decay every edge by 10% in a single scan + atomic batch.\n- let decay = s\n- .edges()\n- .transform_values(&db, |_e, w| Some(w * 0.9))?;\n- db.apply(&decay, Durability::SyncWal)?;\n- println!(\n- \"edge (0->1) after decay: {:.3}\",\n- s.edges().key(&(0, 1)).get(&db)?\n- );\n-\n- Ok(())\n-}\ndiff --git a/durable/src/codec.rs b/durable/src/codec.rs\ndeleted file mode 100644\nindex e9ded9fd10153bb88c34e089e88258c3a0475793..0000000000000000000000000000000000000000\n--- a/durable/src/codec.rs\n+++ /dev/null\n@@ -1,196 +0,0 @@\n-//! Key encoding for durable paths.\n-//!\n-//! Every durable location lowers to a deterministic RocksDB key built from\n-//! *length-prefixed segments*. A segment is `uvarint(len) ++ bytes`, which makes\n-//! segment sequences self-delimiting: no segment can be a prefix of a *different*\n-//! segment, so sibling subtrees never overlap and a parent prefix only ever\n-//! prefixes its own descendants.\n-//!\n-//! Within a location prefix `P` we reserve a one-byte discriminator:\n-//!\n-//! - `P` (exact) → a [`crate::Leaf`] scalar value lives here.\n-//! - `P ++ [DATA] ++ seg` → a child (map entry, struct field, list element).\n-//! - `P ++ [META] ++ seg` → collection metadata (e.g. a list length).\n-//!\n-//! Because children and metadata live *under* `P`, deleting a whole subtree is a\n-//! single RocksDB range delete over `[P, prefix_upper_bound(P))`.\n-\n-/// Discriminator for child data living under a location.\n-pub const DATA: u8 = 0x01;\n-/// Discriminator for collection metadata living under a location.\n-pub const META: u8 = 0x00;\n-\n-/// Append an unsigned LEB128 varint to `out`.\n-pub fn put_uvarint(out: &mut Vec, mut value: u64) {\n- loop {\n- let mut byte = (value & 0x7f) as u8;\n- value >>= 7;\n- if value != 0 {\n- byte |= 0x80;\n- }\n- out.push(byte);\n- if value == 0 {\n- break;\n- }\n- }\n-}\n-\n-/// Decode an unsigned LEB128 varint from the front of `bytes`.\n-///\n-/// Returns the value and the number of bytes consumed, or `None` if the input is\n-/// truncated or overlong.\n-pub fn read_uvarint(bytes: &[u8]) -> Option<(u64, usize)> {\n- let mut result: u64 = 0;\n- let mut shift = 0;\n- for (i, &byte) in bytes.iter().enumerate() {\n- if shift >= 64 {\n- return None;\n- }\n- result |= ((byte & 0x7f) as u64) << shift;\n- if byte & 0x80 == 0 {\n- return Some((result, i + 1));\n- }\n- shift += 7;\n- }\n- None\n-}\n-\n-/// Append a length-prefixed segment to `out`.\n-pub fn put_segment(out: &mut Vec, bytes: &[u8]) {\n- put_uvarint(out, bytes.len() as u64);\n- out.extend_from_slice(bytes);\n-}\n-\n-/// Read one length-prefixed segment from the front of `bytes`.\n-///\n-/// Returns the segment payload and the total number of bytes consumed\n-/// (including the length prefix).\n-pub fn read_segment(bytes: &[u8]) -> Option<(&[u8], usize)> {\n- let (len, header) = read_uvarint(bytes)?;\n- let len = len as usize;\n- let end = header.checked_add(len)?;\n- if end > bytes.len() {\n- return None;\n- }\n- Some((&bytes[header..end], header + len))\n-}\n-\n-/// Build the key for child `seg` under location prefix `parent`.\n-pub fn child_key(parent: &[u8], seg: &[u8]) -> Vec {\n- let mut key = Vec::with_capacity(parent.len() + 2 + seg.len());\n- key.extend_from_slice(parent);\n- key.push(DATA);\n- put_segment(&mut key, seg);\n- key\n-}\n-\n-/// The prefix under which all of `parent`'s child data lives.\n-pub fn child_scan_prefix(parent: &[u8]) -> Vec {\n- let mut key = Vec::with_capacity(parent.len() + 1);\n- key.extend_from_slice(parent);\n- key.push(DATA);\n- key\n-}\n-\n-/// Build a metadata key `name` under location prefix `parent`.\n-pub fn meta_key(parent: &[u8], name: &[u8]) -> Vec {\n- let mut key = Vec::with_capacity(parent.len() + 2 + name.len());\n- key.extend_from_slice(parent);\n- key.push(META);\n- put_segment(&mut key, name);\n- key\n-}\n-\n-/// Smallest key strictly greater than every key prefixed by `prefix`.\n-///\n-/// Returns `None` when `prefix` is empty or all `0xff` (i.e. the range extends to\n-/// the end of the keyspace), in which case callers must fall back to a scan.\n-pub fn prefix_upper_bound(prefix: &[u8]) -> Option> {\n- let mut end = prefix.to_vec();\n- while let Some(last) = end.last_mut() {\n- if *last != 0xff {\n- *last += 1;\n- return Some(end);\n- }\n- end.pop();\n- }\n- None\n-}\n-\n-/// Order-preserving encoding of an `i64` index (used by [`crate::Deque`]).\n-///\n-/// Flipping the sign bit makes the unsigned big-endian byte order match signed\n-/// numeric order, so negative front indices sort before positive ones.\n-pub fn order_i64(index: i64) -> [u8; 8] {\n- ((index as u64) ^ (1u64 << 63)).to_be_bytes()\n-}\n-\n-/// Order-preserving encoding of a `u64` index (used by [`crate::List`]).\n-pub fn order_u64(index: u64) -> [u8; 8] {\n- index.to_be_bytes()\n-}\n-\n-#[cfg(test)]\n-mod tests {\n- use super::*;\n-\n- #[test]\n- fn uvarint_roundtrip() {\n- for value in [0u64, 1, 127, 128, 300, 16384, u32::MAX as u64, u64::MAX] {\n- let mut buf = Vec::new();\n- put_uvarint(&mut buf, value);\n- let (decoded, used) = read_uvarint(&buf).unwrap();\n- assert_eq!(decoded, value);\n- assert_eq!(used, buf.len());\n- }\n- }\n-\n- #[test]\n- fn read_uvarint_rejects_truncated() {\n- assert!(read_uvarint(&[0x80]).is_none());\n- assert!(read_uvarint(&[]).is_none());\n- }\n-\n- #[test]\n- fn segment_roundtrip_and_self_delimiting() {\n- let mut buf = Vec::new();\n- put_segment(&mut buf, b\"alpha\");\n- put_segment(&mut buf, b\"\");\n- put_segment(&mut buf, &[0x00, 0xff, 0x01]);\n-\n- let (a, n1) = read_segment(&buf).unwrap();\n- assert_eq!(a, b\"alpha\");\n- let (b, n2) = read_segment(&buf[n1..]).unwrap();\n- assert_eq!(b, b\"\");\n- let (c, _) = read_segment(&buf[n1 + n2..]).unwrap();\n- assert_eq!(c, &[0x00, 0xff, 0x01]);\n- }\n-\n- #[test]\n- fn segment_no_false_prefix() {\n- // seg(\"a\") must not be a byte-prefix of seg(\"ab\"): length-prefixing guards this.\n- let mut a = Vec::new();\n- put_segment(&mut a, b\"a\");\n- let mut ab = Vec::new();\n- put_segment(&mut ab, b\"ab\");\n- assert!(!ab.starts_with(&a));\n- }\n-\n- #[test]\n- fn upper_bound_basics() {\n- assert_eq!(prefix_upper_bound(&[1, 2, 3]), Some(vec![1, 2, 4]));\n- assert_eq!(prefix_upper_bound(&[1, 2, 0xff]), Some(vec![1, 3]));\n- assert_eq!(prefix_upper_bound(&[0xff, 0xff]), None);\n- assert_eq!(prefix_upper_bound(&[]), None);\n- }\n-\n- #[test]\n- fn order_i64_is_monotonic() {\n- let mut values = [-5i64, -1, 0, 1, 5, i64::MIN, i64::MAX];\n- values.sort();\n- // Encoded byte order must match signed numeric order.\n- for pair in values.windows(2) {\n- assert!(order_i64(pair[0]) < order_i64(pair[1]));\n- }\n- }\n-}\ndiff --git a/durable/src/lib.rs b/durable/src/lib.rs\ndeleted file mode 100644\nindex b58fdcf4d5e85c6f34af10194155effa67ce200b..0000000000000000000000000000000000000000\n--- a/durable/src/lib.rs\n+++ /dev/null\n@@ -1,435 +0,0 @@\n-//! # durable\n-//!\n-//! Deeply nested, precisely updatable RocksDB-backed data structures for Rust,\n-//! built around **paths as data**.\n-//!\n-//! Instead of serializing a big struct into one blob, you describe your data with\n-//! a *schema* of composable types — [`Leaf`], [`Map`], [`List`], [`Deque`],\n-//! [`Sum`], and your own `#[derive(Durable)]` structs — and address any location\n-//! with a typed [`Path`]. A path lowers to a deterministic RocksDB key with no\n-//! I/O, so a mutation touches exactly the keys it names and nothing else.\n-//!\n-//! Terminal operations on a path return reified [`Write`] values (not side\n-//! effects). Compose several into one atomic [`Batch`] and commit them with an\n-//! explicit [`Durability`] policy.\n-//!\n-//! ```\n-//! use durable::{Db, Durable, Durability, Leaf, Map, Sum};\n-//!\n-//! #[derive(Durable)]\n-//! struct Store {\n-//! scores: Map>,\n-//! title: Leaf,\n-//! }\n-//!\n-//! // `#[derive(Durable)]` also generates a `StoreFields` navigator trait,\n-//! // in scope wherever `Store` is.\n-//!\n-//! # fn main() -> durable::Result<()> {\n-//! let dir = tempfile::tempdir().unwrap();\n-//! let db = Db::open(dir.path())?;\n-//!\n-//! let root = Store::root();\n-//! let alice = \"alice\".to_string();\n-//! db.apply(\n-//! &[\n-//! root.scores().key(&alice).add(10), // blind merge, no read\n-//! root.scores().key(&alice).add(5),\n-//! root.title().set(&\"leaderboard\".to_string()),\n-//! ],\n-//! Durability::SyncWal,\n-//! )?;\n-//!\n-//! assert_eq!(root.scores().key(&alice).get(&db)?, 15);\n-//! assert_eq!(root.title().get(&db)?, Some(\"leaderboard\".to_string()));\n-//! # Ok(())\n-//! # }\n-//! ```\n-\n-mod codec;\n-mod path;\n-mod schema;\n-\n-use std::path::Path as FsPath;\n-use std::sync::Arc;\n-\n-use rocksdb::{Options, WriteBatch, WriteOptions, DB as RocksDb};\n-use serde::{de::DeserializeOwned, Serialize};\n-use thiserror::Error;\n-\n-pub use durable_derive::Durable;\n-pub use path::Path;\n-pub use schema::{Deque, Leaf, List, Map, Schema, Sum, Summable};\n-\n-/// Errors returned by durable operations.\n-#[derive(Error, Debug)]\n-pub enum Error {\n- #[error(\"rocksdb error: {0}\")]\n- RocksDb(#[from] rocksdb::Error),\n- #[error(\"serialization error: {0}\")]\n- Serialize(String),\n- #[error(\"deserialization error: {0}\")]\n- Deserialize(String),\n- #[error(\"data corruption: {0}\")]\n- Corruption(String),\n-}\n-\n-/// Result alias used throughout the crate.\n-pub type Result = std::result::Result;\n-\n-/// CBOR-encode a value for leaf storage or key encoding.\n-pub(crate) fn encode_value(value: &T) -> Result> {\n- let mut bytes = Vec::new();\n- ciborium::ser::into_writer(value, &mut bytes).map_err(|e| Error::Serialize(e.to_string()))?;\n- Ok(bytes)\n-}\n-\n-/// CBOR-decode a stored value.\n-pub(crate) fn decode_value(bytes: &[u8]) -> Result {\n- ciborium::de::from_reader(bytes).map_err(|e| Error::Deserialize(e.to_string()))\n-}\n-\n-/// Durability policy for a committed [`Batch`] or [`Db::apply`] call.\n-#[derive(Debug, Clone, Copy, PartialEq, Eq)]\n-pub enum Durability {\n- /// Write through the WAL and fsync it before returning (survives power loss).\n- SyncWal,\n- /// Write through the WAL without forcing an fsync.\n- WalOnly,\n- /// Skip the WAL entirely. Use only for projections rebuildable from another\n- /// durable source of truth.\n- DisableWal,\n-}\n-\n-/// A single reified storage operation.\n-///\n-/// `Op` is the type-erased lowering of a typed terminal operation. It is plain\n-/// data: you can build, inspect, log, and store a list of ops, then apply them.\n-#[derive(Debug, Clone, PartialEq, Eq)]\n-pub enum Op {\n- /// Blind put of a value at an exact key.\n- Put { key: Vec, value: Vec },\n- /// Blind delete of an exact key.\n- Delete { key: Vec },\n- /// Delete every key under `prefix` (a whole subtree, including its root leaf).\n- DeletePrefix { prefix: Vec },\n- /// Blind associative merge (used by [`Sum`]).\n- Merge { key: Vec, value: Vec },\n-}\n-\n-/// A typed, reified mutation produced by a terminal path operation.\n-///\n-/// A `Write` wraps a single [`Op`]. Collect several and hand them to\n-/// [`Db::apply`] (or push them onto a [`Batch`]) to commit atomically.\n-#[derive(Debug, Clone, PartialEq, Eq)]\n-pub struct Write {\n- op: Op,\n-}\n-\n-impl Write {\n- pub(crate) fn new(op: Op) -> Self {\n- Self { op }\n- }\n-\n- /// The underlying reified operation.\n- pub fn op(&self) -> &Op {\n- &self.op\n- }\n-\n- /// Consume the write, yielding its operation.\n- pub fn into_op(self) -> Op {\n- self.op\n- }\n-}\n-\n-/// A handle to an open durable database.\n-///\n-/// Cheap to clone (an `Arc` around the RocksDB handle). Durable assumes a single\n-/// writer process; serialize writes at the application layer.\n-#[derive(Clone)]\n-pub struct Db {\n- inner: Arc,\n-}\n-\n-impl Db {\n- /// Open (or create) a durable database at `path`.\n- pub fn open>(path: P) -> Result {\n- let mut opts = Options::default();\n- opts.create_if_missing(true);\n- opts.set_merge_operator_associative(\"durable.sum\", schema::sum_merge);\n- let db = RocksDb::open(&opts, path)?;\n- Ok(Self {\n- inner: Arc::new(db),\n- })\n- }\n-\n- pub(crate) fn raw(&self) -> &RocksDb {\n- &self.inner\n- }\n-\n- /// Start an atomic batch of writes.\n- pub fn batch(&self) -> Batch {\n- Batch::new(self.clone())\n- }\n-\n- /// Apply reified writes atomically with the given durability policy.\n- pub fn apply(&self, writes: &[Write], durability: Durability) -> Result<()> {\n- let mut batch = self.batch();\n- for write in writes {\n- batch.write(write.clone());\n- }\n- batch.commit_with(durability)\n- }\n-\n- /// Apply a single write with the given durability policy.\n- pub fn run(&self, write: Write, durability: Durability) -> Result<()> {\n- self.apply(std::slice::from_ref(&write), durability)\n- }\n-}\n-\n-/// An atomic batch of writes plus deferred collection appends.\n-///\n-/// Blind writes are recorded immediately. Stateful appends ([`Batch::push_back`]\n-/// / [`Batch::push`]) are resolved at commit time against the current collection\n-/// length, so several appends in one batch land at contiguous indices and the\n-/// whole batch commits as one RocksDB write (one WAL flush).\n-pub struct Batch {\n- db: Db,\n- inner: WriteBatch,\n- prefix_deletes: Vec>,\n- appends: Vec,\n-}\n-\n-pub(crate) enum AppendEnd {\n- /// Append to the tail of a [`List`]; counter is `len`, index is `len`.\n- ListBack,\n- /// Append to the back of a [`Deque`]; counter is `tail`, index is `tail`.\n- DequeBack,\n- /// Push to the front of a [`Deque`]; counter is `head`, index is `head - 1`.\n- DequeFront,\n-}\n-\n-pub(crate) struct PendingAppend {\n- pub(crate) coll_prefix: Vec,\n- pub(crate) end: AppendEnd,\n- pub(crate) value: Vec,\n-}\n-\n-impl Batch {\n- fn new(db: Db) -> Self {\n- Self {\n- db,\n- inner: WriteBatch::default(),\n- prefix_deletes: Vec::new(),\n- appends: Vec::new(),\n- }\n- }\n-\n- /// Record a reified write in this batch.\n- pub fn write(&mut self, write: Write) {\n- match write.into_op() {\n- Op::Put { key, value } => self.inner.put(key, value),\n- Op::Delete { key } => self.inner.delete(key),\n- Op::Merge { key, value } => self.inner.merge(key, value),\n- Op::DeletePrefix { prefix } => self.prefix_deletes.push(prefix),\n- }\n- }\n-\n- /// Record several reified writes.\n- pub fn extend>(&mut self, writes: I) {\n- for write in writes {\n- self.write(write);\n- }\n- }\n-\n- pub(crate) fn raw_put(&mut self, key: Vec, value: Vec) {\n- self.inner.put(key, value);\n- }\n-\n- /// Append `value` to the back of a leaf list (resolved at commit).\n- pub fn push(&mut self, list: &Path>>, value: &T) -> Result<()> {\n- self.appends.push(PendingAppend {\n- coll_prefix: list.prefix().to_vec(),\n- end: AppendEnd::ListBack,\n- value: encode_value(value)?,\n- });\n- Ok(())\n- }\n-\n- /// Append `value` to the back of a leaf deque (resolved at commit).\n- pub fn push_back(\n- &mut self,\n- deque: &Path>>,\n- value: &T,\n- ) -> Result<()> {\n- self.appends.push(PendingAppend {\n- coll_prefix: deque.prefix().to_vec(),\n- end: AppendEnd::DequeBack,\n- value: encode_value(value)?,\n- });\n- Ok(())\n- }\n-\n- /// Push `value` to the front of a leaf deque (resolved at commit).\n- pub fn push_front(\n- &mut self,\n- deque: &Path>>,\n- value: &T,\n- ) -> Result<()> {\n- self.appends.push(PendingAppend {\n- coll_prefix: deque.prefix().to_vec(),\n- end: AppendEnd::DequeFront,\n- value: encode_value(value)?,\n- });\n- Ok(())\n- }\n-\n- /// Commit the batch, fsyncing the WAL (equivalent to\n- /// `commit_with(Durability::SyncWal)`).\n- pub fn commit(self) -> Result<()> {\n- self.commit_with(Durability::SyncWal)\n- }\n-\n- /// Commit the batch with an explicit durability policy.\n- pub fn commit_with(mut self, durability: Durability) -> Result<()> {\n- self.resolve_prefix_deletes()?;\n- self.resolve_appends()?;\n-\n- match durability {\n- Durability::SyncWal => {\n- self.db.raw().write(self.inner)?;\n- self.db.raw().flush_wal(true)?;\n- }\n- Durability::WalOnly => {\n- self.db.raw().write(self.inner)?;\n- }\n- Durability::DisableWal => {\n- let mut opts = WriteOptions::default();\n- opts.disable_wal(true);\n- self.db.raw().write_opt(self.inner, &opts)?;\n- }\n- }\n- Ok(())\n- }\n-\n- fn resolve_prefix_deletes(&mut self) -> Result<()> {\n- for prefix in std::mem::take(&mut self.prefix_deletes) {\n- match codec::prefix_upper_bound(&prefix) {\n- Some(end) => self.inner.delete_range(&prefix, &end),\n- None => {\n- // Range extends to the end of the keyspace: scan and delete.\n- let iter = self.db.raw().iterator(rocksdb::IteratorMode::From(\n- &prefix,\n- rocksdb::Direction::Forward,\n- ));\n- for item in iter {\n- let (key, _) = item?;\n- if !key.starts_with(&prefix) {\n- break;\n- }\n- self.inner.delete(&key);\n- }\n- }\n- }\n- }\n- Ok(())\n- }\n-\n- fn resolve_appends(&mut self) -> Result<()> {\n- use std::collections::HashMap;\n- // Group appends by (collection, end) so contiguous appends get contiguous\n- // indices and each counter is read exactly once.\n- let mut order: Vec<(Vec, u8)> = Vec::new();\n- let mut grouped: HashMap<(Vec, u8), Vec>> = HashMap::new();\n- for append in std::mem::take(&mut self.appends) {\n- let tag = match append.end {\n- AppendEnd::ListBack => 0u8,\n- AppendEnd::DequeBack => 1u8,\n- AppendEnd::DequeFront => 2u8,\n- };\n- let group_key = (append.coll_prefix, tag);\n- let entry = grouped.entry(group_key.clone()).or_insert_with(|| {\n- order.push(group_key);\n- Vec::new()\n- });\n- entry.push(append.value);\n- }\n-\n- for (coll_prefix, tag) in order {\n- let values = grouped.remove(&(coll_prefix.clone(), tag)).unwrap();\n- match tag {\n- 0 => self.resolve_list_back(&coll_prefix, values)?,\n- 1 => self.resolve_deque_end(&coll_prefix, values, true)?,\n- 2 => self.resolve_deque_end(&coll_prefix, values, false)?,\n- _ => unreachable!(),\n- }\n- }\n- Ok(())\n- }\n-\n- fn resolve_list_back(&mut self, coll_prefix: &[u8], values: Vec>) -> Result<()> {\n- let len_key = codec::meta_key(coll_prefix, b\"len\");\n- let mut len = read_u64(&self.db, &len_key)?.unwrap_or(0);\n- for value in values {\n- let elem = codec::child_key(coll_prefix, &codec::order_u64(len));\n- self.inner.put(&elem, &value);\n- len += 1;\n- }\n- self.inner.put(&len_key, len.to_le_bytes());\n- Ok(())\n- }\n-\n- fn resolve_deque_end(\n- &mut self,\n- coll_prefix: &[u8],\n- values: Vec>,\n- back: bool,\n- ) -> Result<()> {\n- let head_key = codec::meta_key(coll_prefix, b\"head\");\n- let tail_key = codec::meta_key(coll_prefix, b\"tail\");\n- let mut head = read_i64(&self.db, &head_key)?.unwrap_or(0);\n- let mut tail = read_i64(&self.db, &tail_key)?.unwrap_or(0);\n- for value in values {\n- if back {\n- let elem = codec::child_key(coll_prefix, &codec::order_i64(tail));\n- self.inner.put(&elem, &value);\n- tail += 1;\n- } else {\n- head -= 1;\n- let elem = codec::child_key(coll_prefix, &codec::order_i64(head));\n- self.inner.put(&elem, &value);\n- }\n- }\n- self.inner.put(&head_key, head.to_le_bytes());\n- self.inner.put(&tail_key, tail.to_le_bytes());\n- Ok(())\n- }\n-}\n-\n-pub(crate) fn read_u64(db: &Db, key: &[u8]) -> Result> {\n- match db.raw().get(key)? {\n- Some(bytes) => {\n- if bytes.len() != 8 {\n- return Err(Error::Corruption(\"expected 8-byte u64 meta\".into()));\n- }\n- let mut buf = [0u8; 8];\n- buf.copy_from_slice(&bytes);\n- Ok(Some(u64::from_le_bytes(buf)))\n- }\n- None => Ok(None),\n- }\n-}\n-\n-pub(crate) fn read_i64(db: &Db, key: &[u8]) -> Result> {\n- match db.raw().get(key)? {\n- Some(bytes) => {\n- if bytes.len() != 8 {\n- return Err(Error::Corruption(\"expected 8-byte i64 meta\".into()));\n- }\n- let mut buf = [0u8; 8];\n- buf.copy_from_slice(&bytes);\n- Ok(Some(i64::from_le_bytes(buf)))\n- }\n- None => Ok(None),\n- }\n-}\ndiff --git a/durable/src/path.rs b/durable/src/path.rs\ndeleted file mode 100644\nindex b2e8bfa4895b1be3a2e44eaf14cda31b443d0b81..0000000000000000000000000000000000000000\n--- a/durable/src/path.rs\n+++ /dev/null\n@@ -1,575 +0,0 @@\n-//! Typed paths: composable, data-only addresses into a durable schema.\n-//!\n-//! A [`Path`] is just a byte prefix plus a phantom schema type. Navigation\n-//! methods are gated by the concrete schema, so only legal steps compile, and\n-//! terminal operations return reified [`Write`]s (for mutations) or read directly\n-//! from a [`Db`].\n-\n-use std::marker::PhantomData;\n-\n-use serde::{de::DeserializeOwned, Serialize};\n-\n-use crate::{\n- codec,\n- schema::{decode_sum, encode_sum, Deque, Leaf, List, Map, Schema, Sum, Summable},\n- decode_value, encode_value, read_i64, read_u64, Db, Error, Op, Result, Write,\n-};\n-\n-/// A typed address into a durable schema.\n-///\n-/// Cheap to clone; carries only the lowered key prefix. Construct the root of a\n-/// schema with [`Path::root`] (typically via the `#[derive(Durable)]`-generated\n-/// `S::root()`), then navigate with schema-specific methods.\n-pub struct Path {\n- prefix: Vec,\n- _schema: PhantomData S>,\n-}\n-\n-impl Clone for Path {\n- fn clone(&self) -> Self {\n- Self {\n- prefix: self.prefix.clone(),\n- _schema: PhantomData,\n- }\n- }\n-}\n-\n-impl std::fmt::Debug for Path {\n- fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {\n- f.debug_struct(\"Path\")\n- .field(\"schema\", &std::any::type_name::())\n- .field(\"prefix\", &self.prefix)\n- .finish()\n- }\n-}\n-\n-impl Path {\n- /// The empty-prefixed root of a schema.\n- ///\n- /// One root schema per database. Navigate from here.\n- pub fn root() -> Self {\n- Self::from_prefix(Vec::new())\n- }\n-\n- /// A root namespaced under `name`, so multiple schemas can share one database.\n- pub fn namespaced(name: &str) -> Self {\n- let mut prefix = Vec::new();\n- codec::put_segment(&mut prefix, name.as_bytes());\n- Self::from_prefix(prefix)\n- }\n-}\n-\n-impl Path {\n- pub(crate) fn from_prefix(prefix: Vec) -> Self {\n- Self {\n- prefix,\n- _schema: PhantomData,\n- }\n- }\n-\n- /// The lowered RocksDB key prefix this path addresses.\n- pub fn prefix(&self) -> &[u8] {\n- &self.prefix\n- }\n-\n- fn child(&self, seg: &[u8]) -> Path {\n- Path::from_prefix(codec::child_key(&self.prefix, seg))\n- }\n-\n- /// Navigate to a `#[derive(Durable)]` struct field. Called by generated code.\n- #[doc(hidden)]\n- pub fn child_field(&self, field_id: u32) -> Path {\n- let mut seg = Vec::new();\n- codec::put_uvarint(&mut seg, field_id as u64);\n- self.child(&seg)\n- }\n-}\n-\n-// ---------------------------------------------------------------------------\n-// Leaf\n-// ---------------------------------------------------------------------------\n-\n-impl Path> {\n- /// Read the value at this leaf, if present.\n- pub fn get(&self, db: &Db) -> Result> {\n- match db.raw().get(&self.prefix)? {\n- Some(bytes) => Ok(Some(decode_value(&bytes)?)),\n- None => Ok(None),\n- }\n- }\n-\n- /// A reified blind write that sets this leaf to `value`.\n- pub fn set(&self, value: &T) -> Write {\n- let bytes = encode_value(value).expect(\"durable: leaf value serialization failed\");\n- Write::new(Op::Put {\n- key: self.prefix.clone(),\n- value: bytes,\n- })\n- }\n-\n- /// A reified blind write that removes this leaf.\n- pub fn delete(&self) -> Write {\n- Write::new(Op::Delete {\n- key: self.prefix.clone(),\n- })\n- }\n-}\n-\n-// ---------------------------------------------------------------------------\n-// Sum\n-// ---------------------------------------------------------------------------\n-\n-impl Path> {\n- /// Read the accumulated value (defaults to zero when absent).\n- pub fn get(&self, db: &Db) -> Result {\n- match db.raw().get(&self.prefix)? {\n- Some(bytes) => decode_sum::(&bytes)\n- .ok_or_else(|| Error::Corruption(\"malformed Sum accumulator\".into())),\n- None => Ok(N::zero()),\n- }\n- }\n-\n- /// A reified blind merge that adds `delta` to the accumulator.\n- ///\n- /// This never reads the current value: it is an O(1) write whose effect is\n- /// resolved lazily by RocksDB's merge operator.\n- pub fn add(&self, delta: N) -> Write {\n- Write::new(Op::Merge {\n- key: self.prefix.clone(),\n- value: encode_sum(delta),\n- })\n- }\n-\n- /// A reified blind write that sets the accumulator to an exact value.\n- pub fn set(&self, value: N) -> Write {\n- Write::new(Op::Put {\n- key: self.prefix.clone(),\n- value: encode_sum(value),\n- })\n- }\n-\n- /// A reified blind write that removes the accumulator.\n- pub fn delete(&self) -> Write {\n- Write::new(Op::Delete {\n- key: self.prefix.clone(),\n- })\n- }\n-}\n-\n-// ---------------------------------------------------------------------------\n-// Map\n-// ---------------------------------------------------------------------------\n-\n-impl Path> {\n- /// Navigate to the sub-schema stored under `key`.\n- pub fn key(&self, key: &K) -> Path {\n- let encoded = encode_value(key).expect(\"durable: map key serialization failed\");\n- self.child(&encoded)\n- }\n-\n- /// A reified write that deletes the entire map (all entries and metadata).\n- pub fn clear(&self) -> Write {\n- Write::new(Op::DeletePrefix {\n- prefix: self.prefix.clone(),\n- })\n- }\n-}\n-\n-impl Path> {\n- /// All keys present in the map, in stored (encoded-byte) order.\n- pub fn keys(&self, db: &Db) -> Result> {\n- let scan = codec::child_scan_prefix(&self.prefix);\n- let iter = db\n- .raw()\n- .iterator(rocksdb::IteratorMode::From(&scan, rocksdb::Direction::Forward));\n- let mut keys = Vec::new();\n- let mut last: Option> = None;\n- for item in iter {\n- let (db_key, _) = item?;\n- if !db_key.starts_with(&scan) {\n- break;\n- }\n- let rest = &db_key[scan.len()..];\n- let (key_seg, _) = codec::read_segment(rest)\n- .ok_or_else(|| Error::Corruption(\"malformed map entry key\".into()))?;\n- if last.as_deref() == Some(key_seg) {\n- continue; // same logical key, deeper sub-key\n- }\n- last = Some(key_seg.to_vec());\n- keys.push(decode_value(key_seg)?);\n- }\n- Ok(keys)\n- }\n-\n- /// The number of distinct keys in the map.\n- pub fn len(&self, db: &Db) -> Result {\n- Ok(self.keys(db)?.len())\n- }\n-\n- /// Whether the map has no entries.\n- pub fn is_empty(&self, db: &Db) -> Result {\n- let scan = codec::child_scan_prefix(&self.prefix);\n- let mut iter = db\n- .raw()\n- .iterator(rocksdb::IteratorMode::From(&scan, rocksdb::Direction::Forward));\n- match iter.next() {\n- Some(item) => {\n- let (db_key, _) = item?;\n- Ok(!db_key.starts_with(&scan))\n- }\n- None => Ok(true),\n- }\n- }\n-\n- /// Whether `key` is present.\n- pub fn contains(&self, db: &Db, key: &K) -> Result {\n- let child = self.key(key);\n- // A present entry has at least one key at-or-under the child prefix.\n- let mut iter = db.raw().iterator(rocksdb::IteratorMode::From(\n- child.prefix(),\n- rocksdb::Direction::Forward,\n- ));\n- match iter.next() {\n- Some(item) => {\n- let (db_key, _) = item?;\n- Ok(db_key.starts_with(child.prefix()))\n- }\n- None => Ok(false),\n- }\n- }\n-\n- /// All keys paired with sub-paths into their values (composable navigation).\n- pub fn entries(&self, db: &Db) -> Result)>> {\n- let keys = self.keys(db)?;\n- Ok(keys\n- .into_iter()\n- .map(|k| {\n- let path = self.key(&k);\n- (k, path)\n- })\n- .collect())\n- }\n-}\n-\n-// Leaf-valued maps gain direct value iteration and bulk transforms.\n-impl Path>> {\n- /// Read the value stored under `key`.\n- pub fn get(&self, db: &Db, key: &K) -> Result> {\n- self.key(key).get(db)\n- }\n-\n- /// All `(key, value)` pairs in stored order.\n- pub fn iter(&self, db: &Db) -> Result> {\n- let scan = codec::child_scan_prefix(&self.prefix);\n- let iter = db\n- .raw()\n- .iterator(rocksdb::IteratorMode::From(&scan, rocksdb::Direction::Forward));\n- let mut out = Vec::new();\n- for item in iter {\n- let (db_key, value) = item?;\n- if !db_key.starts_with(&scan) {\n- break;\n- }\n- let rest = &db_key[scan.len()..];\n- let (key_seg, used) = codec::read_segment(rest)\n- .ok_or_else(|| Error::Corruption(\"malformed map entry key\".into()))?;\n- // Leaf entries are exactly one physical key; reject deeper sub-keys.\n- if used != rest.len() {\n- return Err(Error::Corruption(\"unexpected nested key in leaf map\".into()));\n- }\n- out.push((decode_value(key_seg)?, decode_value(&value)?));\n- }\n- Ok(out)\n- }\n-\n- /// Build reified writes that rewrite each value through `f`.\n- ///\n- /// Returning `Some(new)` sets the value, `None` deletes the entry. This reads\n- /// the map once (a prefix scan) and yields blind writes, so the whole\n- /// transform commits atomically in one batch (e.g. \"decay every edge weight\").\n- pub fn transform_values(\n- &self,\n- db: &Db,\n- mut f: impl FnMut(&K, T) -> Option,\n- ) -> Result> {\n- let mut writes = Vec::new();\n- for (k, v) in self.iter(db)? {\n- let entry = self.key(&k);\n- match f(&k, v) {\n- Some(new) => writes.push(entry.set(&new)),\n- None => writes.push(entry.delete()),\n- }\n- }\n- Ok(writes)\n- }\n-}\n-\n-// Sum-valued maps gain direct accumulator iteration and bulk transforms.\n-impl Path>> {\n- /// Read the accumulator stored under `key` (zero when absent).\n- pub fn get(&self, db: &Db, key: &K) -> Result {\n- self.key(key).get(db)\n- }\n-\n- /// All `(key, value)` accumulator pairs in stored order.\n- pub fn iter(&self, db: &Db) -> Result> {\n- let scan = codec::child_scan_prefix(&self.prefix);\n- let iter = db\n- .raw()\n- .iterator(rocksdb::IteratorMode::From(&scan, rocksdb::Direction::Forward));\n- let mut out = Vec::new();\n- for item in iter {\n- let (db_key, value) = item?;\n- if !db_key.starts_with(&scan) {\n- break;\n- }\n- let rest = &db_key[scan.len()..];\n- let (key_seg, used) = codec::read_segment(rest)\n- .ok_or_else(|| Error::Corruption(\"malformed map entry key\".into()))?;\n- if used != rest.len() {\n- return Err(Error::Corruption(\"unexpected nested key in sum map\".into()));\n- }\n- let n = decode_sum::(&value)\n- .ok_or_else(|| Error::Corruption(\"malformed Sum accumulator\".into()))?;\n- out.push((decode_value(key_seg)?, n));\n- }\n- Ok(out)\n- }\n-\n- /// Build reified writes that rewrite each accumulator through `f`.\n- ///\n- /// `Some(new)` sets the accumulator (blind put), `None` deletes it. Reads the\n- /// map once, yields blind writes — ideal for \"decay every edge weight\".\n- pub fn transform_values(\n- &self,\n- db: &Db,\n- mut f: impl FnMut(&K, N) -> Option,\n- ) -> Result> {\n- let mut writes = Vec::new();\n- for (k, v) in self.iter(db)? {\n- let entry = self.key(&k);\n- match f(&k, v) {\n- Some(new) => writes.push(entry.set(new)),\n- None => writes.push(entry.delete()),\n- }\n- }\n- Ok(writes)\n- }\n-}\n-\n-// ---------------------------------------------------------------------------\n-// List\n-// ---------------------------------------------------------------------------\n-\n-impl Path> {\n- /// Navigate to the element at `index` (no bounds check until read).\n- pub fn at(&self, index: u64) -> Path {\n- self.child(&codec::order_u64(index))\n- }\n-\n- /// The number of elements.\n- pub fn len(&self, db: &Db) -> Result {\n- Ok(read_u64(db, &codec::meta_key(&self.prefix, b\"len\"))?.unwrap_or(0))\n- }\n-\n- /// Whether the list is empty.\n- pub fn is_empty(&self, db: &Db) -> Result {\n- Ok(self.len(db)? == 0)\n- }\n-\n- /// A reified write that deletes the whole list (elements and length).\n- pub fn clear(&self) -> Write {\n- Write::new(Op::DeletePrefix {\n- prefix: self.prefix.clone(),\n- })\n- }\n-}\n-\n-impl Path>> {\n- /// Read the element at `index`.\n- pub fn get(&self, db: &Db, index: u64) -> Result> {\n- if index >= self.len(db)? {\n- return Ok(None);\n- }\n- self.at(index).get(db)\n- }\n-\n- /// Append `value`, returning its index. Commits with `SyncWal`.\n- pub fn push(&self, db: &Db, value: &T) -> Result {\n- let mut batch = db.batch();\n- let index = self.len(db)?;\n- batch.push(self, value)?;\n- batch.commit()?;\n- Ok(index)\n- }\n-\n- /// Remove and return the last element. Commits with `SyncWal`.\n- pub fn pop(&self, db: &Db) -> Result> {\n- let len = self.len(db)?;\n- if len == 0 {\n- return Ok(None);\n- }\n- let last = len - 1;\n- let value = self.at(last).get(db)?;\n- let mut batch = db.batch();\n- batch.write(self.at(last).delete());\n- batch.raw_put(codec::meta_key(&self.prefix, b\"len\"), last.to_le_bytes().to_vec());\n- batch.commit()?;\n- Ok(value)\n- }\n-\n- /// All elements in index order.\n- pub fn iter(&self, db: &Db) -> Result> {\n- let len = self.len(db)?;\n- let mut out = Vec::with_capacity(len as usize);\n- for i in 0..len {\n- match self.at(i).get(db)? {\n- Some(v) => out.push(v),\n- None => return Err(Error::Corruption(\"list element missing below len\".into())),\n- }\n- }\n- Ok(out)\n- }\n-}\n-\n-// ---------------------------------------------------------------------------\n-// Deque\n-// ---------------------------------------------------------------------------\n-\n-impl Path> {\n- fn head(&self, db: &Db) -> Result {\n- Ok(read_i64(db, &codec::meta_key(&self.prefix, b\"head\"))?.unwrap_or(0))\n- }\n-\n- fn tail(&self, db: &Db) -> Result {\n- Ok(read_i64(db, &codec::meta_key(&self.prefix, b\"tail\"))?.unwrap_or(0))\n- }\n-\n- /// The number of elements.\n- pub fn len(&self, db: &Db) -> Result {\n- Ok((self.tail(db)? - self.head(db)?).max(0) as u64)\n- }\n-\n- /// Whether the deque is empty.\n- pub fn is_empty(&self, db: &Db) -> Result {\n- Ok(self.len(db)? == 0)\n- }\n-\n- /// A reified write that deletes the whole deque (elements and metadata).\n- pub fn clear(&self) -> Write {\n- Write::new(Op::DeletePrefix {\n- prefix: self.prefix.clone(),\n- })\n- }\n-}\n-\n-impl Path>> {\n- /// Push to the back. Commits with `SyncWal`.\n- pub fn push_back(&self, db: &Db, value: &T) -> Result<()> {\n- let mut batch = db.batch();\n- batch.push_back(self, value)?;\n- batch.commit()\n- }\n-\n- /// Push to the front. Commits with `SyncWal`.\n- pub fn push_front(&self, db: &Db, value: &T) -> Result<()> {\n- let mut batch = db.batch();\n- batch.push_front(self, value)?;\n- batch.commit()\n- }\n-\n- /// Remove and return the front element. Commits with `SyncWal`.\n- pub fn pop_front(&self, db: &Db) -> Result> {\n- let head = self.head(db)?;\n- let tail = self.tail(db)?;\n- if head >= tail {\n- return Ok(None);\n- }\n- let value = self.child::>(&codec::order_i64(head)).get(db)?;\n- let mut batch = db.batch();\n- batch.write(self.child::>(&codec::order_i64(head)).delete());\n- batch.raw_put(\n- codec::meta_key(&self.prefix, b\"head\"),\n- (head + 1).to_le_bytes().to_vec(),\n- );\n- batch.commit()?;\n- Ok(value)\n- }\n-\n- /// Remove and return the back element. Commits with `SyncWal`.\n- pub fn pop_back(&self, db: &Db) -> Result> {\n- let head = self.head(db)?;\n- let tail = self.tail(db)?;\n- if head >= tail {\n- return Ok(None);\n- }\n- let last = tail - 1;\n- let value = self.child::>(&codec::order_i64(last)).get(db)?;\n- let mut batch = db.batch();\n- batch.write(self.child::>(&codec::order_i64(last)).delete());\n- batch.raw_put(\n- codec::meta_key(&self.prefix, b\"tail\"),\n- last.to_le_bytes().to_vec(),\n- );\n- batch.commit()?;\n- Ok(value)\n- }\n-\n- /// Read the front element without removing it.\n- pub fn front(&self, db: &Db) -> Result> {\n- let head = self.head(db)?;\n- if head >= self.tail(db)? {\n- return Ok(None);\n- }\n- self.child::>(&codec::order_i64(head)).get(db)\n- }\n-\n- /// Read the back element without removing it.\n- pub fn back(&self, db: &Db) -> Result> {\n- let tail = self.tail(db)?;\n- if self.head(db)? >= tail {\n- return Ok(None);\n- }\n- self.child::>(&codec::order_i64(tail - 1)).get(db)\n- }\n-\n- /// All elements from front to back.\n- pub fn iter(&self, db: &Db) -> Result> {\n- let head = self.head(db)?;\n- let tail = self.tail(db)?;\n- let mut out = Vec::with_capacity((tail - head).max(0) as usize);\n- for idx in head..tail {\n- match self.child::>(&codec::order_i64(idx)).get(db)? {\n- Some(v) => out.push(v),\n- None => return Err(Error::Corruption(\"deque element missing in range\".into())),\n- }\n- }\n- Ok(out)\n- }\n-\n- /// Drop elements from the back until the length is at most `max_len`,\n- /// committing with the given durability. A no-op when already short enough.\n- pub fn truncate_back(\n- &self,\n- db: &Db,\n- max_len: u64,\n- durability: crate::Durability,\n- ) -> Result<()> {\n- let head = self.head(db)?;\n- let tail = self.tail(db)?;\n- let len = (tail - head).max(0) as u64;\n- if len <= max_len {\n- return Ok(());\n- }\n- let new_tail = tail - (len - max_len) as i64;\n- let mut batch = db.batch();\n- for idx in new_tail..tail {\n- batch.write(self.child::>(&codec::order_i64(idx)).delete());\n- }\n- batch.raw_put(\n- codec::meta_key(&self.prefix, b\"tail\"),\n- new_tail.to_le_bytes().to_vec(),\n- );\n- batch.commit_with(durability)\n- }\n-}\ndiff --git a/durable/src/schema.rs b/durable/src/schema.rs\ndeleted file mode 100644\nindex f804939aee51683f5773f2415e122ed64ea40020..0000000000000000000000000000000000000000\n--- a/durable/src/schema.rs\n+++ /dev/null\n@@ -1,219 +0,0 @@\n-//! Type-level schema markers.\n-//!\n-//! A *schema* describes the shape of a durable location at the type level. It is\n-//! never instantiated; it only parameterizes a [`crate::Path`] so the compiler\n-//! knows which navigation steps and terminal operations are legal.\n-//!\n-//! - [`Leaf`] — a single CBOR-encoded scalar value.\n-//! - [`Map`] — keys of type `K` to sub-schema `V`.\n-//! - [`List`] — an index-addressed sequence of sub-schema `V`.\n-//! - [`Deque`] — a double-ended queue of sub-schema `V` (O(1) ends).\n-//! - [`Sum`] — a numeric accumulator updated with blind merge writes.\n-//! - any `#[derive(Durable)]` struct — a fixed set of named fields.\n-\n-use std::marker::PhantomData;\n-\n-/// Marker trait implemented by every durable schema.\n-///\n-/// Implemented for [`Leaf`], [`Map`], [`List`], [`Deque`], [`Sum`], and by\n-/// `#[derive(Durable)]` for user structs. It is intentionally minimal; behaviour\n-/// lives on `Path` impls keyed by the concrete schema.\n-pub trait Schema {}\n-\n-/// A single CBOR-encoded scalar value of type `T`.\n-pub struct Leaf(PhantomData);\n-impl Schema for Leaf {}\n-\n-/// A map from keys of type `K` to sub-schema `V`.\n-pub struct Map(PhantomData<(K, V)>);\n-impl Schema for Map {}\n-\n-/// An index-addressed growable sequence of sub-schema `V`.\n-pub struct List(PhantomData);\n-impl Schema for List {}\n-\n-/// A double-ended queue of sub-schema `V` with O(1) push/pop at both ends.\n-pub struct Deque(PhantomData);\n-impl Schema for Deque {}\n-\n-/// A numeric accumulator. Updated with blind, associative merge writes so\n-/// incrementing is O(1) and never reads the current value.\n-pub struct Sum(PhantomData);\n-impl Schema for Sum {}\n-\n-/// Numbers that can back a [`Sum`] accumulator.\n-///\n-/// Stored on disk as `[TAG, b0..b7]`: a one-byte type tag plus the 8-byte\n-/// little-endian payload. The tag lets a single RocksDB merge operator fold\n-/// `f64` and `i64` accumulators correctly.\n-pub trait Summable: Copy + 'static {\n- /// Disk type tag, unique per numeric type.\n- const TAG: u8;\n- /// Additive identity.\n- fn zero() -> Self;\n- /// Combine two values (sum).\n- fn combine(self, other: Self) -> Self;\n- /// Little-endian 8-byte payload.\n- fn to_le_payload(self) -> [u8; 8];\n- /// Decode from a little-endian 8-byte payload.\n- fn from_le_payload(bytes: [u8; 8]) -> Self;\n-}\n-\n-impl Summable for f64 {\n- const TAG: u8 = 0;\n- fn zero() -> Self {\n- 0.0\n- }\n- fn combine(self, other: Self) -> Self {\n- self + other\n- }\n- fn to_le_payload(self) -> [u8; 8] {\n- self.to_le_bytes()\n- }\n- fn from_le_payload(bytes: [u8; 8]) -> Self {\n- f64::from_le_bytes(bytes)\n- }\n-}\n-\n-impl Summable for i64 {\n- const TAG: u8 = 1;\n- fn zero() -> Self {\n- 0\n- }\n- fn combine(self, other: Self) -> Self {\n- self.wrapping_add(other)\n- }\n- fn to_le_payload(self) -> [u8; 8] {\n- self.to_le_bytes()\n- }\n- fn from_le_payload(bytes: [u8; 8]) -> Self {\n- i64::from_le_bytes(bytes)\n- }\n-}\n-\n-impl Summable for u64 {\n- const TAG: u8 = 2;\n- fn zero() -> Self {\n- 0\n- }\n- fn combine(self, other: Self) -> Self {\n- self.wrapping_add(other)\n- }\n- fn to_le_payload(self) -> [u8; 8] {\n- self.to_le_bytes()\n- }\n- fn from_le_payload(bytes: [u8; 8]) -> Self {\n- u64::from_le_bytes(bytes)\n- }\n-}\n-\n-/// Encode a `Summable` to its tagged on-disk form `[TAG, b0..b7]`.\n-pub(crate) fn encode_sum(value: N) -> Vec {\n- let mut out = Vec::with_capacity(9);\n- out.push(N::TAG);\n- out.extend_from_slice(&value.to_le_payload());\n- out\n-}\n-\n-/// Decode a tagged accumulator payload back to `N`, validating the tag.\n-pub(crate) fn decode_sum(bytes: &[u8]) -> Option {\n- if bytes.len() != 9 || bytes[0] != N::TAG {\n- return None;\n- }\n- let mut payload = [0u8; 8];\n- payload.copy_from_slice(&bytes[1..9]);\n- Some(N::from_le_payload(payload))\n-}\n-\n-/// Fold one tagged operand into a running tagged accumulator.\n-///\n-/// Used by the RocksDB merge operator. Operands of mismatched tags are skipped\n-/// rather than panicking, keeping compaction resilient to stray bytes.\n-fn fold_tagged(acc: &mut Option<[u8; 9]>, operand: &[u8]) {\n- if operand.len() != 9 {\n- return;\n- }\n- let tag = operand[0];\n- let mut op_payload = [0u8; 8];\n- op_payload.copy_from_slice(&operand[1..9]);\n-\n- match acc {\n- Some(existing) if existing[0] == tag => {\n- let mut acc_payload = [0u8; 8];\n- acc_payload.copy_from_slice(&existing[1..9]);\n- let combined = match tag {\n- 0 => f64::from_le_bytes(acc_payload)\n- .combine(f64::from_le_bytes(op_payload))\n- .to_le_payload(),\n- 1 => i64::from_le_bytes(acc_payload)\n- .combine(i64::from_le_bytes(op_payload))\n- .to_le_payload(),\n- 2 => u64::from_le_bytes(acc_payload)\n- .combine(u64::from_le_bytes(op_payload))\n- .to_le_payload(),\n- _ => return,\n- };\n- existing[1..9].copy_from_slice(&combined);\n- }\n- Some(_) => {} // tag mismatch: ignore stray operand\n- None => {\n- let mut start = [0u8; 9];\n- start[0] = tag;\n- start[1..9].copy_from_slice(&op_payload);\n- *acc = Some(start);\n- }\n- }\n-}\n-\n-/// Associative merge operator registered on every durable database so that\n-/// [`Sum`] accumulators can be incremented with blind `merge` writes.\n-pub(crate) fn sum_merge(\n- _key: &[u8],\n- existing: Option<&[u8]>,\n- operands: &rocksdb::MergeOperands,\n-) -> Option> {\n- let mut acc: Option<[u8; 9]> = None;\n- if let Some(existing) = existing {\n- if existing.len() == 9 {\n- let mut start = [0u8; 9];\n- start.copy_from_slice(existing);\n- acc = Some(start);\n- }\n- }\n- for operand in operands.iter() {\n- fold_tagged(&mut acc, operand);\n- }\n- acc.map(|bytes| bytes.to_vec())\n-}\n-\n-#[cfg(test)]\n-mod tests {\n- use super::*;\n-\n- #[test]\n- fn sum_roundtrip_tagged() {\n- assert_eq!(decode_sum::(&encode_sum(2.5f64)), Some(2.5));\n- assert_eq!(decode_sum::(&encode_sum(-7i64)), Some(-7));\n- assert_eq!(decode_sum::(&encode_sum(9u64)), Some(9));\n- // Tag mismatch is rejected.\n- assert_eq!(decode_sum::(&encode_sum(2.5f64)), None);\n- }\n-\n- #[test]\n- fn fold_accumulates_same_tag() {\n- let mut acc = None;\n- fold_tagged(&mut acc, &encode_sum(1.5f64));\n- fold_tagged(&mut acc, &encode_sum(2.0f64));\n- let bytes = acc.unwrap();\n- assert_eq!(decode_sum::(&bytes), Some(3.5));\n- }\n-\n- #[test]\n- fn fold_skips_mismatched_tag() {\n- let mut acc = None;\n- fold_tagged(&mut acc, &encode_sum(5i64));\n- fold_tagged(&mut acc, &encode_sum(1.0f64)); // ignored\n- let bytes = acc.unwrap();\n- assert_eq!(decode_sum::(&bytes), Some(5));\n- }\n-}\ndiff --git a/durable/tests/integration.rs b/durable/tests/integration.rs\ndeleted file mode 100644\nindex 6aef7049f5ae442c5153dddb37523308d77503f6..0000000000000000000000000000000000000000\n--- a/durable/tests/integration.rs\n+++ /dev/null\n@@ -1,433 +0,0 @@\n-//! End-to-end tests for the durable paths-as-data API.\n-\n-use durable::{Db, Durability, Durable, Deque, Leaf, List, Map, Op, Sum};\n-use serde::{Deserialize, Serialize};\n-use tempfile::TempDir;\n-\n-#[derive(Serialize, Deserialize, Clone, PartialEq, Debug)]\n-struct Vote {\n- a: String,\n- b: String,\n- ratio: i32,\n-}\n-\n-/// A scope's local ranking state — the kind of thing that used to be one CBOR\n-/// blob, now addressable field-by-field and key-by-key.\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct GroupState {\n- edges: Map<(u32, u32), Sum>,\n- voted_pairs: Map<(u32, u32), Leaf>,\n- recent_votes: Deque>,\n- item_count: Sum,\n-}\n-\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct Store {\n- scopes: Map,\n- nodes: Map>,\n- log: List>,\n-}\n-\n-fn open() -> (TempDir, Db) {\n- let dir = TempDir::new().unwrap();\n- let db = Db::open(dir.path()).unwrap();\n- (dir, db)\n-}\n-\n-#[test]\n-fn leaf_set_get_delete() {\n- let (_dir, db) = open();\n- let root = Store::root();\n- let k = \"reddit.com/r/rust\".to_string();\n-\n- assert_eq!(root.nodes().key(&k).get(&db).unwrap(), None);\n- db.run(root.nodes().key(&k).set(&\"Rust\".to_string()), Durability::SyncWal)\n- .unwrap();\n- assert_eq!(root.nodes().key(&k).get(&db).unwrap(), Some(\"Rust\".to_string()));\n-\n- db.run(root.nodes().key(&k).delete(), Durability::SyncWal).unwrap();\n- assert_eq!(root.nodes().key(&k).get(&db).unwrap(), None);\n-}\n-\n-#[test]\n-fn sum_accumulates_with_blind_merges() {\n- let (_dir, db) = open();\n- let edges = Store::root().scopes().key(&\"s\".to_string()).edges();\n- let e = (3u32, 7u32);\n-\n- // Several blind merges in one atomic batch — no reads involved.\n- db.apply(\n- &[\n- edges.key(&e).add(2.0),\n- edges.key(&e).add(1.0),\n- edges.key(&e).add(0.5),\n- ],\n- Durability::SyncWal,\n- )\n- .unwrap();\n- assert_eq!(edges.key(&e).get(&db).unwrap(), 3.5);\n-\n- // Negative delta decrements; absent key reads as zero.\n- db.run(edges.key(&e).add(-1.5), Durability::SyncWal).unwrap();\n- assert_eq!(edges.key(&e).get(&db).unwrap(), 2.0);\n- assert_eq!(edges.key(&(9, 9)).get(&db).unwrap(), 0.0);\n-}\n-\n-#[test]\n-fn sum_set_then_merge() {\n- let (_dir, db) = open();\n- let count = Store::root().scopes().key(&\"s\".to_string()).item_count();\n- db.run(count.set(10), Durability::SyncWal).unwrap();\n- db.run(count.add(5), Durability::SyncWal).unwrap();\n- assert_eq!(count.get(&db).unwrap(), 15);\n-}\n-\n-#[test]\n-fn reified_writes_are_inspectable_data() {\n- let edges = Store::root().scopes().key(&\"s\".to_string()).edges();\n- let merge = edges.key(&(1u32, 2u32)).add(1.0);\n- assert!(matches!(merge.op(), Op::Merge { .. }));\n-\n- let put = Store::root().nodes().key(&\"x\".to_string()).set(&\"y\".to_string());\n- assert!(matches!(put.op(), Op::Put { .. }));\n-\n- let clear = Store::root().scopes().clear();\n- assert!(matches!(clear.op(), Op::DeletePrefix { .. }));\n-}\n-\n-#[test]\n-fn point_update_touches_only_its_own_key() {\n- let (_dir, db) = open();\n- let root = Store::root();\n- let rust = root.scopes().key(&\"rust\".to_string());\n- let python = root.scopes().key(&\"python\".to_string());\n-\n- // Populate two scopes with several edges and some recent votes.\n- let mut batch = db.batch();\n- for j in 0..5u32 {\n- batch.write(rust.edges().key(&(0, j)).add(j as f64 + 1.0));\n- batch.write(python.edges().key(&(0, j)).add(100.0));\n- }\n- batch\n- .push_back(\n- &rust.recent_votes(),\n- &Vote { a: \"a\".into(), b: \"b\".into(), ratio: 2 },\n- )\n- .unwrap();\n- batch.commit().unwrap();\n-\n- // A single precise update to one edge in `rust`.\n- db.run(rust.edges().key(&(0, 2)).add(10.0), Durability::SyncWal)\n- .unwrap();\n-\n- // Only that edge changed.\n- assert_eq!(rust.edges().key(&(0, 2)).get(&db).unwrap(), 13.0);\n- assert_eq!(rust.edges().key(&(0, 0)).get(&db).unwrap(), 1.0);\n- assert_eq!(rust.edges().key(&(0, 4)).get(&db).unwrap(), 5.0);\n- // The other scope is entirely untouched.\n- for j in 0..5u32 {\n- assert_eq!(python.edges().key(&(0, j)).get(&db).unwrap(), 100.0);\n- }\n- // And the unrelated recent_votes deque is intact.\n- assert_eq!(rust.recent_votes().len(&db).unwrap(), 1);\n-}\n-\n-#[test]\n-fn map_keys_len_contains_entries() {\n- let (_dir, db) = open();\n- let nodes = Store::root().nodes();\n- db.apply(\n- &[\n- nodes.key(&\"a\".to_string()).set(&\"1\".to_string()),\n- nodes.key(&\"b\".to_string()).set(&\"2\".to_string()),\n- nodes.key(&\"c\".to_string()).set(&\"3\".to_string()),\n- ],\n- Durability::SyncWal,\n- )\n- .unwrap();\n-\n- let mut keys = nodes.keys(&db).unwrap();\n- keys.sort();\n- assert_eq!(keys, vec![\"a\".to_string(), \"b\".to_string(), \"c\".to_string()]);\n- assert_eq!(nodes.len(&db).unwrap(), 3);\n- assert!(nodes.contains(&db, &\"b\".to_string()).unwrap());\n- assert!(!nodes.contains(&db, &\"z\".to_string()).unwrap());\n-\n- let mut pairs = nodes.iter(&db).unwrap();\n- pairs.sort();\n- assert_eq!(\n- pairs,\n- vec![\n- (\"a\".to_string(), \"1\".to_string()),\n- (\"b\".to_string(), \"2\".to_string()),\n- (\"c\".to_string(), \"3\".to_string()),\n- ]\n- );\n-}\n-\n-#[test]\n-fn map_keys_dedup_across_nested_subkeys() {\n- // A map whose values are nested structs has many physical keys per logical\n- // key; `keys()`/`len()` must dedup to distinct logical keys.\n- let (_dir, db) = open();\n- let scopes = Store::root().scopes();\n- let rust = scopes.key(&\"rust\".to_string());\n-\n- let mut batch = db.batch();\n- batch.write(rust.edges().key(&(0, 1)).add(1.0));\n- batch.write(rust.edges().key(&(0, 2)).add(1.0));\n- batch.write(rust.item_count().add(3));\n- batch\n- .push_back(&rust.recent_votes(), &Vote { a: \"a\".into(), b: \"b\".into(), ratio: 1 })\n- .unwrap();\n- batch.write(scopes.key(&\"python\".to_string()).item_count().add(1));\n- batch.commit().unwrap();\n-\n- let mut keys = scopes.keys(&db).unwrap();\n- keys.sort();\n- assert_eq!(keys, vec![\"python\".to_string(), \"rust\".to_string()]);\n- assert_eq!(scopes.len(&db).unwrap(), 2);\n-}\n-\n-#[test]\n-fn map_clear_deletes_subtree_only() {\n- let (_dir, db) = open();\n- let rust = Store::root().scopes().key(&\"rust\".to_string());\n-\n- let mut batch = db.batch();\n- batch.write(rust.edges().key(&(0, 1)).add(1.0));\n- batch.write(rust.edges().key(&(0, 2)).add(2.0));\n- batch.write(rust.item_count().add(5));\n- batch.commit().unwrap();\n-\n- // Clear only the edges sub-map.\n- db.run(rust.edges().clear(), Durability::SyncWal).unwrap();\n-\n- assert_eq!(rust.edges().len(&db).unwrap(), 0);\n- assert_eq!(rust.edges().key(&(0, 1)).get(&db).unwrap(), 0.0);\n- // Sibling field under the same scope is untouched.\n- assert_eq!(rust.item_count().get(&db).unwrap(), 5);\n-}\n-\n-#[test]\n-fn transform_values_decays_all_edges_in_one_batch() {\n- let (_dir, db) = open();\n- let edges = Store::root().scopes().key(&\"rust\".to_string()).edges();\n-\n- let mut batch = db.batch();\n- for j in 1..=4u32 {\n- batch.write(edges.key(&(0, j)).add(j as f64 * 10.0));\n- }\n- batch.commit().unwrap();\n-\n- // Decay every edge by half, dropping any that fall to/under 5.0 — built as\n- // reified writes from one scan, applied atomically.\n- let writes = edges\n- .transform_values(&db, |_k, w| {\n- let decayed = w * 0.5;\n- if decayed <= 5.0 {\n- None\n- } else {\n- Some(decayed)\n- }\n- })\n- .unwrap();\n- db.apply(&writes, Durability::SyncWal).unwrap();\n-\n- assert_eq!(edges.key(&(0, 1)).get(&db).unwrap(), 0.0); // 10*0.5=5.0 -> dropped\n- assert_eq!(edges.key(&(0, 2)).get(&db).unwrap(), 10.0);\n- assert_eq!(edges.key(&(0, 3)).get(&db).unwrap(), 15.0);\n- assert_eq!(edges.key(&(0, 4)).get(&db).unwrap(), 20.0);\n- assert_eq!(edges.len(&db).unwrap(), 3);\n-}\n-\n-#[test]\n-fn list_push_pop_iter() {\n- let (_dir, db) = open();\n- let log = Store::root().log();\n-\n- assert!(log.is_empty(&db).unwrap());\n- assert_eq!(log.push(&db, &10).unwrap(), 0);\n- assert_eq!(log.push(&db, &20).unwrap(), 1);\n- assert_eq!(log.push(&db, &30).unwrap(), 2);\n-\n- assert_eq!(log.len(&db).unwrap(), 3);\n- assert_eq!(log.get(&db, 1).unwrap(), Some(20));\n- assert_eq!(log.get(&db, 3).unwrap(), None);\n- assert_eq!(log.iter(&db).unwrap(), vec![10, 20, 30]);\n-\n- assert_eq!(log.pop(&db).unwrap(), Some(30));\n- assert_eq!(log.len(&db).unwrap(), 2);\n- assert_eq!(log.iter(&db).unwrap(), vec![10, 20]);\n-}\n-\n-#[test]\n-fn batched_list_pushes_get_contiguous_indices() {\n- let (_dir, db) = open();\n- let log = Store::root().log();\n-\n- let mut batch = db.batch();\n- batch.push(&log, &1).unwrap();\n- batch.push(&log, &2).unwrap();\n- batch.push(&log, &3).unwrap();\n- batch.commit().unwrap();\n- assert_eq!(log.iter(&db).unwrap(), vec![1, 2, 3]);\n-\n- // A second batch continues from the persisted length.\n- let mut batch = db.batch();\n- batch.push(&log, &4).unwrap();\n- batch.push(&log, &5).unwrap();\n- batch.commit().unwrap();\n- assert_eq!(log.iter(&db).unwrap(), vec![1, 2, 3, 4, 5]);\n- assert_eq!(log.len(&db).unwrap(), 5);\n-}\n-\n-#[test]\n-fn deque_behaves_like_a_double_ended_queue() {\n- let (_dir, db) = open();\n- let dq = Store::root().scopes().key(&\"s\".to_string()).recent_votes();\n- let v = |n: i32| Vote { a: format!(\"a{n}\"), b: format!(\"b{n}\"), ratio: n };\n-\n- dq.push_back(&db, &v(1)).unwrap();\n- dq.push_back(&db, &v(2)).unwrap();\n- dq.push_front(&db, &v(0)).unwrap();\n-\n- assert_eq!(dq.len(&db).unwrap(), 3);\n- assert_eq!(dq.iter(&db).unwrap(), vec![v(0), v(1), v(2)]);\n- assert_eq!(dq.front(&db).unwrap(), Some(v(0)));\n- assert_eq!(dq.back(&db).unwrap(), Some(v(2)));\n-\n- assert_eq!(dq.pop_front(&db).unwrap(), Some(v(0)));\n- assert_eq!(dq.pop_back(&db).unwrap(), Some(v(2)));\n- assert_eq!(dq.iter(&db).unwrap(), vec![v(1)]);\n- assert_eq!(dq.pop_front(&db).unwrap(), Some(v(1)));\n- assert_eq!(dq.pop_front(&db).unwrap(), None);\n- assert!(dq.is_empty(&db).unwrap());\n-}\n-\n-#[test]\n-fn deque_supports_capped_recent_window() {\n- // The motivating use case: keep only the most recent N votes, O(1) per insert.\n- let (_dir, db) = open();\n- let dq = Store::root().scopes().key(&\"s\".to_string()).recent_votes();\n- const CAP: u64 = 3;\n-\n- for n in 0..10 {\n- dq.push_back(&db, &Vote { a: format!(\"{n}\"), b: \"x\".into(), ratio: n })\n- .unwrap();\n- while dq.len(&db).unwrap() > CAP {\n- dq.pop_front(&db).unwrap();\n- }\n- }\n-\n- let kept = dq.iter(&db).unwrap();\n- assert_eq!(kept.len(), 3);\n- assert_eq!(kept.iter().map(|v| v.ratio).collect::>(), vec![7, 8, 9]);\n-}\n-\n-#[test]\n-fn deque_truncate_back_caps_length_keeping_front() {\n- let (_dir, db) = open();\n- let dq = Store::root().scopes().key(&\"s\".to_string()).recent_votes();\n- for n in 0..10 {\n- dq.push_back(&db, &Vote { a: format!(\"{n}\"), b: \"x\".into(), ratio: n })\n- .unwrap();\n- }\n- // Keep only the 3 oldest at front (drop the back/newest beyond cap).\n- dq.truncate_back(&db, 3, Durability::SyncWal).unwrap();\n- let kept = dq.iter(&db).unwrap();\n- assert_eq!(kept.iter().map(|v| v.ratio).collect::>(), vec![0, 1, 2]);\n-\n- // Truncating to a larger-or-equal cap is a no-op.\n- dq.truncate_back(&db, 10, Durability::SyncWal).unwrap();\n- assert_eq!(dq.len(&db).unwrap(), 3);\n-}\n-\n-#[test]\n-fn one_batch_commits_all_or_nothing_and_persists() {\n- let dir = TempDir::new().unwrap();\n- let rust_key = \"rust\".to_string();\n- {\n- let db = Db::open(dir.path()).unwrap();\n- let rust = Store::root().scopes().key(&rust_key);\n- // A \"vote\" as one atomic batch: two edge merges, a pair flag, a recent\n- // vote, and a counter — all distinct keys, one WAL flush.\n- let mut batch = db.batch();\n- batch.write(rust.edges().key(&(0, 1)).add(2.0));\n- batch.write(rust.edges().key(&(1, 0)).add(1.0));\n- batch.write(rust.voted_pairs().key(&(0, 1)).set(&true));\n- batch\n- .push_back(&rust.recent_votes(), &Vote { a: \"0\".into(), b: \"1\".into(), ratio: 2 })\n- .unwrap();\n- batch.write(rust.item_count().add(2));\n- batch.commit().unwrap();\n- }\n-\n- // Reopen: SyncWal data survives.\n- let db = Db::open(dir.path()).unwrap();\n- let rust = Store::root().scopes().key(&rust_key);\n- assert_eq!(rust.edges().key(&(0, 1)).get(&db).unwrap(), 2.0);\n- assert_eq!(rust.edges().key(&(1, 0)).get(&db).unwrap(), 1.0);\n- assert_eq!(rust.voted_pairs().key(&(0, 1)).get(&db).unwrap(), Some(true));\n- assert_eq!(rust.recent_votes().len(&db).unwrap(), 1);\n- assert_eq!(rust.item_count().get(&db).unwrap(), 2);\n-}\n-\n-#[test]\n-fn disable_wal_visible_within_session() {\n- let (_dir, db) = open();\n- let count = Store::root().scopes().key(&\"s\".to_string()).item_count();\n- db.run(count.add(7), Durability::DisableWal).unwrap();\n- assert_eq!(count.get(&db).unwrap(), 7);\n-}\n-\n-#[test]\n-fn wal_only_durability_writes() {\n- let (_dir, db) = open();\n- let node = Store::root().nodes().key(&\"k\".to_string());\n- db.run(node.set(&\"v\".to_string()), Durability::WalOnly).unwrap();\n- assert_eq!(node.get(&db).unwrap(), Some(\"v\".to_string()));\n-}\n-\n-#[test]\n-fn namespaced_roots_do_not_collide() {\n- let (_dir, db) = open();\n- let a = Store::namespaced(\"a\");\n- let b = Store::namespaced(\"b\");\n- db.run(a.nodes().key(&\"k\".to_string()).set(&\"av\".to_string()), Durability::SyncWal)\n- .unwrap();\n- db.run(b.nodes().key(&\"k\".to_string()).set(&\"bv\".to_string()), Durability::SyncWal)\n- .unwrap();\n-\n- assert_eq!(a.nodes().key(&\"k\".to_string()).get(&db).unwrap(), Some(\"av\".to_string()));\n- assert_eq!(b.nodes().key(&\"k\".to_string()).get(&db).unwrap(), Some(\"bv\".to_string()));\n-}\n-\n-#[test]\n-fn persistence_across_reopen_for_all_collection_kinds() {\n- let dir = TempDir::new().unwrap();\n- {\n- let db = Db::open(dir.path()).unwrap();\n- let root = Store::root();\n- let s = root.scopes().key(&\"s\".to_string());\n- db.run(root.nodes().key(&\"n\".to_string()).set(&\"N\".to_string()), Durability::SyncWal)\n- .unwrap();\n- root.log().push(&db, &42).unwrap();\n- db.run(s.edges().key(&(1, 2)).add(9.0), Durability::SyncWal).unwrap();\n- s.recent_votes()\n- .push_back(&db, &Vote { a: \"a\".into(), b: \"b\".into(), ratio: 3 })\n- .unwrap();\n- }\n- let db = Db::open(dir.path()).unwrap();\n- let root = Store::root();\n- let s = root.scopes().key(&\"s\".to_string());\n- assert_eq!(root.nodes().key(&\"n\".to_string()).get(&db).unwrap(), Some(\"N\".to_string()));\n- assert_eq!(root.log().iter(&db).unwrap(), vec![42]);\n- assert_eq!(s.edges().key(&(1, 2)).get(&db).unwrap(), 9.0);\n- assert_eq!(\n- s.recent_votes().front(&db).unwrap(),\n- Some(Vote { a: \"a\".into(), b: \"b\".into(), ratio: 3 })\n- );\n-}\ndiff --git a/durable/tests/proptests.rs b/durable/tests/proptests.rs\ndeleted file mode 100644\nindex 9580a5fa406d29e886ac46c91f5cc723ee21c346..0000000000000000000000000000000000000000\n--- a/durable/tests/proptests.rs\n+++ /dev/null\n@@ -1,103 +0,0 @@\n-//! Property tests: durable collections must behave like their std analogues.\n-\n-use std::collections::{BTreeMap, VecDeque};\n-\n-use durable::{Db, Durability, Durable, Deque, Leaf, List, Map, Sum};\n-use proptest::prelude::*;\n-use tempfile::TempDir;\n-\n-#[derive(Durable)]\n-#[allow(dead_code)]\n-struct Bag {\n- map: Map>,\n- list: List>,\n- deque: Deque>,\n- total: Sum,\n-}\n-\n-fn open() -> (TempDir, Db) {\n- let dir = TempDir::new().unwrap();\n- let db = Db::open(dir.path()).unwrap();\n- (dir, db)\n-}\n-\n-proptest! {\n- #[test]\n- fn map_matches_btreemap(entries in proptest::collection::vec((\".*\", any::()), 0..40)) {\n- let (_dir, db) = open();\n- let map = Bag::root().map();\n- let mut model = BTreeMap::new();\n-\n- let mut batch = db.batch();\n- for (k, v) in &entries {\n- batch.write(map.key(k).set(v));\n- model.insert(k.clone(), *v);\n- }\n- batch.commit_with(Durability::WalOnly).unwrap();\n-\n- prop_assert_eq!(map.len(&db).unwrap(), model.len());\n- for (k, v) in &model {\n- prop_assert_eq!(map.get(&db, k).unwrap(), Some(*v));\n- }\n- let mut got = map.iter(&db).unwrap();\n- got.sort();\n- let mut want: Vec<(String, i64)> = model.into_iter().collect();\n- want.sort();\n- prop_assert_eq!(got, want);\n- }\n-\n- #[test]\n- fn list_roundtrips_in_order(values in proptest::collection::vec(any::(), 0..50)) {\n- let (_dir, db) = open();\n- let list = Bag::root().list();\n- let mut batch = db.batch();\n- for v in &values {\n- batch.push(&list, v).unwrap();\n- }\n- batch.commit_with(Durability::WalOnly).unwrap();\n-\n- prop_assert_eq!(list.len(&db).unwrap(), values.len() as u64);\n- prop_assert_eq!(list.iter(&db).unwrap(), values);\n- }\n-\n- #[test]\n- fn deque_matches_vecdeque(ops in proptest::collection::vec(any::<(bool, i64)>(), 0..60)) {\n- let (_dir, db) = open();\n- let dq = Bag::root().deque();\n- let mut model: VecDeque = VecDeque::new();\n-\n- for (front, v) in &ops {\n- if *front {\n- dq.push_front(&db, v).unwrap();\n- model.push_front(*v);\n- } else {\n- dq.push_back(&db, v).unwrap();\n- model.push_back(*v);\n- }\n- }\n- prop_assert_eq!(dq.len(&db).unwrap(), model.len() as u64);\n- prop_assert_eq!(dq.iter(&db).unwrap(), Vec::from(model.clone()));\n-\n- // Drain alternately from both ends.\n- let mut toggle = true;\n- while !model.is_empty() {\n- if toggle {\n- prop_assert_eq!(dq.pop_front(&db).unwrap(), model.pop_front());\n- } else {\n- prop_assert_eq!(dq.pop_back(&db).unwrap(), model.pop_back());\n- }\n- toggle = !toggle;\n- }\n- prop_assert!(dq.is_empty(&db).unwrap());\n- prop_assert_eq!(dq.pop_front(&db).unwrap(), None);\n- }\n-\n- #[test]\n- fn sum_equals_total_of_deltas(deltas in proptest::collection::vec(-1000i64..1000, 0..50)) {\n- let (_dir, db) = open();\n- let total = Bag::root().total();\n- let writes: Vec<_> = deltas.iter().map(|d| total.add(*d)).collect();\n- db.apply(&writes, Durability::WalOnly).unwrap();\n- prop_assert_eq!(total.get(&db).unwrap(), deltas.iter().sum::());\n- }\n-}\ndiff --git a/scripts/cursor-env-install.sh b/scripts/cursor-env-install.sh\nindex 9aeffdd1164016d5955efc81295f6c7f42b3d925..c4d547a4bdcf39d46cd8929f1b9589537c17dea4 100755\n--- a/scripts/cursor-env-install.sh\n+++ b/scripts/cursor-env-install.sh\n@@ -190,7 +190,6 @@ clojure -P -M\n clojure -M -e \"(com.microsoft.playwright.CLI/main (into-array String [\\\"install\\\" \\\"chromium\\\" \\\"--with-deps\\\"]))\"\n \n # Warm RocksDB + release server link (Clojure tests use release binary).\n-cargo build -p durable --quiet\n cargo build --release --package sorter2-server --quiet\n \n echo \"cursor-env-install: ok (bb=$(bb --version 2>/dev/null || echo missing), CXX=${CXX})\"\ndiff --git a/server/Cargo.toml b/server/Cargo.toml\nindex ad4912791aff59fb1d3293f66ad381ae618cd60b..dfa39beddecfa37dcdeaa602cb30f4b547528fbb 100644\n--- a/server/Cargo.toml\n+++ b/server/Cargo.toml\n@@ -25,7 +25,7 @@ futures-util = { version = \"0.3\", default-features = false, features = [\"std\"] }\n rand = \"0.8\"\n urlencoding = \"2\"\n url = \"2\"\n-durable = { path = \"../durable\" }\n+durable = { git = \"https://github.com/tommy-mor/durable.git\", branch = \"main\" }\n \n [dev-dependencies]\n reqwest = { version = \"0.12\", features = [\"json\"] }\ndiff --git a/server/src/api/ui_html.rs b/server/src/api/ui_html.rs\nindex e6e61815f6a6ddd626f79be81bf70256c047380a..b86581b1f337650564274254d840e8a75b49524d 100644\n--- a/server/src/api/ui_html.rs\n+++ b/server/src/api/ui_html.rs\n@@ -3,14 +3,20 @@ use axum::{\n response::{IntoResponse, Response},\n Form,\n };\n+use axum_extra::extract::cookie::CookieJar;\n use std::collections::HashMap;\n \n use crate::{\n+ auth::{\n+ alias_status_js, alias_redirect_js, config, login_redirect_js, oauth, redirect_js, resolve_vote_actor,\n+ session::{load_valid_session, session_has_pseudonym, session_id_from_jar},\n+ },\n fetch,\n html::{input_panel, js_string_literal, ranking_panel, JsBuilder},\n parser::parse_reddit_url,\n path_types::ItemId,\n state::{parse_item_param, AppState},\n+ storage_schema::pseudonym_owner,\n ui_action::{parse_html_ui_from_form, HtmlUiAction},\n };\n \n@@ -30,8 +36,21 @@ fn parent_from_scope(scope: &str) -> ItemId {\n parse_item_param(scope)\n }\n \n+fn vote_auth_redirect(state: &AppState, jar: &CookieJar) -> Option {\n+ let db = state.projection_store.db();\n+ let session = session_id_from_jar(jar)\n+ .as_deref()\n+ .and_then(|id| load_valid_session(db, id));\n+ match session {\n+ None => Some(login_redirect_js().into_response()),\n+ Some(s) if !session_has_pseudonym(&s) => Some(alias_redirect_js().into_response()),\n+ Some(_) => None,\n+ }\n+}\n+\n pub async fn post_ui_html(\n State(state): State,\n+ jar: CookieJar,\n Form(form): Form>,\n ) -> impl IntoResponse {\n let action = match parse_html_ui_from_form(&form) {\n@@ -48,9 +67,16 @@ pub async fn post_ui_html(\n scope,\n vote_compare,\n } => {\n+ if let Some(resp) = vote_auth_redirect(&state, &jar) {\n+ return resp;\n+ }\n let parent = parent_from_scope(&scope);\n+ let actor = resolve_vote_actor(\n+ state.projection_store.db(),\n+ session_id_from_jar(&jar).as_deref(),\n+ );\n if let Err(e) = state\n- .record_vote(&parent, &a, &b, ratio_left, ratio_right)\n+ .record_vote(&parent, &a, &b, ratio_left, ratio_right, &actor)\n .await\n {\n return ui_js_warn(&e).into_response();\n@@ -72,6 +98,58 @@ pub async fn post_ui_html(\n .morph_selector(\"#ranking-panel\", panel)\n .into_response()\n }\n+ HtmlUiAction::CheckPseudonym { pseudonym } => {\n+ let db = state.projection_store.db();\n+ let session_id = match session_id_from_jar(&jar) {\n+ Some(id) => id,\n+ None => return alias_status_js(\"sign in first\", false).into_response(),\n+ };\n+ let session = match load_valid_session(db, &session_id) {\n+ Some(s) => s,\n+ None => return alias_status_js(\"session expired\", false).into_response(),\n+ };\n+ match oauth::validate_pseudonym(&pseudonym) {\n+ Err(msg) => alias_status_js(msg, false).into_response(),\n+ Ok(name) => match pseudonym_owner(db, &name) {\n+ Ok(None) => alias_status_js(\"available\", true).into_response(),\n+ Ok(Some(owner)) if owner == session.uuid => {\n+ alias_status_js(\"already yours\", true).into_response()\n+ }\n+ Ok(Some(_)) => alias_status_js(\"taken\", false).into_response(),\n+ Err(e) => ui_js_warn(&e.to_string()).into_response(),\n+ },\n+ }\n+ }\n+ HtmlUiAction::ClaimPseudonym {\n+ pseudonym,\n+ return_to,\n+ } => {\n+ let db = state.projection_store.db();\n+ let session_id = match session_id_from_jar(&jar) {\n+ Some(id) => id,\n+ None => return login_redirect_js().into_response(),\n+ };\n+ let session = match load_valid_session(db, &session_id) {\n+ Some(s) => s,\n+ None => return login_redirect_js().into_response(),\n+ };\n+ let name = match oauth::validate_pseudonym(&pseudonym) {\n+ Ok(n) => n,\n+ Err(msg) => return alias_status_js(msg, false).into_response(),\n+ };\n+ if let Ok(Some(owner)) = pseudonym_owner(db, &name) {\n+ if owner != session.uuid {\n+ return alias_status_js(\"taken\", false).into_response();\n+ }\n+ } else if let Err(e) = state.claim_pseudonym(&session.uuid, &name).await {\n+ return ui_js_warn(&e).into_response();\n+ }\n+ if let Err(e) = crate::auth::session::update_session_pseudonym(&db, &session_id, &name)\n+ {\n+ return ui_js_warn(&e).into_response();\n+ }\n+ redirect_js(&config::sanitize_return_to(&return_to)).into_response()\n+ }\n HtmlUiAction::ParseQuery { query } => match parse_reddit_url(&query) {\n Ok(item) => {\n let _ = state.ensure_node(&item).await;\ndiff --git a/server/src/auth/config.rs b/server/src/auth/config.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..a1f042c655bf3e5234eeb87a7d889f64592807fb\n--- /dev/null\n+++ b/server/src/auth/config.rs\n@@ -0,0 +1,9 @@\n+pub const AUTH_RETURN_COOKIE: &str = \"sorter2_auth_return\";\n+\n+pub fn sanitize_return_to(raw: &str) -> String {\n+ let s = raw.trim();\n+ if s.is_empty() || !s.starts_with('/') || s.starts_with(\"//\") {\n+ return \"/\".to_string();\n+ }\n+ s.to_string()\n+}\ndiff --git a/server/src/auth/identity.rs b/server/src/auth/identity.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..20e51ac79ce1a5fd8c5ac2ab16b5799cb9a522d8\n--- /dev/null\n+++ b/server/src/auth/identity.rs\n@@ -0,0 +1,27 @@\n+//! Trust-weight calculation from linked OAuth providers.\n+\n+/// Base weight before any OAuth links.\n+pub const BASE_TRUST_WEIGHT: f64 = 1.0;\n+\n+/// Increment per linked provider (frozen at vote cast time).\n+pub const TRUST_WEIGHT_PER_LINK: f64 = 0.5;\n+\n+pub fn trust_weight_for_link_count(link_count: usize) -> f64 {\n+ BASE_TRUST_WEIGHT + TRUST_WEIGHT_PER_LINK * link_count as f64\n+}\n+\n+pub fn trust_weight_after_link(current: f64) -> f64 {\n+ current + TRUST_WEIGHT_PER_LINK\n+}\n+\n+#[cfg(test)]\n+mod tests {\n+ use super::*;\n+\n+ #[test]\n+ fn trust_weight_scales_with_links() {\n+ assert_eq!(trust_weight_for_link_count(0), 1.0);\n+ assert_eq!(trust_weight_for_link_count(1), 1.5);\n+ assert_eq!(trust_weight_for_link_count(2), 2.0);\n+ }\n+}\ndiff --git a/server/src/auth/mod.rs b/server/src/auth/mod.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..5ed535ba199fa736f0048c32623b14c3b1e5de2d\n--- /dev/null\n+++ b/server/src/auth/mod.rs\n@@ -0,0 +1,412 @@\n+//! GitHub OAuth login, session cookies, and vote actor resolution.\n+\n+pub mod config;\n+pub mod identity;\n+pub mod oauth;\n+pub mod session;\n+\n+use axum::{\n+ extract::{Query, State},\n+ http::StatusCode,\n+ response::{Html, IntoResponse, Redirect, Response},\n+ Form,\n+};\n+use axum_extra::extract::cookie::CookieJar;\n+use maud::{html, Markup};\n+use reqwest::Client;\n+use serde::Deserialize;\n+\n+use crate::{\n+ events::Event,\n+ fetch::now_ms,\n+ form_template::template_json_compact,\n+ html::layout,\n+ state::AppState,\n+ storage_schema::{oauth_link_owner, pseudonym_owner, Store, StoreFields},\n+ ui_action::UI_RPC_FIELD,\n+};\n+\n+pub use session::{resolve_vote_actor, session_id_from_jar, VoteActor};\n+\n+pub fn base_url_from_env(port: u16) -> String {\n+ std::env::var(\"SORTER2_BASE_URL\")\n+ .unwrap_or_else(|_| format!(\"http://127.0.0.1:{port}\"))\n+}\n+\n+fn new_actor_uuid() -> String {\n+ let mut bytes = [0u8; 16];\n+ rand::Rng::fill(&mut rand::thread_rng(), &mut bytes);\n+ format!(\n+ \"{:08x}-{:04x}-{:04x}-{:04x}-{:012x}\",\n+ u32::from_be_bytes([bytes[0], bytes[1], bytes[2], bytes[3]]),\n+ u16::from_be_bytes([bytes[4], bytes[5]]),\n+ u16::from_be_bytes([bytes[6], bytes[7]]) | 0x4000,\n+ u16::from_be_bytes([bytes[8], bytes[9]]) | 0x8000,\n+ u128::from_be_bytes([\n+ 0, 0, bytes[10], bytes[11], bytes[12], bytes[13], bytes[14], bytes[15], 0, 0, 0, 0,\n+ 0, 0, 0, 0,\n+ ]) & 0x0000_FFFF_FFFF_FFFF\n+ )\n+}\n+\n+#[derive(Debug, Deserialize)]\n+pub struct LoginQuery {\n+ #[serde(default)]\n+ pub return_to: Option,\n+}\n+\n+#[derive(Debug, Deserialize)]\n+pub struct GitHubStartQuery {\n+ #[serde(default)]\n+ pub return_to: Option,\n+ #[serde(default)]\n+ pub mock_user: Option,\n+}\n+\n+fn return_from_query_or_jar(jar: &CookieJar, query: Option<&str>) -> String {\n+ if let Some(raw) = query {\n+ return config::sanitize_return_to(raw);\n+ }\n+ jar.get(config::AUTH_RETURN_COOKIE)\n+ .map(|c| config::sanitize_return_to(c.value()))\n+ .unwrap_or_else(|| \"/\".to_string())\n+}\n+\n+fn oauth_providers(base_url: &str, return_to: &str) -> Vec<(&'static str, String)> {\n+ let mut out = Vec::new();\n+ if oauth::GitHubConfig::from_env(base_url).is_some() {\n+ out.push((\n+ \"GitHub\",\n+ format!(\n+ \"/auth/github?return_to={}\",\n+ urlencoding::encode(return_to)\n+ ),\n+ ));\n+ }\n+ out\n+}\n+\n+fn alias_list(db: &durable::Db, uuid: &str) -> Vec {\n+ Store::root()\n+ .user_pseudonyms()\n+ .key(&uuid.to_string())\n+ .iter(db)\n+ .unwrap_or_default()\n+}\n+\n+fn login_body(\n+ session: Option<&session::SessionActor>,\n+ aliases: &[String],\n+ providers: &[(&str, String)],\n+) -> Markup {\n+ html! {\n+ main class=\"panel login-page\" {\n+ div class=\"login-grid\" {\n+ section class=\"login-oauth\" {\n+ h1 { \"sign in\" }\n+ @if providers.is_empty() {\n+ p class=\"muted\" {\n+ \"OAuth is not configured. Set GITHUB_CLIENT_ID and GITHUB_CLIENT_SECRET.\"\n+ }\n+ } @else {\n+ ul class=\"oauth-provider-list\" {\n+ @for (name, href) in providers {\n+ li {\n+ a href=(href) class=\"button oauth-provider\" data-testid=(format!(\"oauth-{}\", name.to_lowercase())) {\n+ (format!(\"Continue with {name}\"))\n+ }\n+ }\n+ }\n+ }\n+ }\n+ @if let Some(actor) = session {\n+ p class=\"muted small\" {\n+ \"session active · weight \" (format!(\"{:.1}\", actor.trust_weight))\n+ }\n+ form method=\"post\" action=\"/auth/logout\" data-navigate=\"full\" {\n+ button type=\"submit\" { \"log out\" }\n+ }\n+ }\n+ }\n+ section class=\"login-aliases\" {\n+ h2 { \"your aliases\" }\n+ ul id=\"alias-list\" class=\"alias-list\" {\n+ @if aliases.is_empty() {\n+ li class=\"muted\" data-testid=\"alias-list-empty\" { \"none yet\" }\n+ } @else {\n+ @for alias in aliases {\n+ li { (alias) }\n+ }\n+ }\n+ }\n+ }\n+ }\n+ p { a href=\"/\" { \"← back\" } }\n+ }\n+ }\n+}\n+\n+pub async fn login_page(\n+ State(state): State,\n+ jar: CookieJar,\n+ Query(query): Query,\n+) -> Response {\n+ let return_to = return_from_query_or_jar(&jar, query.return_to.as_deref());\n+ let jar = jar.add(session::auth_return_cookie_value(&return_to));\n+\n+ let db = state.projection_store.db();\n+ let session = session::session_id_from_jar(&jar)\n+ .as_deref()\n+ .and_then(|id| session::load_session_actor(db, id));\n+ let aliases = session\n+ .as_ref()\n+ .map(|s| alias_list(db, &s.uuid))\n+ .unwrap_or_default();\n+ let providers = oauth_providers(&base_url_from_env(state.cfg.port), &return_to);\n+\n+ let markup = layout(\n+ \"login · sorter2\",\n+ login_body(session.as_ref(), &aliases, &providers),\n+ state.views.get_views(\"/login\"),\n+ );\n+ (jar, Html(markup.into_string())).into_response()\n+}\n+\n+pub async fn alias_page(\n+ State(state): State,\n+ jar: CookieJar,\n+ Query(query): Query,\n+) -> Result {\n+ let return_to = return_from_query_or_jar(&jar, query.return_to.as_deref());\n+ let session_id = session::session_id_from_jar(&jar).ok_or(StatusCode::UNAUTHORIZED)?;\n+ let db = state.projection_store.db();\n+ let session = session::load_valid_session(db, &session_id).ok_or(StatusCode::UNAUTHORIZED)?;\n+ if session::session_has_pseudonym(&session) {\n+ return Ok(Redirect::to(&return_to).into_response());\n+ }\n+\n+ let check_rpc = template_json_compact(&serde_json::json!({\n+ \"action\": \"check_pseudonym\",\n+ \"pseudonym\": {\"$form\": \"pseudonym\"},\n+ }))\n+ .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;\n+ let claim_rpc = template_json_compact(&serde_json::json!({\n+ \"action\": \"claim_pseudonym\",\n+ \"pseudonym\": {\"$form\": \"pseudonym\"},\n+ \"return_to\": return_to,\n+ }))\n+ .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;\n+\n+ let body = html! {\n+ main class=\"panel alias-page\" {\n+ h1 { \"choose alias\" }\n+ p class=\"muted\" { \"pick a unique display name for your votes\" }\n+ form id=\"alias-check-form\" method=\"POST\" action=\"/ui\" {\n+ input type=\"hidden\" name=(UI_RPC_FIELD) value=(check_rpc);\n+ label { \"alias\" }\n+ input type=\"text\" id=\"alias-input\" name=\"pseudonym\" autocomplete=\"off\"\n+ data-testid=\"alias-input\" maxlength=\"64\";\n+ p id=\"alias-status\" class=\"muted\" data-testid=\"alias-status\" { \"type to check availability\" }\n+ }\n+ form id=\"alias-claim-form\" method=\"POST\" action=\"/ui\" {\n+ input type=\"hidden\" name=(UI_RPC_FIELD) value=(claim_rpc);\n+ input type=\"hidden\" name=\"pseudonym\" id=\"alias-claim-field\" value=\"\";\n+ button type=\"submit\" class=\"btn-primary\" data-testid=\"alias-claim\" { \"continue\" }\n+ }\n+ p { a href=\"/login\" { \"← back to login\" } }\n+ }\n+ };\n+\n+ Ok(Html(\n+ layout(\n+ \"choose alias · sorter2\",\n+ body,\n+ state.views.get_views(\"/login/alias\"),\n+ )\n+ .into_string(),\n+ )\n+ .into_response())\n+}\n+\n+pub async fn github_start(\n+ State(state): State,\n+ jar: CookieJar,\n+ Query(query): Query,\n+) -> Result {\n+ let cfg = oauth::GitHubConfig::from_env(&base_url_from_env(state.cfg.port))\n+ .ok_or(StatusCode::SERVICE_UNAVAILABLE)?;\n+ let return_to = return_from_query_or_jar(&jar, query.return_to.as_deref());\n+ let state_token = session::new_oauth_state();\n+ let url = oauth::authorize_url(&cfg, &state_token, query.mock_user.as_deref());\n+ let jar = jar\n+ .add(session::oauth_state_cookie_value(&state_token))\n+ .add(session::auth_return_cookie_value(&return_to));\n+ Ok((jar, Redirect::temporary(&url)).into_response())\n+}\n+\n+#[derive(Debug, Deserialize)]\n+pub struct OAuthCallbackQuery {\n+ pub code: String,\n+ pub state: String,\n+}\n+\n+async fn finish_oauth_login(\n+ state: &AppState,\n+ jar: CookieJar,\n+ provider: &str,\n+ provider_id: String,\n+) -> Result<(CookieJar, String), StatusCode> {\n+ let db = state.projection_store.db();\n+ let return_to = return_from_query_or_jar(&jar, None);\n+\n+ let uuid = match oauth_link_owner(db, provider, &provider_id)\n+ .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?\n+ {\n+ Some(existing) => existing,\n+ None => {\n+ let uuid = new_actor_uuid();\n+ let ts = now_ms();\n+ state\n+ .append_identity_events(vec![\n+ Event::PrincipalCreated {\n+ uuid: uuid.clone(),\n+ ts,\n+ },\n+ Event::OauthLinked {\n+ uuid: uuid.clone(),\n+ provider: provider.to_string(),\n+ provider_id,\n+ ts,\n+ },\n+ ])\n+ .await\n+ .map_err(|e| {\n+ tracing::warn!(err = %e, \"identity event append failed\");\n+ StatusCode::INTERNAL_SERVER_ERROR\n+ })?;\n+ uuid\n+ }\n+ };\n+\n+ let aliases = alias_list(db, &uuid);\n+ let pseudonym = aliases.last().cloned().unwrap_or_default();\n+ let (session_id, _) = session::create_session(db, &uuid, &pseudonym).map_err(|e| {\n+ tracing::warn!(err = %e, \"session create failed\");\n+ StatusCode::INTERNAL_SERVER_ERROR\n+ })?;\n+\n+ let jar = jar\n+ .add(session::session_cookie_value(&session_id))\n+ .add(session::clear_oauth_state_cookie());\n+\n+ let dest = if pseudonym.is_empty() {\n+ format!(\n+ \"/login/alias?return_to={}\",\n+ urlencoding::encode(&return_to)\n+ )\n+ } else {\n+ return_to\n+ };\n+\n+ Ok((jar, dest))\n+}\n+\n+pub async fn github_callback(\n+ State(state): State,\n+ jar: CookieJar,\n+ Query(query): Query,\n+) -> Result {\n+ let cfg = oauth::GitHubConfig::from_env(&base_url_from_env(state.cfg.port))\n+ .ok_or(StatusCode::SERVICE_UNAVAILABLE)?;\n+\n+ let expected_state = session::oauth_state_from_jar(&jar).ok_or(StatusCode::BAD_REQUEST)?;\n+ if expected_state != query.state {\n+ return Err(StatusCode::BAD_REQUEST);\n+ }\n+\n+ let client = Client::builder()\n+ .timeout(std::time::Duration::from_secs(15))\n+ .build()\n+ .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;\n+\n+ let token = oauth::exchange_code(&client, &cfg, &query.code)\n+ .await\n+ .map_err(|e| {\n+ tracing::warn!(err = %e, \"github oauth token exchange failed\");\n+ StatusCode::BAD_GATEWAY\n+ })?;\n+ let user = oauth::fetch_user(&client, &cfg.api_base, &token)\n+ .await\n+ .map_err(|e| {\n+ tracing::warn!(err = %e, \"github user fetch failed\");\n+ StatusCode::BAD_GATEWAY\n+ })?;\n+\n+ let provider = \"github\";\n+ let provider_id = oauth::provider_id(&user);\n+ let (jar, dest) = finish_oauth_login(&state, jar, provider, provider_id).await?;\n+ Ok((jar, Redirect::to(&dest)).into_response())\n+}\n+\n+pub async fn logout(State(state): State, jar: CookieJar) -> impl IntoResponse {\n+ if let Some(session_id) = session::session_id_from_jar(&jar) {\n+ let _ = session::destroy_session(state.projection_store.db(), &session_id);\n+ }\n+ let jar = jar\n+ .add(session::clear_session_cookie())\n+ .add(session::clear_auth_return_cookie());\n+ (jar, Redirect::to(\"/login\"))\n+}\n+\n+#[derive(Deserialize)]\n+pub struct SwitchPseudonymForm {\n+ pseudonym: String,\n+}\n+\n+pub async fn switch_pseudonym(\n+ State(state): State,\n+ jar: CookieJar,\n+ Form(form): Form,\n+) -> Result {\n+ let session_id = session::session_id_from_jar(&jar).ok_or(StatusCode::UNAUTHORIZED)?;\n+ let db = state.projection_store.db();\n+ let actor = session::load_session_actor(db, &session_id).ok_or(StatusCode::UNAUTHORIZED)?;\n+ let owner = pseudonym_owner(db, &form.pseudonym)\n+ .map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?\n+ .ok_or(StatusCode::NOT_FOUND)?;\n+ if owner != actor.uuid {\n+ return Err(StatusCode::FORBIDDEN);\n+ }\n+ session::update_session_pseudonym(db, &session_id, &form.pseudonym).map_err(|_| {\n+ StatusCode::INTERNAL_SERVER_ERROR\n+ })?;\n+ Ok(Redirect::to(\"/login\").into_response())\n+}\n+\n+pub fn redirect_js(path: &str) -> crate::html::JsBuilder {\n+ crate::html::JsBuilder::new().raw(&format!(\n+ \"window.location.href={};\",\n+ crate::html::js_string_literal(path)\n+ ))\n+}\n+\n+pub fn login_redirect_js() -> crate::html::JsBuilder {\n+ crate::html::JsBuilder::new().raw(\n+ \"window.location.href='/login?return_to='+encodeURIComponent(window.location.pathname+window.location.search);\",\n+ )\n+}\n+\n+pub fn alias_redirect_js() -> crate::html::JsBuilder {\n+ crate::html::JsBuilder::new().raw(\n+ \"window.location.href='/login/alias?return_to='+encodeURIComponent(window.location.pathname+window.location.search);\",\n+ )\n+}\n+\n+pub fn alias_status_js(message: &str, ok: bool) -> crate::html::JsBuilder {\n+ let class = if ok { \"alias-ok\" } else { \"alias-bad\" };\n+ crate::html::JsBuilder::new().raw(&format!(\n+ \"var el=document.getElementById('alias-status'); if(el){{ el.textContent={}; el.className={}; }} var cf=document.getElementById('alias-claim-field'); if(cf) cf.value=document.getElementById('alias-input')?.value||'';\",\n+ crate::html::js_string_literal(message),\n+ crate::html::js_string_literal(class),\n+ ))\n+}\ndiff --git a/server/src/auth/oauth.rs b/server/src/auth/oauth.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..b80078fee0c5acd905b9125d29454e46c4e0d066\n--- /dev/null\n+++ b/server/src/auth/oauth.rs\n@@ -0,0 +1,146 @@\n+//! GitHub OAuth (raw reqwest, same style as reddit.rs).\n+\n+use reqwest::Client;\n+use serde::Deserialize;\n+\n+#[derive(Debug, Clone)]\n+pub struct GitHubConfig {\n+ pub client_id: String,\n+ pub client_secret: String,\n+ pub redirect_uri: String,\n+ pub oauth_base: String,\n+ pub api_base: String,\n+}\n+\n+pub fn github_oauth_base() -> String {\n+ std::env::var(\"GITHUB_OAUTH_BASE\").unwrap_or_else(|_| \"https://github.com\".into())\n+}\n+\n+pub fn github_api_base() -> String {\n+ std::env::var(\"GITHUB_API_BASE\").unwrap_or_else(|_| \"https://api.github.com\".into())\n+}\n+\n+impl GitHubConfig {\n+ pub fn from_env(base_url: &str) -> Option {\n+ let client_id = std::env::var(\"GITHUB_CLIENT_ID\").ok()?;\n+ let client_secret = std::env::var(\"GITHUB_CLIENT_SECRET\").ok()?;\n+ if client_id.is_empty() || client_secret.is_empty() {\n+ return None;\n+ }\n+ let oauth_base = github_oauth_base();\n+ let base = base_url.trim_end_matches('/');\n+ Some(Self {\n+ client_id,\n+ client_secret,\n+ redirect_uri: format!(\"{base}/auth/github/callback\"),\n+ oauth_base,\n+ api_base: github_api_base(),\n+ })\n+ }\n+}\n+\n+#[derive(Debug, Deserialize)]\n+struct TokenResponse {\n+ access_token: String,\n+}\n+\n+#[derive(Debug, Deserialize)]\n+pub struct GitHubUser {\n+ pub id: u64,\n+ pub login: String,\n+}\n+\n+pub fn authorize_url(cfg: &GitHubConfig, state: &str, mock_user: Option<&str>) -> String {\n+ let mut url = format!(\n+ \"{}/login/oauth/authorize?client_id={}&redirect_uri={}&scope=read:user&state={}\",\n+ cfg.oauth_base.trim_end_matches('/'),\n+ urlencoding::encode(&cfg.client_id),\n+ urlencoding::encode(&cfg.redirect_uri),\n+ urlencoding::encode(state),\n+ );\n+ if let Some(user) = mock_user {\n+ url.push_str(\"&mock_user=\");\n+ url.push_str(&urlencoding::encode(user));\n+ }\n+ url\n+}\n+\n+pub async fn exchange_code(\n+ client: &Client,\n+ cfg: &GitHubConfig,\n+ code: &str,\n+) -> Result {\n+ let resp = client\n+ .post(format!(\n+ \"{}/login/oauth/access_token\",\n+ cfg.oauth_base.trim_end_matches('/')\n+ ))\n+ .header(\"Accept\", \"application/json\")\n+ .form(&[\n+ (\"client_id\", cfg.client_id.as_str()),\n+ (\"client_secret\", cfg.client_secret.as_str()),\n+ (\"code\", code),\n+ (\"redirect_uri\", cfg.redirect_uri.as_str()),\n+ ])\n+ .send()\n+ .await\n+ .map_err(|e| format!(\"github token request failed: {e}\"))?;\n+\n+ if !resp.status().is_success() {\n+ return Err(format!(\"github token HTTP {}\", resp.status()));\n+ }\n+\n+ let body: TokenResponse = resp\n+ .json()\n+ .await\n+ .map_err(|e| format!(\"github token parse failed: {e}\"))?;\n+ Ok(body.access_token)\n+}\n+\n+pub async fn fetch_user(\n+ client: &Client,\n+ api_base: &str,\n+ access_token: &str,\n+) -> Result {\n+ let resp = client\n+ .get(format!(\"{}/user\", api_base.trim_end_matches('/')))\n+ .header(\"Accept\", \"application/vnd.github+json\")\n+ .header(\"User-Agent\", \"sorter2\")\n+ .bearer_auth(access_token)\n+ .send()\n+ .await\n+ .map_err(|e| format!(\"github user request failed: {e}\"))?;\n+\n+ if !resp.status().is_success() {\n+ return Err(format!(\"github user HTTP {}\", resp.status()));\n+ }\n+\n+ resp.json()\n+ .await\n+ .map_err(|e| format!(\"github user parse failed: {e}\"))\n+}\n+\n+pub fn provider_id(user: &GitHubUser) -> String {\n+ user.id.to_string()\n+}\n+\n+pub fn validate_pseudonym(raw: &str) -> Result {\n+ let trimmed = raw.trim();\n+ if trimmed.is_empty() {\n+ return Err(\"enter a name\");\n+ }\n+ if trimmed.len() > 64 {\n+ return Err(\"too long\");\n+ }\n+ if !trimmed\n+ .chars()\n+ .all(|c| c.is_ascii_alphanumeric() || c == '_' || c == '-')\n+ {\n+ return Err(\"letters, numbers, _ and - only\");\n+ }\n+ Ok(trimmed.to_string())\n+}\n+\n+pub fn sanitize_pseudonym(login: &str) -> String {\n+ validate_pseudonym(login).unwrap_or_else(|_| \"user\".to_string())\n+}\ndiff --git a/server/src/auth/session.rs b/server/src/auth/session.rs\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..09659240b9455c6fca12db5652e1d31cf8c2acfc\n--- /dev/null\n+++ b/server/src/auth/session.rs\n@@ -0,0 +1,240 @@\n+//! Session cookie resolution and durable session CRUD.\n+\n+use axum_extra::extract::cookie::{Cookie, CookieJar, SameSite};\n+use durable::{Db, Durability};\n+use rand::Rng;\n+\n+use crate::{\n+ auth::config::AUTH_RETURN_COOKIE,\n+ fetch::now_ms,\n+ identity::{DEFAULT_ACTOR_UUID, DEFAULT_PSEUDONYM},\n+ storage_dto::{SessionDataV1, SESSION_DATA_VERSION},\n+ storage_schema::{delete_session, load_session, user_trust_weight, write_session},\n+};\n+\n+pub const SESSION_COOKIE: &str = \"sorter2_session\";\n+pub const OAUTH_STATE_COOKIE: &str = \"sorter2_oauth_state\";\n+\n+/// Session lifetime (30 days).\n+pub const SESSION_TTL_MS: i64 = 30 * 24 * 60 * 60 * 1000;\n+\n+pub fn session_has_pseudonym(session: &SessionDataV1) -> bool {\n+ !session.current_pseudonym.trim().is_empty()\n+}\n+\n+pub fn load_valid_session(db: &Db, session_id: &str) -> Option {\n+ let session = load_session(db, session_id).ok()??;\n+ if session.expires_at <= now_ms() {\n+ return None;\n+ }\n+ Some(session)\n+}\n+\n+#[derive(Debug, Clone)]\n+pub struct VoteActor {\n+ pub pseudonym: String,\n+ pub trust_weight: f64,\n+}\n+\n+impl VoteActor {\n+ pub fn anon() -> Self {\n+ Self {\n+ pseudonym: DEFAULT_PSEUDONYM.to_string(),\n+ trust_weight: 1.0,\n+ }\n+ }\n+}\n+\n+#[derive(Debug, Clone)]\n+pub struct SessionActor {\n+ pub session_id: String,\n+ pub uuid: String,\n+ pub pseudonym: String,\n+ pub trust_weight: f64,\n+ pub expires_at: i64,\n+}\n+\n+pub fn new_session_id() -> String {\n+ let mut bytes = [0u8; 32];\n+ rand::thread_rng().fill(&mut bytes);\n+ hex_encode(&bytes)\n+}\n+\n+pub fn new_oauth_state() -> String {\n+ let mut bytes = [0u8; 16];\n+ rand::thread_rng().fill(&mut bytes);\n+ hex_encode(&bytes)\n+}\n+\n+fn hex_encode(bytes: &[u8]) -> String {\n+ bytes.iter().map(|b| format!(\"{b:02x}\")).collect()\n+}\n+\n+pub fn resolve_vote_actor(db: &Db, session_id: Option<&str>) -> VoteActor {\n+ let Some(session_id) = session_id else {\n+ return VoteActor::anon();\n+ };\n+ let Ok(Some(session)) = load_session(db, session_id) else {\n+ return VoteActor::anon();\n+ };\n+ if session.expires_at <= now_ms() {\n+ return VoteActor::anon();\n+ }\n+ let trust_weight = user_trust_weight(db, &session.uuid).unwrap_or(1.0);\n+ VoteActor {\n+ pseudonym: session.current_pseudonym,\n+ trust_weight,\n+ }\n+}\n+\n+pub fn load_session_actor(db: &Db, session_id: &str) -> Option {\n+ let session = load_session(db, session_id).ok()??;\n+ if session.expires_at <= now_ms() {\n+ return None;\n+ }\n+ let trust_weight = user_trust_weight(db, &session.uuid).ok()?;\n+ Some(SessionActor {\n+ session_id: session_id.to_string(),\n+ uuid: session.uuid,\n+ pseudonym: session.current_pseudonym,\n+ trust_weight,\n+ expires_at: session.expires_at,\n+ })\n+}\n+\n+pub fn create_session(\n+ db: &Db,\n+ uuid: &str,\n+ pseudonym: &str,\n+) -> Result<(String, SessionDataV1), String> {\n+ let session_id = new_session_id();\n+ let expires_at = now_ms() + SESSION_TTL_MS;\n+ let data = SessionDataV1 {\n+ version: SESSION_DATA_VERSION,\n+ uuid: uuid.to_string(),\n+ current_pseudonym: pseudonym.to_string(),\n+ expires_at,\n+ };\n+ let mut batch = db.batch();\n+ write_session(&mut batch, &session_id, &data);\n+ batch\n+ .commit_with(Durability::SyncWal)\n+ .map_err(|e| e.to_string())?;\n+ Ok((session_id, data))\n+}\n+\n+pub fn update_session_pseudonym(\n+ db: &Db,\n+ session_id: &str,\n+ pseudonym: &str,\n+) -> Result<(), String> {\n+ let mut session = load_session(db, session_id)\n+ .map_err(|e| e.to_string())?\n+ .ok_or_else(|| \"session not found\".to_string())?;\n+ if session.expires_at <= now_ms() {\n+ return Err(\"session expired\".to_string());\n+ }\n+ session.current_pseudonym = pseudonym.to_string();\n+ let mut batch = db.batch();\n+ write_session(&mut batch, session_id, &session);\n+ batch\n+ .commit_with(Durability::SyncWal)\n+ .map_err(|e| e.to_string())\n+}\n+\n+pub fn destroy_session(db: &Db, session_id: &str) -> Result<(), String> {\n+ let mut batch = db.batch();\n+ delete_session(&mut batch, session_id);\n+ batch\n+ .commit_with(Durability::SyncWal)\n+ .map_err(|e| e.to_string())\n+}\n+\n+pub fn session_cookie_value(session_id: &str) -> Cookie<'static> {\n+ Cookie::build((SESSION_COOKIE, session_id.to_string()))\n+ .http_only(true)\n+ .same_site(SameSite::Lax)\n+ .path(\"/\")\n+ .build()\n+}\n+\n+pub fn clear_session_cookie() -> Cookie<'static> {\n+ Cookie::build((SESSION_COOKIE, \"\"))\n+ .http_only(true)\n+ .same_site(SameSite::Lax)\n+ .path(\"/\")\n+ .removal()\n+ .build()\n+}\n+\n+pub fn oauth_state_cookie_value(state: &str) -> Cookie<'static> {\n+ Cookie::build((OAUTH_STATE_COOKIE, state.to_string()))\n+ .http_only(true)\n+ .same_site(SameSite::Lax)\n+ .path(\"/\")\n+ .build()\n+}\n+\n+pub fn clear_oauth_state_cookie() -> Cookie<'static> {\n+ Cookie::build((OAUTH_STATE_COOKIE, \"\"))\n+ .http_only(true)\n+ .same_site(SameSite::Lax)\n+ .path(\"/\")\n+ .removal()\n+ .build()\n+}\n+\n+pub fn auth_return_cookie_value(return_to: &str) -> Cookie<'static> {\n+ Cookie::build((AUTH_RETURN_COOKIE, return_to.to_string()))\n+ .http_only(true)\n+ .same_site(SameSite::Lax)\n+ .path(\"/\")\n+ .build()\n+}\n+\n+pub fn clear_auth_return_cookie() -> Cookie<'static> {\n+ Cookie::build((AUTH_RETURN_COOKIE, \"\"))\n+ .http_only(true)\n+ .same_site(SameSite::Lax)\n+ .path(\"/\")\n+ .removal()\n+ .build()\n+}\n+\n+pub fn auth_return_from_jar(jar: &CookieJar) -> Option {\n+ jar.get(AUTH_RETURN_COOKIE).map(|c| c.value().to_string())\n+}\n+\n+pub fn session_id_from_jar(jar: &CookieJar) -> Option {\n+ jar.get(SESSION_COOKIE).map(|c| c.value().to_string())\n+}\n+\n+pub fn oauth_state_from_jar(jar: &CookieJar) -> Option {\n+ jar.get(OAUTH_STATE_COOKIE).map(|c| c.value().to_string())\n+}\n+\n+pub fn actor_uuid_for_vote(db: &Db, session_id: Option<&str>) -> String {\n+ let Some(session_id) = session_id else {\n+ return DEFAULT_ACTOR_UUID.to_string();\n+ };\n+ load_session(db, session_id)\n+ .ok()\n+ .flatten()\n+ .filter(|s| s.expires_at > now_ms())\n+ .map(|s| s.uuid)\n+ .unwrap_or_else(|| DEFAULT_ACTOR_UUID.to_string())\n+}\n+\n+#[cfg(test)]\n+mod tests {\n+ use super::*;\n+\n+ #[test]\n+ fn missing_session_falls_back_to_anon() {\n+ let dir = tempfile::tempdir().unwrap();\n+ let db = Db::open(dir.path()).unwrap();\n+ let actor = resolve_vote_actor(&db, None);\n+ assert_eq!(actor.pseudonym, DEFAULT_PSEUDONYM);\n+ assert_eq!(actor.trust_weight, 1.0);\n+ }\n+}\ndiff --git a/server/src/events.rs b/server/src/events.rs\nindex 015208311f6c5c23a0e8aab068d002a69c89c4e1..8d4c1e260f9c034f5f7ec07484dda17efd69ecff 100644\n--- a/server/src/events.rs\n+++ b/server/src/events.rs\n@@ -31,6 +31,9 @@ pub fn event_timestamp(event: &Event) -> i64 {\n match event {\n Event::VoteRecorded { ts, .. } => *ts,\n Event::NodeEnsured { .. } => crate::fetch::now_ms(),\n+ Event::PrincipalCreated { ts, .. } => *ts,\n+ Event::OauthLinked { ts, .. } => *ts,\n+ Event::PseudonymClaimed { ts, .. } => *ts,\n }\n }\n \n@@ -57,6 +60,24 @@ pub enum Event {\n },\n /// Register a node path in the fractal tree (no external fetch).\n NodeEnsured { id: String },\n+\n+ /// New trust anchor (first identity event for a human).\n+ PrincipalCreated { uuid: String, ts: i64 },\n+\n+ /// OAuth provider account linked to an existing UUID.\n+ OauthLinked {\n+ uuid: String,\n+ provider: String,\n+ provider_id: String,\n+ ts: i64,\n+ },\n+\n+ /// Display pseudonym claimed by a UUID (global uniqueness enforced at apply).\n+ PseudonymClaimed {\n+ uuid: String,\n+ pseudonym: String,\n+ ts: i64,\n+ },\n }\n \n impl Event {\ndiff --git a/server/src/html/mod.rs b/server/src/html/mod.rs\nindex 1e2e7a06856d8a62378741aaf5ed94a4ffed337e..46d18b87f1313bf0aeb29955d18f291961057509 100644\n--- a/server/src/html/mod.rs\n+++ b/server/src/html/mod.rs\n@@ -126,7 +126,7 @@ pub fn now_ms() -> i64 {\n t.as_millis() as i64\n }\n \n-fn layout(title: &str, body: Markup, views: u64) -> Markup {\n+pub(crate) fn layout(title: &str, body: Markup, views: u64) -> Markup {\n let ver = asset_version();\n let css_href = format!(\"/static/sorter.css?v={ver}\");\n let js_src = format!(\"/static/sorter_ui.js?v={ver}\");\n@@ -144,6 +144,9 @@ fn layout(title: &str, body: Markup, views: u64) -> Markup {\n @if views > 0 {\n span class=\"view-meta muted\" { (views) \" views\" }\n }\n+ nav class=\"top-nav\" {\n+ a href=\"/login\" { \"login\" }\n+ }\n div id=\"errors\" {}\n (body)\n script src=(js_src) {}\ndiff --git a/server/src/lib.rs b/server/src/lib.rs\nindex da6e33e3dd6d79967bbee7a2708a7d982c93b88c..84f2565b105fb302241b949af64bd5e49916eab2 100644\n--- a/server/src/lib.rs\n+++ b/server/src/lib.rs\n@@ -1,4 +1,5 @@\n pub mod api;\n+pub mod auth;\n pub mod event_log;\n pub mod events;\n pub mod fetch;\n@@ -42,6 +43,12 @@ pub fn create_app(state: AppState) -> Router {\n .route(\"/~/*item_path\", get(crate::html::browse))\n .route(\"/\", get(crate::html::home))\n .route(\"/vote\", get(crate::html::vote::vote_page))\n+ .route(\"/login\", get(crate::auth::login_page))\n+ .route(\"/login/alias\", get(crate::auth::alias_page))\n+ .route(\"/auth/github\", get(crate::auth::github_start))\n+ .route(\"/auth/github/callback\", get(crate::auth::github_callback))\n+ .route(\"/auth/logout\", post(crate::auth::logout))\n+ .route(\"/auth/switch\", post(crate::auth::switch_pseudonym))\n .route(\"/ui\", post(crate::api::ui_html::post_ui_html))\n .with_state(state)\n .layer(TraceLayer::new_for_http())\ndiff --git a/server/src/projection_apply.rs b/server/src/projection_apply.rs\nindex f2859c7e0e718c5bd3d332224eb05d7a8c65a816..4fb4b48ee0a53b7f673fae0c9731eed5acc33332 100644\n--- a/server/src/projection_apply.rs\n+++ b/server/src/projection_apply.rs\n@@ -7,11 +7,13 @@\n use crate::{\n event_log::EventLogError,\n events::{Event, EventRecord},\n- identity::resolve_actor_uuid,\n+ auth::identity::{trust_weight_after_link, BASE_TRUST_WEIGHT},\n path_types::ItemId,\n projection_store::ProjectionStore,\n reducer::VoteData,\n- storage_schema::{ensure_path_writes, vote_writes},\n+ storage_schema::{\n+ ensure_path_writes, oauth_link_key, pseudonym_owner, vote_writes, Store, StoreFields,\n+ },\n };\n \n fn parse_event_id(id: &str) -> Result {\n@@ -72,7 +74,7 @@ pub fn apply_records(\n *trust_weight,\n )\n .ok_or_else(|| EventLogError::Apply(format!(\"invalid vote event: {a} vs {b}\")))?;\n- let actor_uuid = resolve_actor_uuid(db, pseudonym)\n+ let actor_uuid = crate::identity::resolve_actor_uuid(db, pseudonym)\n .map_err(|e| EventLogError::Apply(e))?;\n let parent = parent_from_event_scope(scope);\n vote_writes(&mut batch, &parent, &vote, &actor_uuid)\n@@ -82,6 +84,72 @@ pub fn apply_records(\n let parsed = parse_event_id(id)?;\n ensure_path_writes(&mut batch, &parsed);\n }\n+ Event::PrincipalCreated { uuid, .. } => {\n+ batch.write(\n+ Store::root()\n+ .user_weights()\n+ .key(&uuid.clone())\n+ .set(&BASE_TRUST_WEIGHT),\n+ );\n+ }\n+ Event::OauthLinked {\n+ uuid,\n+ provider,\n+ provider_id,\n+ ..\n+ } => {\n+ let link_key = oauth_link_key(provider, provider_id);\n+ if let Some(existing) = Store::root()\n+ .oauth_links()\n+ .key(&link_key)\n+ .get(db)\n+ .map_err(|e| EventLogError::Apply(e.to_string()))?\n+ {\n+ if existing != *uuid {\n+ return Err(EventLogError::Apply(format!(\n+ \"oauth link {link_key} already owned by {existing}\"\n+ )));\n+ }\n+ } else {\n+ batch.write(Store::root().oauth_links().key(&link_key).set(uuid));\n+ let current = Store::root()\n+ .user_weights()\n+ .key(&uuid.clone())\n+ .get(db)\n+ .map_err(|e| EventLogError::Apply(e.to_string()))?\n+ .unwrap_or(BASE_TRUST_WEIGHT);\n+ batch.write(\n+ Store::root()\n+ .user_weights()\n+ .key(&uuid.clone())\n+ .set(&trust_weight_after_link(current)),\n+ );\n+ }\n+ }\n+ Event::PseudonymClaimed { uuid, pseudonym, .. } => {\n+ if let Some(owner) = pseudonym_owner(db, pseudonym)\n+ .map_err(|e| EventLogError::Apply(e.to_string()))?\n+ {\n+ if owner != *uuid {\n+ return Err(EventLogError::Apply(format!(\n+ \"pseudonym {pseudonym} already claimed by {owner}\"\n+ )));\n+ }\n+ } else {\n+ batch.write(\n+ Store::root()\n+ .pseudonyms()\n+ .key(&pseudonym.clone())\n+ .set(uuid),\n+ );\n+ batch\n+ .push(\n+ &Store::root().user_pseudonyms().key(&uuid.clone()),\n+ &pseudonym.clone(),\n+ )\n+ .map_err(|e| EventLogError::Apply(e.to_string()))?;\n+ }\n+ }\n }\n last_seq = record.seq;\n }\n@@ -93,3 +161,105 @@ pub fn apply_records(\n \n Ok(())\n }\n+\n+#[cfg(test)]\n+mod tests {\n+ use super::*;\n+ use crate::{\n+ events::EventRecord,\n+ identity::resolve_actor_uuid,\n+ projection_store::ProjectionStore,\n+ storage_schema::{oauth_link_owner, user_trust_weight, StoreFields},\n+ };\n+\n+ fn record(seq: u64, event: Event) -> EventRecord {\n+ EventRecord::new(seq, crate::events::event_timestamp(&event), event)\n+ }\n+\n+ #[test]\n+ fn identity_events_project_pseudonym_and_oauth_link() {\n+ let dir = tempfile::tempdir().unwrap();\n+ let db = durable::Db::open(dir.path()).unwrap();\n+ let store = ProjectionStore::from_db(&db).unwrap();\n+ let uuid = \"aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee\";\n+ let ts = 1;\n+\n+ apply_records(\n+ &store,\n+ &[\n+ record(\n+ 1,\n+ Event::PrincipalCreated {\n+ uuid: uuid.into(),\n+ ts,\n+ },\n+ ),\n+ record(\n+ 2,\n+ Event::OauthLinked {\n+ uuid: uuid.into(),\n+ provider: \"github\".into(),\n+ provider_id: \"42\".into(),\n+ ts,\n+ },\n+ ),\n+ record(\n+ 3,\n+ Event::PseudonymClaimed {\n+ uuid: uuid.into(),\n+ pseudonym: \"octocat\".into(),\n+ ts,\n+ },\n+ ),\n+ ],\n+ )\n+ .unwrap();\n+\n+ assert_eq!(\n+ oauth_link_owner(store.db(), \"github\", \"42\").unwrap(),\n+ Some(uuid.to_string())\n+ );\n+ assert_eq!(resolve_actor_uuid(store.db(), \"octocat\").unwrap(), uuid);\n+ assert_eq!(user_trust_weight(store.db(), uuid).unwrap(), 1.5);\n+ let aliases = Store::root()\n+ .user_pseudonyms()\n+ .key(&uuid.to_string())\n+ .iter(store.db())\n+ .unwrap();\n+ assert_eq!(aliases, vec![\"octocat\".to_string()]);\n+ }\n+\n+ #[test]\n+ fn pseudonym_claim_rejects_second_owner() {\n+ let dir = tempfile::tempdir().unwrap();\n+ let db = durable::Db::open(dir.path()).unwrap();\n+ let store = ProjectionStore::from_db(&db).unwrap();\n+\n+ apply_records(\n+ &store,\n+ &[record(\n+ 1,\n+ Event::PseudonymClaimed {\n+ uuid: \"uuid-a\".into(),\n+ pseudonym: \"taken\".into(),\n+ ts: 1,\n+ },\n+ )],\n+ )\n+ .unwrap();\n+\n+ let err = apply_records(\n+ &store,\n+ &[record(\n+ 2,\n+ Event::PseudonymClaimed {\n+ uuid: \"uuid-b\".into(),\n+ pseudonym: \"taken\".into(),\n+ ts: 2,\n+ },\n+ )],\n+ )\n+ .unwrap_err();\n+ assert!(err.to_string().contains(\"already claimed\"));\n+ }\n+}\ndiff --git a/server/src/projection_store.rs b/server/src/projection_store.rs\nindex 8c1a466183a96173fe52b144fb67c2454b715fbd..37b4d1e01ca0aec0b079a936aeeae5306c7d80ca 100644\n--- a/server/src/projection_store.rs\n+++ b/server/src/projection_store.rs\n@@ -18,7 +18,7 @@ use crate::{\n \n const PROJECTION_CURSOR_KEY: &str = \"cursor\";\n const PROJECTION_SCHEMA_KEY: &str = \"schema_version\";\n-const PROJECTION_SCHEMA_VERSION: u64 = 5;\n+const PROJECTION_SCHEMA_VERSION: u64 = 6;\n \n #[derive(Debug, thiserror::Error)]\n pub enum ProjectionStoreError {\ndiff --git a/server/src/state.rs b/server/src/state.rs\nindex dcb82ff4beeaac8f0820de0e0131ca1f6bd81dcc..86949c4d1a40719e3b5bb114387c1fdb8f57d7ca 100644\n--- a/server/src/state.rs\n+++ b/server/src/state.rs\n@@ -227,6 +227,7 @@ impl AppState {\n b: &str,\n ratio_left: i32,\n ratio_right: i32,\n+ actor: &crate::auth::VoteActor,\n ) -> Result<(), String> {\n let ts = crate::html::now_ms();\n let a_raw = a.trim();\n@@ -251,17 +252,49 @@ impl AppState {\n return Err(\"invalid vote: need two distinct items\".to_string());\n }\n \n- let event = Event::vote_recorded(\n+ let event = Event::VoteRecorded {\n ts,\n- a_id.as_str(),\n- b_id.as_str(),\n- left,\n- right,\n- parent.as_str(),\n- );\n+ a: a_id.as_str().to_string(),\n+ b: b_id.as_str().to_string(),\n+ ratio_left: left,\n+ ratio_right: right,\n+ scope: parent.as_str().to_string(),\n+ pseudonym: actor.pseudonym.clone(),\n+ trust_weight: actor.trust_weight,\n+ };\n \n self.journal.append(event).await\n }\n+\n+ /// Append identity events (OAuth link, pseudonym claim, etc.).\n+ pub async fn append_identity_events(&self, events: Vec) -> Result<(), String> {\n+ self.journal.append_many(events).await\n+ }\n+\n+ pub async fn claim_pseudonym(\n+ &self,\n+ uuid: &str,\n+ pseudonym: &str,\n+ ) -> Result<(), String> {\n+ let ts = crate::html::now_ms();\n+ self.journal\n+ .append(Event::PseudonymClaimed {\n+ uuid: uuid.to_string(),\n+ pseudonym: pseudonym.to_string(),\n+ ts,\n+ })\n+ .await\n+ }\n+\n+ /// Session id for the seeded default pseudonym (tests and local dev helpers).\n+ pub fn create_default_session(&self) -> Result {\n+ crate::auth::session::create_session(\n+ self.projection_store.db(),\n+ crate::identity::DEFAULT_ACTOR_UUID,\n+ crate::identity::DEFAULT_PSEUDONYM,\n+ )\n+ .map(|(id, _)| id)\n+ }\n }\n \n #[cfg(test)]\n@@ -497,7 +530,7 @@ mod tests {\n .await;\n \n let err = state\n- .record_vote(&ItemId::root(), \"alpha\", \"beta\", 0, 0)\n+ .record_vote(&ItemId::root(), \"alpha\", \"beta\", 0, 0, &crate::auth::VoteActor::anon())\n .await\n .unwrap_err();\n assert!(err.contains(\"positive preference\"));\n@@ -517,7 +550,14 @@ mod tests {\n .await;\n \n state\n- .record_vote(&ItemId::root(), \"alpha\", \"beta\", 2, 1)\n+ .record_vote(\n+ &ItemId::root(),\n+ \"alpha\",\n+ \"beta\",\n+ 2,\n+ 1,\n+ &crate::auth::VoteActor::anon(),\n+ )\n .await\n .unwrap();\n \n@@ -611,7 +651,14 @@ mod tests {\n };\n let second = AppState::new(cfg).await;\n second\n- .record_vote(&ItemId::root(), \"alpha\", \"gamma\", 3, 1)\n+ .record_vote(\n+ &ItemId::root(),\n+ \"alpha\",\n+ \"gamma\",\n+ 3,\n+ 1,\n+ &crate::auth::VoteActor::anon(),\n+ )\n .await\n .unwrap();\n \ndiff --git a/server/src/storage_dto.rs b/server/src/storage_dto.rs\nindex 22f5ac498ae3a5db4347f3fe94830539c7c01e3e..db8a3094454df66f6d0e0e6c9369017c014784be 100644\n--- a/server/src/storage_dto.rs\n+++ b/server/src/storage_dto.rs\n@@ -13,6 +13,16 @@ use crate::{\n \n pub const VOTE_RECORD_VERSION: u32 = 2;\n pub const ENTITY_DATA_VERSION: u32 = 1;\n+pub const SESSION_DATA_VERSION: u32 = 1;\n+\n+/// Browser session stored in durable (operational; not event-logged).\n+#[derive(Debug, Clone, Serialize, Deserialize)]\n+pub struct SessionDataV1 {\n+ pub version: u32,\n+ pub uuid: String,\n+ pub current_pseudonym: String,\n+ pub expires_at: i64,\n+}\n \n #[derive(Debug, Clone, Serialize, Deserialize)]\n pub struct Versioned {\ndiff --git a/server/src/storage_schema.rs b/server/src/storage_schema.rs\nindex fe2671b876f3df77fdfd402dbc845ff1eb3512cd..b942810afd96d3bf01b7765718303a39be4ef8d2 100644\n--- a/server/src/storage_schema.rs\n+++ b/server/src/storage_schema.rs\n@@ -14,7 +14,7 @@ use crate::{\n reducer::{EntityData, NodeState, ScopeVotes, VoteData, UuidVoteKey, uuid_vote_key},\n storage_dto::{\n decode_entity_data, decode_vote, encode_entity_data, encode_vote, parse_stored_id,\n- StoredEntityDataV1, StoredVoteV1,\n+ SessionDataV1, StoredEntityDataV1, StoredVoteV1,\n },\n };\n \n@@ -34,12 +34,75 @@ pub struct NodeSchema {\n #[allow(dead_code)]\n pub struct Store {\n pub nodes: Map,\n+ pub sessions: Map>,\n+ pub oauth_links: Map>,\n pub pseudonyms: Map>,\n+ pub user_pseudonyms: Map>>,\n+ pub user_weights: Map>,\n pub proj_meta: Map>,\n pub view_counts: Map>,\n pub view_meta: Map>,\n }\n \n+pub fn encode_session(data: &SessionDataV1) -> SessionDataV1 {\n+ data.clone()\n+}\n+\n+pub fn decode_session(data: SessionDataV1) -> SessionDataV1 {\n+ data\n+}\n+\n+pub fn oauth_link_key(provider: &str, provider_id: &str) -> String {\n+ format!(\"{provider}:{provider_id}\")\n+}\n+\n+pub fn user_trust_weight(db: &Db, uuid: &str) -> durable::Result {\n+ Ok(Store::root()\n+ .user_weights()\n+ .key(&uuid.to_string())\n+ .get(db)?\n+ .unwrap_or(1.0))\n+}\n+\n+pub fn load_session(db: &Db, session_id: &str) -> durable::Result> {\n+ Store::root()\n+ .sessions()\n+ .key(&session_id.to_string())\n+ .get(db)\n+}\n+\n+pub fn write_session(batch: &mut Batch, session_id: &str, data: &SessionDataV1) {\n+ batch.write(\n+ Store::root()\n+ .sessions()\n+ .key(&session_id.to_string())\n+ .set(&encode_session(data)),\n+ );\n+}\n+\n+pub fn delete_session(batch: &mut Batch, session_id: &str) {\n+ batch.write(\n+ Store::root()\n+ .sessions()\n+ .key(&session_id.to_string())\n+ .delete(),\n+ );\n+}\n+\n+pub fn pseudonym_owner(db: &Db, pseudonym: &str) -> durable::Result> {\n+ Store::root()\n+ .pseudonyms()\n+ .key(&pseudonym.to_string())\n+ .get(db)\n+}\n+\n+pub fn oauth_link_owner(db: &Db, provider: &str, provider_id: &str) -> durable::Result> {\n+ Store::root()\n+ .oauth_links()\n+ .key(&oauth_link_key(provider, provider_id))\n+ .get(db)\n+}\n+\n pub const RECENT_VOTES_CAP: u64 = 200;\n \n fn id_key(id: &ItemId) -> String {\ndiff --git a/server/src/ui_action.rs b/server/src/ui_action.rs\nindex 227ab3f6e8e9cce1532454450773f65daa24bfe7..a90bf80b191873700e030ee30591576ec6b08cf9 100644\n--- a/server/src/ui_action.rs\n+++ b/server/src/ui_action.rs\n@@ -44,6 +44,14 @@ pub enum HtmlUiAction {\n #[serde(default)]\n kind: FetchTarget,\n },\n+ /// Live alias availability check (alias chooser page).\n+ CheckPseudonym { pseudonym: String },\n+ /// Claim first alias after OAuth, then redirect.\n+ ClaimPseudonym {\n+ pseudonym: String,\n+ #[serde(default)]\n+ return_to: String,\n+ },\n }\n \n #[derive(Debug, Error)]\ndiff --git a/server/static/sorter_ui.js b/server/static/sorter_ui.js\nindex 5d6d765f69f301722e59248c76e919043559fb78..67a6ad7b5a9770290066d648d17feea3716bde54 100644\n--- a/server/static/sorter_ui.js\n+++ b/server/static/sorter_ui.js\n@@ -162,8 +162,36 @@\n update();\n }\n \n+ function initAliasInput() {\n+ var input = document.getElementById('alias-input');\n+ var form = document.getElementById('alias-check-form');\n+ var claimField = document.getElementById('alias-claim-field');\n+ if (!input || !form) return;\n+ var timer;\n+ function syncClaimField() {\n+ if (claimField) claimField.value = input.value || '';\n+ }\n+ function queueCheck() {\n+ syncClaimField();\n+ clearTimeout(timer);\n+ timer = setTimeout(function () {\n+ postUiForm(form);\n+ }, 250);\n+ }\n+ input.addEventListener('input', queueCheck);\n+ syncClaimField();\n+ }\n+\n+ document.addEventListener('input', function (e) {\n+ if (e.target && e.target.id === 'alias-input') {\n+ var claimField = document.getElementById('alias-claim-field');\n+ if (claimField) claimField.value = e.target.value || '';\n+ }\n+ });\n+\n function initSorterUi() {\n initVoteSlider();\n+ initAliasInput();\n document.addEventListener('submit', async function (e) {\n var f = e.target;\n if (!f || f.tagName !== 'FORM') return;\ndiff --git a/server/tests/integration_ui.rs b/server/tests/integration_ui.rs\nindex 40e4cfe1eec36dad29a075fb01aefb4dffd856d0..11ff32489f5feff280657948b1c5a10de5a6d91e 100644\n--- a/server/tests/integration_ui.rs\n+++ b/server/tests/integration_ui.rs\n@@ -3,12 +3,13 @@ use std::net::SocketAddr;\n \n use axum::Router;\n use sorter2_server::{\n+ auth::session::SESSION_COOKIE,\n create_app, create_app_state, path_types::ItemId, state::AppConfig, ui_action::UI_RPC_FIELD,\n };\n use tempfile::TempDir;\n use tokio::net::TcpListener;\n \n-async fn start_test_server() -> (SocketAddr, TempDir) {\n+async fn start_test_server() -> (SocketAddr, TempDir, String) {\n let tmp = TempDir::new().unwrap();\n let data = tmp.path().to_string_lossy().into_owned();\n let cfg = AppConfig {\n@@ -18,6 +19,8 @@ async fn start_test_server() -> (SocketAddr, TempDir) {\n port: 0,\n };\n let state = create_app_state(cfg).await;\n+ let session_id = state.create_default_session().unwrap();\n+ let session_cookie = format!(\"{SESSION_COOKIE}={session_id}\");\n let app: Router = create_app(state);\n \n let listener = TcpListener::bind(\"127.0.0.1:0\").await.unwrap();\n@@ -25,12 +28,12 @@ async fn start_test_server() -> (SocketAddr, TempDir) {\n tokio::spawn(async move {\n axum::serve(listener, app).await.unwrap();\n });\n- (addr, tmp)\n+ (addr, tmp, session_cookie)\n }\n \n #[tokio::test]\n async fn post_ui_vote_compare_morphs_edge_history() {\n- let (addr, _tmp) = start_test_server().await;\n+ let (addr, _tmp, session_cookie) = start_test_server().await;\n let parent = \"reddit.com/r/rust\";\n let a = \"reddit.com/r/rust/comments/aaa/announcing_rust_199\";\n let b = \"reddit.com/r/rust/comments/bbb/what_are_you_working_on\";\n@@ -53,6 +56,7 @@ async fn post_ui_vote_compare_morphs_edge_history() {\n let client = reqwest::Client::new();\n let body = client\n .post(format!(\"http://{addr}/ui\"))\n+ .header(\"Cookie\", &session_cookie)\n .form(&form)\n .send()\n .await\n@@ -85,7 +89,7 @@ async fn post_ui_vote_compare_morphs_edge_history() {\n \n #[tokio::test]\n async fn post_ui_record_vote_morphs_ranking_and_persists() {\n- let (addr, tmp) = start_test_server().await;\n+ let (addr, tmp, session_cookie) = start_test_server().await;\n let rpc = serde_json::json!({\n \"action\": \"record_vote\",\n \"a\": \"alpha\",\n@@ -100,6 +104,7 @@ async fn post_ui_record_vote_morphs_ranking_and_persists() {\n let client = reqwest::Client::new();\n let body = client\n .post(format!(\"http://{addr}/ui\"))\n+ .header(\"Cookie\", &session_cookie)\n .form(&form)\n .send()\n .await\n@@ -139,7 +144,7 @@ async fn post_ui_record_vote_morphs_ranking_and_persists() {\n \n #[tokio::test]\n async fn browse_url_renders_subreddit_page() {\n- let (addr, _tmp) = start_test_server().await;\n+ let (addr, _tmp, _session_cookie) = start_test_server().await;\n let client = reqwest::Client::new();\n let html = client\n .get(format!(\"http://{addr}/~/https://reddit.com/r/rust\"))\n@@ -155,7 +160,7 @@ async fn browse_url_renders_subreddit_page() {\n \n #[tokio::test]\n async fn vote_page_renders_live_ranking_sidebar() {\n- let (addr, _tmp) = start_test_server().await;\n+ let (addr, _tmp, session_cookie) = start_test_server().await;\n let client = reqwest::Client::new();\n let seed_rpc = serde_json::json!({\n \"action\": \"record_vote\",\n@@ -169,6 +174,7 @@ async fn vote_page_renders_live_ranking_sidebar() {\n form.insert(UI_RPC_FIELD.to_string(), seed_rpc);\n client\n .post(format!(\"http://{addr}/ui\"))\n+ .header(\"Cookie\", &session_cookie)\n .form(&form)\n .send()\n .await\n@@ -199,7 +205,7 @@ async fn vote_page_renders_live_ranking_sidebar() {\n \n #[tokio::test]\n async fn post_ui_parse_query_redirects_to_subreddit() {\n- let (addr, _tmp) = start_test_server().await;\n+ let (addr, _tmp, _session_cookie) = start_test_server().await;\n let rpc = serde_json::json!({\n \"action\": \"parse_query\",\n \"query\": \"r/rust\"\ndiff --git a/test/auth_login.clj b/test/auth_login.clj\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..a81e064ebb88fe359c56cf6718880dc1e4c1a7d9\n--- /dev/null\n+++ b/test/auth_login.clj\n@@ -0,0 +1,76 @@\n+(ns test.auth-login\n+ (:require [clojure.string :as str]\n+ [clojure.test :refer [deftest is testing]]\n+ [com.blockether.spel.core :as core]\n+ [com.blockether.spel.locator :as loc]\n+ [com.blockether.spel.page :as page]\n+ [test.support.harness :as harness]\n+ [test.support.seed-auth :as seed-auth]))\n+\n+(defn- move-vote-slider-left [pg]\n+ (page/evaluate pg\n+ \"(() => { const s = document.getElementById('vote-preference-slider'); if (!s) return; s.value = '20'; s.dispatchEvent(new Event('input', { bubbles: true })); })()\"))\n+\n+(defn- type-alias! [pg text]\n+ (page/evaluate pg\n+ (.replace\n+ \"(() => { const i = document.getElementById('alias-input'); const f = document.getElementById('alias-check-form'); if (!i || !f) return;\n+ i.value = __TEXT__;\n+ const cf = document.getElementById('alias-claim-field'); if (cf) cf.value = i.value;\n+ return fetch(f.action, { method: 'POST', credentials: 'same-origin',\n+ headers: { 'Content-Type': 'application/x-www-form-urlencoded' },\n+ body: new URLSearchParams(new FormData(f)).toString() })\n+ .then(function (r) { return r.text(); })\n+ .then(function (t) { eval(t); }); })()\"\n+ \"__TEXT__\"\n+ (pr-str text)))\n+ (Thread/sleep 400))\n+\n+(defn- element-text [pg test-id]\n+ (let [raw (page/evaluate pg\n+ (str \"document.querySelector('[data-testid=\\\"\" test-id \"\\\"]')?.textContent || ''\"))]\n+ (when (string? raw) (str/trim raw))))\n+\n+(defn- wait-for-text [pg test-id text timeout-ms]\n+ (let [deadline (+ (System/currentTimeMillis) timeout-ms)]\n+ (loop []\n+ (let [got (or (element-text pg test-id) \"\")]\n+ (cond\n+ (= got text) got\n+ (< (System/currentTimeMillis) deadline) (do (Thread/sleep 200) (recur))\n+ :else (throw (ex-info \"timeout waiting for text\" {:test-id test-id :want text :got got})))))))\n+\n+(deftest new-user-login-flow-returns-to-vote-pair\n+ (testing \"anonymous vote redirects through OAuth + alias chooser back to the same pair\"\n+ (let [servers (harness/with-auth-servers\n+ (fn [data-dir]\n+ (seed-auth/write-seeder-events! (str data-dir \"/events.jsonl\"))))]\n+ (try\n+ (harness/seed-rust-children! (:app-base servers))\n+ (let [vote-url (seed-auth/seeder-pair-vote-url (:app-base servers))\n+ alias \"newbie-alias\"]\n+ (core/with-testing-page [pg]\n+ (page/navigate pg vote-url)\n+ (page/wait-for-selector pg \"#vote-compare-form\")\n+ (move-vote-slider-left pg)\n+ (loc/click (page/get-by-test-id pg \"vote-post\"))\n+ (page/wait-for-selector pg \"[data-testid=oauth-github]\" {:timeout 15000})\n+ (is (str/includes? (or (element-text pg \"alias-list-empty\") \"\") \"none yet\"))\n+ (loc/click (page/get-by-test-id pg \"oauth-github\"))\n+ (page/wait-for-selector pg \"[data-testid=alias-input]\" {:timeout 15000})\n+ (type-alias! pg \"seeder\")\n+ (wait-for-text pg \"alias-status\" \"taken\" 15000)\n+ (type-alias! pg alias)\n+ (wait-for-text pg \"alias-status\" \"available\" 15000)\n+ (loc/click (page/get-by-test-id pg \"alias-claim\"))\n+ (page/wait-for-selector pg \"#vote-compare-form\" {:timeout 15000})\n+ (is (str/includes? (page/url pg) \"/vote?\"))\n+ (loc/click (page/get-by-test-id pg \"vote-post\"))\n+ (page/wait-for-selector pg \".vote-edge-history-title\" {:timeout 15000})\n+ (let [history (or (element-text pg \"#vote-edge-history-region\") \"\")]\n+ (is (str/includes? history \"votes on this pair\"))\n+ (is (str/includes? history \"3:1\")\n+ \"seeded seeder vote still visible\")\n+ (is (not (str/includes? history \"no votes on this pair yet\"))))))\n+ (finally\n+ ((:stop servers)))))))\ndiff --git a/test/support/harness.clj b/test/support/harness.clj\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..3f05951f418258642dcacb4a10ccccc8bfbe8748\n--- /dev/null\n+++ b/test/support/harness.clj\n@@ -0,0 +1,104 @@\n+(ns test.support.harness\n+ (:require [babashka.process :as process]\n+ [clojure.java.io :as io]\n+ [clojure.string :as str]\n+ [test.support.mock-oauth :as mock-oauth]\n+ [test.support.mock-reddit :as mock-reddit]))\n+\n+(defn repo-root []\n+ (.getCanonicalPath (io/file (System/getProperty \"user.dir\"))))\n+\n+(defn pick-port []\n+ (with-open [s (java.net.ServerSocket. 0)]\n+ (.getLocalPort s)))\n+\n+(defn wait-health [base-url ms]\n+ (let [deadline (+ (System/currentTimeMillis) ms)\n+ url (str base-url \"/healthz\")]\n+ (loop []\n+ (let [resp (try\n+ (process/shell {:out :string :err :string}\n+ \"curl\" \"-sf\" url)\n+ (catch Exception _ nil))]\n+ (if (and resp (zero? (:exit resp)) (= \"ok\" (str/trim (:out resp \"\"))))\n+ true\n+ (if (< (System/currentTimeMillis) deadline)\n+ (do (Thread/sleep 200) (recur))\n+ false))))))\n+\n+(defn curl-fetch-children [base item]\n+ (process/shell {:out :string :err :string}\n+ \"curl\" \"-sfN\" \"--max-time\" \"20\"\n+ \"-X\" \"POST\" (str base \"/ui\")\n+ \"--data-urlencode\"\n+ (str \"__rpc__={\\\"action\\\":\\\"fetch_entity\\\",\\\"item\\\":\\\"\" item\n+ \"\\\",\\\"kind\\\":\\\"children\\\"}\")))\n+\n+(defn app-env\n+ [data-dir app-port oauth-port reddit-port]\n+ (into (into {} (System/getenv))\n+ {\"SORTER2_SKIP_DOTENV\" \"1\"\n+ \"SORTER2_DATA_DIR\" data-dir\n+ \"SORTER2_EVENT_LOG\" (str data-dir \"/events.jsonl\")\n+ \"SORTER2_VIEWS_LOG\" (str data-dir \"/views.jsonl\")\n+ \"PORT\" (str app-port)\n+ \"SORTER2_BASE_URL\" (str \"http://127.0.0.1:\" app-port)\n+ \"GITHUB_CLIENT_ID\" \"test-client\"\n+ \"GITHUB_CLIENT_SECRET\" \"test-secret\"\n+ \"GITHUB_OAUTH_BASE\" (str \"http://127.0.0.1:\" oauth-port)\n+ \"GITHUB_API_BASE\" (str \"http://127.0.0.1:\" oauth-port)\n+ \"REDDIT_API_BASE\" (str \"http://127.0.0.1:\" reddit-port)\n+ \"REDDIT_OAUTH_BASE\" (str \"http://127.0.0.1:\" reddit-port)\n+ \"REDDIT_CLIENT_ID\" \"\"\n+ \"REDDIT_CLIENT_SECRET\" \"\"\n+ \"REDDIT_APP_ID\" \"\"\n+ \"REDDIT_APP_SECRET\" \"\"}))\n+\n+(defn with-auth-servers\n+ \"Start mock Reddit + mock OAuth + release sorter2-server.\n+ `seed-fn` is `(fn [data-dir] ...)` called before the app boots.\n+ Returns `{:stop ... :app-base ...}`.\"\n+ [seed-fn]\n+ (let [root (repo-root)\n+ fixtures (mock-reddit/fixtures-dir root)\n+ data-dir (.getAbsolutePath\n+ (doto (io/file (System/getProperty \"java.io.tmpdir\")\n+ (str \"sorter2-auth-\" (System/currentTimeMillis)))\n+ (.mkdirs)))\n+ reddit-port (pick-port)\n+ oauth-port (pick-port)\n+ app-port (pick-port)\n+ app-base (str \"http://127.0.0.1:\" app-port)\n+ bin (str root \"/target/release/sorter2-server\")\n+ stop-mock-reddit (mock-reddit/start-mock-reddit reddit-port fixtures)\n+ stop-mock-oauth (mock-oauth/start-mock-oauth oauth-port)]\n+ (seed-fn data-dir)\n+ (process/shell {:dir root}\n+ \"cargo\" \"build\" \"--release\" \"--package\" \"sorter2-server\")\n+ (let [proc (process/process {:dir root\n+ :env (app-env data-dir app-port oauth-port reddit-port)\n+ :out :string\n+ :err :string}\n+ bin)]\n+ (when-not (wait-health app-base 25000)\n+ (process/destroy proc)\n+ (stop-mock-oauth)\n+ (stop-mock-reddit)\n+ (throw (ex-info \"app healthz timeout\" {:app-base app-base})))\n+ {:stop (fn []\n+ (process/destroy proc)\n+ (stop-mock-oauth)\n+ (stop-mock-reddit))\n+ :app-base app-base\n+ :data-dir data-dir})))\n+\n+(defn oauth-login-url\n+ [app-base return-to mock-user]\n+ (str app-base \"/auth/github?return_to=\"\n+ (java.net.URLEncoder/encode return-to \"UTF-8\")\n+ \"&mock_user=\" (java.net.URLEncoder/encode mock-user \"UTF-8\")))\n+\n+(defn seed-rust-children! [app-base]\n+ (let [fetch (curl-fetch-children app-base \"reddit.com/r/rust\")]\n+ (when-not (zero? (:exit fetch))\n+ (throw (ex-info \"fetch rust children failed\" {:err (:err fetch)})))))\ndiff --git a/test/support/mock_oauth.clj b/test/support/mock_oauth.clj\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..909d7a7be8b159ea54a122d69c46af3db3318118\n--- /dev/null\n+++ b/test/support/mock_oauth.clj\n@@ -0,0 +1,84 @@\n+(ns test.support.mock-oauth\n+ \"In-process HTTP stub for GitHub OAuth (authorize, token, /user).\"\n+ (:require [clojure.string :as str])\n+ (:import [com.sun.net.httpserver HttpServer HttpHandler HttpExchange]\n+ [java.net InetSocketAddress URLDecoder]))\n+\n+(defn- query-param [query key]\n+ (when query\n+ (some (fn [pair]\n+ (let [[k v] (str/split pair \"=\" 2)]\n+ (when (= k key)\n+ (URLDecoder/decode (or v \"\") \"UTF-8\"))))\n+ (str/split query #\"&\"))))\n+\n+(defn- parse-mock-user [raw]\n+ (let [s (or raw \"1002:newbie\")\n+ [id login] (str/split s #\":\" 2)]\n+ {:id (Long/parseLong id)\n+ :login (or login \"newbie\")}))\n+\n+(defn- send-json [^HttpExchange ex status body]\n+ (let [bytes (.getBytes body \"UTF-8\")]\n+ (.set (.getResponseHeaders ex) \"Content-Type\" \"application/json\")\n+ (.sendResponseHeaders ex status (alength bytes))\n+ (doto (.getResponseBody ex)\n+ (.write bytes)\n+ (.close))))\n+\n+(defn- send-redirect [^HttpExchange ex location]\n+ (.set (.getResponseHeaders ex) \"Location\" location)\n+ (.sendResponseHeaders ex 302 -1)\n+ (.close (.getResponseBody ex)))\n+\n+(defn- read-form-code [^HttpExchange ex]\n+ (let [body (slurp (.getInputStream ex))]\n+ (query-param body \"code\")))\n+\n+(defn- bearer-token [^HttpExchange ex]\n+ (some-> (.getRequestHeaders ex)\n+ (.getFirst \"Authorization\")\n+ (str/replace #\"^[Bb]earer \" \"\")))\n+\n+(defn- parse-token-user [token]\n+ (when (str/starts-with? token \"mock:\")\n+ (parse-mock-user (subs token 5))))\n+\n+(defn start-mock-oauth\n+ \"Start mock GitHub OAuth on `port`. Returns a zero-arg `stop` function.\"\n+ [port]\n+ (let [server (HttpServer/create (InetSocketAddress. \"127.0.0.1\" port) 0)\n+ handler\n+ (proxy [HttpHandler] []\n+ (handle [^HttpExchange exchange]\n+ (let [uri (.getRequestURI exchange)\n+ path (.getPath uri)\n+ query (.getQuery uri)]\n+ (cond\n+ (str/ends-with? path \"/login/oauth/authorize\")\n+ (let [redirect-uri (query-param query \"redirect_uri\")\n+ state (query-param query \"state\")\n+ mock-user (query-param query \"mock_user\")\n+ user (parse-mock-user mock-user)\n+ code (str \"mock:\" (:id user) \":\" (:login user))\n+ loc (str redirect-uri \"?code=\" (java.net.URLEncoder/encode code \"UTF-8\")\n+ \"&state=\" (java.net.URLEncoder/encode state \"UTF-8\"))]\n+ (send-redirect exchange loc))\n+\n+ (str/ends-with? path \"/login/oauth/access_token\")\n+ (let [code (or (read-form-code exchange) \"mock:1002:newbie\")]\n+ (send-json exchange 200 (str \"{\\\"access_token\\\":\\\"\" code \"\\\",\\\"token_type\\\":\\\"bearer\\\"}\")))\n+\n+ (= path \"/user\")\n+ (let [token (bearer-token exchange)\n+ user (or (parse-token-user token) {:id 1002 :login \"newbie\"})]\n+ (send-json exchange 200\n+ (str \"{\\\"id\\\":\" (:id user) \",\\\"login\\\":\\\"\" (:login user) \"\\\"}\")))\n+\n+ :else\n+ (send-json exchange 404 \"{\\\"error\\\":\\\"not found\\\"}\")))))]\n+ (.createContext server \"/\" handler)\n+ (.setExecutor server nil)\n+ (.start server)\n+ (fn stop []\n+ (.stop server 0))))\ndiff --git a/test/support/seed_auth.clj b/test/support/seed_auth.clj\nnew file mode 100644\nindex 0000000000000000000000000000000000000000..88061ceb15daabd2032cc11192ae19c8579cbe5a\n--- /dev/null\n+++ b/test/support/seed_auth.clj\n@@ -0,0 +1,41 @@\n+(ns test.support.seed-auth\n+ \"Append-only event-log seeds for auth integration tests.\")\n+\n+(def seeder-uuid \"00000000-0000-0000-0000-000000000010\")\n+\n+(def rust-scope \"https://reddit.com/r/rust\")\n+(def post-a \"https://reddit.com/r/rust/comments/aaa\")\n+(def post-b \"https://reddit.com/r/rust/comments/bbb\")\n+\n+(defn- esc [s]\n+ (.replace s \"\\\\\" \"\\\\\\\\\"))\n+\n+(defn- line [seq ts event-json]\n+ (str \"{\\\"schema\\\":2,\\\"seq\\\":\" seq \",\\\"ts\\\":\" ts \",\\\"event\\\":\" event-json \"}\" \"\\n\"))\n+\n+(defn seeder-vote-events\n+ \"Events that register a seeder principal and one vote on the rust A/B pair.\"\n+ []\n+ [(line 1 1 (str \"{\\\"type\\\":\\\"principal_created\\\",\\\"uuid\\\":\\\"\" (esc seeder-uuid) \"\\\",\\\"ts\\\":1}\"))\n+ (line 2 2 (str \"{\\\"type\\\":\\\"oauth_linked\\\",\\\"uuid\\\":\\\"\" (esc seeder-uuid)\n+ \"\\\",\\\"provider\\\":\\\"github\\\",\\\"provider_id\\\":\\\"1001\\\",\\\"ts\\\":2}\"))\n+ (line 3 3 (str \"{\\\"type\\\":\\\"pseudonym_claimed\\\",\\\"uuid\\\":\\\"\" (esc seeder-uuid)\n+ \"\\\",\\\"pseudonym\\\":\\\"seeder\\\",\\\"ts\\\":3}\"))\n+ (line 4 4 (str \"{\\\"type\\\":\\\"node_ensured\\\",\\\"id\\\":\\\"\" (esc rust-scope) \"\\\"}\"))\n+ (line 5 5 (str \"{\\\"type\\\":\\\"node_ensured\\\",\\\"id\\\":\\\"\" (esc post-a) \"\\\"}\"))\n+ (line 6 6 (str \"{\\\"type\\\":\\\"node_ensured\\\",\\\"id\\\":\\\"\" (esc post-b) \"\\\"}\"))\n+ (line 7 7 (str \"{\\\"type\\\":\\\"vote_recorded\\\",\\\"ts\\\":7\"\n+ \",\\\"a\\\":\\\"\" (esc post-a) \"\\\",\\\"b\\\":\\\"\" (esc post-b) \"\\\"\"\n+ \",\\\"ratio_left\\\":3,\\\"ratio_right\\\":1\"\n+ \",\\\"scope\\\":\\\"\" (esc rust-scope) \"\\\"\"\n+ \",\\\"pseudonym\\\":\\\"seeder\\\",\\\"trust_weight\\\":1.5}\"))])\n+\n+(defn write-seeder-events!\n+ [event-log-path]\n+ (spit event-log-path (apply str (seeder-vote-events))))\n+\n+(defn seeder-pair-vote-url [app-base]\n+ (str app-base \"/vote?parent=\"\n+ (java.net.URLEncoder/encode rust-scope \"UTF-8\")\n+ \"&left=\" (java.net.URLEncoder/encode post-a \"UTF-8\")\n+ \"&right=\" (java.net.URLEncoder/encode post-b \"UTF-8\")))\ndiff --git a/test/vote_compare.clj b/test/vote_compare.clj\nindex bc0be9f90bb2eccbc756b20cbdf7212ddfed9cbc..8a4f39bb5a1c188b3612e8a5c95aeccd0be9d16d 100644\n--- a/test/vote_compare.clj\n+++ b/test/vote_compare.clj\n@@ -1,101 +1,36 @@\n (ns test.vote-compare\n- (:require [babashka.process :as process]\n- [clojure.java.io :as io]\n- [clojure.string :as str]\n+ (:require [clojure.string :as str]\n [clojure.test :refer [deftest is testing]]\n [com.blockether.spel.core :as core]\n [com.blockether.spel.locator :as loc]\n [com.blockether.spel.page :as page]\n- [test.support.mock-reddit :as mock-reddit])\n- (:import [java.net URLEncoder]))\n-\n-(defn- repo-root []\n- (.getCanonicalPath (io/file (System/getProperty \"user.dir\"))))\n-\n-(defn- pick-port []\n- (with-open [s (java.net.ServerSocket. 0)]\n- (.getLocalPort s)))\n-\n-(defn- wait-health [base-url ms]\n- (let [deadline (+ (System/currentTimeMillis) ms)\n- url (str base-url \"/healthz\")]\n- (loop []\n- (let [resp (try\n- (process/shell {:out :string :err :string}\n- \"curl\" \"-sf\" url)\n- (catch Exception _ nil))]\n- (if (and resp (zero? (:exit resp)) (= \"ok\" (str/trim (:out resp \"\"))))\n- true\n- (if (< (System/currentTimeMillis) deadline)\n- (do (Thread/sleep 200) (recur))\n- false))))))\n-\n-(defn- curl-fetch-children [base item]\n- (process/shell {:out :string :err :string}\n- \"curl\" \"-sfN\" \"--max-time\" \"20\"\n- \"-X\" \"POST\" (str base \"/ui\")\n- \"--data-urlencode\"\n- (str \"__rpc__={\\\"action\\\":\\\"fetch_entity\\\",\\\"item\\\":\\\"\" item\n- \"\\\",\\\"kind\\\":\\\"children\\\"}\")))\n-\n-(defn- vote-page-url [base parent]\n- (str base \"/vote?parent=\"\n- (URLEncoder/encode parent \"UTF-8\")))\n+ [test.support.harness :as harness]\n+ [test.support.seed-auth :as seed-auth]))\n \n (deftest vote-compare-shows-recorded-vote-after-post\n- (testing \"post vote on /vote morphs edge history (mock Reddit children seeded)\"\n- (let [root (repo-root)\n- fixtures (mock-reddit/fixtures-dir root)\n- data-dir (.getAbsolutePath\n- (doto (io/file (System/getProperty \"java.io.tmpdir\")\n- (str \"sorter2-vote-\" (System/currentTimeMillis)))\n- (.mkdirs)))\n- reddit-port (pick-port)\n- app-port (pick-port)\n- reddit-base (str \"http://127.0.0.1:\" reddit-port)\n- app-base (str \"http://127.0.0.1:\" app-port)\n- bin (str root \"/target/release/sorter2-server\")\n- stop-mock (mock-reddit/start-mock-reddit reddit-port fixtures)]\n+ (testing \"post vote on /vote morphs edge history (mock Reddit + auth session)\"\n+ (let [servers (harness/with-auth-servers\n+ (fn [data-dir]\n+ (seed-auth/write-seeder-events! (str data-dir \"/events.jsonl\"))))]\n (try\n- (is (zero? (:exit (process/shell {:dir root}\n- \"cargo\" \"build\" \"--release\" \"--package\" \"sorter2-server\")))\n- \"release build succeeds\")\n- (let [proc (process/process {:dir root\n- :env (into (into {} (System/getenv))\n- {\"SORTER2_SKIP_DOTENV\" \"1\"\n- \"SORTER2_DATA_DIR\" data-dir\n- \"SORTER2_EVENT_LOG\" (str data-dir \"/events.jsonl\")\n- \"PORT\" (str app-port)\n- \"REDDIT_API_BASE\" reddit-base\n- \"REDDIT_OAUTH_BASE\" reddit-base\n- \"REDDIT_CLIENT_ID\" \"\"\n- \"REDDIT_CLIENT_SECRET\" \"\"\n- \"REDDIT_APP_ID\" \"\"\n- \"REDDIT_APP_SECRET\" \"\"})\n- :out :string\n- :err :string}\n- bin)]\n- (try\n- (is (wait-health app-base 20000) \"app healthz\")\n- (let [fetch (curl-fetch-children app-base \"reddit.com/r/rust\")]\n- (is (zero? (:exit fetch)) \"fetch posts via mock Reddit\")\n- (is (str/includes? (:out fetch) \"Idiomorph.morph\")))\n- (core/with-testing-page [pg]\n- (page/navigate pg (vote-page-url app-base \"reddit.com/r/rust\"))\n- (page/wait-for-selector pg \"#vote-compare-form\")\n- (let [before (loc/text-content (page/locator pg \"#vote-edge-history-region\"))]\n- (is (str/includes? before \"no votes on this pair yet\")\n- \"empty edge history before first vote\"))\n+ (harness/seed-rust-children! (:app-base servers))\n+ (let [vote-url (seed-auth/seeder-pair-vote-url (:app-base servers))]\n+ (core/with-testing-page [pg]\n+ (page/navigate pg (harness/oauth-login-url (:app-base servers) \"/\" \"1001:seeder\"))\n+ (page/wait-for-selector pg \".top-nav\" {:timeout 15000})\n+ (page/navigate pg vote-url)\n+ (page/wait-for-selector pg \"#vote-compare-form\")\n+ (let [before (loc/text-content (page/locator pg \"#vote-edge-history-region\"))]\n+ (is (str/includes? before \"votes on this pair\")\n+ \"seeded seeder vote visible before our vote\")\n (loc/click (page/get-by-test-id pg \"vote-post\"))\n (page/wait-for-selector pg \".vote-edge-history-title\")\n (let [after (loc/text-content (page/locator pg \"#vote-edge-history-region\"))]\n (is (str/includes? after \"votes on this pair\")\n \"shows edge history title after vote\")\n- (is (str/includes? after \"1:1\")\n- \"shows submitted ratio after vote (default slider at center)\")\n+ (is (not= before after)\n+ \"edge history updated after authenticated vote\")\n (is (not (str/includes? after \"no votes on this pair yet\"))\n- \"does not revert to empty edge history\")))\n- (finally\n- (process/destroy proc))))\n+ \"does not revert to empty edge history\")))))\n (finally\n- (stop-mock))))))\n+ ((:stop servers)))))))\n","role":"user"}],"model":"~x-ai/grok-latest"}