diff --git a/.sqlx/query-41ef76f85e8b2d37ff906c16973ce08c03219332c86d6f1b7bfa752f0d75275b.json b/.sqlx/query-41ef76f85e8b2d37ff906c16973ce08c03219332c86d6f1b7bfa752f0d75275b.json deleted file mode 100644 index 60d554a..0000000 --- a/.sqlx/query-41ef76f85e8b2d37ff906c16973ce08c03219332c86d6f1b7bfa752f0d75275b.json +++ /dev/null @@ -1,21 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "\n INSERT INTO profiles (did, handle, display_name, description, description_facets, avatar, banner, created_at)\n VALUES ($1, $2, $3, $4, $5, $6, $7, $8)\n ON CONFLICT (did) DO UPDATE SET\n handle = EXCLUDED.handle,\n display_name = EXCLUDED.display_name,\n description = EXCLUDED.description,\n description_facets = EXCLUDED.description_facets,\n avatar = EXCLUDED.avatar,\n banner = EXCLUDED.banner,\n created_at = EXCLUDED.created_at;\n ", - "describe": { - "columns": [], - "parameters": { - "Left": [ - "Text", - "Text", - "Text", - "Text", - "Jsonb", - "Text", - "Text", - "Timestamptz" - ] - }, - "nullable": [] - }, - "hash": "41ef76f85e8b2d37ff906c16973ce08c03219332c86d6f1b7bfa752f0d75275b" -} diff --git a/.sqlx/query-504c321eec8f39add55b4baba18de9e4f0f4e43cebdf62bf010a8a9094decc50.json b/.sqlx/query-504c321eec8f39add55b4baba18de9e4f0f4e43cebdf62bf010a8a9094decc50.json new file mode 100644 index 0000000..62343d5 --- /dev/null +++ b/.sqlx/query-504c321eec8f39add55b4baba18de9e4f0f4e43cebdf62bf010a8a9094decc50.json @@ -0,0 +1,23 @@ +{ + "db_name": "PostgreSQL", + "query": "SELECT uri FROM plays WHERE did = $1 AND cid = $2 LIMIT 1", + "describe": { + "columns": [ + { + "ordinal": 0, + "name": "uri", + "type_info": "Text" + } + ], + "parameters": { + "Left": [ + "Text", + "Text" + ] + }, + "nullable": [ + false + ] + }, + "hash": "504c321eec8f39add55b4baba18de9e4f0f4e43cebdf62bf010a8a9094decc50" +} diff --git a/migrations/20241220000014_deduplicate_plays.sql b/migrations/20241220000014_deduplicate_plays.sql new file mode 100644 index 0000000..98f9eb1 --- /dev/null +++ b/migrations/20241220000014_deduplicate_plays.sql @@ -0,0 +1,23 @@ +-- Deduplicate plays by (did, cid) - keep the record with the earliest rkey +-- This addresses the issue where the same play content (CID) was stored under multiple rkeys + +-- Create temp table with URIs to delete +CREATE TEMP TABLE uris_to_delete AS +WITH duplicates AS ( + SELECT uri, + ROW_NUMBER() OVER (PARTITION BY did, cid ORDER BY rkey ASC) as rn + FROM plays +) +SELECT uri FROM duplicates WHERE rn > 1; + +-- Delete from related tables first +DELETE FROM play_to_artists_extended WHERE play_uri IN (SELECT uri FROM uris_to_delete); +DELETE FROM play_to_artists WHERE play_uri IN (SELECT uri FROM uris_to_delete); +-- Then delete from plays +DELETE FROM plays WHERE uri IN (SELECT uri FROM uris_to_delete); + +-- Drop the temporary table +DROP TABLE uris_to_delete; + +-- Add unique constraint to prevent future duplicates +ALTER TABLE plays ADD CONSTRAINT uq_plays_did_cid UNIQUE (did, cid); diff --git a/services/cadet/src/ingestors/teal/feed_play.rs b/services/cadet/src/ingestors/teal/feed_play.rs index b42e9d8..a7f53bb 100644 --- a/services/cadet/src/ingestors/teal/feed_play.rs +++ b/services/cadet/src/ingestors/teal/feed_play.rs @@ -1490,6 +1490,29 @@ impl PlayIngestor { .as_ref() .map(ToString::to_string); + // Check if we already have this play content (CID) for this user (DID) + // This prevents duplicates when the same content is published under different rkeys + let existing_uri = sqlx::query_scalar!( + "SELECT uri FROM plays WHERE did = $1 AND cid = $2 LIMIT 1", + did, + cid + ) + .fetch_optional(&self.sql) + .await?; + + if let Some(existing) = existing_uri { + if existing != uri { + tracing::debug!( + "Skipping duplicate play for did={} cid={} (existing uri={}, new uri={})", + did, + cid, + existing, + uri + ); + } + return Ok(()); + } + sqlx::query!( r#" INSERT INTO plays (