diff --git a/apps/aqua/src/repos/actor_profile.rs b/apps/aqua/src/repos/actor_profile.rs index b59e343..c13e10b 100644 --- a/apps/aqua/src/repos/actor_profile.rs +++ b/apps/aqua/src/repos/actor_profile.rs @@ -97,11 +97,8 @@ impl ActorProfileRepo for PgDataSource { SELECT record FROM statii WHERE did = p.did - AND COALESCE( - (record->>'expiry')::timestamptz, - (record->>'time')::timestamptz + INTERVAL '10 minutes' - ) > NOW() - ORDER BY (record->>'time')::timestamptz DESC, indexed_at DESC + AND expires_at > NOW() + ORDER BY status_time DESC, indexed_at DESC LIMIT 1 ) s ON TRUE WHERE (p.did = ANY($1)) diff --git a/apps/aqua/src/repos/stats.rs b/apps/aqua/src/repos/stats.rs index 406c0f0..495a6f2 100644 --- a/apps/aqua/src/repos/stats.rs +++ b/apps/aqua/src/repos/stats.rs @@ -1,7 +1,10 @@ use async_trait::async_trait; +use anyhow::anyhow; +use base64::Engine; use jacquard_common::from_json_value; use jacquard_common::types::string::{AtUri, Did}; use serde::{Deserialize, Serialize}; +use sqlx::types::Uuid; use types::fm_teal::alpha::feed::PlayView; use types::fm_teal::alpha::stats::{ArtistView, ReleaseView}; @@ -31,10 +34,75 @@ pub(crate) fn decode_latest_cursor( } pub(crate) fn encode_latest_cursor(cursor: &LatestPlaysCursor) -> anyhow::Result { - Ok(base64::Engine::encode( - &base64::engine::general_purpose::URL_SAFE_NO_PAD, - serde_json::to_vec(cursor)?, - )) + encode_cursor(cursor) +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum StatsPeriod { + SevenDays, + ThirtyDays, +} + +impl StatsPeriod { + fn parse(period: Option<&str>) -> anyhow::Result { + match period.unwrap_or("30days") { + "7days" => Ok(Self::SevenDays), + "30days" => Ok(Self::ThirtyDays), + other => Err(anyhow!("unsupported period: {other}")), + } + } + + fn interval_sql(self) -> &'static str { + match self { + Self::SevenDays => "INTERVAL '7 days'", + Self::ThirtyDays => "INTERVAL '30 days'", + } + } +} + +#[derive(Deserialize, Serialize)] +pub(crate) struct OffsetCursor { + pub offset: i64, +} + +pub struct UserTopArtistsPage { + pub artists: Vec, + pub cursor: Option, +} + +pub struct UserTopReleasesPage { + pub releases: Vec, + pub cursor: Option, +} + +fn normalize_limit(limit: Option) -> i64 { + limit.unwrap_or(50).clamp(1, 100) as i64 +} + +fn decode_offset_cursor(cursor: Option<&str>) -> anyhow::Result { + Ok(cursor + .map(|cursor| decode_cursor::(cursor).map(|cursor| cursor.offset.max(0))) + .transpose()? + .unwrap_or(0)) +} + +fn encode_offset_cursor(offset: i64) -> anyhow::Result { + encode_cursor(&OffsetCursor { offset }) +} + +fn decode_cursor(cursor: &str) -> anyhow::Result +where + T: for<'de> Deserialize<'de>, +{ + let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD.decode(cursor)?; + Ok(serde_json::from_slice(&bytes)?) +} + +fn encode_cursor(cursor: &T) -> anyhow::Result +where + T: Serialize, +{ + Ok(base64::engine::general_purpose::URL_SAFE_NO_PAD.encode(serde_json::to_vec(cursor)?)) } #[async_trait] @@ -43,14 +111,18 @@ pub trait StatsRepo: Send + Sync { async fn get_top_releases(&self, limit: Option) -> anyhow::Result>; async fn get_user_top_artists( &self, - did: &str, + actor: &str, + period: Option<&str>, limit: Option, - ) -> anyhow::Result>; + cursor: Option<&str>, + ) -> anyhow::Result; async fn get_user_top_releases( &self, - did: &str, + actor: &str, + period: Option<&str>, limit: Option, - ) -> anyhow::Result>; + cursor: Option<&str>, + ) -> anyhow::Result; async fn get_latest( &self, limit: Option, @@ -133,85 +205,121 @@ impl StatsRepo for PgDataSource { async fn get_user_top_artists( &self, - did: &str, + actor: &str, + period: Option<&str>, limit: Option, - ) -> anyhow::Result> { - let limit = limit.unwrap_or(50).min(100) as i64; - - let rows = sqlx::query!( + cursor: Option<&str>, + ) -> anyhow::Result { + let did = self.resolve_actor_to_did(actor).await?; + let period = StatsPeriod::parse(period)?; + let limit = normalize_limit(limit); + let offset = decode_offset_cursor(cursor)?; + let query_limit = limit + 1; + let sql = format!( r#" SELECT ae.mbid, ptae.artist_name as name, - COUNT(*) as play_count + COUNT(*)::bigint as play_count FROM plays p INNER JOIN play_to_artists_extended ptae ON p.uri = ptae.play_uri INNER JOIN artists_extended ae ON ptae.artist_id = ae.id WHERE p.did = $1 + AND p.played_time >= NOW() - {} AND ptae.artist_name IS NOT NULL GROUP BY ae.mbid, ptae.artist_name - ORDER BY play_count DESC - LIMIT $2 + ORDER BY play_count DESC, ptae.artist_name ASC, ae.mbid ASC + LIMIT $2 OFFSET $3 "#, - did, - limit - ) + period.interval_sql() + ); + + let rows = sqlx::query_as::<_, (Option, String, i64)>(&sql) + .bind(did) + .bind(query_limit) + .bind(offset) .fetch_all(&self.db) .await?; - let mut result = Vec::with_capacity(rows.len()); - for row in rows { - result.push(ArtistView { - mbid: row.mbid.map(mbid_uri), - name: Some(row.name.into()), - play_count: Some(row.play_count.unwrap_or(0)), + let has_more = rows.len() > limit as usize; + let mut artists = Vec::with_capacity(rows.len().min(limit as usize)); + for (mbid, name, play_count) in rows.into_iter().take(limit as usize) { + artists.push(ArtistView { + mbid: mbid.map(mbid_uri), + name: Some(name.into()), + play_count: Some(play_count), extra_data: Default::default(), }); } - Ok(result) + Ok(UserTopArtistsPage { + artists, + cursor: if has_more { + Some(encode_offset_cursor(offset + limit)?) + } else { + None + }, + }) } async fn get_user_top_releases( &self, - did: &str, + actor: &str, + period: Option<&str>, limit: Option, - ) -> anyhow::Result> { - let limit = limit.unwrap_or(50).min(100) as i64; - - let rows = sqlx::query!( + cursor: Option<&str>, + ) -> anyhow::Result { + let did = self.resolve_actor_to_did(actor).await?; + let period = StatsPeriod::parse(period)?; + let limit = normalize_limit(limit); + let offset = decode_offset_cursor(cursor)?; + let query_limit = limit + 1; + let sql = format!( r#" SELECT p.release_mbid as mbid, p.release_name as name, - COUNT(*) as play_count + COUNT(*)::bigint as play_count FROM plays p WHERE p.did = $1 + AND p.played_time >= NOW() - {} AND p.release_mbid IS NOT NULL AND p.release_name IS NOT NULL GROUP BY p.release_mbid, p.release_name - ORDER BY play_count DESC - LIMIT $2 + ORDER BY play_count DESC, p.release_name ASC, p.release_mbid ASC + LIMIT $2 OFFSET $3 "#, - did, - limit - ) + period.interval_sql() + ); + + let rows = sqlx::query_as::<_, (Option, Option, i64)>(&sql) + .bind(did) + .bind(query_limit) + .bind(offset) .fetch_all(&self.db) .await?; - let mut result = Vec::with_capacity(rows.len()); - for row in rows { - if let (Some(mbid), Some(name)) = (row.mbid, row.name) { - result.push(ReleaseView { + let has_more = rows.len() > limit as usize; + let mut releases = Vec::with_capacity(rows.len().min(limit as usize)); + for (mbid, name, play_count) in rows.into_iter().take(limit as usize) { + if let (Some(mbid), Some(name)) = (mbid, name) { + releases.push(ReleaseView { mbid: Some(mbid_uri(mbid)), name: Some(name.into()), - play_count: Some(row.play_count.unwrap_or(0)), + play_count: Some(play_count), extra_data: Default::default(), }); } } - Ok(result) + Ok(UserTopReleasesPage { + releases, + cursor: if has_more { + Some(encode_offset_cursor(offset + limit)?) + } else { + None + }, + }) } async fn get_latest( @@ -325,9 +433,41 @@ impl StatsRepo for PgDataSource { } } +impl PgDataSource { + async fn resolve_actor_to_did(&self, actor: &str) -> anyhow::Result { + if actor.starts_with("did:") { + return Ok(actor.to_string()); + } + + if let Some(row) = sqlx::query_as::<_, (String,)>( + "SELECT did FROM profiles WHERE LOWER(handle) = LOWER($1) LIMIT 1", + ) + .bind(actor.trim_start_matches("at://")) + .fetch_optional(&self.db) + .await? + { + return Ok(row.0); + } + + let url = format!( + "https://bsky.social/xrpc/com.atproto.identity.resolveHandle?handle={}", + url::form_urlencoded::byte_serialize(actor.as_bytes()).collect::() + ); + let response: serde_json::Value = reqwest::get(&url).await?.json().await?; + response + .get("did") + .and_then(|did| did.as_str()) + .map(ToString::to_string) + .ok_or_else(|| anyhow!("could not resolve actor handle: {actor}")) + } +} + #[cfg(test)] mod tests { - use super::{LatestPlaysCursor, decode_latest_cursor, encode_latest_cursor}; + use super::{ + LatestPlaysCursor, OffsetCursor, StatsPeriod, decode_latest_cursor, decode_offset_cursor, + encode_latest_cursor, encode_offset_cursor, normalize_limit, + }; #[test] fn latest_cursor_round_trips() -> anyhow::Result<()> { @@ -347,4 +487,36 @@ mod tests { fn latest_cursor_rejects_invalid_values() { assert!(decode_latest_cursor(Some("not-a-cursor")).is_err()); } + + #[test] + fn offset_cursor_round_trips() -> anyhow::Result<()> { + let encoded = encode_offset_cursor(100)?; + assert_eq!(decode_offset_cursor(Some(&encoded))?, 100); + Ok(()) + } + + #[test] + fn offset_cursor_defaults_and_clamps_negative_values() -> anyhow::Result<()> { + assert_eq!(decode_offset_cursor(None)?, 0); + let encoded = super::encode_cursor(&OffsetCursor { offset: -10 })?; + assert_eq!(decode_offset_cursor(Some(&encoded))?, 0); + Ok(()) + } + + #[test] + fn stats_period_accepts_only_lexicon_values() { + assert_eq!(StatsPeriod::parse(None).unwrap(), StatsPeriod::ThirtyDays); + assert_eq!( + StatsPeriod::parse(Some("7days")).unwrap(), + StatsPeriod::SevenDays + ); + assert!(StatsPeriod::parse(Some("all")).is_err()); + } + + #[test] + fn stats_limit_is_bounded() { + assert_eq!(normalize_limit(None), 50); + assert_eq!(normalize_limit(Some(-5)), 1); + assert_eq!(normalize_limit(Some(500)), 100); + } } diff --git a/apps/aqua/src/xrpc/stats.rs b/apps/aqua/src/xrpc/stats.rs index 6ebec35..406aaf5 100644 --- a/apps/aqua/src/xrpc/stats.rs +++ b/apps/aqua/src/xrpc/stats.rs @@ -72,12 +72,15 @@ pub async fn get_top_releases( #[derive(Deserialize)] pub struct GetUserTopArtistsQuery { pub actor: String, + pub period: Option, pub limit: Option, + pub cursor: Option, } #[derive(Serialize)] pub struct GetUserTopArtistsResponse { artists: Vec, + cursor: Option, } pub async fn get_user_top_artists( @@ -90,9 +93,18 @@ pub async fn get_user_top_artists( return Err((StatusCode::BAD_REQUEST, "actor is required".to_string())); } - match repo.get_user_top_artists(&query.actor, query.limit).await { - Ok(artists) => Ok(axum::Json(GetUserTopArtistsResponse { - artists: artists.into_static(), + match repo + .get_user_top_artists( + &query.actor, + query.period.as_deref(), + query.limit, + query.cursor.as_deref(), + ) + .await + { + Ok(page) => Ok(axum::Json(GetUserTopArtistsResponse { + artists: page.artists.into_static(), + cursor: page.cursor, })), Err(e) => Err((StatusCode::INTERNAL_SERVER_ERROR, e.to_string())), } @@ -101,12 +113,15 @@ pub async fn get_user_top_artists( #[derive(Deserialize)] pub struct GetUserTopReleasesQuery { pub actor: String, + pub period: Option, pub limit: Option, + pub cursor: Option, } #[derive(Serialize)] pub struct GetUserTopReleasesResponse { releases: Vec, + cursor: Option, } pub async fn get_user_top_releases( @@ -119,9 +134,18 @@ pub async fn get_user_top_releases( return Err((StatusCode::BAD_REQUEST, "actor is required".to_string())); } - match repo.get_user_top_releases(&query.actor, query.limit).await { - Ok(releases) => Ok(axum::Json(GetUserTopReleasesResponse { - releases: releases.into_static(), + match repo + .get_user_top_releases( + &query.actor, + query.period.as_deref(), + query.limit, + query.cursor.as_deref(), + ) + .await + { + Ok(page) => Ok(axum::Json(GetUserTopReleasesResponse { + releases: page.releases.into_static(), + cursor: page.cursor, })), Err(e) => Err((StatusCode::INTERNAL_SERVER_ERROR, e.to_string())), } diff --git a/migrations/20241220000010_status_expiry_indexing.sql b/migrations/20241220000010_status_expiry_indexing.sql new file mode 100644 index 0000000..7d3791a --- /dev/null +++ b/migrations/20241220000010_status_expiry_indexing.sql @@ -0,0 +1,17 @@ +-- Make actor status visibility/indexing explicit. + +ALTER TABLE statii + ADD COLUMN status_time TIMESTAMP WITH TIME ZONE, + ADD COLUMN expires_at TIMESTAMP WITH TIME ZONE; + +UPDATE statii +SET + status_time = (record->>'time')::timestamptz, + expires_at = COALESCE( + (record->>'expiry')::timestamptz, + (record->>'time')::timestamptz + INTERVAL '10 minutes' + ) +WHERE record ? 'time'; + +CREATE INDEX idx_statii_did_visible_latest + ON statii (did, expires_at, status_time DESC, indexed_at DESC); diff --git a/services/cadet/src/ingestors/teal/actor_status.rs b/services/cadet/src/ingestors/teal/actor_status.rs index 1fe563e..cc0b49b 100644 --- a/services/cadet/src/ingestors/teal/actor_status.rs +++ b/services/cadet/src/ingestors/teal/actor_status.rs @@ -1,5 +1,5 @@ use async_trait::async_trait; -use jacquard_common::types::value; +use jacquard_common::types::{string::Datetime, value}; use rocketman::{ingestion::LexiconIngestor, types::event::Event}; use serde_json::Value; use sqlx::PgPool; @@ -25,22 +25,27 @@ impl ActorStatusIngestor { let uri = assemble_at_uri(did, "fm.teal.alpha.actor.status", rkey); let record_json = serde_json::to_value(status)?; + let (status_time, expires_at) = status_times(&status.time, status.expiry.as_ref()); - sqlx::query!( + sqlx::query( r#" - INSERT INTO statii (uri, did, rkey, cid, record) - VALUES ($1, $2, $3, $4, $5) + INSERT INTO statii (uri, did, rkey, cid, record, status_time, expires_at) + VALUES ($1, $2, $3, $4, $5, $6, $7) ON CONFLICT (uri) DO UPDATE SET cid = EXCLUDED.cid, record = EXCLUDED.record, + status_time = EXCLUDED.status_time, + expires_at = EXCLUDED.expires_at, indexed_at = NOW(); "#, - uri, - did, - rkey, - cid, - record_json ) + .bind(uri) + .bind(did) + .bind(rkey) + .bind(cid) + .bind(record_json) + .bind(status_time) + .bind(expires_at) .execute(&self.sql) .await?; @@ -63,6 +68,22 @@ impl ActorStatusIngestor { } } +pub(crate) fn status_times( + time: &Datetime, + expiry: Option<&Datetime>, +) -> (time::OffsetDateTime, time::OffsetDateTime) { + let status_time = datetime_to_time(time); + let expires_at = expiry + .map(datetime_to_time) + .unwrap_or_else(|| status_time + time::Duration::minutes(10)); + (status_time, expires_at) +} + +fn datetime_to_time(datetime: &Datetime) -> time::OffsetDateTime { + time::OffsetDateTime::from_unix_timestamp(datetime.as_ref().timestamp()) + .unwrap_or_else(|_| time::OffsetDateTime::now_utc()) +} + #[async_trait] impl LexiconIngestor for ActorStatusIngestor { async fn ingest(&self, message: Event) -> anyhow::Result<()> { @@ -87,3 +108,31 @@ impl LexiconIngestor for ActorStatusIngestor { Ok(()) } } + +#[cfg(test)] +mod tests { + use jacquard_common::types::string::Datetime; + + use super::status_times; + + fn datetime(value: &str) -> Datetime { + serde_json::from_value(serde_json::json!(value)).expect("valid datetime") + } + + #[test] + fn status_expiry_defaults_to_ten_minutes_after_status_time() { + let time = datetime("2026-06-01T12:00:00Z"); + let (status_time, expires_at) = status_times(&time, None); + + assert_eq!(expires_at - status_time, time::Duration::minutes(10)); + } + + #[test] + fn status_expiry_uses_record_expiry_when_present() { + let time = datetime("2026-06-01T12:00:00Z"); + let expiry = datetime("2026-06-01T12:03:00Z"); + let (status_time, expires_at) = status_times(&time, Some(&expiry)); + + assert_eq!(expires_at - status_time, time::Duration::minutes(3)); + } +} diff --git a/todo.md b/todo.md index a49f2f0..6c5e24d 100644 --- a/todo.md +++ b/todo.md @@ -34,10 +34,10 @@ This file is the working handoff for the Teal-native Teal clone. Keep it updated - [x] Extend CAR import/backfill to process the new profile-status and social record collections, including creates, updates, deletes, and validation failures. - [x] Add database migrations for profile onboarding status, social posts, post reply refs, likes, reposts, playlists, playlist items, badge definitions, badge assignments, rich-text facets, blob CIDs, and derived count/index tables. - [x] Implement delete handling for each new indexed record so primary rows, join rows, counters, and notification rows stay consistent. -- [ ] Finish `fm.teal.alpha.actor.status` indexing semantics: pick the latest status per actor, respect `expiry` with the 10-minute default, omit expired statuses from profile responses, and add regression tests. +- [x] Finish `fm.teal.alpha.actor.status` indexing semantics: pick the latest status per actor, respect `expiry` with the 10-minute default, omit expired statuses from profile responses, and add regression tests. - [x] Index `fm.teal.alpha.actor.profileStatus` so onboarding state can be read from Aqua instead of only direct PDS `getRecord` calls. - [x] Update Aqua `getProfile`/`getProfiles` responses to include indexed actor status and profile onboarding status from the appview. -- [ ] Complete `fm.teal.alpha.stats.getUserTopArtists` and `getUserTopReleases` to honor `period`, `cursor`, handle-to-DID resolution, limit bounds, and lexicon response shapes. +- [x] Complete `fm.teal.alpha.stats.getUserTopArtists` and `getUserTopReleases` to honor `period`, `cursor`, handle-to-DID resolution, limit bounds, and lexicon response shapes. - [ ] Add Aqua social read APIs or lexicons for post feeds, post detail, replies, likes, reposts, playlists, playlist items, badge catalogs, actor badges, and notifications. - [ ] Add Amethyst compose/publish flows for Teal social posts with `trackView`, replies, tags, langs, and `fm.teal.alpha.richtext.facet` mention/link rendering. - [ ] Add Amethyst like and repost actions with optimistic viewer state, counts, undo/delete behavior, and signed-out affordances.