atproto Thingiverse but good
10

Configure Feed

Select the types of activity you want to include in your feed.

PM-51: hot feed pagination

Orual (Jun 29, 2026, 7:06 PM EDT) 83f4ac28 137b59a8

+669 -40
+198
.sqlx/query-5dd94b27dbcc35c0536394c6158cdbc5fb28436947566806bf8f49b19fc6db85.json
··· 1 + { 2 + "db_name": "SQLite", 3 + "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 ?", 4 + "describe": { 5 + "columns": [ 6 + { 7 + "name": "did!", 8 + "ordinal": 0, 9 + "type_info": "Text", 10 + "origin": { 11 + "Table": { 12 + "table": "things", 13 + "name": "did" 14 + } 15 + } 16 + }, 17 + { 18 + "name": "rkey!", 19 + "ordinal": 1, 20 + "type_info": "Text", 21 + "origin": { 22 + "Table": { 23 + "table": "things", 24 + "name": "rkey" 25 + } 26 + } 27 + }, 28 + { 29 + "name": "uri!", 30 + "ordinal": 2, 31 + "type_info": "Text", 32 + "origin": { 33 + "Table": { 34 + "table": "things", 35 + "name": "uri" 36 + } 37 + } 38 + }, 39 + { 40 + "name": "cid!", 41 + "ordinal": 3, 42 + "type_info": "Text", 43 + "origin": { 44 + "Table": { 45 + "table": "things", 46 + "name": "cid" 47 + } 48 + } 49 + }, 50 + { 51 + "name": "name!", 52 + "ordinal": 4, 53 + "type_info": "Text", 54 + "origin": { 55 + "Table": { 56 + "table": "things", 57 + "name": "name" 58 + } 59 + } 60 + }, 61 + { 62 + "name": "summary", 63 + "ordinal": 5, 64 + "type_info": "Text", 65 + "origin": { 66 + "Table": { 67 + "table": "things", 68 + "name": "summary" 69 + } 70 + } 71 + }, 72 + { 73 + "name": "license!", 74 + "ordinal": 6, 75 + "type_info": "Text", 76 + "origin": { 77 + "Table": { 78 + "table": "things", 79 + "name": "license" 80 + } 81 + } 82 + }, 83 + { 84 + "name": "tags_json", 85 + "ordinal": 7, 86 + "type_info": "Text", 87 + "origin": { 88 + "Table": { 89 + "table": "things", 90 + "name": "tags_json" 91 + } 92 + } 93 + }, 94 + { 95 + "name": "cover_json", 96 + "ordinal": 8, 97 + "type_info": "Text", 98 + "origin": { 99 + "Table": { 100 + "table": "things", 101 + "name": "cover_json" 102 + } 103 + } 104 + }, 105 + { 106 + "name": "derived_from_uri", 107 + "ordinal": 9, 108 + "type_info": "Text", 109 + "origin": { 110 + "Table": { 111 + "table": "things", 112 + "name": "derived_from_uri" 113 + } 114 + } 115 + }, 116 + { 117 + "name": "created_at!", 118 + "ordinal": 10, 119 + "type_info": "Integer", 120 + "origin": { 121 + "Table": { 122 + "table": "things", 123 + "name": "created_at" 124 + } 125 + } 126 + }, 127 + { 128 + "name": "indexed_at!", 129 + "ordinal": 11, 130 + "type_info": "Integer", 131 + "origin": { 132 + "Table": { 133 + "table": "things", 134 + "name": "indexed_at" 135 + } 136 + } 137 + }, 138 + { 139 + "name": "record_json", 140 + "ordinal": 12, 141 + "type_info": "Text", 142 + "origin": { 143 + "Table": { 144 + "table": "things", 145 + "name": "record_json" 146 + } 147 + } 148 + }, 149 + { 150 + "name": "like_count!: i64", 151 + "ordinal": 13, 152 + "type_info": "Integer", 153 + "origin": "Expression" 154 + }, 155 + { 156 + "name": "save_count!: i64", 157 + "ordinal": 14, 158 + "type_info": "Integer", 159 + "origin": "Expression" 160 + }, 161 + { 162 + "name": "model_count!: i64", 163 + "ordinal": 15, 164 + "type_info": "Integer", 165 + "origin": "Expression" 166 + }, 167 + { 168 + "name": "part_count!: i64", 169 + "ordinal": 16, 170 + "type_info": "Integer", 171 + "origin": "Expression" 172 + } 173 + ], 174 + "parameters": { 175 + "Right": 2 176 + }, 177 + "nullable": [ 178 + false, 179 + false, 180 + false, 181 + false, 182 + false, 183 + true, 184 + false, 185 + true, 186 + true, 187 + true, 188 + false, 189 + false, 190 + true, 191 + false, 192 + false, 193 + false, 194 + false 195 + ] 196 + }, 197 + "hash": "5dd94b27dbcc35c0536394c6158cdbc5fb28436947566806bf8f49b19fc6db85" 198 + }
+3 -3
.sqlx/query-7141c9ae0cf701b4e9a5b7791b7583737f2d4f66b8f45c4c83ee9e946ecbd2ec.json .sqlx/query-9b2ec255fedbb5d38f44bd5b18fb1ed9cd6eed8a7402347afd8a23a73452008b.json
··· 1 1 { 2 2 "db_name": "SQLite", 3 - "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", 3 + "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 ?", 4 4 "describe": { 5 5 "columns": [ 6 6 { ··· 172 172 } 173 173 ], 174 174 "parameters": { 175 - "Right": 0 175 + "Right": 5 176 176 }, 177 177 "nullable": [ 178 178 false, ··· 194 194 false 195 195 ] 196 196 }, 197 - "hash": "7141c9ae0cf701b4e9a5b7791b7583737f2d4f66b8f45c4c83ee9e946ecbd2ec" 197 + "hash": "9b2ec255fedbb5d38f44bd5b18fb1ed9cd6eed8a7402347afd8a23a73452008b" 198 198 }
+1
AGENTS.md
··· 102 102 ## Code style 103 103 104 104 - **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:")`). 105 + - 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`.
+172 -5
src/appview/tests.rs
··· 15 15 use serde_json::json; 16 16 use sqlx::SqlitePool; 17 17 use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; 18 + use std::collections::HashSet; 18 19 use std::str::FromStr; 19 20 use tower::ServiceExt; 20 21 ··· 126 127 .await 127 128 .unwrap(); 128 129 uri 130 + } 131 + 132 + async fn set_thing_created_at(pool: &SqlitePool, uri: &str, created_at: i64) { 133 + sqlx::query("UPDATE things SET created_at = ? WHERE uri = ?") 134 + .bind(created_at) 135 + .bind(uri) 136 + .execute(pool) 137 + .await 138 + .unwrap(); 129 139 } 130 140 131 141 async fn seed_model( ··· 367 377 } 368 378 369 379 #[tokio::test] 380 + async fn feed_hot_cursor_pins_ranking_time_and_keyset_boundary() { 381 + let state = state().await; 382 + seed_identity(&state.pool, DID_A, "alice.com").await; 383 + let now_a = 1_700_000_000_000_000_000; 384 + let now_b = now_a + 90 * 86_400_000_000_000; 385 + let rows = [ 386 + ("3aaaaaaaaaaaa", "fresh-low", 2, 1), 387 + ("3bbbbbbbbbbbb", "medium-high", 7, 12), 388 + ("3cccccccccccc", "older-high", 48, 100), 389 + ("3dddddddddddd", "newer-medium", 3, 4), 390 + ("3eeeeeeeeeeee", "old-low", 120, 2), 391 + ]; 392 + for (rkey, name, age_hours, likes) in rows { 393 + let uri = seed_thing(&state.pool, DID_A, rkey, name, &[], likes).await; 394 + set_thing_created_at(&state.pool, &uri, now_a - age_hours * 3_600_000_000_000).await; 395 + } 396 + 397 + let expected = views::feed_hot_at(&state, 10, None, None, now_a) 398 + .await 399 + .unwrap() 400 + .items 401 + .into_iter() 402 + .map(|item| item.thing.uri) 403 + .collect::<Vec<_>>(); 404 + assert_eq!(expected.len(), rows.len()); 405 + 406 + let page1 = views::feed_hot_at(&state, 2, None, None, now_a) 407 + .await 408 + .unwrap(); 409 + let cursor = page1 410 + .cursor 411 + .as_ref() 412 + .expect("first hot page should have a cursor"); 413 + assert!( 414 + cursor.as_str().starts_with("hot:v1:"), 415 + "hot cursor must be versioned and opaque" 416 + ); 417 + 418 + let page2 = views::feed_hot_at(&state, 10, Some(cursor.as_str()), None, now_b) 419 + .await 420 + .unwrap(); 421 + let combined = page1 422 + .items 423 + .into_iter() 424 + .chain(page2.items) 425 + .map(|item| item.thing.uri) 426 + .collect::<Vec<_>>(); 427 + let unique = combined.iter().collect::<HashSet<_>>(); 428 + assert_eq!( 429 + unique.len(), 430 + combined.len(), 431 + "hot pagination duplicated an item" 432 + ); 433 + assert_eq!( 434 + combined, expected, 435 + "hot pagination must keep the cursor-pinned order" 436 + ); 437 + } 438 + 439 + #[tokio::test] 440 + async fn feed_hot_cursor_boundary_handles_same_rkey_across_repos() { 441 + let state = state().await; 442 + seed_identity(&state.pool, DID_A, "alice.com").await; 443 + seed_identity(&state.pool, DID_B, "bob.com").await; 444 + let now = 1_700_000_000_000_000_000; 445 + let shared_rkey = "3sameeeeeeeee"; 446 + let uri_a = seed_thing(&state.pool, DID_A, shared_rkey, "alice", &[], 5).await; 447 + let uri_b = seed_thing(&state.pool, DID_B, shared_rkey, "bob", &[], 5).await; 448 + set_thing_created_at(&state.pool, &uri_a, now - 86_400_000_000_000).await; 449 + set_thing_created_at(&state.pool, &uri_b, now - 86_400_000_000_000).await; 450 + 451 + let page1 = views::feed_hot_at(&state, 1, None, None, now) 452 + .await 453 + .unwrap(); 454 + let cursor = page1.cursor.as_deref().expect("first row has another page"); 455 + let page2 = views::feed_hot_at(&state, 1, Some(cursor), None, now) 456 + .await 457 + .unwrap(); 458 + 459 + let uris = page1 460 + .items 461 + .into_iter() 462 + .chain(page2.items) 463 + .map(|item| item.thing.uri) 464 + .collect::<Vec<_>>(); 465 + let unique = uris.iter().collect::<HashSet<_>>(); 466 + assert_eq!(unique.len(), 2, "same-rkey tied rows must not duplicate"); 467 + assert!(uris.iter().any(|uri| uri.as_ref() == uri_a)); 468 + assert!(uris.iter().any(|uri| uri.as_ref() == uri_b)); 469 + } 470 + 471 + #[tokio::test] 472 + async fn feed_hot_zero_engagement_fallback_paginates_after_zero_boundary() { 473 + let state = state().await; 474 + seed_identity(&state.pool, DID_A, "alice.com").await; 475 + let now = 1_700_000_000_000_000_000; 476 + for rkey in [ 477 + "3dddddddddddd", 478 + "3cccccccccccc", 479 + "3bbbbbbbbbbbb", 480 + "3aaaaaaaaaaaa", 481 + ] { 482 + let uri = seed_thing(&state.pool, DID_A, rkey, rkey, &[], 0).await; 483 + set_thing_created_at(&state.pool, &uri, now - 86_400_000_000_000).await; 484 + } 485 + 486 + let page1 = views::feed_hot_at(&state, 2, None, None, now) 487 + .await 488 + .unwrap(); 489 + assert_eq!(page1.items.len(), 2); 490 + assert_eq!(page1.items[0].thing.name.as_str(), "3dddddddddddd"); 491 + assert_eq!(page1.items[1].thing.name.as_str(), "3cccccccccccc"); 492 + let cursor = page1 493 + .cursor 494 + .as_deref() 495 + .expect("zero-score first page has cursor"); 496 + 497 + let page2 = views::feed_hot_at(&state, 2, Some(cursor), None, now) 498 + .await 499 + .unwrap(); 500 + assert_eq!(page2.items.len(), 2); 501 + assert_eq!(page2.items[0].thing.name.as_str(), "3bbbbbbbbbbbb"); 502 + assert_eq!(page2.items[1].thing.name.as_str(), "3aaaaaaaaaaaa"); 503 + assert!(page2.cursor.is_none()); 504 + } 505 + 506 + #[tokio::test] 507 + async fn feed_hot_ranks_older_high_engagement_beyond_previous_candidate_cap() { 508 + let state = state().await; 509 + seed_identity(&state.pool, DID_A, "alice.com").await; 510 + let now = 1_700_000_000_000_000_000; 511 + let older = seed_thing( 512 + &state.pool, 513 + DID_A, 514 + "3aaaaaaaaaaaa", 515 + "older-popular", 516 + &[], 517 + 1_000_000, 518 + ) 519 + .await; 520 + set_thing_created_at(&state.pool, &older, now - 30 * 86_400_000_000_000).await; 521 + 522 + for i in 0..501 { 523 + let rkey = format!("3new{i:010}"); 524 + let uri = seed_thing(&state.pool, DID_A, &rkey, &format!("new-{i}"), &[], 0).await; 525 + set_thing_created_at(&state.pool, &uri, now - 3_600_000_000_000 + i).await; 526 + } 527 + 528 + let feed = views::feed_hot_at(&state, 5, None, None, now) 529 + .await 530 + .unwrap(); 531 + assert_eq!(feed.items[0].thing.uri.as_ref(), older); 532 + assert_eq!(feed.items[0].thing.name.as_str(), "older-popular"); 533 + } 534 + 535 + #[tokio::test] 370 536 async fn search_things_uses_fts() { 371 537 let state = state().await; 372 538 seed_identity(&state.pool, DID_A, "alice.com").await; ··· 383 549 #[tokio::test] 384 550 async fn malformed_ranked_cursor_is_rejected() { 385 551 let state = state().await; 386 - let err = views::search_things(&state, "x", 10, Some("not-a-number"), None) 552 + let search_err = views::search_things(&state, "x", 10, Some("not-a-number"), None) 387 553 .await 388 554 .unwrap_err(); 389 - assert!(matches!(err, AppError::InvalidRequest(_))); 555 + assert!(matches!(search_err, AppError::InvalidRequest(_))); 556 + let hot_err = views::feed_hot(&state, 10, Some("not-a-hot-cursor"), None) 557 + .await 558 + .unwrap_err(); 559 + assert!(matches!(hot_err, AppError::InvalidRequest(_))); 390 560 } 391 561 392 562 #[tokio::test] ··· 556 726 let state = state().await; 557 727 let errs = [ 558 728 views::search_things(&state, "x", 10, Some("-1"), None) 559 - .await 560 - .unwrap_err(), 561 - views::feed_hot(&state, 10, Some("-5"), None) 562 729 .await 563 730 .unwrap_err(), 564 731 views::list_things(
+295 -32
src/appview/views.rs
··· 12 12 //! on the record's TID `rkey` (descending); the returned cursor is the last 13 13 //! item's `rkey`, and the next page uses `rkey < cursor` (a NULL cursor — no 14 14 //! cursor — returns the first page). 15 - //! - Ranked views (`getFeed?algorithm=hot`, `searchThings`) use an offset cursor 16 - //! (stringified); a non-numeric cursor is rejected with 400. Ranking is 17 - //! deterministic: hot scores tie-break on `rkey` desc; search is BM25. 15 + //! - `getFeed?algorithm=hot` uses an opaque versioned keyset cursor over a hot 16 + //! ranking snapshot. The cursor pins the ranking timestamp and last returned 17 + //! `(score, rkey, uri)` boundary so later pages use the same total ordering. 18 + //! - `searchThings` uses a stringified offset cursor for BM25 results; a 19 + //! non-numeric cursor is rejected with 400. 18 20 19 21 use std::cmp::Ordering; 20 22 ··· 33 35 }; 34 36 use sqlx::SqlitePool; 35 37 36 - use super::error::{AppResult, db, internal, invalid_request, not_found}; 38 + use super::error::{AppError, AppResult, db, internal, invalid_request, not_found}; 37 39 use super::state::AppState; 38 40 use crate::oauth::SqliteAuthStore; 41 + 42 + const HOT_POSITIVE_BATCH: i64 = 256; 39 43 40 44 // --------------------------------------------------------------------------- 41 45 // value-construction helpers ··· 669 673 cursor: Option<&str>, 670 674 viewer_did: Option<Did<&str>>, 671 675 ) -> AppResult<FeedView> { 672 - let offset = parse_offset(cursor)?; 673 - // Candidate window: the most recent things. Hot score is computed in Rust so 674 - // the formula (engagement / (age_hours + 2)^1.5) does not depend on SQLite 675 - // math functions, and ranking is deterministic for a fixed `now`. 676 + let default_now_nanos = chrono::Utc::now() 677 + .timestamp_nanos_opt() 678 + .expect("current timestamp must fit in i64 nanoseconds"); 679 + feed_hot_at(state, limit, cursor, viewer_did, default_now_nanos).await 680 + } 681 + 682 + pub(super) async fn feed_hot_at( 683 + state: &AppState, 684 + limit: i64, 685 + cursor: Option<&str>, 686 + viewer_did: Option<Did<&str>>, 687 + default_now_nanos: i64, 688 + ) -> AppResult<FeedView> { 689 + let page_state = parse_hot_cursor(cursor, default_now_nanos)?; 690 + let now = page_state.now_nanos; 691 + let mut scored = hot_positive_candidates(state, limit, &page_state).await?; 692 + 693 + if scored.len() as i64 <= limit { 694 + scored.extend( 695 + hot_zero_engagement_candidates(state, limit + 1 - scored.len() as i64, &page_state) 696 + .await?, 697 + ); 698 + } 699 + 700 + // score desc, then rkey desc and uri desc. rkeys are only per-repo unique, so 701 + // uri completes the total order used by the page boundary. 702 + sort_hot_scores(&mut scored); 703 + scored.truncate(limit as usize + 1); 704 + let cursor = if scored.len() as i64 > limit { 705 + scored 706 + .get((limit - 1) as usize) 707 + .map(|(score, row)| encode_hot_cursor(now, *score, &row.rkey, &row.uri)) 708 + } else { 709 + None 710 + }; 711 + scored.truncate(limit as usize); 712 + let page = scored.into_iter().map(|(_, row)| row).collect(); 713 + feed_view(state, page, cursor, viewer_did).await 714 + } 715 + 716 + async fn hot_positive_candidates( 717 + state: &AppState, 718 + limit: i64, 719 + page_state: &HotPageState, 720 + ) -> AppResult<Vec<(f64, ThingRow)>> { 721 + if !page_state.may_include_positive_scores() { 722 + return Ok(Vec::new()); 723 + } 724 + 725 + let mut scored = Vec::new(); 726 + let mut offset = 0; 727 + loop { 728 + // PM-51 keeps ranking exact without a fixed recent window, but this is 729 + // still a per-request candidate scan. PM-56 tracks replacing it with a 730 + // materialized/indexed hot-feed ranking table. 731 + let rows = db( 732 + sqlx::query_as!( 733 + ThingRow, 734 + r#"SELECT t.did AS "did!", t.rkey AS "rkey!", t.uri AS "uri!", t.cid AS "cid!", 735 + t.name AS "name!", t.summary, t.license AS "license!", t.tags_json, t.cover_json, 736 + t.derived_from_uri, t.created_at AS "created_at!", t.indexed_at AS "indexed_at!", 737 + t.record_json, 738 + COALESCE(cs.like_count, 0) AS "like_count!: i64", 739 + COALESCE(cs.save_count, 0) AS "save_count!: i64", 740 + (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS "model_count!: i64", 741 + (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" 742 + FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri 743 + WHERE COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) > 0 744 + ORDER BY COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) DESC, 745 + t.rkey DESC, t.uri DESC 746 + LIMIT ? OFFSET ?"#, 747 + HOT_POSITIVE_BATCH, 748 + offset, 749 + ) 750 + .fetch_all(&state.pool) 751 + .await, 752 + )?; 753 + if rows.is_empty() { 754 + break; 755 + } 756 + 757 + let batch_len = rows.len() as i64; 758 + let last_engagement = rows 759 + .last() 760 + .map(|row| row.like_count + row.save_count) 761 + .unwrap_or(0); 762 + for row in rows { 763 + let score = hot_score(&row, page_state.now_nanos); 764 + if page_state.is_after_boundary(score, &row.rkey, &row.uri) { 765 + scored.push((score, row)); 766 + } 767 + } 768 + sort_hot_scores(&mut scored); 769 + scored.truncate(limit as usize + 1); 770 + 771 + if batch_len < HOT_POSITIVE_BATCH || can_stop_positive_scan(&scored, limit, last_engagement) 772 + { 773 + break; 774 + } 775 + offset += HOT_POSITIVE_BATCH; 776 + } 777 + Ok(scored) 778 + } 779 + 780 + async fn hot_zero_engagement_candidates( 781 + state: &AppState, 782 + fetch: i64, 783 + page_state: &HotPageState, 784 + ) -> AppResult<Vec<(f64, ThingRow)>> { 785 + if fetch <= 0 || !page_state.may_include_zero_scores() { 786 + return Ok(Vec::new()); 787 + } 788 + let (after_rkey, after_uri) = page_state.zero_score_boundary(); 676 789 let rows = db( 677 790 sqlx::query_as!( 678 791 ThingRow, ··· 685 798 (SELECT COUNT(*) FROM thing_models tm WHERE tm.thing_uri = t.uri) AS "model_count!: i64", 686 799 (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" 687 800 FROM things t LEFT JOIN content_stats cs ON cs.uri = t.uri 688 - ORDER BY t.created_at DESC LIMIT 500"#, 801 + WHERE COALESCE(cs.like_count, 0) + COALESCE(cs.save_count, 0) = 0 802 + AND (? IS NULL OR t.rkey < ? OR (t.rkey = ? AND t.uri < ?)) 803 + ORDER BY t.rkey DESC, t.uri DESC 804 + LIMIT ?"#, 805 + after_rkey, 806 + after_rkey, 807 + after_rkey, 808 + after_uri, 809 + fetch, 689 810 ) 690 811 .fetch_all(&state.pool) 691 812 .await, 692 813 )?; 693 - let now = chrono::Utc::now() 694 - .timestamp_nanos_opt() 695 - .expect("current timestamp must fit in i64 nanoseconds"); 696 - let mut scored: Vec<(f64, ThingRow)> = 697 - rows.into_iter().map(|r| (hot_score(&r, now), r)).collect(); 698 - // score desc, tie-break rkey desc (deterministic, reproducible page boundary). 699 - scored.sort_by(|a, b| { 700 - b.0.partial_cmp(&a.0) 701 - .unwrap_or(Ordering::Equal) 702 - .then_with(|| b.1.rkey.cmp(&a.1.rkey)) 703 - }); 704 - let page: Vec<ThingRow> = scored 705 - .into_iter() 706 - .skip(offset as usize) 707 - .take(limit as usize + 1) 708 - .map(|(_, r)| r) 709 - .collect(); 710 - let cursor = if page.len() as i64 > limit { 711 - Some((offset + limit).to_string()) 712 - } else { 713 - None 814 + Ok(rows.into_iter().map(|row| (0.0, row)).collect()) 815 + } 816 + 817 + fn sort_hot_scores(scored: &mut [(f64, ThingRow)]) { 818 + scored.sort_by(hot_score_order); 819 + } 820 + 821 + fn hot_score_order(a: &(f64, ThingRow), b: &(f64, ThingRow)) -> Ordering { 822 + b.0.partial_cmp(&a.0) 823 + .unwrap_or(Ordering::Equal) 824 + .then_with(|| b.1.rkey.cmp(&a.1.rkey)) 825 + .then_with(|| b.1.uri.cmp(&a.1.uri)) 826 + } 827 + 828 + fn can_stop_positive_scan(scored: &[(f64, ThingRow)], limit: i64, last_engagement: i64) -> bool { 829 + let Some((worst_score, _)) = scored.get(limit as usize) else { 830 + return false; 831 + }; 832 + let max_unseen_score = (last_engagement as f64) / 2.0_f64.powf(1.5); 833 + max_unseen_score < *worst_score 834 + } 835 + 836 + struct HotPageState { 837 + now_nanos: i64, 838 + boundary: Option<HotCursorBoundary>, 839 + } 840 + 841 + struct HotCursorBoundary { 842 + score: f64, 843 + rkey: String, 844 + uri: String, 845 + } 846 + 847 + impl HotPageState { 848 + fn may_include_positive_scores(&self) -> bool { 849 + self.boundary 850 + .as_ref() 851 + .is_none_or(|boundary| boundary.score > 0.0) 852 + } 853 + 854 + fn may_include_zero_scores(&self) -> bool { 855 + self.boundary 856 + .as_ref() 857 + .is_none_or(|boundary| boundary.score >= 0.0) 858 + } 859 + 860 + fn zero_score_boundary(&self) -> (Option<&str>, Option<&str>) { 861 + match &self.boundary { 862 + Some(boundary) if boundary.score == 0.0 => { 863 + (Some(boundary.rkey.as_str()), Some(boundary.uri.as_str())) 864 + } 865 + _ => (None, None), 866 + } 867 + } 868 + 869 + fn is_after_boundary(&self, score: f64, rkey: &str, uri: &str) -> bool { 870 + let Some(boundary) = &self.boundary else { 871 + return true; 872 + }; 873 + score < boundary.score 874 + || (score == boundary.score 875 + && (rkey < boundary.rkey.as_str() 876 + || (rkey == boundary.rkey.as_str() && uri < boundary.uri.as_str()))) 877 + } 878 + } 879 + 880 + fn parse_hot_cursor(cursor: Option<&str>, default_now_nanos: i64) -> AppResult<HotPageState> { 881 + let Some(cursor) = cursor else { 882 + return Ok(HotPageState { 883 + now_nanos: default_now_nanos, 884 + boundary: None, 885 + }); 886 + }; 887 + 888 + // The hot cursor pins the ranking timestamp. Without that snapshot value, 889 + // time-dependent hot scores can shift between requests and make a page 890 + // boundary skip or duplicate rows. 891 + let mut parts = cursor.splitn(6, ':'); 892 + let valid_prefix = parts.next() == Some("hot") && parts.next() == Some("v1"); 893 + let Some(now_nanos) = parts.next() else { 894 + return Err(invalid_hot_cursor()); 895 + }; 896 + let Some(score_bits_hex) = parts.next() else { 897 + return Err(invalid_hot_cursor()); 898 + }; 899 + let Some(rkey) = parts.next() else { 900 + return Err(invalid_hot_cursor()); 901 + }; 902 + let Some(uri) = parts.next() else { 903 + return Err(invalid_hot_cursor()); 714 904 }; 715 - let page = page.into_iter().take(limit as usize).collect(); 716 - feed_view(state, page, cursor, viewer_did).await 905 + if !valid_prefix || rkey.is_empty() || uri.is_empty() { 906 + return Err(invalid_hot_cursor()); 907 + } 908 + let now_nanos = now_nanos.parse().map_err(|_| invalid_hot_cursor())?; 909 + let score_bits = u64::from_str_radix(score_bits_hex, 16).map_err(|_| invalid_hot_cursor())?; 910 + let score = f64::from_bits(score_bits); 911 + if !(score.is_finite() && score >= 0.0) { 912 + return Err(invalid_hot_cursor()); 913 + } 914 + Ok(HotPageState { 915 + now_nanos, 916 + boundary: Some(HotCursorBoundary { 917 + score, 918 + rkey: rkey.to_owned(), 919 + uri: uri.to_owned(), 920 + }), 921 + }) 922 + } 923 + 924 + fn encode_hot_cursor(now_nanos: i64, score: f64, rkey: &str, uri: &str) -> String { 925 + format!("hot:v1:{now_nanos}:{:016x}:{rkey}:{uri}", score.to_bits()) 926 + } 927 + 928 + fn invalid_hot_cursor() -> AppError { 929 + invalid_request("cursor must be an opaque hot feed cursor") 717 930 } 718 931 719 932 fn hot_score(row: &ThingRow, now_nanos: i64) -> f64 { 720 933 let age_hours = ((now_nanos - row.created_at).max(0) as f64) / 3_600_000_000_000.0; 721 934 let engagement = row.like_count as f64 + row.save_count as f64; 722 935 engagement / (age_hours + 2.0).powf(1.5) 936 + } 937 + 938 + #[cfg(test)] 939 + mod hot_tests { 940 + use super::*; 941 + 942 + fn row(rkey: &str, uri: &str) -> ThingRow { 943 + ThingRow { 944 + did: "did:plc:test".to_owned(), 945 + rkey: rkey.to_owned(), 946 + uri: uri.to_owned(), 947 + cid: "cid".to_owned(), 948 + name: "test".to_owned(), 949 + summary: None, 950 + license: "CC-BY-4.0".to_owned(), 951 + tags_json: None, 952 + cover_json: None, 953 + derived_from_uri: None, 954 + created_at: 0, 955 + indexed_at: 0, 956 + record_json: None, 957 + like_count: 0, 958 + save_count: 0, 959 + model_count: 0, 960 + part_count: 0, 961 + } 962 + } 963 + 964 + #[test] 965 + fn positive_scan_bound_stops_only_when_unseen_scores_cannot_enter_page() { 966 + let page = vec![ 967 + (10.0, row("3c", "at://did:plc:c/thing/3c")), 968 + (9.0, row("3b", "at://did:plc:b/thing/3b")), 969 + ]; 970 + assert!(can_stop_positive_scan(&page, 1, 1)); 971 + assert!(!can_stop_positive_scan(&page, 1, 100)); 972 + assert!(!can_stop_positive_scan(&page[..1], 1, 1)); 973 + } 974 + 975 + #[test] 976 + fn rejects_malformed_hot_cursor_score_bits() { 977 + let bits = |s: f64| format!("{:016x}", s.to_bits()); 978 + let cursor = |score_bits: String| format!("hot:v1:100:{score_bits}:rkey:at://did:plc:x/x"); 979 + assert!(parse_hot_cursor(Some(&cursor(bits(f64::NAN))), 0).is_err()); 980 + assert!(parse_hot_cursor(Some(&cursor(bits(f64::INFINITY))), 0).is_err()); 981 + assert!(parse_hot_cursor(Some(&cursor(bits(f64::NEG_INFINITY))), 0).is_err()); 982 + assert!(parse_hot_cursor(Some(&cursor(bits(-1.0))), 0).is_err()); 983 + assert!(parse_hot_cursor(Some(&cursor(bits(0.0))), 0).is_ok()); 984 + assert!(parse_hot_cursor(Some(&cursor(bits(12.5))), 0).is_ok()); 985 + } 723 986 } 724 987 725 988 pub(super) async fn search_things(