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!("", p.id, ago); + let thread_attr = p + .thread + .as_deref() + .map(|thread| format!(" thread=\"{thread}\"")) + .unwrap_or_default(); + println!( + "", + p.id, ago, p.room, thread_attr + ); print!("{}", p.body); if !p.body.ends_with('\n') { println!(); } println!(""); 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, + requested_since: Option, + 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, 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 = 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.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.principal == viewer).then_some((index, ing.ts)) }) - }) - .filter(|id| !reduced.redacted_posts.contains(*id)) - .collect(); - let total = matching.len(); - let posts: Vec = 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(); - let line = line_ok(RpcResult::Feed(FeedResponse { - delegate: None, + }); + let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT); + let line = line_ok(RpcResult::Feed(rpc_feed( + &reduced, + &viewer, + None, since, - posts, - total, - })); + implicit_anchor, + limit, + ))); drop(reduced); line } diff --git a/server/tests/integration_rooms.rs b/server/tests/integration_rooms.rs index a3b378267e090a016a7e663f103a0dd9f482690b..e623523ba352ce456674a399542c689505760b78 100644 --- a/server/tests/integration_rooms.rs +++ b/server/tests/integration_rooms.rs @@ -1,6 +1,9 @@ mod support; use slug_types::room_route_segment; +use slugsocial_server::events::{ + AgentBound, Event, GrantAdded, GrantRevoked, Ingest, RoomCreated, ThreadCapability, +}; use support::*; #[tokio::test] @@ -473,6 +476,247 @@ async fn test_feed_without_delegate_uses_principal_last_post_including_delegate( ); } +#[tokio::test] +async fn test_feed_uses_delegate_ingest_position_when_multi_user_timestamps_collide() { + let (addr, _tmp, _log, state, _handle) = create_test_server_with_state().await; + seed_test_identity(&state, "bob", "bobtok", "bobsecret").await; + let client = reqwest::Client::new(); + let alice_delegate = + "00000000-0000-0000-0000-0000000000c1:feedtest:local/alice-model"; + let bob_delegate = + "00000000-0000-0000-0000-0000000000c2:feedtest:local/bob-model"; + + { + let mut reduced = state.reduced.write().await; + for (agent, username) in [ + (alice_delegate, "testuser"), + (bob_delegate, "bob"), + ] { + reduced.apply_event(Event::AgentBound(AgentBound { + ts: 1, + agent: agent.to_string(), + username: username.to_string(), + })); + } + reduced.apply_event(Event::Ingest(Ingest { + ts: 100, + id: "alice-anchor".into(), + raw: "alice anchor".into(), + principal: "testuser".into(), + delegate: Some(alice_delegate.into()), + room_id: "public".into(), + thread_tag: "feed-order".into(), + })); + reduced.apply_event(Event::Ingest(Ingest { + ts: 100, + id: "bob-same-millisecond".into(), + raw: "bob same-millisecond change".into(), + principal: "bob".into(), + delegate: Some(bob_delegate.into()), + room_id: "public".into(), + thread_tag: "feed-order".into(), + })); + // Event-log order remains authoritative even if the wall clock moves backwards. + reduced.apply_event(Event::Ingest(Ingest { + ts: 99, + id: "bob-clock-rollback".into(), + raw: "bob change after clock rollback".into(), + principal: "bob".into(), + delegate: Some(bob_delegate.into()), + room_id: "public".into(), + thread_tag: "feed-order".into(), + })); + } + + let feed = rpc_batch( + &client, + addr, + Some(&test_bearer()), + serde_json::json!([{ + "GetFeed": { "delegate": alice_delegate, "limit": 20 } + }]), + ) + .await; + let posts = feed["results"][0]["result"]["Feed"]["posts"] + .as_array() + .unwrap(); + let ids: Vec<&str> = posts.iter().filter_map(|p| p["id"].as_str()).collect(); + assert_eq!( + ids, + ["bob-clock-rollback", "bob-same-millisecond"], + "implicit feed cutoff must use append position, not timestamp" + ); +} + +#[tokio::test] +async fn test_feed_multi_user_private_room_visibility_and_revoked_anchor_access() { + let (addr, _tmp, _log, state, _handle) = create_test_server_with_state().await; + seed_test_identity(&state, "bob", "bobtok", "bobsecret").await; + seed_test_identity(&state, "carol", "caroltok", "carolsecret").await; + let client = reqwest::Client::new(); + let alice_delegate = + "00000000-0000-0000-0000-0000000000d1:feedtest:local/alice-model"; + let bob_delegate = + "00000000-0000-0000-0000-0000000000d2:feedtest:local/bob-model"; + let room_alice = "alice01/alice-room"; + let room_bob = "bob0001/bob-room"; + + { + let mut reduced = state.reduced.write().await; + for (agent, username) in [ + (alice_delegate, "testuser"), + (bob_delegate, "bob"), + ] { + reduced.apply_event(Event::AgentBound(AgentBound { + ts: 1, + agent: agent.to_string(), + username: username.to_string(), + })); + } + for (room_id, slug, owner) in [ + (room_alice, "alice-room", "testuser"), + (room_bob, "bob-room", "bob"), + ] { + reduced.apply_event(Event::RoomCreated(RoomCreated { + ts: 2, + room_id: room_id.to_string(), + slug: slug.to_string(), + owner: owner.to_string(), + })); + reduced.apply_event(Event::GrantAdded(GrantAdded { + ts: 2, + room_id: room_id.to_string(), + username: owner.to_string(), + capabilities: vec![ThreadCapability::View], + granted_by: owner.to_string(), + })); + } + + let mut add_ingest = + |ts, id: &str, raw: &str, principal: &str, delegate: Option<&str>, room: &str| { + reduced.apply_event(Event::Ingest(Ingest { + ts, + id: id.into(), + raw: raw.into(), + principal: principal.into(), + delegate: delegate.map(str::to_string), + room_id: room.into(), + thread_tag: "feed-permissions".into(), + })); + }; + add_ingest(10, "prehistory", "must remain before alice anchor", "carol", None, "public"); + add_ingest( + 11, + "alice-private-anchor", + "alice last posted here", + "testuser", + Some(alice_delegate), + room_alice, + ); + add_ingest( + 12, + "bob-public-anchor", + "bob last posted here", + "bob", + Some(bob_delegate), + "public", + ); + add_ingest(13, "alice-room-change", "visible only to alice", "carol", None, room_alice); + add_ingest(14, "bob-room-change", "visible only to bob", "carol", None, room_bob); + add_ingest(15, "public-change", "visible to everyone", "carol", None, "public"); + + reduced.apply_event(Event::GrantRevoked(GrantRevoked { + ts: 16, + room_id: room_alice.into(), + username: "testuser".into(), + capabilities: vec![ThreadCapability::View], + revoked_by: "testuser".into(), + })); + } + + let alice_feed = rpc_batch( + &client, + addr, + Some(&test_bearer()), + serde_json::json!([{ + "GetFeed": { "delegate": alice_delegate, "limit": 20 } + }]), + ) + .await; + let alice_posts = alice_feed["results"][0]["result"]["Feed"]["posts"] + .as_array() + .unwrap(); + let alice_ids: Vec<&str> = alice_posts + .iter() + .filter_map(|p| p["id"].as_str()) + .collect(); + assert_eq!( + alice_ids, + ["public-change", "bob-public-anchor"], + "revoked private content must be hidden without moving the delegate anchor backwards" + ); + assert!(!alice_ids.contains(&"prehistory")); + assert!( + alice_posts + .iter() + .all(|post| post["room"].as_str() == Some("public")), + "private posts must not leak and every feed post must identify its room" + ); + + let bob_bearer = test_bearer_for("bobtok", "bobsecret"); + let bob_feed = rpc_batch( + &client, + addr, + Some(&bob_bearer), + serde_json::json!([{ + "GetFeed": { "delegate": bob_delegate, "limit": 20 } + }]), + ) + .await; + let bob_ids: Vec<&str> = bob_feed["results"][0]["result"]["Feed"]["posts"] + .as_array() + .unwrap() + .iter() + .filter_map(|p| p["id"].as_str()) + .collect(); + assert_eq!(bob_ids, ["public-change", "bob-room-change"]); + assert_eq!( + bob_feed["results"][0]["result"]["Feed"]["posts"][1]["room"], + room_bob + ); + + // Restoring View exposes only changes after the same stable delegate anchor. + { + let mut reduced = state.reduced.write().await; + reduced.apply_event(Event::GrantAdded(GrantAdded { + ts: 17, + room_id: room_alice.into(), + username: "testuser".into(), + capabilities: vec![ThreadCapability::View], + granted_by: "testuser".into(), + })); + } + let restored = rpc_batch( + &client, + addr, + Some(&test_bearer()), + serde_json::json!([{ + "GetFeed": { "delegate": alice_delegate, "limit": 20 } + }]), + ) + .await; + let restored_ids: Vec<&str> = restored["results"][0]["result"]["Feed"]["posts"] + .as_array() + .unwrap() + .iter() + .filter_map(|p| p["id"].as_str()) + .collect(); + assert_eq!( + restored_ids, + ["public-change", "alice-room-change", "bob-public-anchor"] + ); +} + #[tokio::test] async fn test_private_room_thread_urls_use_t_segment() { let (addr, _tmp, _log, _handle) = create_test_server().await; diff --git a/server/tests/support/mod.rs b/server/tests/support/mod.rs index 4a620eaa875e7a1145f2a4e82cc4ff2be21d5345..a4320e991763617c1a760efad4b621977e2b74d0 100644 --- a/server/tests/support/mod.rs +++ b/server/tests/support/mod.rs @@ -20,8 +20,10 @@ pub fn sha256_hex(s: &str) -> String { /// Fixed bearer for integration tests (`TokenIssued` seeded into reducer in `create_test_server`). pub fn test_bearer() -> String { - let token_id = "testtok"; - let secret = "secret"; + test_bearer_for("testtok", "secret") +} + +pub fn test_bearer_for(token_id: &str, secret: &str) -> String { format!("slug_{token_id}_{secret}") } @@ -71,19 +73,27 @@ pub async fn rpc_batch( } pub async fn seed_test_token(state: &AppState) { + seed_test_identity(state, "testuser", "testtok", "secret").await; +} + +/// Add a distinct principal and bearer to a running integration-test server. +pub async fn seed_test_identity( + state: &AppState, + username: &str, + token_id: &str, + secret: &str, +) { let registered = Event::UserRegistered(UserRegistered { ts: 0, - username: "testuser".to_string(), + username: username.to_string(), provider: "test".to_string(), - provider_id: "testuser".to_string(), + provider_id: username.to_string(), }); - let token_id = "testtok"; - let secret = "secret"; let salt = "salt"; let token_hash = sha256_hex(&format!("{salt}:{secret}")); let ev = Event::TokenIssued(TokenIssued { ts: 0, - username: "testuser".to_string(), + username: username.to_string(), token_id: token_id.to_string(), token_hash, salt: salt.to_string(), diff --git a/types/src/lib.rs b/types/src/lib.rs index 211493935d19582607f5c87fb492faf47bdc6f53..cd6f94c60a4e1c3bfbce13bc0404b803f5f57c6d 100644 --- a/types/src/lib.rs +++ b/types/src/lib.rs @@ -264,6 +264,9 @@ pub struct FeedResponse { pub struct FeedPost { pub ts: i64, pub id: String, + /// Permission scope containing the post: `"public"` or a private room id. + #[serde(default)] + pub room: String, /// Primary thread tag (without #), if the ingest declared one. #[serde(skip_serializing_if = "Option::is_none")] pub thread: Option,