diff --git a/apps/amethyst/app/(tabs)/index.tsx b/apps/amethyst/app/(tabs)/index.tsx index 822bd96..3b31c6c 100644 --- a/apps/amethyst/app/(tabs)/index.tsx +++ b/apps/amethyst/app/(tabs)/index.tsx @@ -3,6 +3,11 @@ import { useEffect, useState } from "react"; import { ActivityIndicator, ScrollView, View } from "react-native"; import { Redirect, Stack, useRouter } from "expo-router"; import ActorView from "@/components/actor/actorView"; +import { + LEGACY_PROFILE_STATUS_COLLECTION, + readRepoRecordWithLegacyFallback, + STABLE_PROFILE_STATUS_COLLECTION, +} from "@/lib/atp/onboardingRecords"; import { useStore } from "@/stores/mainStore"; import { Record as ProfileStatusRecord } from "@teal/lexicons/src/types/fm/teal/actor/profileStatus"; @@ -16,7 +21,8 @@ export default function Screen() { const agent = useStore((state) => state.pdsAgent); const profile = useStore((state) => state.profiles[agent?.did ?? ""]); const tealDid = useStore((state) => state.tealDid); - const [profileStatus, setProfileStatus] = useState(null); + const [profileStatus, setProfileStatus] = + useState(null); const [statusLoading, setStatusLoading] = useState(true); useEffect(() => { @@ -26,14 +32,15 @@ export default function Screen() { try { if (!agent) return; - const res = await agent.call("com.atproto.repo.getRecord", { - repo: agent.did, - collection: "fm.teal.actor.profileStatus", - rkey: "self", - }); + const result = + await readRepoRecordWithLegacyFallback( + agent, + STABLE_PROFILE_STATUS_COLLECTION, + LEGACY_PROFILE_STATUS_COLLECTION, + ); if (isMounted) { - setProfileStatus(res.data.value as ProfileStatusRecord); + setProfileStatus(result?.record ?? null); } } catch (error) { if (isMounted) { @@ -65,7 +72,10 @@ export default function Screen() { return ; } - if (!statusLoading && (!profileStatus || profileStatus.completedOnboarding === "none")) { + if ( + !statusLoading && + (!profileStatus || profileStatus.completedOnboarding === "none") + ) { return ( diff --git a/apps/amethyst/app/onboarding/index.tsx b/apps/amethyst/app/onboarding/index.tsx index 824a219..2da1cfe 100644 --- a/apps/amethyst/app/onboarding/index.tsx +++ b/apps/amethyst/app/onboarding/index.tsx @@ -4,7 +4,15 @@ import { SafeAreaView } from "react-native-safe-area-context"; import { useRouter } from "expo-router"; import ProgressDots from "@/components/onboarding/progressDots"; import { Text } from "@/components/ui/text"; // Your UI components - +import getImageCdnLink from "@/lib/atp/getImageCdnLink"; +import { + getBlobHash, + LEGACY_PROFILE_COLLECTION, + LEGACY_PROFILE_STATUS_COLLECTION, + readRepoRecordWithLegacyFallback, + STABLE_PROFILE_COLLECTION, + STABLE_PROFILE_STATUS_COLLECTION, +} from "@/lib/atp/onboardingRecords"; import { useStore } from "@/stores/mainStore"; import { Record as ProfileRecord } from "@teal/lexicons/src/types/fm/teal/actor/profile"; @@ -32,30 +40,33 @@ export default function OnboardingPage() { const [submissionStep, setSubmissionStep] = useState(0); - // Profile status hooks - must be at top level - const [profileStatus, setProfileStatus] = useState(null); + const [profileStatus, setProfileStatus] = + useState(null); const [statusLoading, setStatusLoading] = useState(true); + const [profileLoading, setProfileLoading] = useState(true); const router = useRouter(); const agent = useStore((store) => store.pdsAgent); const profile = useStore((store) => store.profiles); - // Check profile status + // Read both records with a legacy fallback so an existing user is never + // treated as a new user just because the namespace changed. React.useEffect(() => { const checkProfileStatus = async () => { if (!agent) return; try { - const res = await agent.call("com.atproto.repo.getRecord", { - repo: agent.did, - collection: "fm.teal.actor.profileStatus", - rkey: "self", - }); - setProfileStatus(res.data.value as ProfileStatusRecord); - } catch { - // If no record exists, user hasn't completed onboarding + const result = + await readRepoRecordWithLegacyFallback( + agent, + STABLE_PROFILE_STATUS_COLLECTION, + LEGACY_PROFILE_STATUS_COLLECTION, + ); + setProfileStatus(result?.record ?? null); + } catch (error) { + console.error("Error fetching profile status:", error); setProfileStatus(null); } finally { setStatusLoading(false); @@ -65,7 +76,47 @@ export default function OnboardingPage() { checkProfileStatus(); }, [agent]); - const handleImageSelectionComplete = (avatar: string | undefined, banner: string | undefined) => { + React.useEffect(() => { + const loadProfile = async () => { + if (!agent) return; + + try { + const result = await readRepoRecordWithLegacyFallback( + agent, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ); + if (result) { + setDisplayName(result.record.displayName ?? ""); + setDescription(result.record.description ?? ""); + + const avatarHash = getBlobHash(result.record.avatar); + const bannerHash = getBlobHash(result.record.banner); + setAvatarUri( + avatarHash + ? (getImageCdnLink({ did: agent.did!, hash: avatarHash }) ?? "") + : "", + ); + setBannerUri( + bannerHash + ? (getImageCdnLink({ did: agent.did!, hash: bannerHash }) ?? "") + : "", + ); + } + } catch (error) { + console.error("Error fetching user profile:", error); + } finally { + setProfileLoading(false); + } + }; + + loadProfile(); + }, [agent]); + + const handleImageSelectionComplete = ( + avatar: string | undefined, + banner: string | undefined, + ) => { setAvatarUri(avatar ?? ""); setBannerUri(banner ?? ""); onComplete({ displayName, description }, avatar, banner); @@ -94,13 +145,18 @@ export default function OnboardingPage() { let currentUser: ProfileRecord | undefined; let cid: string | undefined; try { - const res = await agent.call("com.atproto.repo.getRecord", { - repo: agent.did, - collection: "fm.teal.actor.profile", - rkey: "self", - }); - currentUser = res.data.value; - cid = res.data.cid; + const result = await readRepoRecordWithLegacyFallback( + agent, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ); + currentUser = result?.record; + // A legacy record must be copied into the stable collection, not + // updated in place. Its blob refs are retained below. + cid = + result?.collection === STABLE_PROFILE_COLLECTION + ? result.cid + : undefined; } catch (error) { console.error("Error fetching user profile:", error); } @@ -114,7 +170,8 @@ export default function OnboardingPage() { const data = await fetch(newAvatarUri).then((r) => r.blob()); const fileType = newAvatarUri.split(";")[0].split(":")[1]; const blob = new Blob([data], { type: fileType }); - newAvatarBlob = (await agent.uploadBlob(blob, { encoding: fileType })).data.blob as unknown as ProfileRecord["avatar"]; + newAvatarBlob = (await agent.uploadBlob(blob, { encoding: fileType })) + .data.blob as unknown as ProfileRecord["avatar"]; } } if (newBannerUri) { @@ -123,40 +180,41 @@ export default function OnboardingPage() { const data = await fetch(newBannerUri).then((r) => r.blob()); const fileType = newBannerUri.split(";")[0].split(":")[1]; const blob = new Blob([data], { type: fileType }); - newBannerBlob = (await agent.uploadBlob(blob, { encoding: fileType })).data.blob as unknown as ProfileRecord["banner"]; + newBannerBlob = (await agent.uploadBlob(blob, { encoding: fileType })) + .data.blob as unknown as ProfileRecord["banner"]; } } setSubmissionStep(4); - let record: ProfileRecord = { + const record: ProfileRecord = { + ...currentUser, + $type: STABLE_PROFILE_COLLECTION, displayName: updatedProfile.displayName, description: updatedProfile.description, avatar: newAvatarBlob, banner: newBannerBlob, }; - let post; - if (cid) { - post = await agent.call( + await agent.call( "com.atproto.repo.putRecord", {}, { repo: agent.did, - collection: "fm.teal.actor.profile", + collection: STABLE_PROFILE_COLLECTION, rkey: "self", record, swapRecord: cid, }, ); } else { - post = await agent.call( + await agent.call( "com.atproto.repo.createRecord", {}, { repo: agent.did, - collection: "fm.teal.actor.profile", + collection: STABLE_PROFILE_COLLECTION, rkey: "self", record, }, @@ -165,6 +223,7 @@ export default function OnboardingPage() { // Update profile status to mark onboarding as completed const profileStatusRecord: ProfileStatusRecord = { + $type: STABLE_PROFILE_STATUS_COLLECTION, completedOnboarding: "profileOnboarding", createdAt: new Date().toISOString(), updatedAt: new Date().toISOString(), @@ -176,7 +235,7 @@ export default function OnboardingPage() { {}, { repo: agent.did, - collection: "fm.teal.actor.profileStatus", + collection: STABLE_PROFILE_STATUS_COLLECTION, rkey: "self", record: profileStatusRecord, }, @@ -189,7 +248,7 @@ export default function OnboardingPage() { {}, { repo: agent.did, - collection: "fm.teal.actor.profileStatus", + collection: STABLE_PROFILE_STATUS_COLLECTION, rkey: "self", record: { ...profileStatusRecord, @@ -214,7 +273,7 @@ export default function OnboardingPage() { return
Loading...
; } - if (statusLoading) { + if (statusLoading || profileLoading) { return ( diff --git a/apps/amethyst/components/actor/actorView.tsx b/apps/amethyst/components/actor/actorView.tsx index 0cb9645..a434463 100644 --- a/apps/amethyst/components/actor/actorView.tsx +++ b/apps/amethyst/components/actor/actorView.tsx @@ -5,6 +5,11 @@ import { Avatar, AvatarFallback, AvatarImage } from "@/components/ui/avatar"; import { Button } from "@/components/ui/button"; import { Text } from "@/components/ui/text"; import getImageCdnLink from "@/lib/atp/getImageCdnLink"; +import { + LEGACY_PROFILE_COLLECTION, + readRepoRecordWithLegacyFallback, + STABLE_PROFILE_COLLECTION, +} from "@/lib/atp/onboardingRecords"; import { Icon } from "@/lib/icons/iconWithClassName"; import { useStore } from "@/stores/mainStore"; import { Agent } from "@atproto/api"; @@ -84,13 +89,18 @@ export default function ActorView({ actorDid, pdsAgent }: ActorViewProps) { let currentUser: ProfileRecord | undefined; let cid: string | undefined; try { - const res = await pdsAgent.call("com.atproto.repo.getRecord", { - repo: pdsAgent.did, - collection: "fm.teal.actor.profile", - rkey: "self", - }); - currentUser = res.data.value; - cid = res.data.cid; + const result = await readRepoRecordWithLegacyFallback( + pdsAgent, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ); + currentUser = result?.record; + // Legacy records are copied to the stable collection on the next edit; + // keep their blob refs and never try to update the old collection. + cid = + result?.collection === STABLE_PROFILE_COLLECTION + ? result.cid + : undefined; } catch (error) { console.error("Error fetching user profile:", error); } @@ -104,7 +114,9 @@ export default function ActorView({ actorDid, pdsAgent }: ActorViewProps) { const data = await fetch(newAvatarUri).then((r) => r.blob()); const fileType = newAvatarUri.split(";")[0].split(":")[1]; const blob = new Blob([data], { type: fileType }); - newAvatarBlob = (await pdsAgent.uploadBlob(blob, { encoding: fileType })).data.blob as unknown as ProfileRecord["avatar"]; + newAvatarBlob = ( + await pdsAgent.uploadBlob(blob, { encoding: fileType }) + ).data.blob as unknown as ProfileRecord["avatar"]; } } if (newBannerUri) { @@ -112,11 +124,15 @@ export default function ActorView({ actorDid, pdsAgent }: ActorViewProps) { const data = await fetch(newBannerUri).then((r) => r.blob()); const fileType = newBannerUri.split(";")[0].split(":")[1]; const blob = new Blob([data], { type: fileType }); - newBannerBlob = (await pdsAgent.uploadBlob(blob, { encoding: fileType })).data.blob as unknown as ProfileRecord["banner"]; + newBannerBlob = ( + await pdsAgent.uploadBlob(blob, { encoding: fileType }) + ).data.blob as unknown as ProfileRecord["banner"]; } } - let record: ProfileRecord = { + const record: ProfileRecord = { + ...currentUser, + $type: STABLE_PROFILE_COLLECTION, displayName: updatedProfile.displayName, description: updatedProfile.description, avatar: newAvatarBlob, @@ -131,7 +147,7 @@ export default function ActorView({ actorDid, pdsAgent }: ActorViewProps) { {}, { repo: pdsAgent.did, - collection: "fm.teal.actor.profile", + collection: STABLE_PROFILE_COLLECTION, rkey: "self", record, swapRecord: cid, @@ -143,7 +159,7 @@ export default function ActorView({ actorDid, pdsAgent }: ActorViewProps) { {}, { repo: pdsAgent.did, - collection: "fm.teal.actor.profile", + collection: STABLE_PROFILE_COLLECTION, rkey: "self", record, }, diff --git a/apps/amethyst/lib/__tests__/onboardingRecords.test.ts b/apps/amethyst/lib/__tests__/onboardingRecords.test.ts new file mode 100644 index 0000000..00b7232 --- /dev/null +++ b/apps/amethyst/lib/__tests__/onboardingRecords.test.ts @@ -0,0 +1,87 @@ +import { + getBlobHash, + LEGACY_PROFILE_COLLECTION, + readRepoRecordWithLegacyFallback, + STABLE_PROFILE_COLLECTION, +} from "../atp/onboardingRecords"; + +describe("readRepoRecordWithLegacyFallback", () => { + it("uses the stable record when it exists", async () => { + const call = jest.fn().mockResolvedValue({ + data: { cid: "stable-cid", value: { displayName: "Stable" } }, + }); + const agent = { call, did: "did:plc:test" }; + + const result = await readRepoRecordWithLegacyFallback( + agent as never, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ); + + expect(result).toEqual({ + collection: STABLE_PROFILE_COLLECTION, + cid: "stable-cid", + record: { displayName: "Stable" }, + }); + expect(call).toHaveBeenCalledTimes(1); + expect(call).toHaveBeenCalledWith("com.atproto.repo.getRecord", { + repo: "did:plc:test", + collection: STABLE_PROFILE_COLLECTION, + rkey: "self", + }); + }); + + it("falls back to the legacy record when stable is missing", async () => { + const call = jest + .fn() + .mockRejectedValueOnce({ error: "RecordNotFound", message: "missing" }) + .mockResolvedValueOnce({ + data: { cid: "legacy-cid", value: { displayName: "Legacy" } }, + }); + const agent = { call, did: "did:plc:test" }; + + const result = await readRepoRecordWithLegacyFallback( + agent as never, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ); + + expect(result).toEqual({ + collection: LEGACY_PROFILE_COLLECTION, + cid: "legacy-cid", + record: { displayName: "Legacy" }, + }); + expect(call).toHaveBeenNthCalledWith(2, "com.atproto.repo.getRecord", { + repo: "did:plc:test", + collection: LEGACY_PROFILE_COLLECTION, + rkey: "self", + }); + }); + + it("returns null when neither namespace has a record", async () => { + const call = jest + .fn() + .mockRejectedValue({ error: "RecordNotFound", message: "missing" }); + const agent = { call, did: "did:plc:test" }; + + await expect( + readRepoRecordWithLegacyFallback( + agent as never, + STABLE_PROFILE_COLLECTION, + LEGACY_PROFILE_COLLECTION, + ), + ).resolves.toBeNull(); + }); +}); + +describe("getBlobHash", () => { + it("extracts hashes from legacy repo blob refs", () => { + expect( + getBlobHash({ ref: { $link: "bafkreilegacy" }, mimeType: "image/jpeg" }), + ).toBe("bafkreilegacy"); + }); + + it("keeps string hashes compatible with appview profile responses", () => { + expect(getBlobHash("bafkreiappview")).toBe("bafkreiappview"); + }); +}); diff --git a/apps/amethyst/lib/atp/onboardingRecords.ts b/apps/amethyst/lib/atp/onboardingRecords.ts new file mode 100644 index 0000000..b372b91 --- /dev/null +++ b/apps/amethyst/lib/atp/onboardingRecords.ts @@ -0,0 +1,92 @@ +export const STABLE_PROFILE_COLLECTION = "fm.teal.actor.profile"; +export const LEGACY_PROFILE_COLLECTION = "fm.teal.alpha.actor.profile"; +export const STABLE_PROFILE_STATUS_COLLECTION = "fm.teal.actor.profileStatus"; +export const LEGACY_PROFILE_STATUS_COLLECTION = + "fm.teal.alpha.actor.profileStatus"; + +export interface RepoRecordAgent { + call: ( + methodNsid: string, + params?: Record, + ) => Promise<{ data: { cid?: string; value: unknown } }>; + did?: string; +} + +export interface RepoRecord { + collection: string; + cid?: string; + record: T; +} + +/** + * Read a self record from the stable namespace, falling back to the old + * namespace for users whose repositories have not been migrated yet. + */ +export async function readRepoRecordWithLegacyFallback( + agent: RepoRecordAgent, + stableCollection: string, + legacyCollection: string, +): Promise | null> { + try { + return await readRepoRecord(agent, stableCollection); + } catch (stableError) { + try { + return await readRepoRecord(agent, legacyCollection); + } catch (legacyError) { + if (isRecordNotFound(stableError) || isRecordNotFound(legacyError)) { + return null; + } + + throw stableError; + } + } +} + +function readRepoRecord( + agent: RepoRecordAgent, + collection: string, +): Promise> { + return agent + .call("com.atproto.repo.getRecord", { + repo: agent.did, + collection, + rkey: "self", + }) + .then((response) => ({ + collection, + cid: response.data.cid, + record: response.data.value as T, + })); +} + +function isRecordNotFound(error: unknown): boolean { + if (!error || typeof error !== "object") return false; + + const details = error as { error?: unknown; message?: unknown }; + const errorCode = typeof details.error === "string" ? details.error : ""; + const message = typeof details.message === "string" ? details.message : ""; + + return ( + errorCode === "RecordNotFound" || + /(?:record(?:.*)?(?:not found|does not exist)|(?:not found|does not exist).*record|could not find record)/i.test( + message, + ) + ); +} + +export function getBlobHash(blob: unknown): string | undefined { + if (typeof blob === "string") return blob; + if (!blob || typeof blob !== "object") return undefined; + + const ref = (blob as { ref?: unknown }).ref; + if (typeof ref === "string") return ref; + if (ref && typeof ref === "object") { + const link = (ref as { $link?: unknown }).$link; + if (typeof link === "string") return link; + + const stringified = String(ref); + if (stringified !== "[object Object]") return stringified; + } + + return undefined; +} diff --git a/services/cadet/src/ingestors/car/README.md b/services/cadet/src/ingestors/car/README.md index a9789e9..fc24cfa 100644 --- a/services/cadet/src/ingestors/car/README.md +++ b/services/cadet/src/ingestors/car/README.md @@ -73,6 +73,10 @@ The CAR importer automatically detects and processes these Teal record types: - **`fm.teal.actor.profile`** - User profile data - **`fm.teal.actor.status`** - User status updates +Historical CAR files using the legacy `fm.teal.alpha.*` collections are also +accepted. They are dispatched through the stable ingestors and stored with +stable `fm.teal.*` AT URIs. + Records are processed using the same logic as real-time Jetstream ingestion, ensuring data consistency. ## Architecture diff --git a/services/cadet/src/ingestors/car/car_import.rs b/services/cadet/src/ingestors/car/car_import.rs index a254b6c..7c5ee52 100644 --- a/services/cadet/src/ingestors/car/car_import.rs +++ b/services/cadet/src/ingestors/car/car_import.rs @@ -52,6 +52,59 @@ pub struct ExtractedRecord { pub data: serde_json::Value, } +const STABLE_PLAY_COLLECTION: &str = "fm.teal.feed.play"; +const STABLE_PROFILE_COLLECTION: &str = "fm.teal.actor.profile"; +const STABLE_STATUS_COLLECTION: &str = "fm.teal.actor.status"; +const LEGACY_PLAY_COLLECTION: &str = "fm.teal.alpha.feed.play"; +const LEGACY_PROFILE_COLLECTION: &str = "fm.teal.alpha.actor.profile"; +const LEGACY_STATUS_COLLECTION: &str = "fm.teal.alpha.actor.status"; + +/// Return the stable collection handled by the existing Teal ingestor. +/// +/// Historical CAR files contain records under the alpha namespace. Their +/// records have the same shape as the stable records, so they can be passed +/// through the stable ingestors after the collection is normalized. +fn stable_collection_for(collection: &str) -> Option<&'static str> { + match collection { + STABLE_PLAY_COLLECTION | LEGACY_PLAY_COLLECTION => Some(STABLE_PLAY_COLLECTION), + STABLE_PROFILE_COLLECTION | LEGACY_PROFILE_COLLECTION => Some(STABLE_PROFILE_COLLECTION), + STABLE_STATUS_COLLECTION | LEGACY_STATUS_COLLECTION => Some(STABLE_STATUS_COLLECTION), + _ => None, + } +} + +fn is_teal_record_key(key: &str) -> bool { + let Some((collection, rkey)) = key.rsplit_once('/') else { + return false; + }; + + !rkey.is_empty() && stable_collection_for(collection).is_some() +} + +fn parse_teal_key(key: &str) -> Option<(String, String)> { + if !is_teal_record_key(key) { + return None; + } + + let (collection, rkey) = key.rsplit_once('/')?; + 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) +} + /// CAR Import Ingestor handles importing Teal records from CAR files using atmst pub struct CarImportIngestor { sql: PgPool, @@ -159,9 +212,9 @@ impl CarImportIngestor { match result { Ok((key, record_cid)) => { // Check if this is a Teal record based on the key pattern - if self.is_teal_record_key(&key) { + if is_teal_record_key(&key) { info!("🎵 Found Teal record: {} -> {}", key, record_cid); - if let Some((collection, rkey)) = self.parse_teal_key(&key) { + if let Some((collection, rkey)) = parse_teal_key(&key) { info!(" Collection: {}, rkey: {}", collection, rkey); // Get the actual record data using the CID match self.get_record_data(&record_cid, car_importer).await { @@ -247,8 +300,8 @@ impl CarImportIngestor { "🔄 Processing {} record: {}", record.collection, record.rkey ); - match record.collection.as_str() { - "fm.teal.feed.play" => { + match stable_collection_for(&record.collection) { + Some(STABLE_PLAY_COLLECTION) => { info!(" 📀 Processing play record..."); let result = self .process_play_record(&record.data, did, &record.rkey) @@ -260,7 +313,7 @@ impl CarImportIngestor { } result } - "fm.teal.actor.profile" => { + Some(STABLE_PROFILE_COLLECTION) => { info!(" 👤 Processing profile record..."); let result = self .process_profile_record(&record.data, did, &record.rkey) @@ -272,7 +325,7 @@ impl CarImportIngestor { } result } - "fm.teal.actor.status" => { + Some(STABLE_STATUS_COLLECTION) => { info!(" 📢 Processing status record..."); let result = self .process_status_record(&record.data, did, &record.rkey) @@ -291,29 +344,14 @@ impl CarImportIngestor { } } - /// Check if a key represents a Teal record - fn is_teal_record_key(&self, key: &str) -> bool { - key.starts_with("fm.teal.") && key.contains("/") - } - - /// Parse a Teal MST key to extract collection and rkey - fn parse_teal_key(&self, key: &str) -> Option<(String, String)> { - if let Some(slash_pos) = key.rfind('/') { - let collection = key[..slash_pos].to_string(); - let rkey = key[slash_pos + 1..].to_string(); - Some((collection, rkey)) - } else { - None - } - } - /// Process a play record using the existing PlayIngestor async fn process_play_record(&self, data: &Value, did: &str, rkey: &str) -> Result<()> { + let data = normalize_legacy_record_type(data); let play_record: types::fm_teal::feed::play::Play = - value::from_json_value::(data.clone())?; + value::from_json_value::(data)?; let play_ingestor = super::super::teal::feed_play::PlayIngestor::new(self.sql.clone()); - let uri = super::super::teal::assemble_at_uri(did, "fm.teal.feed.play", rkey); + let uri = super::super::teal::assemble_at_uri(did, STABLE_PLAY_COLLECTION, rkey); play_ingestor .insert_play( @@ -334,8 +372,9 @@ impl CarImportIngestor { /// Process a profile record using the existing ActorProfileIngestor async fn process_profile_record(&self, data: &Value, did: &str, _rkey: &str) -> Result<()> { + let data = normalize_legacy_record_type(data); let profile_record: types::fm_teal::actor::profile::Profile = - value::from_json_value::(data.clone())?; + value::from_json_value::(data)?; let profile_ingestor = super::super::teal::actor_profile::ActorProfileIngestor::new(self.sql.clone()); @@ -353,8 +392,9 @@ impl CarImportIngestor { /// Process a status record using the existing ActorStatusIngestor async fn process_status_record(&self, data: &Value, did: &str, rkey: &str) -> Result<()> { + let data = normalize_legacy_record_type(data); let status_record: types::fm_teal::actor::status::Status = - value::from_json_value::(data.clone())?; + value::from_json_value::(data)?; let status_ingestor = super::super::teal::actor_status::ActorStatusIngestor::new(self.sql.clone()); @@ -652,32 +692,70 @@ mod tests { #[test] fn test_parse_teal_key() { - // This test doesn't need a database connection or async - let key = "fm.teal.feed.play/3k2akjdlkjsf"; - - // Test the parsing logic directly - if let Some(slash_pos) = key.rfind('/') { - let collection = key[..slash_pos].to_string(); - let rkey = key[slash_pos + 1..].to_string(); + for collection in [STABLE_PLAY_COLLECTION, LEGACY_PLAY_COLLECTION] { + let key = format!("{collection}/3k2akjdlkjsf"); + let (parsed_collection, rkey) = parse_teal_key(&key).expect("valid Teal key"); - assert_eq!(collection, "fm.teal.feed.play"); + assert_eq!(parsed_collection, collection); assert_eq!(rkey, "3k2akjdlkjsf"); - } else { - panic!("Should have found slash in key"); } } #[test] fn test_is_teal_record_key() { - // Test the logic directly without needing an ingestor instance - fn is_teal_record_key(key: &str) -> bool { - key.starts_with("fm.teal.") && key.contains("/") - } - assert!(is_teal_record_key("fm.teal.feed.play/abc123")); assert!(is_teal_record_key("fm.teal.actor.profile/def456")); + assert!(is_teal_record_key("fm.teal.alpha.feed.play/abc123")); + assert!(is_teal_record_key("fm.teal.alpha.actor.profile/def456")); + assert!(is_teal_record_key("fm.teal.alpha.actor.status/ghi789")); assert!(!is_teal_record_key("app.bsky.feed.post/xyz789")); assert!(!is_teal_record_key("fm.teal.feed.play")); // No rkey + assert!(!is_teal_record_key("fm.teal.feed.play/")); // Empty rkey + } + + #[test] + fn test_legacy_collections_route_to_stable_ingestors() { + let collections = [ + ( + LEGACY_PLAY_COLLECTION, + STABLE_PLAY_COLLECTION, + "3k2akjdlkjsf", + ), + (LEGACY_PROFILE_COLLECTION, STABLE_PROFILE_COLLECTION, "self"), + (LEGACY_STATUS_COLLECTION, STABLE_STATUS_COLLECTION, "self"), + ]; + + for (legacy_collection, stable_collection, rkey) in collections { + let key = format!("{legacy_collection}/{rkey}"); + let (parsed_collection, parsed_rkey) = + parse_teal_key(&key).expect("legacy Teal record should be retained"); + + assert_eq!(parsed_collection, legacy_collection); + assert_eq!(parsed_rkey, rkey); + assert_eq!( + stable_collection_for(&parsed_collection), + Some(stable_collection) + ); + assert_eq!( + crate::ingestors::teal::assemble_at_uri( + "did:plc:historical", + stable_collection_for(&parsed_collection).unwrap(), + &parsed_rkey, + ), + format!("at://did:plc:historical/{stable_collection}/{rkey}"), + ); + } + } + + #[test] + fn test_normalize_legacy_record_type() { + let data = serde_json::json!({ + "$type": "fm.teal.alpha.actor.status", + "time": "2024-01-01T00:00:00Z" + }); + + let normalized = normalize_legacy_record_type(&data); + assert_eq!(normalized["$type"], "fm.teal.actor.status"); } #[test] diff --git a/services/cadet/src/ingestors/teal/feed_play.rs b/services/cadet/src/ingestors/teal/feed_play.rs index ab20e8b..7a47b7e 100644 --- a/services/cadet/src/ingestors/teal/feed_play.rs +++ b/services/cadet/src/ingestors/teal/feed_play.rs @@ -1600,7 +1600,11 @@ impl LexiconIngestor for PlayIngestor { // TODO: verify cid self.insert_play( &record, - &assemble_at_uri(&message.did, &commit.collection, &commit.rkey), + &assemble_at_uri( + &message.did, + crate::ingestors::teal::canonical_collection(&commit.collection), + &commit.rkey, + ), cid, &message.did, &commit.rkey, @@ -1610,7 +1614,12 @@ impl LexiconIngestor for PlayIngestor { } } else { println!("{}: Message {} deleted", message.did, commit.rkey); - self.remove_play(&message.did).await?; + let uri = assemble_at_uri( + &message.did, + crate::ingestors::teal::canonical_collection(&commit.collection), + &commit.rkey, + ); + self.remove_play(&uri).await?; } } else { return Err(anyhow!("Message has no commit")); diff --git a/services/cadet/src/ingestors/teal/mod.rs b/services/cadet/src/ingestors/teal/mod.rs index 5f402ac..daa0b7c 100644 --- a/services/cadet/src/ingestors/teal/mod.rs +++ b/services/cadet/src/ingestors/teal/mod.rs @@ -2,6 +2,42 @@ pub mod actor_profile; pub mod actor_status; pub mod feed_play; +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"; + +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"; + +const COLLECTION_ALIASES: [(&str, &str); 3] = [ + (ALPHA_FEED_PLAY, STABLE_FEED_PLAY), + (ALPHA_ACTOR_PROFILE, STABLE_ACTOR_PROFILE), + (ALPHA_ACTOR_STATUS, STABLE_ACTOR_STATUS), +]; + +pub fn canonical_collection(collection: &str) -> &str { + COLLECTION_ALIASES + .iter() + .find_map(|(alias, stable)| (*alias == collection).then_some(*stable)) + .unwrap_or(collection) +} + +pub fn wanted_collections() -> Vec { + [ + STABLE_FEED_PLAY, + STABLE_ACTOR_PROFILE, + STABLE_ACTOR_STATUS, + ALPHA_FEED_PLAY, + ALPHA_ACTOR_PROFILE, + ALPHA_ACTOR_STATUS, + "com.atproto.repo.importRepo", + ] + .into_iter() + .map(str::to_string) + .collect() +} + /// Parses an AT uri into parts: /// did/handle, collection, rkey // fn parse_at_parts(aturi: &str) -> (&str, Option<&str>, Option<&str>) { @@ -15,3 +51,35 @@ pub mod feed_play; pub fn assemble_at_uri(did: &str, collection: &str, rkey: &str) -> String { format!("at://{did}/{collection}/{rkey}") } + +#[cfg(test)] +mod tests { + use super::{canonical_collection, wanted_collections}; + + #[test] + fn canonicalizes_alpha_collections_to_stable_names() { + for (alpha, stable) in [ + ("fm.teal.alpha.feed.play", "fm.teal.feed.play"), + ("fm.teal.alpha.actor.profile", "fm.teal.actor.profile"), + ("fm.teal.alpha.actor.status", "fm.teal.actor.status"), + ] { + assert_eq!(canonical_collection(alpha), stable); + } + } + + #[test] + fn wanted_collections_include_stable_and_alpha_records() { + let wanted = wanted_collections(); + + for collection in [ + "fm.teal.feed.play", + "fm.teal.actor.profile", + "fm.teal.actor.status", + "fm.teal.alpha.feed.play", + "fm.teal.alpha.actor.profile", + "fm.teal.alpha.actor.status", + ] { + assert!(wanted.iter().any(|wanted| wanted == collection)); + } + } +} diff --git a/services/cadet/src/main.rs b/services/cadet/src/main.rs index 7594abf..b942549 100644 --- a/services/cadet/src/main.rs +++ b/services/cadet/src/main.rs @@ -48,41 +48,46 @@ async fn main() { .expect("Could not get PostgreSQL pool"); let opts = JetstreamOptions::builder() - .wanted_collections( - [ - "fm.teal.feed.play", - "fm.teal.actor.profile", - "fm.teal.actor.status", - "com.atproto.repo.importRepo", - ] - .iter() - .map(|collection| collection.to_string()) - .collect(), - ) + .wanted_collections(ingestors::teal::wanted_collections()) .build(); let jetstream = JetstreamConnection::new(opts); let mut ingestors: HashMap> = HashMap::new(); - ingestors.insert( - "fm.teal.feed.play".to_string(), - Box::new(ingestors::teal::feed_play::PlayIngestor::new(pool.clone())), - ); + for collection in [ + ingestors::teal::STABLE_FEED_PLAY, + ingestors::teal::ALPHA_FEED_PLAY, + ] { + ingestors.insert( + collection.to_string(), + Box::new(ingestors::teal::feed_play::PlayIngestor::new(pool.clone())), + ); + } - ingestors.insert( - "fm.teal.actor.profile".to_string(), - Box::new(ingestors::teal::actor_profile::ActorProfileIngestor::new( - pool.clone(), - )), - ); + for collection in [ + ingestors::teal::STABLE_ACTOR_PROFILE, + ingestors::teal::ALPHA_ACTOR_PROFILE, + ] { + ingestors.insert( + collection.to_string(), + Box::new(ingestors::teal::actor_profile::ActorProfileIngestor::new( + pool.clone(), + )), + ); + } - ingestors.insert( - "fm.teal.actor.status".to_string(), - Box::new(ingestors::teal::actor_status::ActorStatusIngestor::new( - pool.clone(), - )), - ); + for collection in [ + ingestors::teal::STABLE_ACTOR_STATUS, + ingestors::teal::ALPHA_ACTOR_STATUS, + ] { + ingestors.insert( + collection.to_string(), + Box::new(ingestors::teal::actor_status::ActorStatusIngestor::new( + pool.clone(), + )), + ); + } ingestors.insert( "com.atproto.repo.importRepo".to_string(),