//! Tests for bulk_copy SQL queries //! //! These tests validate that the SQL queries used by bulk_copy functions //! are syntactically correct against the current database schema. //! #![allow(clippy::unseparated_literal_suffix)] //! IMPORTANT: These tests run against the parakeet_test database and verify //! that all SQL operations compile and execute without schema errors. mod common; use common::*; use consumer::db::bulk_copy::{self, queries}; use consumer::database_writer::bulk_types::{FollowCopyData, BlockCopyData, PostCopyData, PostLikeCopyData, RepostCopyData}; use consumer::db::ensure_actor_id; use parakeet_db::types::ActorStatus; use chrono::Utc; use eyre::WrapErr; #[tokio::test] async fn test_insert_posts_from_staging_query() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction to create staging table let txn = conn.transaction().await?; // Create staging table using extracted function queries::create_posts_staging_table(&txn).await?; // Test with empty staging table (valid SQL, should return no results) let result = queries::insert_posts_from_staging(&txn).await; assert!( result.is_ok(), "insert_posts_from_staging query should be valid SQL: {:?}", result.err() ); // Should return empty vec since staging table is empty assert_eq!(result.unwrap().len(), 0); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_posts_from_staging_with_data() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction let txn = conn.transaction().await?; // Create test actor using source code function let actor_id = ensure_actor_id( &txn, "did:plc:test123", Some(&ActorStatus::Active), None, Utc::now(), ).await?; // Use production bulk_copy function with real data type let post_data = PostCopyData { actor_id, rkey: 1234567890, cid: vec![1u8, 2, 3, 4], content_compressed: Some(vec![10u8, 20, 30]), langs: vec!["en".to_string(), "es".to_string()], tags: vec!["test".to_string()], parent_post_actor_id: None, parent_post_rkey: None, root_post_actor_id: None, root_post_rkey: None, embed_type: Some("images".to_string()), embed_subtype: None, violates_threadgate: false, tokens: vec!["hello".to_string(), "world".to_string()], 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 result = bulk_copy::copy_posts(&txn, vec![post_data]).await; assert!( result.is_ok(), "copy_posts should work with type casts: {:?}", result.err() ); // Should return 1 inserted post let posts = result.unwrap(); assert_eq!(posts.len(), 1); assert_eq!(posts[0].rkey, 1234567890); assert!(posts[0].post_uri.contains("did:plc:test123")); assert_eq!(posts[0].actor_id, actor_id); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_post_likes_from_staging_query() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction to create staging table let txn = conn.transaction().await?; // Create staging table using extracted function queries::create_post_likes_staging_table(&txn).await?; // Test with empty staging table (valid SQL, should return empty vec) let result = bulk_copy::copy_post_likes(&txn, vec![]).await; assert!( result.is_ok(), "copy_post_likes query should be valid SQL: {:?}", result.err() ); // Should return empty vec since no data provided assert_eq!(result.unwrap().len(), 0); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_reposts_from_staging_query() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction to create staging table let txn = conn.transaction().await?; // Create staging table using extracted function queries::create_reposts_staging_table(&txn).await?; // Test with empty staging table (valid SQL, should return no results) let result = queries::insert_reposts_from_staging(&txn).await; assert!( result.is_ok(), "insert_reposts_from_staging query should be valid SQL: {:?}", result.err() ); // Should return empty vec since staging table is empty assert_eq!(result.unwrap().len(), 0); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_follows_from_staging_with_data() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction let txn = conn.transaction().await?; // Create test actors using source code function let actor_id = ensure_actor_id( &txn, "did:plc:follower123", Some(&ActorStatus::Active), None, Utc::now(), ).await?; let subject_actor_id = ensure_actor_id( &txn, "did:plc:subject456", Some(&ActorStatus::Active), None, Utc::now(), ).await?; // Use production bulk_copy function with real data type let follow_data = FollowCopyData { actor_id, rkey: 1234567890, subject_actor_id, created_at: Utc::now(), }; let result = bulk_copy::copy_follows(&txn, vec![follow_data]).await; assert!( result.is_ok(), "copy_follows should work: {:?}", result.err() ); // Should return 1 inserted follow let follows = result.unwrap(); assert_eq!(follows.len(), 1); assert_eq!(follows[0].actor_id, actor_id); assert_eq!(follows[0].subject_actor_id, subject_actor_id); assert_eq!(follows[0].subject_did, "did:plc:subject456"); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_blocks_from_staging_with_data() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction let txn = conn.transaction().await?; // Create test actors using source code function let actor_id = ensure_actor_id( &txn, "did:plc:blocker123", Some(&ActorStatus::Active), None, Utc::now(), ).await?; let subject_actor_id = ensure_actor_id( &txn, "did:plc:blocked456", Some(&ActorStatus::Active), None, Utc::now(), ).await?; // Use production bulk_copy function with real data type let block_data = BlockCopyData { actor_id, rkey: 9876543210, subject_actor_id, created_at: Utc::now(), }; let result = bulk_copy::copy_blocks(&txn, vec![block_data]).await; assert!( result.is_ok(), "copy_blocks should work: {:?}", result.err() ); // Should return 1 inserted block let blocks = result.unwrap(); assert_eq!(blocks.len(), 1); assert_eq!(blocks[0].actor_id, actor_id); assert_eq!(blocks[0].subject_actor_id, subject_actor_id); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_post_likes_from_staging_with_data() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction let txn = conn.transaction().await?; // Create test actors using source code function let actor_id = ensure_actor_id( &txn, "did:plc:liker123", Some(&ActorStatus::Active), None, Utc::now(), ).await?; let post_author_id = ensure_actor_id( &txn, "did:plc:postauthor456", Some(&ActorStatus::Active), None, Utc::now(), ).await?; // Create a test post using production bulk_copy function let post_data = PostCopyData { actor_id: post_author_id, rkey: 9999999999, cid: vec![5u8, 6, 7, 8], 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 = posts[0].actor_id; let post_rkey = posts[0].rkey; // Use production bulk_copy function for post likes let like_data = PostLikeCopyData { actor_id, rkey: 1234567890, post_actor_id, post_rkey, via_repost_actor_id: None, via_repost_rkey: None, }; let result = bulk_copy::copy_post_likes(&txn, vec![like_data]).await; assert!( result.is_ok(), "copy_post_likes should work with type casts: {:?}", result.err() ); // Should return 1 inserted like let likes = result.unwrap(); assert_eq!(likes.len(), 1); assert_eq!(likes[0].actor_id, actor_id); assert_eq!(likes[0].post_actor_id, post_actor_id); assert_eq!(likes[0].post_rkey, post_rkey); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_reposts_from_staging_with_data() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction let txn = conn.transaction().await?; // Create test actors using source code function let reposter_id = ensure_actor_id( &txn, "did:plc:reposter123", Some(&ActorStatus::Active), None, Utc::now(), ).await?; let post_author_id = ensure_actor_id( &txn, "did:plc:origauthor456", Some(&ActorStatus::Active), None, Utc::now(), ).await?; // Create a test post using production bulk_copy function let post_data = PostCopyData { actor_id: post_author_id, rkey: 8888888888, cid: vec![9u8, 10, 11, 12], 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 = posts[0].actor_id; let post_rkey = posts[0].rkey; // Use production bulk_copy function for reposts let repost_data = RepostCopyData { actor_id: reposter_id, rkey: 5555555555, post_actor_id, post_rkey, cid: vec![13u8, 14, 15, 16], created_at: Utc::now(), via_repost_actor_id: None, via_repost_rkey: None, }; let result = bulk_copy::copy_reposts(&txn, vec![repost_data]).await; assert!( result.is_ok(), "copy_reposts should work with URI generation: {:?}", result.err() ); // Should return 1 inserted repost let reposts = result.unwrap(); assert_eq!(reposts.len(), 1); assert_eq!(reposts[0].post_actor_id, post_actor_id); assert_eq!(reposts[0].post_rkey, post_rkey); assert!(reposts[0].post_uri.contains("did:plc:origauthor456")); assert!(reposts[0].post_uri.contains("at://")); assert!(reposts[0].post_uri.contains("/app.bsky.feed.post/")); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_posts_with_video_embed_and_captions() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; // Begin transaction let txn = conn.transaction().await?; // Create test actor let actor_id = ensure_actor_id( &txn, "did:plc:videotest", Some(&ActorStatus::Active), None, Utc::now(), ).await?; // Create VideoEmbed with all 3 captions to test serialization let video_embed = parakeet_db::composite::VideoEmbed { mime_type: parakeet_db::types::VideoMimeType::Mp4, cid: vec![1u8, 2, 3, 4, 5, 6, 7, 8], alt: Some("Test video".to_string()), width: Some(1920), height: Some(1080), caption_1: Some(parakeet_db::composite::VideoCaption { lang: parakeet_db::types::LanguageCode::En, mime_type: parakeet_db::types::CaptionMimeType::Vtt, cid: vec![11u8, 12, 13, 14], }), caption_2: Some(parakeet_db::composite::VideoCaption { lang: parakeet_db::types::LanguageCode::Es, mime_type: parakeet_db::types::CaptionMimeType::Vtt, cid: vec![21u8, 22, 23, 24], }), caption_3: Some(parakeet_db::composite::VideoCaption { lang: parakeet_db::types::LanguageCode::Fr, mime_type: parakeet_db::types::CaptionMimeType::Vtt, cid: vec![31u8, 32, 33, 34], }), }; let post_data = PostCopyData { actor_id, rkey: 9876543210, cid: vec![99u8, 98, 97, 96], content_compressed: Some(vec![5u8, 6, 7]), langs: vec!["en".to_string()], tags: vec![], parent_post_actor_id: None, parent_post_rkey: None, root_post_actor_id: None, root_post_rkey: None, embed_type: Some("video".to_string()), embed_subtype: None, violates_threadgate: false, tokens: vec!["video".to_string(), "test".to_string()], ext_embed: None, video_embed: Some(video_embed), 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, }; // This test validates that VideoEmbed serialization with captions works correctly // Previously failed with "wrong number of columns: 5, expected 8" error let result = bulk_copy::copy_posts(&txn, vec![post_data]).await; assert!( result.is_ok(), "copy_posts should work with VideoEmbed containing captions: {:?}", result.err() ); // Should return 1 inserted post let posts = result.unwrap(); assert_eq!(posts.len(), 1); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_posts_with_image_embeds() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let actor_id = ensure_actor_id(&txn, "did:plc:imagetest", Some(&ActorStatus::Active), None, Utc::now()).await?; let post_data = PostCopyData { actor_id, rkey: 1111111111, cid: vec![1u8, 2, 3, 4], content_compressed: Some(vec![5u8, 6, 7]), langs: vec!["en".to_string()], tags: vec![], parent_post_actor_id: None, parent_post_rkey: None, root_post_actor_id: None, root_post_rkey: None, embed_type: Some("images".to_string()), embed_subtype: None, violates_threadgate: false, tokens: vec![], ext_embed: None, video_embed: None, image_1: Some(parakeet_db::composite::ImageEmbed { mime_type: parakeet_db::types::ImageMimeType::Jpeg, cid: vec![11u8, 12, 13, 14], alt: Some("Test image".to_string()), width: Some(1920), height: Some(1080), }), image_2: Some(parakeet_db::composite::ImageEmbed { mime_type: parakeet_db::types::ImageMimeType::Png, cid: vec![21u8, 22, 23, 24], alt: None, width: Some(800), height: Some(600), }), 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 result = bulk_copy::copy_posts(&txn, vec![post_data]).await; assert!(result.is_ok(), "copy_posts should work with ImageEmbed: {:?}", result.err()); assert_eq!(result.unwrap().len(), 1); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_posts_with_external_embed() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let actor_id = ensure_actor_id(&txn, "did:plc:exttest", Some(&ActorStatus::Active), None, Utc::now()).await?; let post_data = PostCopyData { actor_id, rkey: 2222222222, cid: vec![1u8, 2, 3, 4], content_compressed: Some(vec![5u8, 6, 7]), langs: vec!["en".to_string()], tags: vec![], parent_post_actor_id: None, parent_post_rkey: None, root_post_actor_id: None, root_post_rkey: None, embed_type: Some("external".to_string()), embed_subtype: None, violates_threadgate: false, tokens: vec![], ext_embed: Some(parakeet_db::composite::ExtEmbed { uri: "https://example.com".to_string(), title: Some("Example Site".to_string()), description: Some("Test description".to_string()), thumb_mime_type: Some(parakeet_db::types::ImageMimeType::Jpeg), thumb_cid: Some(vec![31u8, 32, 33, 34]), }), 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 result = bulk_copy::copy_posts(&txn, vec![post_data]).await; assert!(result.is_ok(), "copy_posts should work with ExtEmbed: {:?}", result.err()); assert_eq!(result.unwrap().len(), 1); txn.rollback().await?; Ok(()) } #[tokio::test] async fn test_insert_posts_with_facets() -> eyre::Result<()> { let pool = test_pool(); let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let txn = conn.transaction().await?; let actor_id = ensure_actor_id(&txn, "did:plc:facettest", Some(&ActorStatus::Active), None, Utc::now()).await?; let mention_actor_id = ensure_actor_id(&txn, "did:plc:mentioned", Some(&ActorStatus::Active), None, Utc::now()).await?; let post_data = PostCopyData { actor_id, rkey: 3333333333, cid: vec![1u8, 2, 3, 4], content_compressed: Some(vec![5u8, 6, 7]), langs: vec!["en".to_string()], tags: vec!["test".to_string()], 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: Some(parakeet_db::composite::FacetEmbed { facet_type: parakeet_db::types::FacetType::Link, index_start: 0, index_end: 10, link_uri: Some("https://example.com".to_string()), mention_actor_id: None, tag: None, }), facet_2: Some(parakeet_db::composite::FacetEmbed { facet_type: parakeet_db::types::FacetType::Mention, index_start: 11, index_end: 20, link_uri: None, mention_actor_id: Some(mention_actor_id), tag: None, }), facet_3: Some(parakeet_db::composite::FacetEmbed { facet_type: parakeet_db::types::FacetType::Tag, index_start: 21, index_end: 30, link_uri: None, mention_actor_id: None, tag: Some("test".to_string()), }), facet_4: None, facet_5: None, facet_6: None, facet_7: None, facet_8: None, mentions: Some(vec![Some(mention_actor_id)]), }; let result = bulk_copy::copy_posts(&txn, vec![post_data]).await; assert!(result.is_ok(), "copy_posts should work with FacetEmbed: {:?}", result.err()); assert_eq!(result.unwrap().len(), 1); txn.rollback().await?; Ok(()) }