Side A is a large mechanical refactor (REST endpoints collapsed into an RPC batch dispatcher) with lots of churn but no new correctness guarantee — it's mostly moving code around. Side B fixes a real, subtle bug (feed cutoff using wall-clock timestamps instead of durable append order, causing missed/duplicated posts under concurrent writes or clock rollback) and adds permission-aware multi-room visibility with targeted regression tests, delivering concrete lasting correctness value in a small, precise diff.
constitution · epochs · watch · epoch 3
c_2595b6007624 (tommy-mor) vs c_0a9a8eab32ba (tommy-mor)
download prompt · raw event · cmp_4859231b2767ad
council reasoning
A is a foundational redesign: batch RPC replaces the REST surface, events/reducer separate room scope from forum thread tags (RoomCreated, ScopeId::Room, ingests_by_scope_thread), and the CLI/tests are rewired to that model. B is a high-quality but narrow correctness fix—stable feed anchors on ingest index plus permission filtering—with strong tests, yet it only hardens one path after A’s architecture is in place.
Side A replaces many ad hoc REST endpoints with a unified RPC batch API, restructures the CLI around scoped public/private rooms, introduces room-aware routing in events and the reducer, and factors ingest validation into a reusable module. Although it is a large refactor, it establishes a lasting architectural foundation and migrates existing functionality, whereas Side B is a focused improvement that makes feed catch-up stable under timestamp collisions and permission-aware by anchoring to durable ingest order and adding targeted tests.
sides
A — c_2595b6007624 (tommy-mor)
message
[96b6da05] rpc + reducer changes first pass
diff preview
diff --git a/cli/src/main.rs b/cli/src/main.rs
index 630c5dea1f78c0ec9bc53e6b96234a0dc75bb705..8d0442959f4332bafe499a2a8cdf364731a06871 100644
--- a/cli/src/main.rs
+++ b/cli/src/main.rs
@@ -21,8 +21,9 @@ struct Cli {
cmd: Option<Command>,
}
+/// Commands scoped to a room (`public` or `shortid/slug`).
#[derive(Subcommand, Debug)]
-enum Command {
+enum ScopedCmd {
/// Browse the garden (ontology) — light mode, ranked by votes
Garden {
#[command(subcommand)]
@@ -174,6 +175,23 @@ enum Command {
#[arg(long)]
json: bool,
},
+}
+
+#[derive(Subcommand, Debug)]
+enum Command {
+ /// Public site (same as room `public`)
+ Public {
+ #[command(subcommand)]
+ sub: ScopedCmd,
+ },
+ /// Private room id (`shortid/slug` from `room create`)
+ Private {
+ /// Room id, e.g. `a1b2c3d/my-project`
+ #[arg(value_name = "ROOM_ID")]
+ room: String,
+ #[command(subcommand)]
+ sub: ScopedCmd,
+ },
/// Show all activity since you last posted (global feed)
///
@@ -628,6 +646,37 @@ fn http_client() -> Result<reqwest::Client> {
.build()?)
}
+async fn send_rpc(
+ client: &reqwest::Client,
+ base: &str,
+ bearer: Option<&str>,
+ commands: Vec<RpcCommand>,
+) -> Result<RpcBatchResponse> {
+ let url = format!("{}/api/v0/rpc", base.trim_end_matches('/'));
+ let mut req = client.post(url).json(&RpcBatch(commands));
+ if let Some(b) = bearer {
+ req = req.header("Authorization", format!("Bearer {}", b));
+ }
+ let resp = req.send().await?;
+ let status = resp.status();
+ let text = resp.text().await.unwrap_or_default();
+ if !status.is_success() {
+ return Err(anyhow!("rpc HTTP {}: {}", status, text.trim()));
+ }
+ serde_json::from_str(&text).map_err(|e| anyhow!("rpc response: {e}"))
+}
+
+fn rpc_line_ok(line: &RpcLine) -> Result<&RpcResult> {
+ if !line.ok {
+ let mut m = line.error.clone().unwrap_or_else(|| "rpc error".into());
+ if let Some(h) = &line.hint {
+ m.push_str(&format!("\nhint: {h}"));
+ }
+ return Err(anyhow!(m));
+ }
+ line.result.as_ref().ok_or_else(|| anyhow!("rpc missing result"))
+}
+
/// Normalize ontology path for API. Accepts path with or without ~/ (shell expands ~ to $HOME).
/// Returns a bare slug path (e.g. `languages/python`) with no leading `/` or `~/`.
/// Call `ontology_path_for_api_query` before sending `item=` / `parent=` params so the server
@@ -729,231 +778,262 @@ fn write_secret_file(name: &str, contents: &str) -> Result<()> {
Ok(())
}
-#[tokio::main]
-async fn main() -> Result<()> {
- let Cli { cmd, server } = Cli::parse();
-
- // If no command provided, print the guide
- let Some(cmd) = cmd else {
- print!("{}", include_str!("../GUIDE.sorter"));
- return Ok(());
- };
-
- let base = server.trim_end_matches('/');
-
- match cmd {
- Command::Healthz { json } => {
- let client = http_client()?;
- let url = format!("{base}/healthz");
- let body = client.get(url).send().await?.text().await?;
- if json {
- // Wrap plain text response in a JSON object
- println!("{}", serde_json::json!({ "ok": true, "body": body.trim() }));
- } else {
- println!("{body}");
- }
- }
-
- Command::Search { query, json } => {
- let client = http_client()?;
- let url = format!("{base}/api/v0/search?q={}", urlencoding::encode(&query));
- let resp: slug_types::SearchResponse = expect_json(client.get(url).send().await?).await?;
- if json {
- println!("{}", serde_json::to_string_pretty(&resp)?);
- } else {
- if !resp.items.is_empty() {
- println!("items ({})", resp.items.len());
- for item in &resp.items {
- print!(" {}", item.path);
- if let Some(body) = &item.body {
- let first_line = body.lines().next().unwrap_or("").trim();
- if !first_line.is_empty() {
- print!(" {}", first_line);
+async fn run_scoped(base: &str, room: &str, sub: ScopedCmd) -> Result<()> {
+ let room = room.trim();
+ let client = http_client()?;
+ match sub {
+ ScopedCmd::Garden { sub } => match sub {
+ GardenCmd::Tree { json } => {
+ let batch = send_rpc(&client, base, None, vec![RpcCommand::GetLeaves { room: room.to_string() }]).await?;
+ match rpc_line_ok(&batch.results[0])? {
+ RpcResult::Leaves(resp) => {
+ if json {
+ println!("{}", serde_json::to_string_pretty(&resp)?);
+ } else {
+ for p in &resp.paths {
+ println!("~/{}", p);
}
}
- println!();
- }
- }
- if !resp.threads.is_empty() {
- if !resp.items.is_empty() { println!(); }
- println!("threads ({})", resp.threads.len());
- let now_ms = std::time::SystemTime::now()
- .duration_since(std::time::UNIX_EPOCH)
- .unwrap_or_default()
- .as_millis() as i64;
- for t in &resp.threads {
- println!(" {} {}n {}", t.tag, t.post_count, slug_types::timeago::timeago(now_ms, t.last_activity));
- }
- }
- if !resp.posts.is_empty() {
- if !resp.items.is_empty() || !resp.threads.is_empty() { println!(); }
- println!("posts ({})", resp.posts.len());
- let now_ms = std::time::SystemTime::now()
- .duration_since(std::time::UNIX_EPOCH)
- .unwrap_or_default()
- .as_millis() as i64;
- for p in &resp.posts {
- let first_line = p.snippet.lines().next().unwrap_or("").trim();
- println!(" {} · {} {}", p.thread, slug_types::timeago::timeago(now_ms, p.ts), first_line);
- }
- }
- if resp.items.is_empty() && resp.threads.is_empty() && resp.posts.is_empty() {
- println!("no results");
- }
- }
- }
-
- Command::Garden { sub } => match sub {
- GardenCmd::Tree { json } => {
- let client = http_client()?;
- let url = format!("{base}/api/v0/leaves");
- let builder = client.get(url);
- let resp: LeavesResponse = expect_json(builder.send().await?).await?;
- if json {
- println!("{}", serde_json::to_string_pretty(&resp)?);
- } else {
- for p in &resp.paths {
- println!("~/{}", p);
}
+ _ => return Err(anyhow!("unexpected RPC result")),
}
}
-
GardenCmd::Body { path, json, full } => {
let path = normalize_ontology_path_input(&path).map_err(anyhow::Error::msg)?;
let item_q = ontology_path_for_api_query(&path);
- let client = http_client()?;
- let mut url = format!("{base}/api/v0/item?item={}", urlencoding::encode(&item_q));
- if full {
- url.push_str("&full=true");
- }
- let builder = client.get(url);
- let resp: ItemResponse = expect_json(builder.send().await?).await?;
- if json {
- println!("{}", serde_json::to_string_pretty(&resp)?);
- } else {
- print_item_response(&resp);
+ let batch = send_rpc(
+ &client,
+ base,
+ None,
+ vec![RpcCommand::GetGardenItem {
+ room: room.to_string(),
+ item_path: item_q,
+ full: Some(full),
+ }],
+ )
+ .await?;
+ match rpc_line_ok(&batch.results[0])? {
+ RpcResult::GardenItem(resp) => {
+ if json {
+ println!("{}", serde_json::to_string_pretty(&resp)?);
+ } else {
+ print_item_response(&resp);
+ }
+ }
+ _ => return Err(anyhow!("unexpected RPC result")),
}
}
-
GardenCmd::Children { paths, depth, json } => {
let paths: Vec<String> = paths
.iter()
.map(|p| normalize_ontology_path_input(p).map_err(anyhow::Error::msg))
.collect::<Result<Vec<_>>>()?;
- let client = http_client()?;
let parent_param = paths
.iter()
.map(|p| ontology_path_for_api_query(p))
.collect::<Vec<_>>()
.join(",");
- let mut url = format!("{base}/api/v0/rank?parent={}", urlencoding::encode(&parent_param));
- if let Some(d) = depth {
- url.push_str(&format!("&depth={d}"));
- }
- let builder = client.get(url);
- let resp: RankResponse = expect_json(builder.send().await?).await?;
-
- if json {
- println!("{}", serde_json::to_string_pretty(&resp)?);
- } else {
- print_rank_response(&resp);
+ let batch = send_rpc(
+ &client,
+ base,
+ None,
+ vec![RpcCommand::GetGardenRank {
+ room: room.to_string(),
+ parent_path: parent_param,
+ depth,
+ offset: None,
+ limit: None,
+ percent: None,
+ }],
+ )
+ .await?;
+ match rpc_line_ok(&batch.results[0])? {
+ RpcResult::GardenRank(resp) => {
+ if json {
+ println!("{}", serde_json::to_string_pretty(&resp)?);
+ } else {
+ print_rank_response(&resp);
+ }
+ }
+ _ => return Err(anyhow!("unexpected RPC result")),
}
}
-
GardenCmd::Pair { path, json } => {
let path = normalize_ontology_path_input(&path).map_err(anyhow::Error::msg)?;
let parent_q = ontology_path_for_api_query(&path);
- let client = http_client()?;
- let url = format!("{base}/api/v0/pair?parent={}", urlencoding::encode(&parent_q));
- let builder = client.get(url);
- let resp: PairResponse = expect_json(builder.send().await?).await?;
- if json {
- println!("{}", serde_json::to_string_pretty(&resp)?);
- } else {
- print_pair_response(&resp);
+ let batch = send_rpc(
+ &client,
+ base,
+ None,
+ vec![RpcCommand::GetPair {
+ room: room.to_string(),
+ parent_path: parent_q,
+ }]
… preview truncated; 237,234 characters omittedB — c_0a9a8eab32ba (tommy-mor)
message
[c94456ff] Make feed catch-up stable and permission-aware Anchor implicit feeds to durable ingest order and cover multi-user private-room visibility so concurrent posts are not missed or leaked. Co-authored-by: Cursor <cursoragent@cursor.com>
diff preview
diff --git a/cli/src/main.rs b/cli/src/main.rs
index abb5a55b49f60fe28fbfd4ec02715cb94ea0b4ec..c4f1494df3aedbd8b883aea6249579aa8336abfe 100644
--- a/cli/src/main.rs
+++ b/cli/src/main.rs
@@ -878,6 +878,30 @@ mod tests {
"graph: 4 items, 3/6 pairs (50.0% density), 1 component, connected"
);
}
+
+ #[test]
+ fn feed_without_since_uses_logged_in_delegate_from_env() {
+ let key = "SLUG_DELEGATE";
+ let previous = std::env::var_os(key);
+ let expected = "00000000-0000-0000-0000-0000000000ee:test:local/model";
+ std::env::set_var(key, expected);
+
+ let cli = Cli::try_parse_from(["slugsocial", "feed"]).expect("parse feed");
+
+ match previous {
+ Some(value) => std::env::set_var(key, value),
+ None => std::env::remove_var(key),
+ }
+ match cli.cmd {
+ Some(Command::Feed {
+ delegate, since, ..
+ }) => {
+ assert_eq!(delegate.as_deref(), Some(expected));
+ assert!(since.is_none());
+ }
+ _ => panic!("expected feed command"),
+ }
+ }
}
async fn run_scoped(base: &str, room: &str, sub: ScopedCmd) -> Result<()> {
@@ -1626,7 +1650,15 @@ async fn run() -> Result<()> {
} else {
for p in &resp.posts {
let ago = slug_types::timeago::timeago(now_ms, p.ts);
- println!("<post id=\"{}\" ts=\"{}\">", p.id, ago);
+ let thread_attr = p
+ .thread
+ .as_deref()
+ .map(|thread| format!(" thread=\"{thread}\""))
+ .unwrap_or_default();
+ println!(
+ "<post id=\"{}\" ts=\"{}\" room=\"{}\"{}>",
+ p.id, ago, p.room, thread_attr
+ );
print!("{}", p.body);
if !p.body.ends_with('\n') { println!(); }
println!("</post>");
diff --git a/server/src/api/rpc.rs b/server/src/api/rpc.rs
index afc4f95bef160c1e38ecff2c096d6440cd94e2b3..46d748f918d9b225acc4ecedfe5a1793089407b5 100644
--- a/server/src/api/rpc.rs
+++ b/server/src/api/rpc.rs
@@ -72,6 +72,80 @@ fn can_view_scope(reduced: &ReducerState, scope: &ScopeId, principal: Option<&st
}
}
+/// Build a feed in durable ingest order.
+///
+/// An implicit feed boundary is an ingest position, not only its millisecond timestamp. Two users
+/// can post in the same millisecond, and wall-clock timestamps can move backwards during replay.
+/// Explicit `since` remains a timestamp query for API compatibility, but scans the whole ordered
+/// ledger rather than assuming timestamps are monotonic.
+fn rpc_feed(
+ reduced: &ReducerState,
+ viewer: &str,
+ delegate: Option<String>,
+ requested_since: Option<i64>,
+ implicit_anchor: Option<(usize, i64)>,
+ limit: usize,
+) -> FeedResponse {
+ let since = requested_since.or_else(|| implicit_anchor.map(|(_, ts)| ts));
+ let implicit_anchor_index = requested_since
+ .is_none()
+ .then(|| implicit_anchor.map(|(index, _)| index))
+ .flatten();
+
+ let matching: Vec<&str> = reduced
+ .ingests_ordered
+ .iter()
+ .enumerate()
+ .rev()
+ .filter(|(index, id)| {
+ reduced.ingests_by_id.get(id.as_str()).is_some_and(|ing| {
+ match requested_since {
+ Some(cutoff) => ing.ts > cutoff,
+ None => implicit_anchor_index.is_none_or(|anchor| *index > anchor),
+ }
+ })
+ })
+ .map(|(_, id)| id.as_str())
+ .filter(|id| {
+ reduced.ingests_by_id.get(*id).is_some_and(|ing| {
+ let scope = scope_from_room_wire(&ing.room_id);
+ can_view_scope(reduced, &scope, Some(viewer))
+ })
+ })
+ .filter(|id| !reduced.redacted_posts.contains(*id))
+ .collect();
+
+ let total = matching.len();
+ let posts = matching
+ .into_iter()
+ .take(limit)
+ .filter_map(|id| reduced.ingests_by_id.get(id))
+ .map(|ing| {
+ let scope = scope_from_room_wire(&ing.room_id);
+ let thread_post_index = reduced.try_thread_post_index_chronological(
+ &scope,
+ &ing.thread_tag,
+ &ing.id,
+ );
+ FeedPost {
+ ts: ing.ts,
+ id: ing.id.clone(),
+ room: ing.room_id.clone(),
+ thread: Some(ing.thread_tag.clone()),
+ thread_post_index,
+ body: ing.raw.clone(),
+ }
+ })
+ .collect();
+
+ FeedResponse {
+ delegate,
+ since,
+ posts,
+ total,
+ }
+}
+
fn principal_from_optional_bearer(headers: &HeaderMap, reduced: &ReducerState) -> Result<Option<String>, RpcErr> {
if headers.contains_key(axum::http::header::AUTHORIZATION) {
verify_bearer_principal(headers, reduced)
@@ -1506,60 +1580,27 @@ pub async fn handle_rpc_batch(
Some("this delegate is not bound to your signed-in account".into()),
)
} else {
- let since_default = reduced
+ let implicit_anchor = reduced
.ingests_ordered
.iter()
+ .enumerate()
.rev()
- .filter_map(|id| reduced.ingests_by_id.get(id))
- .find(|ing| {
- if ing.delegate.as_deref() != Some(delegate_stored.as_str()) {
- return false;
- }
- let scope = scope_from_room_wire(&ing.room_id);
- can_view_scope(&reduced, &scope, Some(viewer.as_str()))
- })
- .map(|ing| ing.ts);
- let since = since.or(since_default);
- let cutoff = since.unwrap_or(0);
- let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
- let matching: Vec<&str> = reduced.ingests_ordered.iter().rev()
- .map(|id| id.as_str())
- .take_while(|id| reduced.ingests_by_id.get(*id).is_some_and(|ing| ing.ts > cutoff))
- .filter(|id| {
- reduced.ingests_by_id.get(*id).is_some_and(|ing| {
- let scope = scope_from_room_wire(&ing.room_id);
- can_view_scope(&reduced, &scope, Some(viewer.as_str()))
+ .find_map(|(index, id)| {
+ reduced.ingests_by_id.get(id).and_then(|ing| {
+ (ing.delegate.as_deref()
+ == Some(delegate_stored.as_str()))
+ .then_some((index, ing.ts))
})
- })
- .filter(|id| !reduced.redacted_posts.contains(*id))
- .collect();
- let total = matching.len();
- let posts: Vec<FeedPost> = matching.into_iter()
- .take(limit)
- .filter_map(|id| reduced.ingests_by_id.get(id))
- .map(|ing| {
- let scope = scope_from_room_wire(&ing.room_id);
- let thread_post_index = reduced
- .try_thread_post_index_chronological(
- &scope,
- &ing.thread_tag,
- &ing.id,
- );
- FeedPost {
- ts: ing.ts,
- id: ing.id.clone(),
- thread: Some(ing.thread_tag.clone()),
- thread_post_index,
- body: ing.raw.clone(),
- }
- })
- .collect();
- line_ok(RpcResult::Feed(FeedResponse {
- delegate: Some(delegate_stored),
+ });
+ let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
+ line_ok(RpcResult::Feed(rpc_feed(
+ &reduced,
+ &viewer,
+ Some(delegate_stored),
since,
- posts,
- total,
- }))
+ implicit_anchor,
+ limit,
+ )))
};
drop(reduced);
line
@@ -1567,60 +1608,25 @@ pub async fn handle_rpc_batch(
None => {
// Session catch-up: last time *you* posted anything (delegate or not), so revisiting
// an old chat with only a token still gets a sane cutoff.
- let since_default = reduced
+ let implicit_anchor = reduced
.ingests_ordered
.iter()
+ .enumerate()
.rev()
- .filter_map(|id| reduced.ingests_by_id.get(id))
- .find(|ing| {
- if ing.principal != viewer {
- return false;
- }
- let scope = scope_from_room_wire(&ing.room_id);
- can_view_scope(&reduced, &scope, Some(viewer.as_str()))
- })
- .map(|ing| ing.ts);
- let since = since.or(since_default);
- let cutoff = since.unwrap_or(0);
- let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
- let matching: Vec<&str> = reduced.ingests_ordered.iter().rev()
- .map(|id| id.as_str())
- .take_while(|id| reduced.inges
… preview truncated; 14,758 characters omittedHardlinks — judgments / attempts / prompt
judgments
attempts
Prompt text is loaded only by the download route.