//! Tests for bulk_resolve SQL queries //! //! These tests validate that the SQL queries used by bulk_resolve functions //! are syntactically correct against the current database schema. #![allow(clippy::unseparated_literal_suffix)] mod common; use chrono::Utc; use common::*; use consumer::db::{bulk_resolve, ensure_actor_id}; use consumer::database_writer::bulk_types::PostCopyData; use consumer::db::bulk_copy; use parakeet_db::types::ActorStatus; use eyre::WrapErr; #[tokio::test] async fn test_resolve_repost_uris_query() -> eyre::Result<()> { let pool = test_pool(); let conn = pool.get().await.wrap_err("Failed to get connection")?; // Test with empty arrays (valid SQL, should return no results) let dids: Vec = vec![]; let rkeys: Vec = vec![]; let result = bulk_resolve::queries::resolve_repost_uris(&conn, &dids, &rkeys).await; assert!( result.is_ok(), "resolve_repost_uris query should be valid SQL: {:?}", result.err() ); Ok(()) } #[tokio::test] async fn test_resolve_repost_uris_with_data_query() -> eyre::Result<()> { let pool = test_pool(); let conn = pool.get().await.wrap_err("Failed to get connection")?; // Test with sample data (won't find anything, but SQL should be valid) let dids: Vec = vec!["did:plc:test1".to_string(), "did:plc:test2".to_string()]; let rkeys: Vec = vec![1234567890, 9876543210]; let result = bulk_resolve::queries::resolve_repost_uris(&conn, &dids, &rkeys).await; assert!( result.is_ok(), "resolve_repost_uris query with data should be valid SQL: {:?}", result.err() ); // Should return empty vec since test data doesn't exist assert_eq!(result.unwrap().len(), 0); Ok(()) } // NOTE: test_create_repost_stubs_query removed - we no longer create repost stubs with natural keys // ============================================================================ // Actor DID Resolution Tests // ============================================================================ #[tokio::test] async fn test_resolve_actor_dids_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let dids: Vec<&str> = vec![]; let result = bulk_resolve::resolve_actor_dids_bulk(&txn, &dids).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_actor_dids_bulk_existing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Create test actors let actor1_id = ensure_actor_id(&txn, "did:plc:test1", Some(&ActorStatus::Active), None, Utc::now()).await?; let actor2_id = ensure_actor_id(&txn, "did:plc:test2", Some(&ActorStatus::Active), None, Utc::now()).await?; // Resolve them let dids = vec!["did:plc:test1", "did:plc:test2"]; let result = bulk_resolve::resolve_actor_dids_bulk(&txn, &dids).await?; assert_eq!(result.len(), 2); assert_eq!(result.get("did:plc:test1"), Some(&actor1_id)); assert_eq!(result.get("did:plc:test2"), Some(&actor2_id)); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_actor_dids_bulk_missing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Create one actor let actor1_id = ensure_actor_id(&txn, "did:plc:test1", Some(&ActorStatus::Active), None, Utc::now()).await?; // Try to resolve two actors (one exists, one doesn't) let dids = vec!["did:plc:test1", "did:plc:missing"]; let result = bulk_resolve::resolve_actor_dids_bulk(&txn, &dids).await?; assert_eq!(result.len(), 1, "Should only return existing actor"); assert_eq!(result.get("did:plc:test1"), Some(&actor1_id)); assert_eq!(result.get("did:plc:missing"), None); txn.rollback().await?; Ok(()) } // ============================================================================ // Actor Stub Creation Tests // ============================================================================ #[tokio::test] async fn test_create_actor_stubs_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let dids: Vec<&str> = vec![]; let result = bulk_resolve::create_actor_stubs_bulk(&txn, &dids).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_create_actor_stubs_bulk_new() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Create stubs for new actors let dids = vec!["did:plc:stub1", "did:plc:stub2"]; let result = bulk_resolve::create_actor_stubs_bulk(&txn, &dids).await?; assert_eq!(result.len(), 2); assert!(result.contains_key("did:plc:stub1")); assert!(result.contains_key("did:plc:stub2")); // Verify actors exist in database let check = bulk_resolve::resolve_actor_dids_bulk(&txn, &dids).await?; assert_eq!(check.len(), 2); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_create_actor_stubs_bulk_idempotent() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let dids = vec!["did:plc:stub1"]; // Create once let result1 = bulk_resolve::create_actor_stubs_bulk(&txn, &dids).await?; assert_eq!(result1.len(), 1); let actor_id1 = *result1.get("did:plc:stub1").unwrap(); // Create again (ON CONFLICT DO NOTHING means no rows returned) let result2 = bulk_resolve::create_actor_stubs_bulk(&txn, &dids).await?; assert_eq!(result2.len(), 0, "Second call should return empty due to ON CONFLICT DO NOTHING"); // Verify actor still exists with same ID let check = bulk_resolve::resolve_actor_dids_bulk(&txn, &dids).await?; assert_eq!(check.get("did:plc:stub1"), Some(&actor_id1), "Actor should still exist with same ID"); txn.rollback().await?; Ok(()) } // ============================================================================ // Post URI Resolution Tests // ============================================================================ #[tokio::test] async fn test_resolve_post_uris_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let uris: Vec<&str> = vec![]; let result = bulk_resolve::resolve_post_uris_bulk(&txn, &uris).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_post_uris_bulk_existing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Create test actor and posts let actor_id = ensure_actor_id(&txn, "did:plc:author", Some(&ActorStatus::Active), None, Utc::now()).await?; let post1 = PostCopyData { actor_id, rkey: 1746247680000000000, // Jan 15 2024, properly formatted as TID (timestamp << 10) cid: vec![1u8, 2, 3, 4], content_compressed: None, langs: vec![], tags: vec![], parent_post_actor_id: None, parent_post_rkey: None, root_post_actor_id: None, root_post_rkey: None, embed_type: None, embed_subtype: None, violates_threadgate: false, tokens: vec![], ext_embed: None, video_embed: None, image_1: None, image_2: None, image_3: None, image_4: None, embedded_post_actor_id: None, embedded_post_rkey: None, record_detached: None, facet_1: None, facet_2: None, facet_3: None, facet_4: None, facet_5: None, facet_6: None, facet_7: None, facet_8: None, mentions: None, }; let posts = bulk_copy::copy_posts(&txn, vec![post1]).await?; let (post_actor_id, post_rkey) = (posts[0].actor_id, posts[0].rkey); let post_uri = &posts[0].post_uri; // Resolve the post let uris = vec![post_uri.as_str()]; let result = bulk_resolve::resolve_post_uris_bulk(&txn, &uris).await?; assert_eq!(result.len(), 1); assert_eq!(result.get(post_uri.as_str()), Some(&(post_actor_id, post_rkey))); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_post_uris_bulk_missing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // With natural keys, we can resolve URIs even if actors/posts don't exist // The function creates actor stubs if needed and returns natural keys let uris = vec!["at://did:plc:missing/app.bsky.feed.post/3l6kdoqxe7k2a"]; let result = bulk_resolve::resolve_post_uris_bulk(&txn, &uris).await?; assert_eq!(result.len(), 1, "Should return natural key even if post doesn't exist"); assert!(result.contains_key("at://did:plc:missing/app.bsky.feed.post/3l6kdoqxe7k2a")); txn.rollback().await?; Ok(()) } // ============================================================================ // Post Stub Creation Tests // ============================================================================ // NOTE: test_create_post_stubs_bulk tests removed - we no longer create post stubs with natural keys // ============================================================================ // Resolve and Ensure Posts Tests // ============================================================================ #[tokio::test] async fn test_resolve_and_ensure_posts_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let uris_with_cids: Vec<(&str, &str)> = vec![]; let result = bulk_resolve::resolve_and_ensure_posts_bulk(&txn, &uris_with_cids).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_and_ensure_posts_bulk_existing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Create test actor and post let actor_id = ensure_actor_id(&txn, "did:plc:author", Some(&ActorStatus::Active), None, Utc::now()).await?; let post_data = PostCopyData { actor_id, rkey: 1746247680000000000, // Jan 15 2024, properly formatted as TID (timestamp << 10) cid: vec![1u8, 2, 3, 4], content_compressed: None, langs: vec![], tags: vec![], parent_post_actor_id: None, parent_post_rkey: None, root_post_actor_id: None, root_post_rkey: None, embed_type: None, embed_subtype: None, violates_threadgate: false, tokens: vec![], ext_embed: None, video_embed: None, image_1: None, image_2: None, image_3: None, image_4: None, embedded_post_actor_id: None, embedded_post_rkey: None, record_detached: None, facet_1: None, facet_2: None, facet_3: None, facet_4: None, facet_5: None, facet_6: None, facet_7: None, facet_8: None, mentions: None, }; let posts = bulk_copy::copy_posts(&txn, vec![post_data]).await?; let (post_actor_id, post_rkey) = (posts[0].actor_id, posts[0].rkey); let post_uri = &posts[0].post_uri; // Resolve and ensure (should find existing) let cid_str = test_cid().to_string(); let uris_with_cids = vec![(post_uri.as_str(), cid_str.as_str())]; let result = bulk_resolve::resolve_and_ensure_posts_bulk(&txn, &uris_with_cids).await?; assert_eq!(result.len(), 1); assert_eq!(result.get(post_uri.as_str()), Some(&(post_actor_id, post_rkey))); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_and_ensure_posts_bulk_creates_stubs() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Ensure actor exists ensure_actor_id(&txn, "did:plc:author", Some(&ActorStatus::Active), None, Utc::now()).await?; // Resolve and ensure for a post that doesn't exist (should create stub) let cid_str = test_cid().to_string(); let uris_with_cids = vec![ ("at://did:plc:author/app.bsky.feed.post/3l6kdoqxe7k2a", cid_str.as_str()), ]; let result = bulk_resolve::resolve_and_ensure_posts_bulk(&txn, &uris_with_cids).await?; assert_eq!(result.len(), 1); assert!(result.contains_key("at://did:plc:author/app.bsky.feed.post/3l6kdoqxe7k2a")); txn.rollback().await?; Ok(()) } // ============================================================================ // Feedgen URI Resolution Tests // ============================================================================ #[tokio::test] async fn test_resolve_feedgen_uris_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let uris: Vec<&str> = vec![]; let result = bulk_resolve::resolve_feedgen_uris_bulk(&txn, &uris).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_feedgen_uris_bulk_missing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Try to resolve feedgen that doesn't exist let uris = vec!["at://did:plc:feedgen/app.bsky.feed.generator/test"]; let result = bulk_resolve::resolve_feedgen_uris_bulk(&txn, &uris).await?; assert_eq!(result.len(), 0, "Missing feedgen should not be in result"); txn.rollback().await?; Ok(()) } // ============================================================================ // Labeler DID Resolution Tests // ============================================================================ #[tokio::test] async fn test_resolve_labeler_dids_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let dids: Vec<&str> = vec![]; let result = bulk_resolve::resolve_labeler_dids_bulk(&txn, &dids).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_labeler_dids_bulk_missing() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Try to resolve labeler that doesn't exist let dids = vec!["did:plc:labeler"]; let result = bulk_resolve::resolve_labeler_dids_bulk(&txn, &dids).await?; assert_eq!(result.len(), 0, "Missing labeler should not be in result"); txn.rollback().await?; Ok(()) } // ============================================================================ // Resolve and Ensure Reposts Tests // ============================================================================ #[tokio::test] async fn test_resolve_and_ensure_reposts_bulk_empty() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let repost_data: Vec<(&str, &str, &str, &str)> = vec![]; let result = bulk_resolve::resolve_and_ensure_reposts_bulk(&txn, &repost_data).await?; assert_eq!(result.len(), 0, "Empty input should return empty map"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_resolve_and_ensure_reposts_bulk_creates_stubs() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; // Ensure author actor exists ensure_actor_id(&txn, "did:plc:author", Some(&ActorStatus::Active), None, Utc::now()).await?; // Resolve and ensure a repost that doesn't exist (should create stub) let cid_str = test_cid().to_string(); let repost_data = vec![ ( "at://did:plc:author/app.bsky.feed.repost/3l6kdoqxe7k2a", cid_str.as_str(), "at://did:plc:author/app.bsky.feed.post/3l6kdoqxe7k2b", cid_str.as_str() ), ]; let result = bulk_resolve::resolve_and_ensure_reposts_bulk(&txn, &repost_data).await?; assert_eq!(result.len(), 1); assert!(result.contains_key("at://did:plc:author/app.bsky.feed.repost/3l6kdoqxe7k2a")); txn.rollback().await?; Ok(()) }