diff --git a/.sqlx/query-5dd94b27dbcc35c0536394c6158cdbc5fb28436947566806bf8f49b19fc6db85.json b/.sqlx/query-5dd94b27dbcc35c0536394c6158cdbc5fb28436947566806bf8f49b19fc6db85.json new file mode 100644 index 0000000..84fe451 --- /dev/null +++ b/.sqlx/query-5dd94b27dbcc35c0536394c6158cdbc5fb28436947566806bf8f49b19fc6db85.json @@ -0,0 +1,198 @@ +{ + "db_name": "SQLite", + "query": "SELECT t.did AS \"did!\", t.rkey AS \"rkey!\", t.uri AS \"uri!\", t.cid AS \"cid!\",\n t.name AS \"name!\", t.summary, t.license AS \"license!\", t.tags_json, t.cover_json,\n t.derived_from_uri, t.created_at AS \"created_at!\", t.indexed_at AS \"indexed_at!\",\n t.record_json,\n COALESCE(cs.like_count, 0) AS \"like_count!: i64\",\n COALESCE(cs.save_count, 0) AS \"save_count!: i64\",\n (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS \"model_count!: i64\",\n (SELECT COUNT(*) FROM thing_models tm JOIN model_parts mp ON mp.model_uri = tm.model_uri WHERE tm.thing_uri = t.uri) AS \"part_count!: i64\"\n FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri\n WHERE COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) > 0\n ORDER BY COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) DESC,\n t.rkey DESC, t.uri DESC\n LIMIT ? OFFSET ?", + "describe": { + "columns": [ + { + "name": "did!", + "ordinal": 0, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "did" + } + } + }, + { + "name": "rkey!", + "ordinal": 1, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "rkey" + } + } + }, + { + "name": "uri!", + "ordinal": 2, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "uri" + } + } + }, + { + "name": "cid!", + "ordinal": 3, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "cid" + } + } + }, + { + "name": "name!", + "ordinal": 4, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "name" + } + } + }, + { + "name": "summary", + "ordinal": 5, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "summary" + } + } + }, + { + "name": "license!", + "ordinal": 6, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "license" + } + } + }, + { + "name": "tags_json", + "ordinal": 7, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "tags_json" + } + } + }, + { + "name": "cover_json", + "ordinal": 8, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "cover_json" + } + } + }, + { + "name": "derived_from_uri", + "ordinal": 9, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "derived_from_uri" + } + } + }, + { + "name": "created_at!", + "ordinal": 10, + "type_info": "Integer", + "origin": { + "Table": { + "table": "things", + "name": "created_at" + } + } + }, + { + "name": "indexed_at!", + "ordinal": 11, + "type_info": "Integer", + "origin": { + "Table": { + "table": "things", + "name": "indexed_at" + } + } + }, + { + "name": "record_json", + "ordinal": 12, + "type_info": "Text", + "origin": { + "Table": { + "table": "things", + "name": "record_json" + } + } + }, + { + "name": "like_count!: i64", + "ordinal": 13, + "type_info": "Integer", + "origin": "Expression" + }, + { + "name": "save_count!: i64", + "ordinal": 14, + "type_info": "Integer", + "origin": "Expression" + }, + { + "name": "model_count!: i64", + "ordinal": 15, + "type_info": "Integer", + "origin": "Expression" + }, + { + "name": "part_count!: i64", + "ordinal": 16, + "type_info": "Integer", + "origin": "Expression" + } + ], + "parameters": { + "Right": 2 + }, + "nullable": [ + false, + false, + false, + false, + false, + true, + false, + true, + true, + true, + false, + false, + true, + false, + false, + false, + false + ] + }, + "hash": "5dd94b27dbcc35c0536394c6158cdbc5fb28436947566806bf8f49b19fc6db85" +} diff --git a/.sqlx/query-7141c9ae0cf701b4e9a5b7791b7583737f2d4f66b8f45c4c83ee9e946ecbd2ec.json b/.sqlx/query-9b2ec255fedbb5d38f44bd5b18fb1ed9cd6eed8a7402347afd8a23a73452008b.json similarity index 92% rename from .sqlx/query-7141c9ae0cf701b4e9a5b7791b7583737f2d4f66b8f45c4c83ee9e946ecbd2ec.json rename to .sqlx/query-9b2ec255fedbb5d38f44bd5b18fb1ed9cd6eed8a7402347afd8a23a73452008b.json index 73d7ea2..7496c23 100644 --- a/.sqlx/query-7141c9ae0cf701b4e9a5b7791b7583737f2d4f66b8f45c4c83ee9e946ecbd2ec.json +++ b/.sqlx/query-9b2ec255fedbb5d38f44bd5b18fb1ed9cd6eed8a7402347afd8a23a73452008b.json @@ -1,6 +1,6 @@ { "db_name": "SQLite", - "query": "SELECT t.did AS \"did!\", t.rkey AS \"rkey!\", t.uri AS \"uri!\", t.cid AS \"cid!\",\n t.name AS \"name!\", t.summary, t.license AS \"license!\", t.tags_json, t.cover_json,\n t.derived_from_uri, t.created_at AS \"created_at!\", t.indexed_at AS \"indexed_at!\",\n t.record_json,\n COALESCE(cs.like_count, 0) AS \"like_count!: i64\",\n COALESCE(cs.save_count, 0) AS \"save_count!: i64\",\n (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS \"model_count!: i64\",\n (SELECT COUNT(*) FROM thing_models tm JOIN model_parts mp ON mp.model_uri = tm.model_uri WHERE tm.thing_uri = t.uri) AS \"part_count!: i64\"\n FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri\n ORDER BY t.created_at DESC LIMIT 500", + "query": "SELECT t.did AS \"did!\", t.rkey AS \"rkey!\", t.uri AS \"uri!\", t.cid AS \"cid!\",\n t.name AS \"name!\", t.summary, t.license AS \"license!\", t.tags_json, t.cover_json,\n t.derived_from_uri, t.created_at AS \"created_at!\", t.indexed_at AS \"indexed_at!\",\n t.record_json,\n COALESCE(cs.like_count, 0) AS \"like_count!: i64\",\n COALESCE(cs.save_count, 0) AS \"save_count!: i64\",\n (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS \"model_count!: i64\",\n (SELECT COUNT(*) FROM thing_models tm JOIN model_parts mp ON mp.model_uri = tm.model_uri WHERE tm.thing_uri = t.uri) AS \"part_count!: i64\"\n FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri\n WHERE COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) = 0\n AND (? IS NULL OR t.rkey < ? OR (t.rkey = ? AND t.uri < ?))\n ORDER BY t.rkey DESC, t.uri DESC\n LIMIT ?", "describe": { "columns": [ { @@ -172,7 +172,7 @@ } ], "parameters": { - "Right": 0 + "Right": 5 }, "nullable": [ false, @@ -194,5 +194,5 @@ false ] }, - "hash": "7141c9ae0cf701b4e9a5b7791b7583737f2d4f66b8f45c4c83ee9e946ecbd2ec" + "hash": "9b2ec255fedbb5d38f44bd5b18fb1ed9cd6eed8a7402347afd8a23a73452008b" } diff --git a/AGENTS.md b/AGENTS.md index 137091d..754ee2e 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -102,3 +102,4 @@ The server keeps the browser single-origin: `space.polymodel.*` reads are served ## Code style - **Prefer branded types over bare `String`/`&str`.** Any use of a bare `String` or `&str` where a wrapping, possibly-validated branded type exists (`Did`, `Handle`, `AtUri`, `Nsid`, `Cid`, `AtIdentifier`, …) is a bug to fix immediately unless clearly rebutted — e.g. binding to SQLite (TEXT columns), human-facing error/log messages, or where the surrounding error structure already carries the needed context. Thread the brand from the request/parse site down to the system boundary (DB bind, serialization) and convert with `.as_ref()`/`.as_str()` only there. Don't reconstruct a brand you already have (e.g. don't `Did::new_owned(s)` from a string you got by stringifying a `Did`), and prefer matching on a branded enum's variants over sniffing its serialized form (`ident.starts_with("did:")`). +- Production SQL must use SQLx compile-time checked macros (`query!`, `query_as!`) rather than runtime `query`/`query_as` escape hatches. When changing production SQL or migrations, regenerate and commit the `.sqlx` cache with `just sqlx-prepare`. diff --git a/src/appview/tests.rs b/src/appview/tests.rs index 8879437..a0a6e34 100644 --- a/src/appview/tests.rs +++ b/src/appview/tests.rs @@ -15,6 +15,7 @@ use polymodel_api::space_polymodel::actor::ProfileViewUnion; use serde_json::json; use sqlx::SqlitePool; use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; +use std::collections::HashSet; use std::str::FromStr; use tower::ServiceExt; @@ -128,6 +129,15 @@ async fn seed_thing( uri } +async fn set_thing_created_at(pool: &SqlitePool, uri: &str, created_at: i64) { + sqlx::query("UPDATE things SET created_at = ? WHERE uri = ?") + .bind(created_at) + .bind(uri) + .execute(pool) + .await + .unwrap(); +} + async fn seed_model( pool: &SqlitePool, did: &str, @@ -366,6 +376,162 @@ async fn feed_hot_is_deterministic_with_tied_scores() { assert_eq!(feed.items[0].thing.name.as_str(), "b"); } +#[tokio::test] +async fn feed_hot_cursor_pins_ranking_time_and_keyset_boundary() { + let state = state().await; + seed_identity(&state.pool, DID_A, "alice.com").await; + let now_a = 1_700_000_000_000_000_000; + let now_b = now_a + 90 * 86_400_000_000_000; + let rows = [ + ("3aaaaaaaaaaaa", "fresh-low", 2, 1), + ("3bbbbbbbbbbbb", "medium-high", 7, 12), + ("3cccccccccccc", "older-high", 48, 100), + ("3dddddddddddd", "newer-medium", 3, 4), + ("3eeeeeeeeeeee", "old-low", 120, 2), + ]; + for (rkey, name, age_hours, likes) in rows { + let uri = seed_thing(&state.pool, DID_A, rkey, name, &[], likes).await; + set_thing_created_at(&state.pool, &uri, now_a - age_hours * 3_600_000_000_000).await; + } + + let expected = views::feed_hot_at(&state, 10, None, None, now_a) + .await + .unwrap() + .items + .into_iter() + .map(|item| item.thing.uri) + .collect::>(); + assert_eq!(expected.len(), rows.len()); + + let page1 = views::feed_hot_at(&state, 2, None, None, now_a) + .await + .unwrap(); + let cursor = page1 + .cursor + .as_ref() + .expect("first hot page should have a cursor"); + assert!( + cursor.as_str().starts_with("hot:v1:"), + "hot cursor must be versioned and opaque" + ); + + let page2 = views::feed_hot_at(&state, 10, Some(cursor.as_str()), None, now_b) + .await + .unwrap(); + let combined = page1 + .items + .into_iter() + .chain(page2.items) + .map(|item| item.thing.uri) + .collect::>(); + let unique = combined.iter().collect::>(); + assert_eq!( + unique.len(), + combined.len(), + "hot pagination duplicated an item" + ); + assert_eq!( + combined, expected, + "hot pagination must keep the cursor-pinned order" + ); +} + +#[tokio::test] +async fn feed_hot_cursor_boundary_handles_same_rkey_across_repos() { + let state = state().await; + seed_identity(&state.pool, DID_A, "alice.com").await; + seed_identity(&state.pool, DID_B, "bob.com").await; + let now = 1_700_000_000_000_000_000; + let shared_rkey = "3sameeeeeeeee"; + let uri_a = seed_thing(&state.pool, DID_A, shared_rkey, "alice", &[], 5).await; + let uri_b = seed_thing(&state.pool, DID_B, shared_rkey, "bob", &[], 5).await; + set_thing_created_at(&state.pool, &uri_a, now - 86_400_000_000_000).await; + set_thing_created_at(&state.pool, &uri_b, now - 86_400_000_000_000).await; + + let page1 = views::feed_hot_at(&state, 1, None, None, now) + .await + .unwrap(); + let cursor = page1.cursor.as_deref().expect("first row has another page"); + let page2 = views::feed_hot_at(&state, 1, Some(cursor), None, now) + .await + .unwrap(); + + let uris = page1 + .items + .into_iter() + .chain(page2.items) + .map(|item| item.thing.uri) + .collect::>(); + let unique = uris.iter().collect::>(); + assert_eq!(unique.len(), 2, "same-rkey tied rows must not duplicate"); + assert!(uris.iter().any(|uri| uri.as_ref() == uri_a)); + assert!(uris.iter().any(|uri| uri.as_ref() == uri_b)); +} + +#[tokio::test] +async fn feed_hot_zero_engagement_fallback_paginates_after_zero_boundary() { + let state = state().await; + seed_identity(&state.pool, DID_A, "alice.com").await; + let now = 1_700_000_000_000_000_000; + for rkey in [ + "3dddddddddddd", + "3cccccccccccc", + "3bbbbbbbbbbbb", + "3aaaaaaaaaaaa", + ] { + let uri = seed_thing(&state.pool, DID_A, rkey, rkey, &[], 0).await; + set_thing_created_at(&state.pool, &uri, now - 86_400_000_000_000).await; + } + + let page1 = views::feed_hot_at(&state, 2, None, None, now) + .await + .unwrap(); + assert_eq!(page1.items.len(), 2); + assert_eq!(page1.items[0].thing.name.as_str(), "3dddddddddddd"); + assert_eq!(page1.items[1].thing.name.as_str(), "3cccccccccccc"); + let cursor = page1 + .cursor + .as_deref() + .expect("zero-score first page has cursor"); + + let page2 = views::feed_hot_at(&state, 2, Some(cursor), None, now) + .await + .unwrap(); + assert_eq!(page2.items.len(), 2); + assert_eq!(page2.items[0].thing.name.as_str(), "3bbbbbbbbbbbb"); + assert_eq!(page2.items[1].thing.name.as_str(), "3aaaaaaaaaaaa"); + assert!(page2.cursor.is_none()); +} + +#[tokio::test] +async fn feed_hot_ranks_older_high_engagement_beyond_previous_candidate_cap() { + let state = state().await; + seed_identity(&state.pool, DID_A, "alice.com").await; + let now = 1_700_000_000_000_000_000; + let older = seed_thing( + &state.pool, + DID_A, + "3aaaaaaaaaaaa", + "older-popular", + &[], + 1_000_000, + ) + .await; + set_thing_created_at(&state.pool, &older, now - 30 * 86_400_000_000_000).await; + + for i in 0..501 { + let rkey = format!("3new{i:010}"); + let uri = seed_thing(&state.pool, DID_A, &rkey, &format!("new-{i}"), &[], 0).await; + set_thing_created_at(&state.pool, &uri, now - 3_600_000_000_000 + i).await; + } + + let feed = views::feed_hot_at(&state, 5, None, None, now) + .await + .unwrap(); + assert_eq!(feed.items[0].thing.uri.as_ref(), older); + assert_eq!(feed.items[0].thing.name.as_str(), "older-popular"); +} + #[tokio::test] async fn search_things_uses_fts() { let state = state().await; @@ -383,10 +549,14 @@ async fn search_things_uses_fts() { #[tokio::test] async fn malformed_ranked_cursor_is_rejected() { let state = state().await; - let err = views::search_things(&state, "x", 10, Some("not-a-number"), None) + let search_err = views::search_things(&state, "x", 10, Some("not-a-number"), None) .await .unwrap_err(); - assert!(matches!(err, AppError::InvalidRequest(_))); + assert!(matches!(search_err, AppError::InvalidRequest(_))); + let hot_err = views::feed_hot(&state, 10, Some("not-a-hot-cursor"), None) + .await + .unwrap_err(); + assert!(matches!(hot_err, AppError::InvalidRequest(_))); } #[tokio::test] @@ -558,9 +728,6 @@ async fn negative_offset_cursor_is_rejected() { views::search_things(&state, "x", 10, Some("-1"), None) .await .unwrap_err(), - views::feed_hot(&state, 10, Some("-5"), None) - .await - .unwrap_err(), views::list_things( &state, "at://did:plc:l/space.polymodel.graph.list/self", diff --git a/src/appview/views.rs b/src/appview/views.rs index 8117dd1..72dabba 100644 --- a/src/appview/views.rs +++ b/src/appview/views.rs @@ -12,9 +12,11 @@ //! on the record's TID `rkey` (descending); the returned cursor is the last //! item's `rkey`, and the next page uses `rkey < cursor` (a NULL cursor — no //! cursor — returns the first page). -//! - Ranked views (`getFeed?algorithm=hot`, `searchThings`) use an offset cursor -//! (stringified); a non-numeric cursor is rejected with 400. Ranking is -//! deterministic: hot scores tie-break on `rkey` desc; search is BM25. +//! - `getFeed?algorithm=hot` uses an opaque versioned keyset cursor over a hot +//! ranking snapshot. The cursor pins the ranking timestamp and last returned +//! `(score, rkey, uri)` boundary so later pages use the same total ordering. +//! - `searchThings` uses a stringified offset cursor for BM25 results; a +//! non-numeric cursor is rejected with 400. use std::cmp::Ordering; @@ -33,10 +35,12 @@ use polymodel_api::space_polymodel::library::{ }; use sqlx::SqlitePool; -use super::error::{AppResult, db, internal, invalid_request, not_found}; +use super::error::{AppError, AppResult, db, internal, invalid_request, not_found}; use super::state::AppState; use crate::oauth::SqliteAuthStore; +const HOT_POSITIVE_BATCH: i64 = 256; + // --------------------------------------------------------------------------- // value-construction helpers // --------------------------------------------------------------------------- @@ -669,10 +673,119 @@ pub(super) async fn feed_hot( cursor: Option<&str>, viewer_did: Option>, ) -> AppResult { - let offset = parse_offset(cursor)?; - // Candidate window: the most recent things. Hot score is computed in Rust so - // the formula (engagement / (age_hours + 2)^1.5) does not depend on SQLite - // math functions, and ranking is deterministic for a fixed `now`. + let default_now_nanos = chrono::Utc::now() + .timestamp_nanos_opt() + .expect("current timestamp must fit in i64 nanoseconds"); + feed_hot_at(state, limit, cursor, viewer_did, default_now_nanos).await +} + +pub(super) async fn feed_hot_at( + state: &AppState, + limit: i64, + cursor: Option<&str>, + viewer_did: Option>, + default_now_nanos: i64, +) -> AppResult { + let page_state = parse_hot_cursor(cursor, default_now_nanos)?; + let now = page_state.now_nanos; + let mut scored = hot_positive_candidates(state, limit, &page_state).await?; + + if scored.len() as i64 <= limit { + scored.extend( + hot_zero_engagement_candidates(state, limit + 1 - scored.len() as i64, &page_state) + .await?, + ); + } + + // score desc, then rkey desc and uri desc. rkeys are only per-repo unique, so + // uri completes the total order used by the page boundary. + sort_hot_scores(&mut scored); + scored.truncate(limit as usize + 1); + let cursor = if scored.len() as i64 > limit { + scored + .get((limit - 1) as usize) + .map(|(score, row)| encode_hot_cursor(now, *score, &row.rkey, &row.uri)) + } else { + None + }; + scored.truncate(limit as usize); + let page = scored.into_iter().map(|(_, row)| row).collect(); + feed_view(state, page, cursor, viewer_did).await +} + +async fn hot_positive_candidates( + state: &AppState, + limit: i64, + page_state: &HotPageState, +) -> AppResult> { + if !page_state.may_include_positive_scores() { + return Ok(Vec::new()); + } + + let mut scored = Vec::new(); + let mut offset = 0; + loop { + // PM-51 keeps ranking exact without a fixed recent window, but this is + // still a per-request candidate scan. PM-56 tracks replacing it with a + // materialized/indexed hot-feed ranking table. + let rows = db( + sqlx::query_as!( + ThingRow, + r#"SELECT t.did AS "did!", t.rkey AS "rkey!", t.uri AS "uri!", t.cid AS "cid!", + t.name AS "name!", t.summary, t.license AS "license!", t.tags_json, t.cover_json, + t.derived_from_uri, t.created_at AS "created_at!", t.indexed_at AS "indexed_at!", + t.record_json, + COALESCE(cs.like_count, 0) AS "like_count!: i64", + COALESCE(cs.save_count, 0) AS "save_count!: i64", + (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS "model_count!: i64", + (SELECT COUNT(*) FROM thing_models tm JOIN model_parts mp ON mp.model_uri = tm.model_uri WHERE tm.thing_uri = t.uri) AS "part_count!: i64" + FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri + WHERE COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) > 0 + ORDER BY COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) DESC, + t.rkey DESC, t.uri DESC + LIMIT ? OFFSET ?"#, + HOT_POSITIVE_BATCH, + offset, + ) + .fetch_all(&state.pool) + .await, + )?; + if rows.is_empty() { + break; + } + + let batch_len = rows.len() as i64; + let last_engagement = rows + .last() + .map(|row| row.like_count + row.save_count) + .unwrap_or(0); + for row in rows { + let score = hot_score(&row, page_state.now_nanos); + if page_state.is_after_boundary(score, &row.rkey, &row.uri) { + scored.push((score, row)); + } + } + sort_hot_scores(&mut scored); + scored.truncate(limit as usize + 1); + + if batch_len < HOT_POSITIVE_BATCH || can_stop_positive_scan(&scored, limit, last_engagement) + { + break; + } + offset += HOT_POSITIVE_BATCH; + } + Ok(scored) +} + +async fn hot_zero_engagement_candidates( + state: &AppState, + fetch: i64, + page_state: &HotPageState, +) -> AppResult> { + if fetch <= 0 || !page_state.may_include_zero_scores() { + return Ok(Vec::new()); + } + let (after_rkey, after_uri) = page_state.zero_score_boundary(); let rows = db( sqlx::query_as!( ThingRow, @@ -685,35 +798,135 @@ pub(super) async fn feed_hot( (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS "model_count!: i64", (SELECT COUNT(*) FROM thing_models tm JOIN model_parts mp ON mp.model_uri = tm.model_uri WHERE tm.thing_uri = t.uri) AS "part_count!: i64" FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri - ORDER BY t.created_at DESC LIMIT 500"#, + WHERE COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) = 0 + AND (? IS NULL OR t.rkey < ? OR (t.rkey = ? AND t.uri < ?)) + ORDER BY t.rkey DESC, t.uri DESC + LIMIT ?"#, + after_rkey, + after_rkey, + after_rkey, + after_uri, + fetch, ) .fetch_all(&state.pool) .await, )?; - let now = chrono::Utc::now() - .timestamp_nanos_opt() - .expect("current timestamp must fit in i64 nanoseconds"); - let mut scored: Vec<(f64, ThingRow)> = - rows.into_iter().map(|r| (hot_score(&r, now), r)).collect(); - // score desc, tie-break rkey desc (deterministic, reproducible page boundary). - scored.sort_by(|a, b| { - b.0.partial_cmp(&a.0) - .unwrap_or(Ordering::Equal) - .then_with(|| b.1.rkey.cmp(&a.1.rkey)) - }); - let page: Vec = scored - .into_iter() - .skip(offset as usize) - .take(limit as usize + 1) - .map(|(_, r)| r) - .collect(); - let cursor = if page.len() as i64 > limit { - Some((offset + limit).to_string()) - } else { - None + Ok(rows.into_iter().map(|row| (0.0, row)).collect()) +} + +fn sort_hot_scores(scored: &mut [(f64, ThingRow)]) { + scored.sort_by(hot_score_order); +} + +fn hot_score_order(a: &(f64, ThingRow), b: &(f64, ThingRow)) -> Ordering { + b.0.partial_cmp(&a.0) + .unwrap_or(Ordering::Equal) + .then_with(|| b.1.rkey.cmp(&a.1.rkey)) + .then_with(|| b.1.uri.cmp(&a.1.uri)) +} + +fn can_stop_positive_scan(scored: &[(f64, ThingRow)], limit: i64, last_engagement: i64) -> bool { + let Some((worst_score, _)) = scored.get(limit as usize) else { + return false; }; - let page = page.into_iter().take(limit as usize).collect(); - feed_view(state, page, cursor, viewer_did).await + let max_unseen_score = (last_engagement as f64) / 2.0_f64.powf(1.5); + max_unseen_score < *worst_score +} + +struct HotPageState { + now_nanos: i64, + boundary: Option, +} + +struct HotCursorBoundary { + score: f64, + rkey: String, + uri: String, +} + +impl HotPageState { + fn may_include_positive_scores(&self) -> bool { + self.boundary + .as_ref() + .is_none_or(|boundary| boundary.score > 0.0) + } + + fn may_include_zero_scores(&self) -> bool { + self.boundary + .as_ref() + .is_none_or(|boundary| boundary.score >= 0.0) + } + + fn zero_score_boundary(&self) -> (Option<&str>, Option<&str>) { + match &self.boundary { + Some(boundary) if boundary.score == 0.0 => { + (Some(boundary.rkey.as_str()), Some(boundary.uri.as_str())) + } + _ => (None, None), + } + } + + fn is_after_boundary(&self, score: f64, rkey: &str, uri: &str) -> bool { + let Some(boundary) = &self.boundary else { + return true; + }; + score < boundary.score + || (score == boundary.score + && (rkey < boundary.rkey.as_str() + || (rkey == boundary.rkey.as_str() && uri < boundary.uri.as_str()))) + } +} + +fn parse_hot_cursor(cursor: Option<&str>, default_now_nanos: i64) -> AppResult { + let Some(cursor) = cursor else { + return Ok(HotPageState { + now_nanos: default_now_nanos, + boundary: None, + }); + }; + + // The hot cursor pins the ranking timestamp. Without that snapshot value, + // time-dependent hot scores can shift between requests and make a page + // boundary skip or duplicate rows. + let mut parts = cursor.splitn(6, ':'); + let valid_prefix = parts.next() == Some("hot") && parts.next() == Some("v1"); + let Some(now_nanos) = parts.next() else { + return Err(invalid_hot_cursor()); + }; + let Some(score_bits_hex) = parts.next() else { + return Err(invalid_hot_cursor()); + }; + let Some(rkey) = parts.next() else { + return Err(invalid_hot_cursor()); + }; + let Some(uri) = parts.next() else { + return Err(invalid_hot_cursor()); + }; + if !valid_prefix || rkey.is_empty() || uri.is_empty() { + return Err(invalid_hot_cursor()); + } + let now_nanos = now_nanos.parse().map_err(|_| invalid_hot_cursor())?; + let score_bits = u64::from_str_radix(score_bits_hex, 16).map_err(|_| invalid_hot_cursor())?; + let score = f64::from_bits(score_bits); + if !(score.is_finite() && score >= 0.0) { + return Err(invalid_hot_cursor()); + } + Ok(HotPageState { + now_nanos, + boundary: Some(HotCursorBoundary { + score, + rkey: rkey.to_owned(), + uri: uri.to_owned(), + }), + }) +} + +fn encode_hot_cursor(now_nanos: i64, score: f64, rkey: &str, uri: &str) -> String { + format!("hot:v1:{now_nanos}:{:016x}:{rkey}:{uri}", score.to_bits()) +} + +fn invalid_hot_cursor() -> AppError { + invalid_request("cursor must be an opaque hot feed cursor") } fn hot_score(row: &ThingRow, now_nanos: i64) -> f64 { @@ -722,6 +935,56 @@ fn hot_score(row: &ThingRow, now_nanos: i64) -> f64 { engagement / (age_hours + 2.0).powf(1.5) } +#[cfg(test)] +mod hot_tests { + use super::*; + + fn row(rkey: &str, uri: &str) -> ThingRow { + ThingRow { + did: "did:plc:test".to_owned(), + rkey: rkey.to_owned(), + uri: uri.to_owned(), + cid: "cid".to_owned(), + name: "test".to_owned(), + summary: None, + license: "CC-BY-4.0".to_owned(), + tags_json: None, + cover_json: None, + derived_from_uri: None, + created_at: 0, + indexed_at: 0, + record_json: None, + like_count: 0, + save_count: 0, + model_count: 0, + part_count: 0, + } + } + + #[test] + fn positive_scan_bound_stops_only_when_unseen_scores_cannot_enter_page() { + let page = vec![ + (10.0, row("3c", "at://did:plc:c/thing/3c")), + (9.0, row("3b", "at://did:plc:b/thing/3b")), + ]; + assert!(can_stop_positive_scan(&page, 1, 1)); + assert!(!can_stop_positive_scan(&page, 1, 100)); + assert!(!can_stop_positive_scan(&page[..1], 1, 1)); + } + + #[test] + fn rejects_malformed_hot_cursor_score_bits() { + let bits = |s: f64| format!("{:016x}", s.to_bits()); + let cursor = |score_bits: String| format!("hot:v1:100:{score_bits}:rkey:at://did:plc:x/x"); + assert!(parse_hot_cursor(Some(&cursor(bits(f64::NAN))), 0).is_err()); + assert!(parse_hot_cursor(Some(&cursor(bits(f64::INFINITY))), 0).is_err()); + assert!(parse_hot_cursor(Some(&cursor(bits(f64::NEG_INFINITY))), 0).is_err()); + assert!(parse_hot_cursor(Some(&cursor(bits(-1.0))), 0).is_err()); + assert!(parse_hot_cursor(Some(&cursor(bits(0.0))), 0).is_ok()); + assert!(parse_hot_cursor(Some(&cursor(bits(12.5))), 0).is_ok()); + } +} + pub(super) async fn search_things( state: &AppState, q: &str,