diff --git a/Cargo.toml b/Cargo.toml index 17ba4cb..5cfe408 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,9 +20,9 @@ sqlx = { version = "0.8", features = [ "chrono", "tls-rustls", ] } -jacquard-common = { version = "0.12.0-beta.2", features = ["reqwest-client"] } -jacquard-derive = { version = "0.12.0-beta.2" } -jacquard-lexicon = { version = "0.12.0-beta.2" } +jacquard-common = { version = "0.12.1", features = ["reqwest-client"] } +jacquard-derive = { version = "0.12.1" } +jacquard-lexicon = { version = "0.12.1" } serde = { version = "1.0", features = ["derive"] } anyhow = "1.0" serde_json = "1.0" diff --git a/apps/amethyst/lib/__tests__/onboardingRecords.test.ts b/apps/amethyst/lib/__tests__/onboardingRecords.test.ts index 9ce8ca7..65c522a 100644 --- a/apps/amethyst/lib/__tests__/onboardingRecords.test.ts +++ b/apps/amethyst/lib/__tests__/onboardingRecords.test.ts @@ -2,18 +2,34 @@ import { getBlobHash, LEGACY_PROFILE_COLLECTION, readRepoRecordWithLegacyFallback, + type RepoRecordAgent, STABLE_PROFILE_COLLECTION, } from "../atp/onboardingRecords"; +type MockRepoRecordAgent = RepoRecordAgent & { + call: jest.MockedFunction; +}; + +function createAgent(): MockRepoRecordAgent { + return { + call: jest.fn< + ReturnType, + Parameters + >(), + did: "did:plc:test", + }; +} + describe("readRepoRecordWithLegacyFallback", () => { it("uses the stable record when it exists", async () => { - const call = jest.fn().mockResolvedValue({ + const agent = createAgent(); + const { call } = agent; + call.mockResolvedValue({ data: { cid: "stable-cid", value: { displayName: "Stable" } }, }); - const agent = { call, did: "did:plc:test" }; const result = await readRepoRecordWithLegacyFallback( - agent as never, + agent, STABLE_PROFILE_COLLECTION, LEGACY_PROFILE_COLLECTION, ); @@ -32,16 +48,16 @@ describe("readRepoRecordWithLegacyFallback", () => { }); it("falls back to the legacy record when stable is missing", async () => { - const call = jest - .fn() + const agent = createAgent(); + const { call } = agent; + call .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, + agent, STABLE_PROFILE_COLLECTION, LEGACY_PROFILE_COLLECTION, ); @@ -59,18 +75,23 @@ describe("readRepoRecordWithLegacyFallback", () => { }); 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" }; + const agent = createAgent(); + const { call } = agent; + call.mockRejectedValue({ error: "RecordNotFound", message: "missing" }); await expect( readRepoRecordWithLegacyFallback( - agent as never, + agent, STABLE_PROFILE_COLLECTION, LEGACY_PROFILE_COLLECTION, ), ).resolves.toBeNull(); + expect(call).toHaveBeenCalledTimes(2); + expect(call).toHaveBeenNthCalledWith(2, "com.atproto.repo.getRecord", { + repo: "did:plc:test", + collection: LEGACY_PROFILE_COLLECTION, + rkey: "self", + }); }); it("propagates a legacy read failure when the stable record is missing", async () => { @@ -78,15 +99,15 @@ describe("readRepoRecordWithLegacyFallback", () => { error: "NetworkError", message: "legacy read failed", }; - const call = jest - .fn() + const agent = createAgent(); + const { call } = agent; + call .mockRejectedValueOnce({ error: "RecordNotFound", message: "missing" }) .mockRejectedValueOnce(legacyError); - const agent = { call, did: "did:plc:test" }; await expect( readRepoRecordWithLegacyFallback( - agent as never, + agent, STABLE_PROFILE_COLLECTION, LEGACY_PROFILE_COLLECTION, ), @@ -98,15 +119,15 @@ describe("readRepoRecordWithLegacyFallback", () => { error: "NetworkError", message: "stable read failed", }; - const call = jest - .fn() + const agent = createAgent(); + const { call } = agent; + call .mockRejectedValueOnce(stableError) .mockRejectedValueOnce({ error: "RecordNotFound", message: "missing" }); - const agent = { call, did: "did:plc:test" }; await expect( readRepoRecordWithLegacyFallback( - agent as never, + agent, STABLE_PROFILE_COLLECTION, LEGACY_PROFILE_COLLECTION, ), diff --git a/services/cadet/src/ingestors/car/car_import.rs b/services/cadet/src/ingestors/car/car_import.rs index de03018..0582fe4 100644 --- a/services/cadet/src/ingestors/car/car_import.rs +++ b/services/cadet/src/ingestors/car/car_import.rs @@ -43,6 +43,7 @@ use redis::AsyncCommands; use rocketman::{ingestion::LexiconIngestor, types::event::Event}; use serde_json::Value; use sqlx::PgPool; +use std::collections::HashMap; use tracing::{info, warn}; /// Helper struct for extracted records @@ -88,6 +89,49 @@ fn stable_collection_for(collection: &str) -> Option<&'static str> { } } +fn is_legacy_collection(collection: &str) -> bool { + collection.starts_with("fm.teal.alpha.") +} + +/// Drop duplicate stable/alpha records from a repository snapshot. +/// +/// Records are duplicates when their canonical collection, rkey, and +/// normalized JSON are identical. Stable records win when both namespace +/// variants are present. +fn deduplicate_records(records: Vec) -> Vec { + let mut deduplicated: Vec<(ExtractedRecord, Value)> = Vec::with_capacity(records.len()); + let mut positions_by_key: HashMap<(&'static str, String), Vec> = HashMap::new(); + + for record in records { + let normalized_data = normalize_legacy_record_type(&record.data); + let Some(canonical_collection) = stable_collection_for(&record.collection) else { + deduplicated.push((record, normalized_data)); + continue; + }; + let key = (canonical_collection, record.rkey.clone()); + let duplicate_index = positions_by_key.get(&key).and_then(|positions| { + positions + .iter() + .copied() + .find(|&index| deduplicated[index].1 == normalized_data) + }); + + if let Some(index) = duplicate_index { + if is_legacy_collection(&deduplicated[index].0.collection) + && !is_legacy_collection(&record.collection) + { + deduplicated[index] = (record, normalized_data); + } + } else { + let index = deduplicated.len(); + deduplicated.push((record, normalized_data)); + positions_by_key.entry(key).or_default().push(index); + } + } + + deduplicated.into_iter().map(|(record, _)| record).collect() +} + /// CAR Import Ingestor handles importing Teal records from CAR files using atmst pub struct CarImportIngestor { sql: PgPool, @@ -151,8 +195,14 @@ impl CarImportIngestor { let records = self .extract_records_from_mst(&mst, &data_importer, did) .await?; + let extracted_count = records.len(); + let records = deduplicate_records(records); - info!("Extracted {} records from MST", records.len()); + info!( + "Extracted {} records from MST ({} after namespace deduplication)", + extracted_count, + records.len() + ); // Process each record through the appropriate ingestor let mut processed_count = 0_usize; @@ -879,6 +929,72 @@ mod tests { assert!(!is_teal_record_key("fm.teal.alpha.feed.play")); // No rkey } + #[test] + fn test_deduplicate_identical_alpha_and_stable_records_prefers_stable() { + let record_fields = serde_json::json!({ + "trackName": "Same recording", + "playedTime": "2024-01-01T00:00:00Z" + }); + let mut legacy_data = record_fields.clone(); + legacy_data["$type"] = serde_json::json!("fm.teal.alpha.feed.play"); + let mut stable_data = record_fields; + stable_data["$type"] = serde_json::json!("fm.teal.feed.play"); + + let records = deduplicate_records(vec![ + ExtractedRecord { + collection: "fm.teal.alpha.feed.play".to_string(), + rkey: "same-rkey".to_string(), + cid: "legacy-cid".to_string(), + data: legacy_data, + }, + ExtractedRecord { + collection: "fm.teal.feed.play".to_string(), + rkey: "same-rkey".to_string(), + cid: "stable-cid".to_string(), + data: stable_data, + }, + ]); + + assert_eq!(records.len(), 1); + assert_eq!(records[0].collection, "fm.teal.feed.play"); + assert_eq!(records[0].cid, "stable-cid"); + } + + #[test] + fn test_deduplicate_compares_only_matching_collection_and_rkey() { + let records = deduplicate_records(vec![ + ExtractedRecord { + collection: "fm.teal.alpha.actor.status".to_string(), + rkey: "same-rkey".to_string(), + cid: "status-cid".to_string(), + data: serde_json::json!({ + "$type": "fm.teal.alpha.actor.status", + "time": "2024-01-01T00:00:00Z" + }), + }, + ExtractedRecord { + collection: "fm.teal.feed.play".to_string(), + rkey: "same-rkey".to_string(), + cid: "play-cid".to_string(), + data: serde_json::json!({ + "$type": "fm.teal.feed.play", + "time": "2024-01-01T00:00:00Z" + }), + }, + ExtractedRecord { + collection: "fm.teal.actor.status".to_string(), + rkey: "different-rkey".to_string(), + cid: "other-status-cid".to_string(), + data: serde_json::json!({ + "$type": "fm.teal.actor.status", + "time": "2024-01-01T00:00:00Z" + }), + }, + ]); + + assert_eq!(records.len(), 3); + } + #[test] fn test_ipld_to_json_conversion() { // Test IPLD to JSON conversion logic directly diff --git a/todo.md b/todo.md index 55dda3e..ca8fca2 100644 --- a/todo.md +++ b/todo.md @@ -18,6 +18,7 @@ Last synced with GitHub and Linear issues: 2026-06-14. - Dark mode preview refresh (2026-08-11): added the cohesive dark palette across Amethyst shells, feeds, social surfaces, profile/stats views, manual stamping, onboarding, web chrome, and native tab colors. Rebuilt the named Cloudflare preview from `codex/manual-listens` after cherry-picking `07e4af5`; `https://sigilyph.teal.fm` metadata and latest-listens XRPC verification pass. - Manual theme toggle (2026-08-11): the Teal shell now exposes an accessible light/dark toggle on desktop and mobile; Settings retains the System option, and web theme choices persist via localStorage. - Focused code-health pass (2026-07-10): fixed `did:web` path resolution and resolver error handling in Cadet, avoided 10-second Postgres connection churn in Aqua/Cadet/Satellite, returned proper 404s for missing feed plays, and preserved delimiter-containing artist names in Satellite's latest-play response. Verified with targeted Rust checks and Cadet resolver tests. +- Review follow-up (2026-08-13): indexed CAR namespace deduplication by canonical collection/rkey, strengthened onboarding legacy-fallback tests with typed repo-agent mocks, and aligned workspace Jacquard dependencies with the 0.12.1 generator. Aqua actor-feed pagination and Cadet play-link cleanup were already implemented and required no changes. - Focused Amethyst code-health pass (2026-07-10): moved color-scheme initialization out of render, made Escape handling safe on native clients, kept independent home and right-rail requests visible when a sibling request fails, and migrated linting to ESLint 9 flat config while ignoring generated Expo output. Verified with TypeScript, focused lint, Jest, and a web export built against `https://sigilyph.teal.fm`. - Dependency refresh (2026-07-10): updated the Rust lockfile, root and standalone lexicon CLI pnpm locks, Expo SDK 57/RN 0.86, AT Protocol clients and lexicon generator, plus current compatible workspace tooling. Regenerated lexicons now normalize TypeScript relative imports for Metro; Amethyst record creation supplies required `$type` fields. Verified with offline Rust tests/checks, TypeScript, Jest, and full workspace builds. - Cadet live-ingestion recovery (2026-07-11): the preview consumer was repeatedly stalling because every incoming play synchronously refreshed four materialized views. The Compose Cadet service now enables its existing deferred-refresh mode so Jetstream events can drain without blocking for several seconds per play. Follow-up: add a periodic materialized-view refresh path before relying on live aggregate counts.