diff --git a/consumer/tests/actor_operations_test.rs b/consumer/tests/actor_operations_test.rs --- a/consumer/tests/actor_operations_test.rs +++ b/consumer/tests/actor_operations_test.rs @@ -24,7 +24,6 @@ async fn test_profile_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -33,7 +32,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:profiletest1") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:profiletest1") .await .wrap_err("Failed to ensure actor")?; @@ -95,7 +94,6 @@ async fn test_profile_upsert_update() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -104,7 +102,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:profiletest2") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:profiletest2") .await .wrap_err("Failed to ensure actor")?; @@ -176,7 +174,6 @@ async fn test_profile_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -185,7 +182,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:profiletest3") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:profiletest3") .await .wrap_err("Failed to ensure actor")?; @@ -242,7 +239,6 @@ async fn test_status_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -251,7 +247,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:statustest1") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:statustest1") .await .wrap_err("Failed to ensure actor")?; @@ -295,7 +291,6 @@ async fn test_status_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -304,7 +299,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:statustest2") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:statustest2") .await .wrap_err("Failed to ensure actor")?; @@ -355,7 +350,6 @@ async fn test_chat_decl_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -364,7 +358,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:chattest1") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:chattest1") .await .wrap_err("Failed to ensure actor")?; @@ -403,7 +397,6 @@ async fn test_chat_decl_upsert_update() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -412,7 +405,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:chattest2") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:chattest2") .await .wrap_err("Failed to ensure actor")?; @@ -460,7 +453,6 @@ async fn test_chat_decl_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -469,7 +461,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:chattest3") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:chattest3") .await .wrap_err("Failed to ensure actor")?; @@ -517,7 +509,6 @@ async fn test_notif_decl_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -526,7 +517,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:notiftest1") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:notiftest1") .await .wrap_err("Failed to ensure actor")?; @@ -569,7 +560,6 @@ async fn test_notif_decl_upsert_update() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -578,7 +568,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:notiftest2") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:notiftest2") .await .wrap_err("Failed to ensure actor")?; @@ -630,7 +620,6 @@ async fn test_notif_decl_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -639,7 +628,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:notiftest3") + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:notiftest3") .await .wrap_err("Failed to ensure actor")?; @@ -694,7 +683,6 @@ async fn test_actor_set_repo_state() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -703,7 +691,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:repostate1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:repostate1") .await .wrap_err("Failed to ensure actor")?; @@ -777,7 +765,6 @@ async fn test_actor_get_repo_rev() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -786,7 +773,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:repostate2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:repostate2") .await .wrap_err("Failed to ensure actor")?; diff --git a/consumer/tests/cleanup_worker_test.rs b/consumer/tests/cleanup_worker_test.rs --- a/consumer/tests/cleanup_worker_test.rs +++ b/consumer/tests/cleanup_worker_test.rs @@ -13,7 +13,7 @@ mod common; use chrono::{Duration, Utc}; -use common::{test_cid, test_pool, test_redis}; +use common::{test_cid, test_pool}; use consumer::database_writer::EventSource; use consumer::db::operations::feed; use consumer::types::records::{AppBskyEmbed, AppBskyEmbedRecord, AppBskyFeedLike, AppBskyFeedPost, AppBskyFeedRepost, ReplyRef}; @@ -24,7 +24,6 @@ /// Note: Caller must ensure actor exists via get_actor_id() before calling this async fn create_post_with_tid( tx: &C, - redis: &mut redis::aio::MultiplexedConnection, did: &str, tid: &str, text: &str, @@ -41,7 +40,7 @@ }; let uri = format!("at://{}/app.bsky.feed.post/{}", did, tid); - feed::post_insert(tx, redis, &uri, did, test_cid(), post, EventSource::Jetstream) + feed::post_insert(tx, &uri, did, test_cid(), post, EventSource::Jetstream) .await .wrap_err("Failed to insert post")?; @@ -52,7 +51,6 @@ /// Note: Caller must ensure actor exists via get_actor_id() before calling this async fn create_like_with_rkey( tx: &C, - redis: &mut redis::aio::MultiplexedConnection, did: &str, rkey: &str, subject_uri: &str, @@ -66,7 +64,7 @@ created_at: Utc::now(), }; - feed::like_insert(tx, redis, rkey, did, test_cid(), like, EventSource::Jetstream) + feed::like_insert(tx, rkey, did, test_cid(), like, EventSource::Jetstream) .await .wrap_err("Failed to insert like")?; @@ -77,7 +75,6 @@ /// Note: Caller must ensure actor exists via get_actor_id() before calling this async fn create_repost_with_rkey( tx: &C, - redis: &mut redis::aio::MultiplexedConnection, did: &str, rkey: &str, subject_uri: &str, @@ -91,7 +88,7 @@ created_at: Utc::now(), }; - feed::repost_insert(tx, redis, rkey, did, test_cid(), repost, EventSource::Jetstream) + feed::repost_insert(tx, rkey, did, test_cid(), repost, EventSource::Jetstream) .await .wrap_err("Failed to insert repost")?; @@ -106,26 +103,25 @@ async fn test_cleanup_deletes_old_likes() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actors and a post to like - feed::get_actor_id(&tx, &mut redis, "did:plc:poster").await?; - feed::get_actor_id(&tx, &mut redis, "did:plc:liker").await?; + feed::get_actor_id(&tx, "did:plc:poster").await?; + feed::get_actor_id(&tx, "did:plc:liker").await?; - let post_uri = create_post_with_tid(&tx, &mut redis, "did:plc:poster", "3l7mkz4lmk235", "Test post").await?; + let post_uri = create_post_with_tid(&tx, "did:plc:poster", "3l7mkz4lmk235", "Test post").await?; // Create old like (35 days old) - should be deleted let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - create_like_with_rkey(&tx, &mut redis, "did:plc:liker", &old_tid, &post_uri).await?; + create_like_with_rkey(&tx, "did:plc:liker", &old_tid, &post_uri).await?; // Create recent like (5 days old) - should be kept let recent_time = Utc::now() - Duration::days(5); let recent_tid = parakeet_db::tid_util::timestamp_to_tid(recent_time); - create_like_with_rkey(&tx, &mut redis, "did:plc:liker", &recent_tid, &post_uri).await?; + create_like_with_rkey(&tx, "did:plc:liker", &recent_tid, &post_uri).await?; // Verify both likes exist let like_count_before: i64 = tx @@ -171,26 +167,25 @@ async fn test_cleanup_deletes_old_reposts() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actors and a post to repost - feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor").await?; - feed::get_actor_id(&tx, &mut redis, "did:plc:reposter").await?; + feed::get_actor_id(&tx, "did:plc:postauthor").await?; + feed::get_actor_id(&tx, "did:plc:reposter").await?; - let post_uri = create_post_with_tid(&tx, &mut redis, "did:plc:postauthor", "3l7mkz4lmk236", "Post to repost").await?; + let post_uri = create_post_with_tid(&tx, "did:plc:postauthor", "3l7mkz4lmk236", "Post to repost").await?; // Create old repost (35 days old) - should be deleted let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - create_repost_with_rkey(&tx, &mut redis, "did:plc:reposter", &old_tid, &post_uri).await?; + create_repost_with_rkey(&tx, "did:plc:reposter", &old_tid, &post_uri).await?; // Create recent repost (5 days old) - should be kept let recent_time = Utc::now() - Duration::days(5); let recent_tid = parakeet_db::tid_util::timestamp_to_tid(recent_time); - create_repost_with_rkey(&tx, &mut redis, "did:plc:reposter", &recent_tid, &post_uri).await?; + create_repost_with_rkey(&tx, "did:plc:reposter", &recent_tid, &post_uri).await?; // Verify both reposts exist let repost_count_before: i64 = tx @@ -236,18 +231,17 @@ async fn test_cleanup_marks_old_unreferenced_posts() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actor - feed::get_actor_id(&tx, &mut redis, "did:plc:oldposter").await?; + feed::get_actor_id(&tx, "did:plc:oldposter").await?; // Create old post (35 days old) with NO recent references - should be marked pruned let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - create_post_with_tid(&tx, &mut redis, "did:plc:oldposter", &old_tid, "Old unreferenced post").await?; + create_post_with_tid(&tx, "did:plc:oldposter", &old_tid, "Old unreferenced post").await?; // Verify post status is 'complete' before cleanup let status_before: String = tx @@ -333,24 +327,23 @@ async fn test_cleanup_preserves_posts_with_recent_likes() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actors - feed::get_actor_id(&tx, &mut redis, "did:plc:oldposter2").await?; - feed::get_actor_id(&tx, &mut redis, "did:plc:liker2").await?; + feed::get_actor_id(&tx, "did:plc:oldposter2").await?; + feed::get_actor_id(&tx, "did:plc:liker2").await?; // Create old post (35 days old) let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - let post_uri = create_post_with_tid(&tx, &mut redis, "did:plc:oldposter2", &old_tid, "Old post with recent like").await?; + let post_uri = create_post_with_tid(&tx, "did:plc:oldposter2", &old_tid, "Old post with recent like").await?; // Add a RECENT like (5 days old) - this should prevent pruning let recent_time = Utc::now() - Duration::days(5); let recent_like_tid = parakeet_db::tid_util::timestamp_to_tid(recent_time); - create_like_with_rkey(&tx, &mut redis, "did:plc:liker2", &recent_like_tid, &post_uri).await?; + create_like_with_rkey(&tx, "did:plc:liker2", &recent_like_tid, &post_uri).await?; // Run cleanup query to mark old posts as pruned (30 days) let cutoff_time = Utc::now() - Duration::days(30); @@ -424,24 +417,23 @@ async fn test_cleanup_preserves_posts_with_recent_reposts() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actors - feed::get_actor_id(&tx, &mut redis, "did:plc:oldposter3").await?; - feed::get_actor_id(&tx, &mut redis, "did:plc:reposter3").await?; + feed::get_actor_id(&tx, "did:plc:oldposter3").await?; + feed::get_actor_id(&tx, "did:plc:reposter3").await?; // Create old post (35 days old) let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - let post_uri = create_post_with_tid(&tx, &mut redis, "did:plc:oldposter3", &old_tid, "Old post with recent repost").await?; + let post_uri = create_post_with_tid(&tx, "did:plc:oldposter3", &old_tid, "Old post with recent repost").await?; // Add a RECENT repost (5 days old) - this should prevent pruning let recent_time = Utc::now() - Duration::days(5); let recent_repost_tid = parakeet_db::tid_util::timestamp_to_tid(recent_time); - create_repost_with_rkey(&tx, &mut redis, "did:plc:reposter3", &recent_repost_tid, &post_uri).await?; + create_repost_with_rkey(&tx, "did:plc:reposter3", &recent_repost_tid, &post_uri).await?; // Run cleanup query to mark old posts as pruned (30 days) let cutoff_time = Utc::now() - Duration::days(30); @@ -515,19 +507,18 @@ async fn test_cleanup_preserves_posts_with_recent_replies() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actors - feed::get_actor_id(&tx, &mut redis, "did:plc:oldposter4").await?; - feed::get_actor_id(&tx, &mut redis, "did:plc:replier4").await?; + feed::get_actor_id(&tx, "did:plc:oldposter4").await?; + feed::get_actor_id(&tx, "did:plc:replier4").await?; // Create old parent post (35 days old) let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - let parent_uri = create_post_with_tid(&tx, &mut redis, "did:plc:oldposter4", &old_tid, "Old parent post").await?; + let parent_uri = create_post_with_tid(&tx, "did:plc:oldposter4", &old_tid, "Old parent post").await?; // Add a RECENT reply (5 days old) - this should prevent pruning of parent let recent_time = Utc::now() - Duration::days(5); @@ -554,7 +545,7 @@ }; let reply_uri = format!("at://did:plc:replier4/app.bsky.feed.post/{}", recent_tid); - feed::post_insert(&tx, &mut redis, &reply_uri, "did:plc:replier4", test_cid(), reply_post, EventSource::Jetstream).await?; + feed::post_insert(&tx, &reply_uri, "did:plc:replier4", test_cid(), reply_post, EventSource::Jetstream).await?; // Run cleanup query to mark old posts as pruned (30 days) let cutoff_time = Utc::now() - Duration::days(30); @@ -628,19 +619,18 @@ async fn test_cleanup_preserves_posts_with_recent_quotes() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actors - feed::get_actor_id(&tx, &mut redis, "did:plc:oldposter5").await?; - feed::get_actor_id(&tx, &mut redis, "did:plc:quoter5").await?; + feed::get_actor_id(&tx, "did:plc:oldposter5").await?; + feed::get_actor_id(&tx, "did:plc:quoter5").await?; // Create old quoted post (35 days old) let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - let quoted_uri = create_post_with_tid(&tx, &mut redis, "did:plc:oldposter5", &old_tid, "Old quoted post").await?; + let quoted_uri = create_post_with_tid(&tx, "did:plc:oldposter5", &old_tid, "Old quoted post").await?; // Add a RECENT quote post (5 days old) - this should prevent pruning of quoted post let recent_time = Utc::now() - Duration::days(5); @@ -658,7 +648,7 @@ }; let quote_uri = format!("at://did:plc:quoter5/app.bsky.feed.post/{}", recent_tid); - feed::post_insert(&tx, &mut redis, "e_uri, "did:plc:quoter5", test_cid(), quote_post, EventSource::Jetstream).await?; + feed::post_insert(&tx, "e_uri, "did:plc:quoter5", test_cid(), quote_post, EventSource::Jetstream).await?; // Add the embed (quote) let embed = AppBskyEmbed::Record(AppBskyEmbedRecord { @@ -668,7 +658,7 @@ }, }); - feed::post_embed_insert(&tx, &mut redis, "e_uri, embed, Utc::now(), EventSource::Jetstream).await?; + feed::post_embed_insert(&tx, "e_uri, embed, Utc::now(), EventSource::Jetstream).await?; // Run cleanup query to mark old posts as pruned (30 days) let cutoff_time = Utc::now() - Duration::days(30); @@ -742,18 +732,17 @@ async fn test_cleanup_preserves_forbidden_posts() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actor - feed::get_actor_id(&tx, &mut redis, "did:plc:forbiddenposter").await?; + feed::get_actor_id(&tx, "did:plc:forbiddenposter").await?; // Create old post (35 days old) let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - create_post_with_tid(&tx, &mut redis, "did:plc:forbiddenposter", &old_tid, "Forbidden post").await?; + create_post_with_tid(&tx, "did:plc:forbiddenposter", &old_tid, "Forbidden post").await?; // Manually mark it as 'forbidden' (moderation action) tx.execute( @@ -841,18 +830,17 @@ async fn test_cleanup_deletes_pruned_posts() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = test_pool(); - let mut redis = test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Setup: Create actor - feed::get_actor_id(&tx, &mut redis, "did:plc:prunedposter").await?; + feed::get_actor_id(&tx, "did:plc:prunedposter").await?; // Create old post let old_time = Utc::now() - Duration::days(35); let old_tid = parakeet_db::tid_util::timestamp_to_tid(old_time); - create_post_with_tid(&tx, &mut redis, "did:plc:prunedposter", &old_tid, "Post to prune").await?; + create_post_with_tid(&tx, "did:plc:prunedposter", &old_tid, "Post to prune").await?; // Manually mark it as 'pruned' tx.execute( diff --git a/consumer/tests/common.rs b/consumer/tests/common.rs --- a/consumer/tests/common.rs +++ b/consumer/tests/common.rs @@ -1,10 +1,9 @@ //! Common test utilities for database integration tests //! -//! This module provides helpers for setting up test databases, transactions, -//! and mock Redis connections for testing SQL operations. +//! This module provides helpers for setting up test databases and transactions +//! for testing SQL operations. use deadpool_postgres::{Config, ManagerConfig, Pool, RecyclingMethod, Runtime}; -use redis::aio::MultiplexedConnection; use tokio_postgres::NoTls; /// Creates a connection pool to the parakeet_test database @@ -69,25 +68,6 @@ .expect("Failed to create test database pool") } -/// Creates a test Redis connection -/// -/// This creates a real Redis connection for testing. If Redis is not available, -/// tests will fail gracefully. -/// -/// # Environment Variables -/// - `TEST_REDIS_URL`: Redis URL (default: redis://127.0.0.1:6379) -#[allow(dead_code)] -pub async fn test_redis() -> MultiplexedConnection { - let redis_url = - std::env::var("TEST_REDIS_URL").unwrap_or_else(|_| "redis://127.0.0.1:6379".to_string()); - - let client = redis::Client::open(redis_url).expect("Failed to create Redis client"); - - client - .get_multiplexed_async_connection() - .await - .expect("Failed to connect to Redis") -} /// Helper to begin a test transaction that will automatically rollback /// diff --git a/consumer/tests/community_starterpack_operations_test.rs b/consumer/tests/community_starterpack_operations_test.rs --- a/consumer/tests/community_starterpack_operations_test.rs +++ b/consumer/tests/community_starterpack_operations_test.rs @@ -23,7 +23,6 @@ #[tokio::test] async fn test_bookmark_upsert_post_subject() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -31,10 +30,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:bookmarkuser1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:bookmarkuser1") .await .wrap_err("Failed to ensure bookmark actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postauthor1") .await .wrap_err("Failed to ensure post author")?; @@ -50,7 +49,6 @@ created_at: chrono::Utc::now()}; consumer::db::operations::feed::post_insert( &tx, - &mut redis, "at://did:plc:postauthor1/app.bsky.feed.post/3l7mkz4lmk22q", "did:plc:postauthor1", test_cid(), @@ -65,7 +63,7 @@ tags: vec!["favorite".to_string(), "tech".to_string()], created_at: Utc::now()}; - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:bookmarkuser1").await?; + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:bookmarkuser1").await?; let result = community::bookmark_upsert(&tx, actor_id, "3l7mkz4lmk22r", bookmark).await; @@ -100,7 +98,6 @@ #[tokio::test] async fn test_bookmark_upsert_idempotent() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -108,10 +105,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors and create post - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:bookmarkuser2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:bookmarkuser2") .await .wrap_err("Failed to ensure bookmark actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postauthor2") .await .wrap_err("Failed to ensure post author")?; @@ -126,7 +123,6 @@ created_at: chrono::Utc::now()}; consumer::db::operations::feed::post_insert( &tx, - &mut redis, "at://did:plc:postauthor2/app.bsky.feed.post/3l7mkz4lmk22s", "did:plc:postauthor2", test_cid(), @@ -141,7 +137,7 @@ tags: vec![], created_at: Utc::now()}; - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:bookmarkuser2").await?; + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:bookmarkuser2").await?; // Insert bookmark first time let result1 = community::bookmark_upsert( &tx, @@ -182,7 +178,6 @@ #[tokio::test] async fn test_bookmark_delete() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -190,10 +185,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors and create post - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:bookmarkuser3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:bookmarkuser3") .await .wrap_err("Failed to ensure bookmark actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postauthor3") .await .wrap_err("Failed to ensure post author")?; @@ -208,7 +203,6 @@ created_at: chrono::Utc::now()}; consumer::db::operations::feed::post_insert( &tx, - &mut redis, "at://did:plc:postauthor3/app.bsky.feed.post/3l7mkz4lmk22u", "did:plc:postauthor3", test_cid(), @@ -223,7 +217,7 @@ tags: vec!["delete-me".to_string()], created_at: Utc::now()}; - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:bookmarkuser3").await?; + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:bookmarkuser3").await?; community::bookmark_upsert(&tx, actor_id, "3l7mkz4lmk22v", bookmark) .await .wrap_err("INSERT bookmark for deletion test")?; @@ -256,7 +250,6 @@ #[tokio::test] async fn test_starter_pack_upsert_insert() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -264,7 +257,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:starterpackowner1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:starterpackowner1") .await .wrap_err("Failed to ensure starterpack owner")?; @@ -279,7 +272,7 @@ created_at: Utc::now()}; graph::list_upsert( - &tx, &mut redis, + &tx, "at://did:plc:starterpackowner1/app.bsky.graph.list/3l7mkz4lmk22w", "did:plc:starterpackowner1", test_cid(), @@ -304,7 +297,7 @@ let at_uri = "at://did:plc:starterpackowner1/app.bsky.graph.starterpack/3l7mkz4lmk22x"; let result = starter_pack::starter_pack_upsert( - &tx, &mut redis, + &tx, at_uri, "did:plc:starterpackowner1", test_cid(), @@ -351,7 +344,6 @@ #[tokio::test] async fn test_starter_pack_upsert_update() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -359,7 +351,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:starterpackowner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:starterpackowner2") .await .wrap_err("Failed to ensure starterpack owner")?; @@ -374,7 +366,7 @@ created_at: Utc::now()}; graph::list_upsert( - &tx, &mut redis, + &tx, "at://did:plc:starterpackowner2/app.bsky.graph.list/3l7mkz4lmk22y", "did:plc:starterpackowner2", test_cid(), @@ -395,7 +387,7 @@ created_at: Utc::now()}; let result1 = starter_pack::starter_pack_upsert( - &tx, &mut redis, + &tx, at_uri, "did:plc:starterpackowner2", test_cid(), @@ -415,7 +407,7 @@ created_at: Utc::now()}; let result2 = starter_pack::starter_pack_upsert( - &tx, &mut redis, + &tx, at_uri, "did:plc:starterpackowner2", test_cid(), @@ -457,7 +449,6 @@ #[tokio::test] async fn test_starter_pack_delete() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -465,7 +456,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:starterpackowner3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:starterpackowner3") .await .wrap_err("Failed to ensure starterpack owner")?; @@ -480,7 +471,7 @@ created_at: Utc::now()}; graph::list_upsert( - &tx, &mut redis, + &tx, "at://did:plc:starterpackowner3/app.bsky.graph.list/3l7mkz4lmk232", "did:plc:starterpackowner3", test_cid(), @@ -501,7 +492,7 @@ created_at: Utc::now()}; starter_pack::starter_pack_upsert( - &tx, &mut redis, + &tx, at_uri, "did:plc:starterpackowner3", test_cid(), @@ -534,7 +525,6 @@ #[tokio::test] async fn test_starter_pack_with_feeds_insertion() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -542,13 +532,13 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:starterpackowner4") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:starterpackowner4") .await .wrap_err("Failed to ensure starterpack owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedowner") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedowner") .await .wrap_err("Failed to ensure feed owner")?; - let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:web:feedservice.example") + let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:web:feedservice.example") .await .wrap_err("Failed to ensure feed service")?; @@ -634,7 +624,7 @@ }; graph::list_upsert( - &tx, &mut redis, + &tx, "at://did:plc:starterpackowner4/app.bsky.graph.list/3l7mkz4lmk234", "did:plc:starterpackowner4", test_cid(), @@ -666,7 +656,7 @@ let at_uri = "at://did:plc:starterpackowner4/app.bsky.graph.starterpack/3l7mkz4lmk235"; let result = starter_pack::starter_pack_upsert( - &tx, &mut redis, + &tx, at_uri, "did:plc:starterpackowner4", test_cid(), diff --git a/consumer/tests/feed_operations_test.rs b/consumer/tests/feed_operations_test.rs --- a/consumer/tests/feed_operations_test.rs +++ b/consumer/tests/feed_operations_test.rs @@ -27,7 +27,6 @@ async fn test_post_insert_normal() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -36,7 +35,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:poster1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:poster1") .await .wrap_err("Failed to ensure actor")?; @@ -54,7 +53,6 @@ let at_uri = "at://did:plc:poster1/app.bsky.feed.post/3l7mkz4lmk234"; let result = feed::post_insert( &tx, - &mut redis, at_uri, "did:plc:poster1", test_cid(), @@ -100,7 +98,6 @@ async fn test_post_insert_with_reply() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -109,10 +106,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootauthor") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:replier1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:replier1") .await .wrap_err("Failed to ensure actor")?; @@ -128,7 +125,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:rootauthor/app.bsky.feed.post/3l7mkz4lmk235", "did:plc:rootauthor", test_cid(), @@ -158,7 +154,6 @@ let at_uri = "at://did:plc:replier1/app.bsky.feed.post/3l7mkz4lmk236"; let result = feed::post_insert( &tx, - &mut redis, at_uri, "did:plc:replier1", test_cid(), @@ -198,7 +193,6 @@ async fn test_post_insert_with_facets() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -207,10 +201,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postauthor") .await .wrap_err("Failed to ensure post author")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:mentioned") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:mentioned") .await .wrap_err("Failed to ensure mentioned actor")?; @@ -263,7 +257,6 @@ let at_uri = "at://did:plc:postauthor/app.bsky.feed.post/3l7mkz4lmk25a"; let result = feed::post_insert( &tx, - &mut redis, at_uri, "did:plc:postauthor", test_cid(), @@ -371,7 +364,6 @@ async fn test_post_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -380,7 +372,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:poster2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:poster2") .await .wrap_err("Failed to ensure actor")?; @@ -398,7 +390,6 @@ let at_uri = "at://did:plc:poster2/app.bsky.feed.post/3l7mkz4lmk237"; feed::post_insert( &tx, - &mut redis, at_uri, "did:plc:poster2", test_cid(), @@ -450,7 +441,6 @@ async fn test_like_insert_post_subject() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -459,10 +449,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postauthor") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:liker") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:liker") .await .wrap_err("Failed to ensure actor")?; @@ -478,7 +468,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:postauthor/app.bsky.feed.post/3l7mkz4lmk23a", "did:plc:postauthor", test_cid(), @@ -499,7 +488,6 @@ // Insert the like let result = feed::like_insert( &tx, - &mut redis, "3l7mkz4lmk23y", "did:plc:liker", test_cid(), @@ -540,7 +528,6 @@ async fn test_like_insert_stub_subject() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -549,7 +536,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:liker2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:liker2") .await .wrap_err("Failed to ensure actor")?; @@ -564,7 +551,6 @@ // Insert the like (should succeed and create stub post) let result = feed::like_insert( &tx, - &mut redis, "3l7mkz4lmk23z", "did:plc:liker2", test_cid(), @@ -618,7 +604,6 @@ async fn test_like_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -627,10 +612,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postauthor2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postauthor2") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:liker3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:liker3") .await .wrap_err("Failed to ensure actor")?; @@ -646,7 +631,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:postauthor2/app.bsky.feed.post/3l7mkz4lmk23c", "did:plc:postauthor2", test_cid(), @@ -666,7 +650,6 @@ feed::like_insert( &tx, - &mut redis, "3l7mkz4lmk242", "did:plc:liker3", test_cid(), @@ -717,7 +700,6 @@ async fn test_repost_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -726,10 +708,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:repostedauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:repostedauthor") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:reposter1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:reposter1") .await .wrap_err("Failed to ensure actor")?; @@ -745,7 +727,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:repostedauthor/app.bsky.feed.post/3l7mkz4lmk23d", "did:plc:repostedauthor", test_cid(), @@ -765,7 +746,6 @@ let result = feed::repost_insert( &tx, - &mut redis, "3l7mkz4lmk243", "did:plc:reposter1", test_cid(), @@ -813,7 +793,6 @@ async fn test_repost_insert_stub_post() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -822,7 +801,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:reposter2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:reposter2") .await .wrap_err("Failed to ensure actor")?; @@ -836,7 +815,6 @@ let result = feed::repost_insert( &tx, - &mut redis, "3l7mkz4lmk244", "did:plc:reposter2", test_cid(), @@ -890,7 +868,6 @@ async fn test_repost_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -899,10 +876,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:repostedauthor2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:repostedauthor2") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:reposter3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:reposter3") .await .wrap_err("Failed to ensure actor")?; @@ -918,7 +895,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:repostedauthor2/app.bsky.feed.post/3l7mkz4lmk23f", "did:plc:repostedauthor2", test_cid(), @@ -937,7 +913,6 @@ feed::repost_insert( &tx, - &mut redis, "3l7mkz4lmk245", "did:plc:reposter3", test_cid(), @@ -984,7 +959,6 @@ async fn test_repost_insert_idempotent() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -993,10 +967,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:repostedauthor3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:repostedauthor3") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:reposter4") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:reposter4") .await .wrap_err("Failed to ensure actor")?; @@ -1012,7 +986,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:repostedauthor3/app.bsky.feed.post/3l7mkz4lmk23g", "did:plc:repostedauthor3", test_cid(), @@ -1032,7 +1005,6 @@ // Insert first time let result1 = feed::repost_insert( &tx, - &mut redis, "3l7mkz4lmk246", "did:plc:reposter4", test_cid(), @@ -1052,7 +1024,6 @@ let result2 = feed::repost_insert( &tx, - &mut redis, "3l7mkz4lmk246", "did:plc:reposter4", test_cid(), @@ -1076,7 +1047,6 @@ async fn test_postgate_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1085,7 +1055,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postgateowner1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postgateowner1") .await .wrap_err("Failed to ensure actor")?; @@ -1101,7 +1071,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:postgateowner1/app.bsky.feed.post/3l7mkz4lmk23h", "did:plc:postgateowner1", test_cid(), @@ -1147,7 +1116,6 @@ async fn test_postgate_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1156,7 +1124,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:postgateowner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:postgateowner2") .await .wrap_err("Failed to ensure actor")?; @@ -1172,7 +1140,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:postgateowner2/app.bsky.feed.post/3l7mkz4lmk23k", "did:plc:postgateowner2", test_cid(), @@ -1228,7 +1195,6 @@ async fn test_threadgate_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1237,7 +1203,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:threadgateowner1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:threadgateowner1") .await .wrap_err("Failed to ensure actor")?; @@ -1253,7 +1219,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:threadgateowner1/app.bsky.feed.post/3l7mkz4lmk23n", "did:plc:threadgateowner1", test_cid(), @@ -1313,7 +1278,6 @@ async fn test_threadgate_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1322,7 +1286,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:threadgateowner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:threadgateowner2") .await .wrap_err("Failed to ensure actor")?; @@ -1338,7 +1302,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:threadgateowner2/app.bsky.feed.post/3l7mkz4lmk23r", "did:plc:threadgateowner2", test_cid(), @@ -1394,7 +1357,6 @@ async fn test_post_embed_insert_images() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1403,7 +1365,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:embedimages1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:embedimages1") .await .wrap_err("Failed to ensure actor")?; @@ -1419,7 +1381,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:embedimages1/app.bsky.feed.post/3l7mkz4lmk23t", "did:plc:embedimages1", test_cid(), @@ -1453,7 +1414,7 @@ ]}); let post_uri = "at://did:plc:embedimages1/app.bsky.feed.post/3l7mkz4lmk23t"; - let result = feed::post_embed_insert(&tx, &mut redis, post_uri, embed, Utc::now(),EventSource::Jetstream).await; + let result = feed::post_embed_insert(&tx, post_uri, embed, Utc::now(),EventSource::Jetstream).await; assert!( result.is_ok(), @@ -1497,7 +1458,6 @@ async fn test_post_embed_insert_video() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1506,7 +1466,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:embedvideo1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:embedvideo1") .await .wrap_err("Failed to ensure actor")?; @@ -1522,7 +1482,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:embedvideo1/app.bsky.feed.post/3l7mkz4lmk23u", "did:plc:embedvideo1", test_cid(), @@ -1558,7 +1517,7 @@ ])}); let post_uri = "at://did:plc:embedvideo1/app.bsky.feed.post/3l7mkz4lmk23u"; - let result = feed::post_embed_insert(&tx, &mut redis, post_uri, embed, Utc::now(),EventSource::Jetstream).await; + let result = feed::post_embed_insert(&tx, post_uri, embed, Utc::now(),EventSource::Jetstream).await; assert!( result.is_ok(), @@ -1601,7 +1560,6 @@ async fn test_post_embed_insert_external() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1610,7 +1568,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:embedext1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:embedext1") .await .wrap_err("Failed to ensure actor")?; @@ -1626,7 +1584,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:embedext1/app.bsky.feed.post/3l7mkz4lmk23v", "did:plc:embedext1", test_cid(), @@ -1648,7 +1605,7 @@ size: 8192})}}); let post_uri = "at://did:plc:embedext1/app.bsky.feed.post/3l7mkz4lmk23v"; - let result = feed::post_embed_insert(&tx, &mut redis, post_uri, embed, Utc::now(),EventSource::Jetstream).await; + let result = feed::post_embed_insert(&tx, post_uri, embed, Utc::now(),EventSource::Jetstream).await; assert!( result.is_ok(), @@ -1683,7 +1640,6 @@ async fn test_post_embed_insert_record() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1692,10 +1648,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:quoter1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:quoter1") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:quotedauthor1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:quotedauthor1") .await .wrap_err("Failed to ensure actor")?; @@ -1711,7 +1667,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:quotedauthor1/app.bsky.feed.post/3l7mkz4lmk23w", "did:plc:quotedauthor1", test_cid(), @@ -1733,7 +1688,6 @@ created_at: Utc::now()}; feed::post_insert( &tx, - &mut redis, "at://did:plc:quoter1/app.bsky.feed.post/3l7mkz4lmk23x", "did:plc:quoter1", test_cid(), @@ -1750,7 +1704,7 @@ cid: test_cid()}}); let post_uri = "at://did:plc:quoter1/app.bsky.feed.post/3l7mkz4lmk23x"; - let result = feed::post_embed_insert(&tx, &mut redis, post_uri, embed, Utc::now(),EventSource::Jetstream).await; + let result = feed::post_embed_insert(&tx, post_uri, embed, Utc::now(),EventSource::Jetstream).await; assert!( result.is_ok(), diff --git a/consumer/tests/feedgen_labeler_operations_test.rs b/consumer/tests/feedgen_labeler_operations_test.rs --- a/consumer/tests/feedgen_labeler_operations_test.rs +++ b/consumer/tests/feedgen_labeler_operations_test.rs @@ -25,7 +25,6 @@ #[tokio::test] async fn test_feedgen_upsert_insert() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -33,10 +32,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist first (owner and service) - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedgenowner1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedgenowner1") .await .wrap_err("Failed to ensure owner actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedservice1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedservice1") .await .wrap_err("Failed to ensure service actor")?; @@ -56,7 +55,7 @@ let at_uri = "at://did:plc:feedgenowner1/app.bsky.feed.generator/3k2a7abc123"; - let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedservice1").await?; + let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedservice1").await?; let result = feed::feedgen_upsert(&tx, at_uri, "did:plc:feedgenowner1", test_cid(), service_actor_id, feedgen).await; assert!( @@ -109,7 +108,6 @@ #[tokio::test] async fn test_feedgen_upsert_update() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -117,10 +115,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedgenowner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedgenowner2") .await .wrap_err("Failed to ensure owner actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedservice2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedservice2") .await .wrap_err("Failed to ensure service actor")?; @@ -138,7 +136,7 @@ content_mode: Some("contentModeUnspecified".to_string()), created_at: Utc::now()}; - let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedservice2").await?; + let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedservice2").await?; let result1 = feed::feedgen_upsert(&tx, at_uri, "did:plc:feedgenowner2", test_cid(), service_actor_id, feedgen1).await; assert!(result1.unwrap(), "First insert should return true"); @@ -199,7 +197,6 @@ #[tokio::test] async fn test_feedgen_delete() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -207,10 +204,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedgenowner3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedgenowner3") .await .wrap_err("Failed to ensure owner actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedservice3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedservice3") .await .wrap_err("Failed to ensure service actor")?; @@ -228,7 +225,7 @@ content_mode: None, created_at: Utc::now()}; - let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:feedservice3").await?; + let (service_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:feedservice3").await?; feed::feedgen_upsert(&tx, at_uri, "did:plc:feedgenowner3", test_cid(), service_actor_id, feedgen) .await .wrap_err_with(|| format!("INSERT feedgen: {}", at_uri))?; @@ -260,7 +257,6 @@ #[tokio::test] async fn test_labeler_upsert_insert() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -268,7 +264,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:labeler1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:labeler1") .await .wrap_err("Failed to ensure labeler actor")?; @@ -297,7 +293,7 @@ subject_collections: Some(vec!["app.bsky.feed.post".to_string()]), created_at: Utc::now()}; - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:labeler1").await?; + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:labeler1").await?; let result = labeler::labeler_upsert(&tx, actor_id, "did:plc:labeler1", test_cid(), labeler).await; assert_eq!( result.wrap_err("Operation failed")?, @@ -390,7 +386,6 @@ #[tokio::test] async fn test_labeler_upsert_update() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -398,7 +393,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:labeler2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:labeler2") .await .wrap_err("Failed to ensure labeler actor")?; @@ -428,7 +423,7 @@ subject_collections: None, created_at: Utc::now()}; - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:labeler2").await?; + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:labeler2").await?; let result1 = labeler::labeler_upsert(&tx, actor_id, "did:plc:labeler2", test_cid(), labeler1).await; assert_eq!(result1.unwrap(), 0, "labeler_upsert always returns 0"); @@ -552,7 +547,6 @@ #[tokio::test] async fn test_labeler_delete() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn .transaction() @@ -560,7 +554,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:labeler3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:labeler3") .await .wrap_err("Failed to ensure labeler actor")?; @@ -583,7 +577,7 @@ subject_collections: None, created_at: Utc::now()}; - let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:labeler3").await?; + let (actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:labeler3").await?; labeler::labeler_upsert(&tx, actor_id, "did:plc:labeler3", test_cid(), labeler) .await .wrap_err_with(|| format!("INSERT labeler: {}", at_uri))?; diff --git a/consumer/tests/graph_operations_test.rs b/consumer/tests/graph_operations_test.rs --- a/consumer/tests/graph_operations_test.rs +++ b/consumer/tests/graph_operations_test.rs @@ -22,7 +22,6 @@ async fn test_follow_insert_creates_actors_and_follow() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; // Begin transaction - will rollback automatically let mut conn = pool.get().await.wrap_err("Failed to get connection")?; @@ -32,10 +31,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:user123") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:user123") .await .wrap_err("Failed to ensure actor")?; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:target123") + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:target123") .await .wrap_err("Failed to ensure actor")?; @@ -48,7 +47,6 @@ // Call the actual SQL function let result = graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk245", // rkey "did:plc:user123", // follower DID subject_actor_id, @@ -107,7 +105,6 @@ async fn test_follow_insert_duplicate_is_idempotent() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -116,10 +113,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:user456") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:user456") .await .wrap_err("Failed to ensure actor")?; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:target456") + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:target456") .await .wrap_err("Failed to ensure actor")?; @@ -135,7 +132,6 @@ // Insert first time let result1 = graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk246", "did:plc:user456", subject_actor_id, @@ -152,7 +148,6 @@ }; let result2 = graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk246", "did:plc:user456", subject_actor_id, @@ -172,7 +167,6 @@ async fn test_follow_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -181,10 +175,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:user789") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:user789") .await .wrap_err("Failed to ensure actor")?; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:target789") + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:target789") .await .wrap_err("Failed to ensure actor")?; @@ -196,7 +190,6 @@ graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk247", "did:plc:user789", subject_actor_id, @@ -244,7 +237,6 @@ async fn test_block_insert_creates_both_actors() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -253,10 +245,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:blocker123") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:blocker123") .await .wrap_err("Failed to ensure actor")?; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:blocked123") + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:blocked123") .await .wrap_err("Failed to ensure actor")?; @@ -267,7 +259,6 @@ let result = graph::block_insert( &tx, - &mut redis, "3l7mkz4lmk24c", "did:plc:blocker123", subject_actor_id, @@ -327,7 +318,6 @@ async fn test_list_upsert_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -336,7 +326,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner") .await .wrap_err("Failed to ensure actor")?; @@ -351,7 +341,7 @@ }; let at_uri = "at://did:plc:listowner/app.bsky.graph.list/3l7mkz4lmk24a"; - let result = graph::list_upsert(&tx, &mut redis, at_uri, "did:plc:listowner", test_cid(), list).await; + let result = graph::list_upsert(&tx, at_uri, "did:plc:listowner", test_cid(), list).await; assert!( result.is_ok(), @@ -381,7 +371,6 @@ async fn test_list_upsert_update() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -390,7 +379,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner2") .await .wrap_err("Failed to ensure actor")?; @@ -406,7 +395,7 @@ }; let at_uri = "at://did:plc:listowner2/app.bsky.graph.list/3l7mkz4lmk24b"; - let result1 = graph::list_upsert(&tx, &mut redis, at_uri, "did:plc:listowner2", test_cid(), list1).await; + let result1 = graph::list_upsert(&tx, at_uri, "did:plc:listowner2", test_cid(), list1).await; assert!(result1.unwrap(), "First insert should return true"); @@ -421,7 +410,7 @@ created_at: Utc::now(), }; - let result2 = graph::list_upsert(&tx, &mut redis, at_uri, "did:plc:listowner2", test_cid(), list2).await; + let result2 = graph::list_upsert(&tx, at_uri, "did:plc:listowner2", test_cid(), list2).await; assert!(!result2.unwrap(), "Update should return false"); @@ -444,7 +433,6 @@ async fn test_list_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -453,7 +441,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner3") .await .wrap_err("Failed to ensure actor")?; @@ -469,7 +457,7 @@ }; let at_uri = "at://did:plc:listowner3/app.bsky.graph.list/3l7mkz4lmk24c"; - graph::list_upsert(&tx, &mut redis, at_uri, "did:plc:listowner3", test_cid(), list) + graph::list_upsert(&tx, at_uri, "did:plc:listowner3", test_cid(), list) .await .wrap_err("Failed to insert list")?; @@ -508,7 +496,6 @@ async fn test_verification_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -517,10 +504,10 @@ .wrap_err("Failed to start transaction")?; // Ensure both verifier and subject actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:verifier1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:verifier1") .await .wrap_err("Failed to ensure verifier actor")?; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject1") + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject1") .await .wrap_err("Failed to ensure subject actor")?; @@ -587,7 +574,6 @@ async fn test_verification_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -596,10 +582,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:verifier2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:verifier2") .await .wrap_err("Failed to ensure verifier actor")?; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject2") + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject2") .await .wrap_err("Failed to ensure subject actor")?; @@ -651,7 +637,6 @@ async fn test_list_item_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -660,10 +645,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner4") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner4") .await .wrap_err("Failed to ensure list owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject4") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject4") .await .wrap_err("Failed to ensure subject")?; @@ -679,7 +664,7 @@ }; let list_uri = "at://did:plc:listowner4/app.bsky.graph.list/3l7mkz4lmk24d"; - graph::list_upsert(&tx, &mut redis, list_uri, "did:plc:listowner4", test_cid(), list) + graph::list_upsert(&tx, list_uri, "did:plc:listowner4", test_cid(), list) .await .wrap_err("Failed to create list")?; @@ -691,8 +676,8 @@ }; let item_uri = "at://did:plc:listowner4/app.bsky.graph.listitem/3l7mkz4lmk24g"; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject4").await?; - let result = graph::list_item_insert(&tx, &mut redis, item_uri, subject_actor_id, test_cid(), list_item).await; + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject4").await?; + let result = graph::list_item_insert(&tx, item_uri, subject_actor_id, test_cid(), list_item).await; assert!( result.is_ok(), @@ -736,7 +721,6 @@ async fn test_list_item_insert_creates_stub_list() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -745,10 +729,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner5") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner5") .await .wrap_err("Failed to ensure list owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject5") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject5") .await .wrap_err("Failed to ensure subject")?; @@ -763,8 +747,8 @@ }; let item_uri = "at://did:plc:listowner5/app.bsky.graph.listitem/3l7mkz4lmk24i"; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject5").await?; - let result = graph::list_item_insert(&tx, &mut redis, item_uri, subject_actor_id, test_cid(), list_item).await; + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject5").await?; + let result = graph::list_item_insert(&tx, item_uri, subject_actor_id, test_cid(), list_item).await; assert!( result.is_ok(), @@ -813,7 +797,6 @@ async fn test_list_item_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -822,10 +805,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner6") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner6") .await .wrap_err("Failed to ensure list owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject6") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject6") .await .wrap_err("Failed to ensure subject")?; @@ -841,7 +824,7 @@ }; let list_uri = "at://did:plc:listowner6/app.bsky.graph.list/3l7mkz4lmk24j"; - graph::list_upsert(&tx, &mut redis, list_uri, "did:plc:listowner6", test_cid(), list) + graph::list_upsert(&tx, list_uri, "did:plc:listowner6", test_cid(), list) .await .wrap_err("Failed to create list")?; @@ -853,8 +836,8 @@ }; let item_uri = "at://did:plc:listowner6/app.bsky.graph.listitem/3l7mkz4lmk24k"; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:subject6").await?; - graph::list_item_insert(&tx, &mut redis, item_uri, subject_actor_id, test_cid(), list_item) + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:subject6").await?; + graph::list_item_insert(&tx, item_uri, subject_actor_id, test_cid(), list_item) .await .wrap_err("Failed to insert list item")?; @@ -893,7 +876,6 @@ async fn test_list_block_insert() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -902,10 +884,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - blocker and list owner - let (blocker_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:blocker7") + let (blocker_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:blocker7") .await .wrap_err("Failed to ensure blocker")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner7") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner7") .await .wrap_err("Failed to ensure list owner")?; @@ -921,7 +903,7 @@ }; let list_uri = "at://did:plc:listowner7/app.bsky.graph.list/3l7mkz4lmk24l"; - graph::list_upsert(&tx, &mut redis, list_uri, "did:plc:listowner7", test_cid(), list) + graph::list_upsert(&tx, list_uri, "did:plc:listowner7", test_cid(), list) .await .wrap_err("Failed to create list")?; @@ -977,7 +959,6 @@ async fn test_list_block_delete() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -986,10 +967,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - let (blocker_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:blocker8") + let (blocker_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:blocker8") .await .wrap_err("Failed to ensure blocker")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner8") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner8") .await .wrap_err("Failed to ensure list owner")?; @@ -1005,7 +986,7 @@ }; let list_uri = "at://did:plc:listowner8/app.bsky.graph.list/3l7mkz4lmk24n"; - graph::list_upsert(&tx, &mut redis, list_uri, "did:plc:listowner8", test_cid(), list) + graph::list_upsert(&tx, list_uri, "did:plc:listowner8", test_cid(), list) .await .wrap_err("Failed to create list")?; diff --git a/consumer/tests/notification_test.rs b/consumer/tests/notification_test.rs --- a/consumer/tests/notification_test.rs +++ b/consumer/tests/notification_test.rs @@ -24,7 +24,6 @@ async fn test_is_thread_muted_not_muted() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -33,7 +32,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:user1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:user1") .await .wrap_err("Failed to ensure actor")?; @@ -58,7 +57,6 @@ async fn test_is_thread_muted_muted() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -67,10 +65,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:user2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:user2") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootpostauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootpostauthor") .await .wrap_err("Failed to ensure actor")?; @@ -88,7 +86,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:rootpostauthor/app.bsky.feed.post/3l7mkz4lmk2aa", "did:plc:rootpostauthor", test_cid(), @@ -156,7 +153,6 @@ async fn test_reply_chain_walker_single_level() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -165,10 +161,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootauthor") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:replier") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:replier") .await .wrap_err("Failed to ensure actor")?; @@ -186,7 +182,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:rootauthor/app.bsky.feed.post/3l7mkz4lmk2ab", "did:plc:rootauthor", test_cid(), @@ -219,7 +214,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:replier/app.bsky.feed.post/3l7mkz4lmk2ac", "did:plc:replier", test_cid(), @@ -258,7 +252,6 @@ async fn test_reply_chain_walker_multi_level() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -267,16 +260,16 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:root") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:root") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:replier1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:replier1") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:replier2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:replier2") .await .wrap_err("Failed to ensure actor")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:replier3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:replier3") .await .wrap_err("Failed to ensure actor")?; @@ -294,7 +287,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:root/app.bsky.feed.post/3l7mkz4lmk2ad", "did:plc:root", test_cid(), @@ -327,7 +319,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:replier1/app.bsky.feed.post/3l7mkz4lmk2ae", "did:plc:replier1", test_cid(), @@ -360,7 +351,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:replier2/app.bsky.feed.post/3l7mkz4lmk2af", "did:plc:replier2", test_cid(), @@ -393,7 +383,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:replier3/app.bsky.feed.post/3l7mkz4lmk2ag", "did:plc:replier3", test_cid(), @@ -455,7 +444,6 @@ async fn test_reply_chain_walker_post_not_found() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn diff --git a/consumer/tests/threadgate_enforcement_test.rs b/consumer/tests/threadgate_enforcement_test.rs --- a/consumer/tests/threadgate_enforcement_test.rs +++ b/consumer/tests/threadgate_enforcement_test.rs @@ -24,8 +24,6 @@ async fn test_threadgate_enforcement_no_threadgate() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -34,10 +32,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootauthor") .await .wrap_err("Failed to ensure root author")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:replier") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:replier") .await .wrap_err("Failed to ensure replier")?; @@ -55,7 +53,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:rootauthor/app.bsky.feed.post/3l7mkz4lmk2ba", "did:plc:rootauthor", test_cid(), @@ -87,8 +84,6 @@ async fn test_threadgate_enforcement_same_author() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -97,7 +92,7 @@ .wrap_err("Failed to start transaction")?; // Ensure actor exists - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:sameauthor") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:sameauthor") .await .wrap_err("Failed to ensure actor")?; @@ -115,7 +110,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:sameauthor/app.bsky.feed.post/3l7mkz4lmk2bb", "did:plc:sameauthor", test_cid(), @@ -161,8 +155,6 @@ async fn test_threadgate_enforcement_empty_allow_list() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -171,10 +163,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:threadowner") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:threadowner") .await .wrap_err("Failed to ensure thread owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:blocker") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:blocker") .await .wrap_err("Failed to ensure blocker")?; @@ -192,7 +184,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:threadowner/app.bsky.feed.post/3l7mkz4lmk2bd", "did:plc:threadowner", test_cid(), @@ -241,8 +232,6 @@ async fn test_threadgate_enforcement_following_rule_allows() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -251,10 +240,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:follower1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:follower1") .await .wrap_err("Failed to ensure follower")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:followed1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:followed1") .await .wrap_err("Failed to ensure followed")?; @@ -263,10 +252,9 @@ subject: "did:plc:followed1".to_string(), created_at: Utc::now(), }; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:followed1").await?; + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:followed1").await?; graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk2bf", "did:plc:follower1", subject_actor_id, @@ -290,7 +278,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:followed1/app.bsky.feed.post/3l7mkz4lmk2bg", "did:plc:followed1", test_cid(), @@ -342,8 +329,6 @@ async fn test_threadgate_enforcement_following_rule_blocks() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -352,10 +337,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:notfollowed") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:notfollowed") .await .wrap_err("Failed to ensure not followed")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:stranger1") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:stranger1") .await .wrap_err("Failed to ensure stranger")?; @@ -373,7 +358,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:notfollowed/app.bsky.feed.post/3l7mkz4lmk2bi", "did:plc:notfollowed", test_cid(), @@ -422,8 +406,6 @@ async fn test_threadgate_enforcement_follower_rule_allows() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -432,10 +414,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootauthor2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootauthor2") .await .wrap_err("Failed to ensure root author")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:follower2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:follower2") .await .wrap_err("Failed to ensure follower")?; @@ -444,10 +426,9 @@ subject: "did:plc:rootauthor2".to_string(), created_at: Utc::now(), }; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootauthor2").await?; + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootauthor2").await?; graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk2bk", "did:plc:follower2", subject_actor_id, @@ -471,7 +452,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:rootauthor2/app.bsky.feed.post/3l7mkz4lmk2bl", "did:plc:rootauthor2", test_cid(), @@ -520,8 +500,6 @@ async fn test_threadgate_enforcement_follower_rule_blocks() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -530,10 +508,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:rootauthor3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:rootauthor3") .await .wrap_err("Failed to ensure root author")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:nonfollower") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:nonfollower") .await .wrap_err("Failed to ensure non-follower")?; @@ -551,7 +529,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:rootauthor3/app.bsky.feed.post/3l7mkz4lmk2bn", "did:plc:rootauthor3", test_cid(), @@ -600,8 +577,6 @@ async fn test_threadgate_enforcement_mention_rule_allows() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -610,10 +585,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:mentioner") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:mentioner") .await .wrap_err("Failed to ensure mentioner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:mentioned") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:mentioned") .await .wrap_err("Failed to ensure mentioned")?; @@ -641,7 +616,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:mentioner/app.bsky.feed.post/3l7mkz4lmk2bp", "did:plc:mentioner", test_cid(), @@ -690,8 +664,6 @@ async fn test_threadgate_enforcement_mention_rule_blocks() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -700,10 +672,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:mentioner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:mentioner2") .await .wrap_err("Failed to ensure mentioner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:notmentioned") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:notmentioned") .await .wrap_err("Failed to ensure not mentioned")?; @@ -721,7 +693,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:mentioner2/app.bsky.feed.post/3l7mkz4lmk2br", "did:plc:mentioner2", test_cid(), @@ -770,8 +741,6 @@ async fn test_threadgate_enforcement_list_rule_allows() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -780,10 +749,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner") .await .wrap_err("Failed to ensure list owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listmember") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listmember") .await .wrap_err("Failed to ensure list member")?; @@ -799,7 +768,7 @@ }; graph::list_upsert( - &tx, &mut redis, + &tx, "at://did:plc:listowner/app.bsky.graph.list/3l7mkz4lmk2bt", "did:plc:listowner", test_cid(), @@ -814,9 +783,9 @@ subject: "did:plc:listmember".to_string(), created_at: Utc::now(), }; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listmember").await?; + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listmember").await?; graph::list_item_insert( - &tx, &mut redis, + &tx, "at://did:plc:listowner/app.bsky.graph.listitem/3l7mkz4lmk2bu", subject_actor_id, test_cid(), @@ -839,7 +808,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:listowner/app.bsky.feed.post/3l7mkz4lmk2bv", "did:plc:listowner", test_cid(), @@ -890,8 +858,6 @@ async fn test_threadgate_enforcement_list_rule_blocks() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -900,10 +866,10 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:listowner2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:listowner2") .await .wrap_err("Failed to ensure list owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:notinlist") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:notinlist") .await .wrap_err("Failed to ensure not in list")?; @@ -919,7 +885,7 @@ }; graph::list_upsert( - &tx, &mut redis, + &tx, "at://did:plc:listowner2/app.bsky.graph.list/3l7mkz4lmk2bx", "did:plc:listowner2", test_cid(), @@ -942,7 +908,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:listowner2/app.bsky.feed.post/3l7mkz4lmk2by", "did:plc:listowner2", test_cid(), @@ -993,8 +958,6 @@ async fn test_threadgate_enforcement_multiple_rules() -> eyre::Result<()> { common::ensure_test_db_ready().await; let pool = common::test_pool(); - let mut redis = common::test_redis().await; - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn @@ -1003,13 +966,13 @@ .wrap_err("Failed to start transaction")?; // Ensure actors exist - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:multiowner") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:multiowner") .await .wrap_err("Failed to ensure owner")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:follower3") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:follower3") .await .wrap_err("Failed to ensure follower")?; - consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:stranger2") + consumer::db::operations::feed::get_actor_id(&tx, "did:plc:stranger2") .await .wrap_err("Failed to ensure stranger")?; @@ -1018,10 +981,9 @@ subject: "did:plc:multiowner".to_string(), created_at: Utc::now(), }; - let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, &mut redis, "did:plc:multiowner").await?; + let (subject_actor_id, _, _) = consumer::db::operations::feed::get_actor_id(&tx, "did:plc:multiowner").await?; graph::follow_insert( &tx, - &mut redis, "3l7mkz4lmk2ca", "did:plc:follower3", subject_actor_id, @@ -1045,7 +1007,6 @@ feed::post_insert( &tx, - &mut redis, "at://did:plc:multiowner/app.bsky.feed.post/3l7mkz4lmk2cb", "did:plc:multiowner", test_cid(), diff --git a/consumer/tests/workers_test.rs b/consumer/tests/workers_test.rs --- a/consumer/tests/workers_test.rs +++ b/consumer/tests/workers_test.rs @@ -24,7 +24,6 @@ #[tokio::test] async fn test_bulk_ensure_actors_new() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -41,7 +40,6 @@ #[tokio::test] async fn test_bulk_ensure_actors_existing() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -64,7 +62,6 @@ #[tokio::test] async fn test_bulk_ensure_actors_mixed() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -91,7 +88,6 @@ #[tokio::test] async fn test_bulk_ensure_actors_empty() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -113,7 +109,6 @@ #[tokio::test] async fn test_backfill_update_actor_status() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -157,7 +152,6 @@ #[tokio::test] async fn test_backfill_update_actor_status_nonexistent() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -188,7 +182,6 @@ #[tokio::test] async fn test_backfill_mark_processing_new() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -220,7 +213,6 @@ #[tokio::test] async fn test_backfill_mark_processing_existing() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; @@ -261,7 +253,6 @@ #[tokio::test] async fn test_get_pinned_post_uri_none() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let conn = pool.get().await.wrap_err("Failed to get connection")?; let result = workers::get_pinned_post_uri(&conn, "did:plc:no_profile").await; @@ -282,12 +273,11 @@ #[tokio::test] async fn test_get_pinned_post_uri_exists() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Create actor using production function - let (actor_id, _, _) = operations::feed::get_actor_id(&tx, &mut redis, "did:plc:with_pinned") + let (actor_id, _, _) = operations::feed::get_actor_id(&tx, "did:plc:with_pinned") .await .wrap_err("Failed to ensure actor")?; @@ -303,7 +293,7 @@ tags: None, created_at: Utc::now(), }; - operations::feed::post_insert(&tx, &mut redis, post_uri, "did:plc:with_pinned", test_cid(), post,EventSource::Jetstream) + operations::feed::post_insert(&tx, post_uri, "did:plc:with_pinned", test_cid(), post,EventSource::Jetstream) .await .wrap_err("Failed to insert post")?; @@ -350,7 +340,6 @@ #[tokio::test] async fn test_get_recent_post_uris_empty() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let conn = pool.get().await.wrap_err("Failed to get connection")?; let result = workers::get_recent_post_uris(&conn, "did:plc:no_posts").await; @@ -373,12 +362,11 @@ #[tokio::test] async fn test_get_recent_post_uris_ordered() -> eyre::Result<()> { let pool = test_pool(); - let mut redis = common::test_redis().await; let mut conn = pool.get().await.wrap_err("Failed to get connection")?; let tx = conn.transaction().await.wrap_err("Failed to start transaction")?; // Create actor using production function - operations::feed::get_actor_id(&tx, &mut redis, "did:plc:with_posts") + operations::feed::get_actor_id(&tx, "did:plc:with_posts") .await .wrap_err("Failed to ensure actor")?; @@ -397,7 +385,7 @@ tags: None, created_at, }; - operations::feed::post_insert(&tx, &mut redis, &post_uri, "did:plc:with_posts", test_cid(), post,EventSource::Jetstream) + operations::feed::post_insert(&tx, &post_uri, "did:plc:with_posts", test_cid(), post,EventSource::Jetstream) .await .wrap_err("Failed to insert post")?; } diff --git a/consumer/src/db/backfill_jobs.rs b/consumer/src/db/backfill_jobs.rs --- a/consumer/src/db/backfill_jobs.rs +++ b/consumer/src/db/backfill_jobs.rs @@ -236,7 +236,7 @@ #[cfg(test)] mod tests { - use super::*; + #[test] fn test_backfill_jobs_module_compiles() { diff --git a/consumer/src/db/constellation_enrichment_queue.rs b/consumer/src/db/constellation_enrichment_queue.rs --- a/consumer/src/db/constellation_enrichment_queue.rs +++ b/consumer/src/db/constellation_enrichment_queue.rs @@ -156,7 +156,7 @@ #[cfg(test)] mod tests { - use super::*; + #[test] fn test_constellation_enrichment_queue_module_compiles() { diff --git a/consumer/src/db/cursors.rs b/consumer/src/db/cursors.rs --- a/consumer/src/db/cursors.rs +++ b/consumer/src/db/cursors.rs @@ -45,7 +45,7 @@ #[cfg(test)] mod tests { - use super::*; + #[test] fn test_cursor_module_compiles() { diff --git a/consumer/src/db/fetch_queue.rs b/consumer/src/db/fetch_queue.rs --- a/consumer/src/db/fetch_queue.rs +++ b/consumer/src/db/fetch_queue.rs @@ -163,7 +163,7 @@ #[cfg(test)] mod tests { - use super::*; + #[test] fn test_fetch_queue_module_compiles() { diff --git a/consumer/src/db/handle_resolution_queue.rs b/consumer/src/db/handle_resolution_queue.rs --- a/consumer/src/db/handle_resolution_queue.rs +++ b/consumer/src/db/handle_resolution_queue.rs @@ -134,7 +134,7 @@ #[cfg(test)] mod tests { - use super::*; + #[test] fn test_handle_resolution_queue_module_compiles() { diff --git a/consumer/src/db/notifications.rs b/consumer/src/db/notifications.rs --- a/consumer/src/db/notifications.rs +++ b/consumer/src/db/notifications.rs @@ -127,7 +127,7 @@ #[cfg(test)] mod tests { - use super::*; + #[test] fn test_notifications_module_compiles() {