diff --git a/.sqlx/query-41b7d971ace6a57750c7a4d3fbf19d213124144d9ecb22aa2302d675a80ce7c4.json b/.sqlx/query-41b7d971ace6a57750c7a4d3fbf19d213124144d9ecb22aa2302d675a80ce7c4.json new file mode 100644 index 0000000..aca0fda --- /dev/null +++ b/.sqlx/query-41b7d971ace6a57750c7a4d3fbf19d213124144d9ecb22aa2302d675a80ce7c4.json @@ -0,0 +1,42 @@ +{ + "db_name": "PostgreSQL", + "query": "\n SELECT avatar, did, display_name, handle\n FROM profiles\n WHERE display_name ILIKE '%' || $1 || '%' ESCAPE '!'\n OR description ILIKE '%' || $1 || '%' ESCAPE '!'\n OR handle ILIKE '%' || $1 || '%' ESCAPE '!'\n ORDER BY display_name NULLS LAST, did\n LIMIT $2 OFFSET $3\n ", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "avatar", + "type_info": "Text" + }, + { + "ordinal": 1, + "name": "did", + "type_info": "Text" + }, + { + "ordinal": 2, + "name": "display_name", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "handle", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Text", + "Int8", + "Int8" + ] + }, + "nullable": [ + true, + false, + true, + true + ] + }, + "hash": "41b7d971ace6a57750c7a4d3fbf19d213124144d9ecb22aa2302d675a80ce7c4" +} diff --git a/.sqlx/query-7e4ab6be5c601dab08817293924e681b97669c1d6d33455cc0f7f33ae107aad5.json b/.sqlx/query-7e4ab6be5c601dab08817293924e681b97669c1d6d33455cc0f7f33ae107aad5.json new file mode 100644 index 0000000..25db21f --- /dev/null +++ b/.sqlx/query-7e4ab6be5c601dab08817293924e681b97669c1d6d33455cc0f7f33ae107aad5.json @@ -0,0 +1,41 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT p.avatar, p.did, p.display_name, p.handle\n FROM profiles p\n WHERE (p.did = ANY($1)) OR (p.handle = ANY($2))\n ORDER BY p.display_name NULLS LAST, p.did", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "avatar", + "type_info": "Text" + }, + { + "ordinal": 1, + "name": "did", + "type_info": "Text" + }, + { + "ordinal": 2, + "name": "display_name", + "type_info": "Text" + }, + { + "ordinal": 3, + "name": "handle", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "TextArray", + "TextArray" + ] + }, + "nullable": [ + true, + false, + true, + true + ] + }, + "hash": "7e4ab6be5c601dab08817293924e681b97669c1d6d33455cc0f7f33ae107aad5" +} diff --git a/apps/amethyst/components/actor/actorView.tsx b/apps/amethyst/components/actor/actorView.tsx index a434463..c6c9ef1 100644 --- a/apps/amethyst/components/actor/actorView.tsx +++ b/apps/amethyst/components/actor/actorView.tsx @@ -52,7 +52,7 @@ export default function ActorView({ actorDid, pdsAgent }: ActorViewProps) { { headers: { "atproto-proxy": tealDid + "#teal_fm_appview" } }, ); if (isMounted) { - setProfile(res.data["actor"] as GetProfileOutputSchema["actor"]); + setProfile(res.data.actor as GetProfileOutputSchema["actor"]); } } catch (error) { console.error("Error fetching profile:", error); diff --git a/apps/amethyst/lib/__tests__/onboardingRecords.test.ts b/apps/amethyst/lib/__tests__/onboardingRecords.test.ts index 00b7232..9ce8ca7 100644 --- a/apps/amethyst/lib/__tests__/onboardingRecords.test.ts +++ b/apps/amethyst/lib/__tests__/onboardingRecords.test.ts @@ -72,6 +72,46 @@ describe("readRepoRecordWithLegacyFallback", () => { ), ).resolves.toBeNull(); }); + + it("propagates a legacy read failure when the stable record is missing", async () => { + const legacyError = { + error: "NetworkError", + message: "legacy read failed", + }; + const call = jest + .fn() + .mockRejectedValueOnce({ error: "RecordNotFound", message: "missing" }) + .mockRejectedValueOnce(legacyError); + const agent = { call, did: "did:plc:test" }; + + await expect( + readRepoRecordWithLegacyFallback( + agent as never, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ), + ).rejects.toBe(legacyError); + }); + + it("propagates a stable read failure when the legacy record is missing", async () => { + const stableError = { + error: "NetworkError", + message: "stable read failed", + }; + const call = jest + .fn() + .mockRejectedValueOnce(stableError) + .mockRejectedValueOnce({ error: "RecordNotFound", message: "missing" }); + const agent = { call, did: "did:plc:test" }; + + await expect( + readRepoRecordWithLegacyFallback( + agent as never, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ), + ).rejects.toBe(stableError); + }); }); describe("getBlobHash", () => { diff --git a/apps/amethyst/lib/atp/onboardingRecords.ts b/apps/amethyst/lib/atp/onboardingRecords.ts index b372b91..8d3b044 100644 --- a/apps/amethyst/lib/atp/onboardingRecords.ts +++ b/apps/amethyst/lib/atp/onboardingRecords.ts @@ -33,11 +33,14 @@ export async function readRepoRecordWithLegacyFallback( try { return await readRepoRecord(agent, legacyCollection); } catch (legacyError) { - if (isRecordNotFound(stableError) || isRecordNotFound(legacyError)) { + const stableNotFound = isRecordNotFound(stableError); + const legacyNotFound = isRecordNotFound(legacyError); + + if (stableNotFound && legacyNotFound) { return null; } - throw stableError; + throw stableNotFound ? legacyError : stableError; } } } diff --git a/apps/amethyst/stores/authenticationSlice.tsx b/apps/amethyst/stores/authenticationSlice.tsx index d117cef..28d5f40 100644 --- a/apps/amethyst/stores/authenticationSlice.tsx +++ b/apps/amethyst/stores/authenticationSlice.tsx @@ -192,7 +192,7 @@ export const createAuthenticationSlice: StateCreator = ( ) .then((profile) => { console.log(profile); - return profile.data.agent || null; + return profile.data.actor || null; }); set({ diff --git a/apps/aqua/src/repos/actor_profile.rs b/apps/aqua/src/repos/actor_profile.rs index 642d8d5..239be32 100644 --- a/apps/aqua/src/repos/actor_profile.rs +++ b/apps/aqua/src/repos/actor_profile.rs @@ -3,7 +3,7 @@ use jacquard_common::from_json_value; use serde_json::Value; use types::{ app_bsky::richtext::facet::Facet, - fm_teal::actor::{ProfileView, StatusView}, + fm_teal::actor::{MiniProfileView, ProfileView, StatusView}, }; use super::{pg::PgDataSource, utc_to_atrium_datetime}; @@ -15,6 +15,16 @@ pub trait ActorProfileRepo { &self, identities: &[String], ) -> anyhow::Result>; + async fn get_multiple_actor_mini_profiles( + &self, + identities: &[String], + ) -> anyhow::Result>; + async fn search_actor_profiles( + &self, + query: &str, + limit: i64, + offset: i64, + ) -> anyhow::Result>; } pub struct PgProfileRepoRows { @@ -28,6 +38,25 @@ pub struct PgProfileRepoRows { pub status: Option, } +pub struct PgMiniProfileRepoRows { + pub avatar: Option, + pub did: Option, + pub display_name: Option, + pub handle: Option, +} + +impl From for MiniProfileView { + fn from(row: PgMiniProfileRepoRows) -> Self { + Self { + avatar: row.avatar.map(Into::into), + did: row.did.map(Into::into), + display_name: row.display_name.map(Into::into), + handle: row.handle.map(Into::into), + extra_data: Default::default(), + } + } +} + impl From for ProfileView { fn from(row: PgProfileRepoRows) -> Self { Self { @@ -96,4 +125,68 @@ impl ActorProfileRepo for PgDataSource { .await?; Ok(profiles.into_iter().map(|p| p.into()).collect()) } + + async fn get_multiple_actor_mini_profiles( + &self, + identities: &[String], + ) -> anyhow::Result> { + let (dids, handles) = split_identities(identities); + let profiles = sqlx::query_as!( + PgMiniProfileRepoRows, + "SELECT p.avatar, p.did, p.display_name, p.handle + FROM profiles p + WHERE (p.did = ANY($1)) OR (p.handle = ANY($2)) + ORDER BY p.display_name NULLS LAST, p.did", + &dids, + &handles, + ) + .fetch_all(&self.db) + .await?; + + Ok(profiles.into_iter().map(Into::into).collect()) + } + + async fn search_actor_profiles( + &self, + query: &str, + limit: i64, + offset: i64, + ) -> anyhow::Result> { + let query = query + .replace('!', "!!") + .replace('%', "!%") + .replace('_', "!_"); + let profiles = sqlx::query_as!( + PgMiniProfileRepoRows, + r#" + SELECT avatar, did, display_name, handle + FROM profiles + WHERE display_name ILIKE '%' || $1 || '%' ESCAPE '!' + OR description ILIKE '%' || $1 || '%' ESCAPE '!' + OR handle ILIKE '%' || $1 || '%' ESCAPE '!' + ORDER BY display_name NULLS LAST, did + LIMIT $2 OFFSET $3 + "#, + query, + limit, + offset, + ) + .fetch_all(&self.db) + .await?; + + Ok(profiles.into_iter().map(Into::into).collect()) + } +} + +fn split_identities(identities: &[String]) -> (Vec, Vec) { + let mut dids = Vec::new(); + let mut handles = Vec::new(); + for identity in identities { + if identity.starts_with("did:") { + dids.push(identity.clone()); + } else { + handles.push(identity.clone()); + } + } + (dids, handles) } diff --git a/apps/aqua/src/xrpc/actor.rs b/apps/aqua/src/xrpc/actor.rs index 67282af..7b6ddfb 100644 --- a/apps/aqua/src/xrpc/actor.rs +++ b/apps/aqua/src/xrpc/actor.rs @@ -1,14 +1,15 @@ use crate::ctx::Context; -use axum::{Extension, http::StatusCode, response::IntoResponse, routing::get}; +use axum::{http::StatusCode, response::IntoResponse, routing::get, Extension}; use jacquard_common::IntoStatic; use serde::{Deserialize, Serialize}; -use types::fm_teal::actor::ProfileView; +use types::fm_teal::actor::{MiniProfileView, ProfileView}; // mount actor routes pub fn actor_routes() -> axum::Router { axum::Router::new() .route("/fm.teal.actor.getProfile", get(get_actor)) .route("/fm.teal.actor.getProfiles", get(get_actors)) + .route("/fm.teal.actor.searchActors", get(search_actors)) } #[derive(Deserialize)] @@ -18,7 +19,7 @@ pub struct GetProfileQuery { #[derive(Serialize)] pub struct GetProfileResponse { - profile: ProfileView, + actor: ProfileView, } pub async fn get_actor( @@ -26,24 +27,91 @@ pub async fn get_actor( axum::extract::Query(query): axum::extract::Query, ) -> Result { let repo = &ctx.db; // assuming ctx.db is Box - let identity = &query.actor; + let identity = match query.actor.as_deref() { + Some(identity) => identity, + None => return Err((StatusCode::BAD_REQUEST, "actor is required".to_string())), + }; - if identity.is_none() { - return Err((StatusCode::BAD_REQUEST, "actor is required".to_string())); - } - - match repo - .get_actor_profile(identity.as_ref().expect("actor is not none").as_str()) - .await - { + match repo.get_actor_profile(identity).await { Ok(Some(profile)) => Ok(axum::Json(GetProfileResponse { - profile: profile.into_static(), + actor: profile.into_static(), })), Ok(None) => Err((StatusCode::NOT_FOUND, "Profile not found".to_string())), Err(e) => Err((StatusCode::INTERNAL_SERVER_ERROR, e.to_string())), } } +#[derive(Deserialize)] +pub struct SearchActorsQuery { + pub q: String, + pub limit: Option, + pub cursor: Option, +} + +#[derive(Serialize)] +pub struct SearchActorsResponse { + actors: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + cursor: Option, +} + +pub async fn search_actors( + Extension(ctx): Extension, + axum::extract::Query(query): axum::extract::Query, +) -> Result { + if query.q.trim().is_empty() { + return Err((StatusCode::BAD_REQUEST, "q is required".to_string())); + } + + let limit = query.limit.unwrap_or(25); + if !(1..=25).contains(&limit) { + return Err(( + StatusCode::BAD_REQUEST, + "limit must be between 1 and 25".to_string(), + )); + } + let offset = query + .cursor + .as_deref() + .unwrap_or("0") + .parse::() + .map_err(|_| { + ( + StatusCode::BAD_REQUEST, + "cursor must be a non-negative integer".to_string(), + ) + })?; + if offset < 0 { + return Err(( + StatusCode::BAD_REQUEST, + "cursor must be a non-negative integer".to_string(), + )); + } + + let next_offset = offset.checked_add(limit).ok_or_else(|| { + ( + StatusCode::BAD_REQUEST, + "cursor is out of range".to_string(), + ) + })?; + let mut actors = ctx + .db + .search_actor_profiles(query.q.trim(), limit + 1, offset) + .await + .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; + let cursor = if actors.len() > limit as usize { + actors.truncate(limit as usize); + Some(next_offset.to_string()) + } else { + None + }; + + Ok(axum::Json(SearchActorsResponse { + actors: actors.into_static(), + cursor, + })) +} + #[derive(Deserialize)] pub struct GetProfilesQuery { pub actors: Vec, @@ -51,7 +119,7 @@ pub struct GetProfilesQuery { #[derive(Serialize)] pub struct GetProfilesResponse { - profiles: Vec, + actors: Vec, } pub async fn get_actors( @@ -65,9 +133,9 @@ pub async fn get_actors( return Err((StatusCode::BAD_REQUEST, "actor is required".to_string())); } - match repo.get_multiple_actor_profiles(actor).await { - Ok(profiles) => Ok(axum::Json(GetProfilesResponse { - profiles: profiles.into_static(), + match repo.get_multiple_actor_mini_profiles(actor).await { + Ok(actors) => Ok(axum::Json(GetProfilesResponse { + actors: actors.into_static(), })), Err(e) => Err((StatusCode::INTERNAL_SERVER_ERROR, e.to_string())), } diff --git a/apps/aqua/src/xrpc/feed.rs b/apps/aqua/src/xrpc/feed.rs index b81c598..72aff6e 100644 --- a/apps/aqua/src/xrpc/feed.rs +++ b/apps/aqua/src/xrpc/feed.rs @@ -1,5 +1,5 @@ use crate::ctx::Context; -use axum::{Extension, http::StatusCode, response::IntoResponse, routing::get}; +use axum::{http::StatusCode, response::IntoResponse, routing::get, Extension}; use jacquard_common::IntoStatic; use serde::{Deserialize, Serialize}; use types::fm_teal::feed::PlayView; @@ -9,11 +9,14 @@ pub fn feed_routes() -> axum::Router { axum::Router::new() .route("/fm.teal.feed.getPlay", get(get_feed_play)) .route("/fm.teal.feed.getPlays", get(get_feed_plays)) + .route("/fm.teal.feed.getActorFeed", get(get_actor_feed)) } #[derive(Deserialize)] pub struct GetFeedPlayQuery { - pub identity: Option, + #[serde(rename = "authorDID")] + pub author_did: Option, + pub rkey: Option, } #[derive(Serialize)] @@ -26,16 +29,20 @@ pub async fn get_feed_play( axum::extract::Query(query): axum::extract::Query, ) -> Result { let repo = &ctx.db; - let identity = &query.identity; + let (author_did, rkey) = match (query.author_did.as_deref(), query.rkey.as_deref()) { + (Some(author_did), Some(rkey)) if !author_did.is_empty() && !rkey.is_empty() => { + (author_did, rkey) + } + _ => { + return Err(( + StatusCode::BAD_REQUEST, + "authorDID and rkey are required".to_string(), + )); + } + }; + let identity = format!("at://{author_did}/fm.teal.feed.play/{rkey}"); - if identity.is_none() { - return Err((StatusCode::BAD_REQUEST, "identity is required".to_string())); - } - - match repo - .get_feed_play(identity.as_ref().expect("identity is not none").as_str()) - .await - { + match repo.get_feed_play(&identity).await { Ok(Some(play)) => Ok(axum::Json(GetFeedPlayResponse { play: play.into_static(), })), @@ -75,3 +82,54 @@ pub async fn get_feed_plays( Err(e) => Err((StatusCode::INTERNAL_SERVER_ERROR, e.to_string())), } } + +#[derive(Deserialize)] +pub struct GetActorFeedQuery { + #[serde(rename = "authorDID")] + pub author_did: String, + pub cursor: Option, + pub limit: Option, +} + +#[derive(Serialize)] +pub struct GetActorFeedResponse { + plays: Vec, +} + +pub async fn get_actor_feed( + Extension(ctx): Extension, + axum::extract::Query(query): axum::extract::Query, +) -> Result { + if !query.author_did.starts_with("did:") { + return Err(( + StatusCode::BAD_REQUEST, + "authorDID must be a DID".to_string(), + )); + } + + let limit = query.limit.unwrap_or(20); + if !(1..=50).contains(&limit) { + return Err(( + StatusCode::BAD_REQUEST, + "limit must be between 1 and 50".to_string(), + )); + } + let limit = limit as usize; + let offset = query + .cursor + .as_deref() + .unwrap_or("0") + .parse::() + .map_err(|_| (StatusCode::BAD_REQUEST, "cursor is invalid".to_string()))?; + + match ctx + .db + .get_feed_plays_for_profile(std::slice::from_ref(&query.author_did)) + .await + { + Ok(plays) => Ok(axum::Json(GetActorFeedResponse { + plays: plays.into_iter().skip(offset).take(limit).collect(), + })), + Err(e) => Err((StatusCode::INTERNAL_SERVER_ERROR, e.to_string())), + } +} diff --git a/services/cadet/src/ingestors/car/car_import.rs b/services/cadet/src/ingestors/car/car_import.rs index 3327cce..4dafbdd 100644 --- a/services/cadet/src/ingestors/car/car_import.rs +++ b/services/cadet/src/ingestors/car/car_import.rs @@ -31,6 +31,7 @@ //! and use the original rkey from the AT Protocol MST structure. use crate::ingestors::car::jobs::{queue_keys, CarImportJob}; +use crate::ingestors::teal::normalize_legacy_record_type; use crate::redis_client::RedisClient; use anyhow::{anyhow, Result}; use async_trait::async_trait; @@ -90,21 +91,6 @@ fn parse_teal_key(key: &str) -> Option<(String, String)> { Some((collection.to_string(), rkey.to_string())) } -fn normalize_legacy_record_type(data: &Value) -> Value { - let Value::Object(object) = data else { - return data.clone(); - }; - - let mut normalized = object.clone(); - if let Some(Value::String(record_type)) = normalized.get_mut("$type") { - if let Some(stable_type) = record_type.strip_prefix("fm.teal.alpha.") { - *record_type = format!("fm.teal.{stable_type}"); - } - } - - Value::Object(normalized) -} - fn is_legacy_collection(collection: &str) -> bool { matches!( collection, diff --git a/services/cadet/src/ingestors/teal/actor_profile.rs b/services/cadet/src/ingestors/teal/actor_profile.rs index ff4fbee..b1a5aa3 100644 --- a/services/cadet/src/ingestors/teal/actor_profile.rs +++ b/services/cadet/src/ingestors/teal/actor_profile.rs @@ -4,6 +4,7 @@ use rocketman::{ingestion::LexiconIngestor, types::event::Event}; use serde_json::Value; use sqlx::PgPool; +use crate::ingestors::teal::normalize_legacy_record_type; use crate::resolve::resolve_identity; pub struct ActorProfileIngestor { @@ -86,10 +87,7 @@ impl LexiconIngestor for ActorProfileIngestor { async fn ingest(&self, message: Event) -> anyhow::Result<()> { if let Some(commit) = &message.commit { if let Some(ref record) = &commit.record { - let record: types::fm_teal::actor::profile::Profile = - value::from_json_value::( - record.clone(), - )?; + let record = parse_profile_record(record)?; if let Some(ref commit) = message.commit { if let Some(ref _cid) = commit.cid { // TODO: verify cid @@ -106,3 +104,52 @@ impl LexiconIngestor for ActorProfileIngestor { Ok(()) } } + +/// Normalize legacy namespace tags before deserializing the generated record type. +fn parse_profile_record(record: &Value) -> anyhow::Result { + Ok(value::from_json_value::< + types::fm_teal::actor::profile::Profile, + >(normalize_legacy_record_type(record))?) +} + +#[cfg(test)] +mod tests { + use super::parse_profile_record; + use crate::ingestors::teal::{ALPHA_ACTOR_PROFILE, STABLE_ACTOR_PROFILE}; + use rocketman::types::event::Event; + use serde_json::{json, Value}; + + fn direct_record(operation: &str, collection: &str) -> Value { + let event: Event = serde_json::from_value(json!({ + "did": "did:plc:test", + "kind": "commit", + "commit": { + "rev": "3ltest", + "operation": operation, + "collection": collection, + "rkey": "self", + "record": {"$type": collection, "displayName": "Test User"}, + "cid": "bafytest" + } + })) + .expect("valid direct event"); + + event.commit.expect("commit event").record.expect("record") + } + + #[test] + fn direct_create_and_update_records_accept_stable_and_legacy_types() { + for operation in ["create", "update"] { + for record_type in [STABLE_ACTOR_PROFILE, ALPHA_ACTOR_PROFILE] { + let record = direct_record(operation, record_type); + let parsed = parse_profile_record(&record) + .unwrap_or_else(|error| panic!("{operation} record should parse: {error}")); + + assert_eq!( + parsed.display_name.as_ref().map(|name| name.as_str()), + Some("Test User") + ); + } + } + } +} diff --git a/services/cadet/src/ingestors/teal/actor_status.rs b/services/cadet/src/ingestors/teal/actor_status.rs index 6f085e7..df8ba63 100644 --- a/services/cadet/src/ingestors/teal/actor_status.rs +++ b/services/cadet/src/ingestors/teal/actor_status.rs @@ -4,7 +4,7 @@ use rocketman::{ingestion::LexiconIngestor, types::event::Event}; use serde_json::Value; use sqlx::PgPool; -use crate::ingestors::teal::assemble_at_uri; +use crate::ingestors::teal::{assemble_at_uri, normalize_legacy_record_type}; pub struct ActorStatusIngestor { sql: PgPool, @@ -68,10 +68,7 @@ impl LexiconIngestor for ActorStatusIngestor { async fn ingest(&self, message: Event) -> anyhow::Result<()> { if let Some(commit) = &message.commit { if let Some(ref record) = &commit.record { - let record: types::fm_teal::actor::status::Status = - value::from_json_value::( - record.clone(), - )?; + let record = parse_status_record(record)?; if let Some(ref cid) = commit.cid { self.insert_status(&message.did, &commit.rkey, cid, &record) @@ -87,3 +84,56 @@ impl LexiconIngestor for ActorStatusIngestor { Ok(()) } } + +/// Normalize legacy namespace tags before deserializing the generated record type. +fn parse_status_record(record: &Value) -> anyhow::Result { + Ok(value::from_json_value::< + types::fm_teal::actor::status::Status, + >(normalize_legacy_record_type(record))?) +} + +#[cfg(test)] +mod tests { + use super::parse_status_record; + use crate::ingestors::teal::{ALPHA_ACTOR_STATUS, STABLE_ACTOR_STATUS}; + use rocketman::types::event::Event; + use serde_json::{json, Value}; + + fn direct_record(operation: &str, collection: &str) -> Value { + let event: Event = serde_json::from_value(json!({ + "did": "did:plc:test", + "kind": "commit", + "commit": { + "rev": "3ltest", + "operation": operation, + "collection": collection, + "rkey": "self", + "record": { + "$type": collection, + "time": "2024-01-01T00:00:00Z", + "item": { + "artists": [{"artistName": "Test Artist"}], + "trackName": "Test Song" + } + }, + "cid": "bafytest" + } + })) + .expect("valid direct event"); + + event.commit.expect("commit event").record.expect("record") + } + + #[test] + fn direct_create_and_update_records_accept_stable_and_legacy_types() { + for operation in ["create", "update"] { + for record_type in [STABLE_ACTOR_STATUS, ALPHA_ACTOR_STATUS] { + let record = direct_record(operation, record_type); + let parsed = parse_status_record(&record) + .unwrap_or_else(|error| panic!("{operation} record should parse: {error}")); + + assert_eq!(parsed.item.track_name.as_str(), "Test Song"); + } + } + } +} diff --git a/services/cadet/src/ingestors/teal/feed_play.rs b/services/cadet/src/ingestors/teal/feed_play.rs index 7a47b7e..6276439 100644 --- a/services/cadet/src/ingestors/teal/feed_play.rs +++ b/services/cadet/src/ingestors/teal/feed_play.rs @@ -12,7 +12,7 @@ use serde_json::Value; use sqlx::{types::Uuid, PgPool}; use unicode_normalization::UnicodeNormalization; -use super::assemble_at_uri; +use super::{assemble_at_uri, normalize_legacy_record_type}; #[derive(Debug, Clone)] struct FuzzyMatchCandidate { @@ -1591,10 +1591,7 @@ impl LexiconIngestor for PlayIngestor { async fn ingest(&self, message: Event) -> anyhow::Result<()> { if let Some(commit) = &message.commit { if let Some(ref record) = &commit.record { - let record: types::fm_teal::feed::play::Play = - value::from_json_value::( - record.clone(), - )?; + let record = parse_play_record(record)?; if let Some(ref commit) = message.commit { if let Some(ref cid) = commit.cid { // TODO: verify cid @@ -1627,3 +1624,49 @@ impl LexiconIngestor for PlayIngestor { Ok(()) } } + +/// Normalize legacy namespace tags before deserializing the generated record type. +fn parse_play_record(record: &Value) -> anyhow::Result { + Ok(value::from_json_value::( + normalize_legacy_record_type(record), + )?) +} + +#[cfg(test)] +mod tests { + use super::parse_play_record; + use crate::ingestors::teal::{ALPHA_FEED_PLAY, STABLE_FEED_PLAY}; + use rocketman::types::event::Event; + use serde_json::{json, Value}; + + fn direct_record(operation: &str, collection: &str) -> Value { + let event: Event = serde_json::from_value(json!({ + "did": "did:plc:test", + "kind": "commit", + "commit": { + "rev": "3ltest", + "operation": operation, + "collection": collection, + "rkey": "3ltest", + "record": {"$type": collection, "trackName": "Test Song"}, + "cid": "bafytest" + } + })) + .expect("valid direct event"); + + event.commit.expect("commit event").record.expect("record") + } + + #[test] + fn direct_create_and_update_records_accept_stable_and_legacy_types() { + for operation in ["create", "update"] { + for record_type in [STABLE_FEED_PLAY, ALPHA_FEED_PLAY] { + let record = direct_record(operation, record_type); + let parsed = parse_play_record(&record) + .unwrap_or_else(|error| panic!("{operation} record should parse: {error}")); + + assert_eq!(parsed.track_name.as_str(), "Test Song"); + } + } + } +} diff --git a/services/cadet/src/ingestors/teal/mod.rs b/services/cadet/src/ingestors/teal/mod.rs index daa0b7c..0213b6b 100644 --- a/services/cadet/src/ingestors/teal/mod.rs +++ b/services/cadet/src/ingestors/teal/mod.rs @@ -2,6 +2,8 @@ pub mod actor_profile; pub mod actor_status; pub mod feed_play; +use serde_json::Value; + pub const STABLE_FEED_PLAY: &str = "fm.teal.feed.play"; pub const STABLE_ACTOR_PROFILE: &str = "fm.teal.actor.profile"; pub const STABLE_ACTOR_STATUS: &str = "fm.teal.actor.status"; @@ -10,6 +12,22 @@ pub const ALPHA_FEED_PLAY: &str = "fm.teal.alpha.feed.play"; pub const ALPHA_ACTOR_PROFILE: &str = "fm.teal.alpha.actor.profile"; pub const ALPHA_ACTOR_STATUS: &str = "fm.teal.alpha.actor.status"; +/// Normalize the root record namespace used by historical Teal records. +pub fn normalize_legacy_record_type(data: &Value) -> Value { + let Value::Object(object) = data else { + return data.clone(); + }; + + let mut normalized = object.clone(); + if let Some(Value::String(record_type)) = normalized.get_mut("$type") { + if let Some(stable_type) = record_type.strip_prefix("fm.teal.alpha.") { + *record_type = format!("fm.teal.{stable_type}"); + } + } + + Value::Object(normalized) +} + const COLLECTION_ALIASES: [(&str, &str); 3] = [ (ALPHA_FEED_PLAY, STABLE_FEED_PLAY), (ALPHA_ACTOR_PROFILE, STABLE_ACTOR_PROFILE), diff --git a/tools/lexicon-cli/src/commands/validate.ts b/tools/lexicon-cli/src/commands/validate.ts index 5d25c9a..c9bc4b1 100644 --- a/tools/lexicon-cli/src/commands/validate.ts +++ b/tools/lexicon-cli/src/commands/validate.ts @@ -46,7 +46,7 @@ async function validateTypeScriptGeneration(workspaceRoot: string) { const sourceFiles = await glob('**/*.json', { cwd: lexiconsPath }); for (const sourceFile of sourceFiles) { - const namespace = sourceFile + const namespace: string = sourceFile .replace(/\.json$/, '') .split('/') .flatMap((segment) => segment.split('.'))