//! §11.6 D3 integration test: group-shared durable queue store. //! //! Proves the competing-consumer primitive: N messages enqueued into one //! group-keyed store are claimed EXACTLY ONCE across concurrent claimers (the //! `FOR UPDATE SKIP LOCKED` guarantee that makes per-descendant materialized //! consumers safe). Also checks ack removes a row and nack re-defers it. //! //! Deterministic: drives `GroupQueueRepo` directly (no dispatcher). Skips when //! `DATABASE_URL` is unset. #![allow(clippy::too_many_lines, clippy::many_single_char_names)] use std::collections::HashSet; use picloud_manager_core::group_dead_letter_repo::{ GroupDeadLetterRepo, PostgresGroupDeadLetterRepo, }; use picloud_manager_core::group_queue_repo::{ GroupQueueRepo, NewGroupQueueMessage, PostgresGroupQueueRepo, }; use picloud_shared::GroupId; use sqlx::postgres::PgPoolOptions; use sqlx::PgPool; use uuid::Uuid; async fn pool_or_skip() -> Option { let Ok(url) = std::env::var("DATABASE_URL") else { picloud_test_support::abort_if_db_required("group_queue"); eprintln!("group_queue: DATABASE_URL unset — skipping"); return None; }; let pool = PgPoolOptions::new() .max_connections(6) .connect(&url) .await .expect("connect"); sqlx::migrate!("./migrations") .run(&pool) .await .expect("migrate"); Some(pool) } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn competing_consumers_claim_each_message_exactly_once() { let Some(pool) = pool_or_skip().await else { return; }; let sfx = Uuid::new_v4().simple().to_string(); let g: (Uuid,) = sqlx::query_as("INSERT INTO groups (slug, name) VALUES ($1, $1) RETURNING id") .bind(format!("gq-g-{sfx}")) .fetch_one(&pool) .await .unwrap(); let group = GroupId::from(g.0); let repo = PostgresGroupQueueRepo::new(pool.clone()); // Enqueue 50 distinct messages into the shared `tasks` queue. let n_msgs: usize = 50; for i in 0..n_msgs { repo.enqueue(NewGroupQueueMessage { group_id: group, collection: "tasks".into(), payload: serde_json::json!({ "i": i }), deliver_after: None, max_attempts: 3, enqueued_by_principal: None, }) .await .unwrap(); } // Four concurrent "consumers" claim until the queue is drained, ack-ing each. // Every claimed payload id must be unique — no message delivered twice. let claim_loop = |repo: PostgresGroupQueueRepo, group: GroupId| async move { let mut got: Vec = Vec::new(); while let Some(msg) = repo.claim(group, "tasks").await.unwrap() { got.push(msg.payload["i"].as_i64().unwrap()); assert!(repo.ack(msg.id, msg.claim_token).await.unwrap(), "ack ok"); } got }; let (a, b, c, d) = tokio::join!( claim_loop(PostgresGroupQueueRepo::new(pool.clone()), group), claim_loop(PostgresGroupQueueRepo::new(pool.clone()), group), claim_loop(PostgresGroupQueueRepo::new(pool.clone()), group), claim_loop(PostgresGroupQueueRepo::new(pool.clone()), group), ); let mut all: Vec = Vec::new(); all.extend(a); all.extend(b); all.extend(c); all.extend(d); let unique: HashSet = all.iter().copied().collect(); assert_eq!( all.len(), n_msgs, "every message delivered exactly once (no dupes)" ); assert_eq!(unique.len(), n_msgs, "all distinct ids covered"); assert_eq!( repo.depth(group, "tasks").await.unwrap(), 0, "queue drained" ); // nack re-defers a message (a later claim gets it back). let id = repo .enqueue(NewGroupQueueMessage { group_id: group, collection: "tasks".into(), payload: serde_json::json!({ "i": 999 }), deliver_after: None, max_attempts: 3, enqueued_by_principal: None, }) .await .unwrap(); let msg = repo.claim(group, "tasks").await.unwrap().expect("claimed"); assert_eq!(msg.id, id); repo.nack(msg.id, msg.claim_token, chrono::Duration::milliseconds(0)) .await .unwrap(); // After nack (0ms delay) it is claimable again. let again = repo.claim(group, "tasks").await.unwrap().expect("re-claim"); assert_eq!(again.id, id, "nacked message is re-delivered"); assert_eq!(again.attempt, 2, "attempt incremented on re-claim"); repo.ack(again.id, again.claim_token).await.unwrap(); // Cleanup (messages cascade on group delete anyway). let _ = sqlx::query("DELETE FROM group_queue_messages WHERE group_id = $1") .bind(g.0) .execute(&pool) .await; let _ = sqlx::query("DELETE FROM groups WHERE id = $1") .bind(g.0) .execute(&pool) .await; } #[tokio::test] async fn dead_letter_moves_an_exhausted_message_to_the_group_store() { // §11.6 D3: an exhausted shared-queue message is MOVED to the group dead- // letter store (not dropped) — the row leaves the live queue and appears in // `group_dead_letters`, operator-visible via GroupDeadLetterRepo. let Some(pool) = pool_or_skip().await else { return; }; let sfx = Uuid::new_v4().simple().to_string(); let g: (Uuid,) = sqlx::query_as("INSERT INTO groups (slug, name) VALUES ($1, $1) RETURNING id") .bind(format!("gdl-g-{sfx}")) .fetch_one(&pool) .await .unwrap(); let group = GroupId::from(g.0); let queue = PostgresGroupQueueRepo::new(pool.clone()); let dlq = PostgresGroupDeadLetterRepo::new(pool.clone()); let id = queue .enqueue(NewGroupQueueMessage { group_id: group, collection: "jobs".into(), payload: serde_json::json!({ "task": "boom" }), deliver_after: None, max_attempts: 1, enqueued_by_principal: None, }) .await .unwrap(); let msg = queue.claim(group, "jobs").await.unwrap().expect("claimed"); assert_eq!(msg.id, id); // Exhausted → dead-letter (mirrors the dispatcher's q_terminal shared arm). let dl_id = queue .dead_letter( msg.id, msg.claim_token, group, "jobs", None, None, msg.attempt, msg.enqueued_at, "handler exploded", ) .await .expect("dead_letter"); // The live queue no longer holds it. assert_eq!( queue.depth(group, "jobs").await.unwrap(), 0, "the dead-lettered message left the live queue" ); // It is preserved + operator-visible in the group dead-letter store. let rows = dlq.list_for_group(group, false, 100).await.unwrap(); assert_eq!(rows.len(), 1, "exactly one dead-letter recorded"); let row = &rows[0]; assert_eq!(row.id.into_inner(), dl_id.into_inner()); assert_eq!(row.collection, "jobs"); assert_eq!(row.source, "queue"); assert_eq!(row.last_error, "handler exploded"); assert_eq!(row.payload["message"]["task"], "boom"); assert!( row.resolved_at.is_none(), "a fresh dead-letter is unresolved" ); // A claim_token mismatch cannot dead-letter a re-claimed message: enqueue // another, claim it, then try to DL with a bogus token → the source-row // fetch finds nothing (RowNotFound), leaving the live queue intact. let id2 = queue .enqueue(NewGroupQueueMessage { group_id: group, collection: "jobs".into(), payload: serde_json::json!({ "task": "keep" }), deliver_after: None, max_attempts: 1, enqueued_by_principal: None, }) .await .unwrap(); let m2 = queue .claim(group, "jobs") .await .unwrap() .expect("claimed 2"); let bogus = Uuid::new_v4(); assert!( queue .dead_letter( m2.id, bogus, group, "jobs", None, None, 1, m2.enqueued_at, "x" ) .await .is_err(), "a claim-token mismatch must not dead-letter the message" ); assert_eq!( queue.depth(group, "jobs").await.unwrap(), 1, "the mismatched message is still live" ); let _ = id2; let _ = sqlx::query("DELETE FROM groups WHERE id = $1") .bind(g.0) .execute(&pool) .await; }