use crate::common;
async fn explain(pool: &sqlx::PgPool, statement: &str) -> String {
let lines: Vec<String> = sqlx::query_scalar(sqlx::AssertSqlSafe(format!(
"EXPLAIN (ANALYZE, BUFFERS, FORMAT TEXT) {statement}"
)))
.fetch_all(pool)
.await
.unwrap();
lines.join("\n")
}
const SEED_ROWS: u32 = 100_000;
async fn seed(pool: &sqlx::PgPool, n: u32) {
sqlx::query(
"INSERT INTO outbox (message_id, message_type, message_version, conversation_id, \
content_type, payload, available_at, created_at) \
SELECT uuidv7(), 'orders.created', 1, uuidv7(), 'application/json', '{}'::bytea, \
now() - (g || ' seconds')::interval, now() - (g || ' seconds')::interval \
FROM generate_series(1, $1::bigint) AS g",
)
.bind(i64::from(n))
.execute(pool)
.await
.unwrap();
sqlx::query(
"UPDATE outbox SET dead_at = now() - interval '1 hour', dead_reason = 'permanent_error' \
WHERE id IN (SELECT id FROM outbox ORDER BY id LIMIT 1000)",
)
.execute(pool)
.await
.unwrap();
sqlx::query("VACUUM (ANALYZE) outbox")
.execute(pool)
.await
.unwrap();
}
static SHARED_SEEDED_POOL: tokio::sync::OnceCell<sqlx::PgPool> = tokio::sync::OnceCell::const_new();
async fn shared_seeded_pool() -> sqlx::PgPool {
SHARED_SEEDED_POOL
.get_or_init(|| async {
let pool = common::fresh_db().await;
seed(&pool, SEED_ROWS).await;
pool
})
.await
.clone()
}
async fn list_dead_uses_ix_outbox_dead_cursor_with_no_sort() {
let pool = shared_seeded_pool().await;
let plan = explain(
&pool,
"SELECT id, message_id, message_type, message_version, \
correlation_id, conversation_id, causation_id, request_id, \
content_type, payload, tenant_id, expires_at, ordering_key, \
metadata, headers, metadata_version, \
created_at, available_at, \
attempts, locked_by, claim_token, \
published_at, dead_at, dead_reason, last_error \
FROM outbox \
WHERE dead_at IS NOT NULL \
AND (NULL::text IS NULL OR message_type = NULL) \
AND (NULL::text IS NULL OR tenant_id = NULL) \
AND (NULL::timestamptz IS NULL OR dead_at < NULL) \
AND (NULL::timestamptz IS NULL OR (dead_at, id) > (NULL, NULL::uuid)) \
ORDER BY dead_at ASC, id ASC \
LIMIT 100",
)
.await;
assert!(
plan.contains("ix_outbox_dead_cursor"),
"expected ix_outbox_dead_cursor to back list_dead's keyset listing; plan:\n{plan}"
);
assert!(
!plan.contains("Sort"),
"list_dead must read in death-time order directly off the index, never via a Sort node; \
plan:\n{plan}"
);
assert!(
!plan.contains("Seq Scan on outbox"),
"list_dead must never fall back to a sequential scan; plan:\n{plan}"
);
assert!(
plan.contains("Limit"),
"list_dead must still be bounded by LIMIT; plan:\n{plan}"
);
}
async fn dead_retention_purge_sub_select_uses_ix_outbox_dead_cursor() {
let pool = shared_seeded_pool().await;
let plan = explain(
&pool,
"SELECT id FROM outbox \
WHERE dead_at IS NOT NULL \
AND dead_at < now() - interval '30 days' \
LIMIT 1000",
)
.await;
assert!(
plan.contains("ix_outbox_dead_cursor"),
"expected ix_outbox_dead_cursor to back the dead-retention purge sub-select; plan:\n{plan}"
);
}
async fn claim_still_plans_on_ix_outbox_claimable() {
let pool = shared_seeded_pool().await;
let plan = explain(
&pool,
"SELECT id FROM outbox \
WHERE published_at IS NULL AND dead_at IS NULL \
AND available_at <= now() \
AND (expires_at IS NULL OR expires_at > now()) \
ORDER BY available_at, id \
LIMIT 50 \
FOR UPDATE SKIP LOCKED",
)
.await;
assert!(
plan.contains("ix_outbox_claimable"),
"expected ix_outbox_claimable to still back the claim scan; plan:\n{plan}"
);
assert!(
!plan.contains("Seq Scan"),
"the claim must never fall back to a sequential scan; plan:\n{plan}"
);
assert!(
!plan.contains("Sort"),
"id is a key column of ix_outbox_claimable, so the scan must already be ordered, with no \
separate Sort node; plan:\n{plan}"
);
}
async fn mixed_seeded_pool() -> sqlx::PgPool {
let pool = common::fresh_db().await;
let per_bucket: u32 = 20_000;
sqlx::query(
"INSERT INTO outbox (message_id, message_type, message_version, conversation_id, \
content_type, payload, available_at, created_at) \
SELECT uuidv7(), 'orders.created', 1, uuidv7(), 'application/json', '{}'::bytea, \
now() - (g || ' seconds')::interval, now() - (g || ' seconds')::interval \
FROM generate_series(1, $1::bigint) AS g",
)
.bind(i64::from(per_bucket))
.execute(&pool)
.await
.unwrap();
sqlx::query(
"INSERT INTO outbox (message_id, message_type, message_version, conversation_id, \
content_type, payload, available_at, created_at, locked_by) \
SELECT uuidv7(), 'orders.created', 1, uuidv7(), 'application/json', '{}'::bytea, \
now() + interval '1 hour', now(), 'perf-worker' \
FROM generate_series(1, $1::bigint) AS g",
)
.bind(i64::from(per_bucket))
.execute(&pool)
.await
.unwrap();
sqlx::query(
"INSERT INTO outbox (message_id, message_type, message_version, conversation_id, \
content_type, payload, available_at, created_at, published_at) \
SELECT uuidv7(), 'orders.created', 1, uuidv7(), 'application/json', '{}'::bytea, \
now(), now(), now() \
FROM generate_series(1, $1::bigint) AS g",
)
.bind(i64::from(per_bucket))
.execute(&pool)
.await
.unwrap();
sqlx::query(
"INSERT INTO outbox (message_id, message_type, message_version, conversation_id, \
content_type, payload, available_at, created_at, dead_at, \
dead_reason) \
SELECT uuidv7(), 'orders.created', 1, uuidv7(), 'application/json', '{}'::bytea, \
now(), now(), now(), 'permanent_error' \
FROM generate_series(1, $1::bigint) AS g",
)
.bind(i64::from(per_bucket))
.execute(&pool)
.await
.unwrap();
sqlx::query(
"INSERT INTO outbox (message_id, message_type, message_version, conversation_id, \
content_type, payload, available_at, created_at, expires_at) \
SELECT uuidv7(), 'orders.created', 1, uuidv7(), 'application/json', '{}'::bytea, \
now(), now(), now() - interval '1 second' \
FROM generate_series(1, $1::bigint) AS g",
)
.bind(i64::from(per_bucket))
.execute(&pool)
.await
.unwrap();
sqlx::query("VACUUM (ANALYZE) outbox")
.execute(&pool)
.await
.unwrap();
pool
}
async fn stats_pending_and_oldest_pending_are_index_only_scans_with_no_heap_fetches() {
let pool = mixed_seeded_pool().await;
let pending_plan = explain(
&pool,
"SELECT count(*) FROM outbox \
WHERE published_at IS NULL AND dead_at IS NULL \
AND available_at <= now() \
AND (expires_at IS NULL OR expires_at > now())",
)
.await;
assert!(
pending_plan.contains("Index Only Scan"),
"expected stats' pending subquery to be an Index Only Scan; plan:\n{pending_plan}"
);
assert!(
pending_plan.contains("Heap Fetches: 0"),
"expected zero heap fetches on a vacuumed table; plan:\n{pending_plan}"
);
let oldest_plan = explain(
&pool,
"SELECT available_at FROM outbox \
WHERE published_at IS NULL AND dead_at IS NULL \
AND available_at <= now() \
AND (expires_at IS NULL OR expires_at > now()) \
ORDER BY available_at, id LIMIT 1",
)
.await;
assert!(
oldest_plan.contains("Index Only Scan"),
"expected stats' oldest_pending_available_at subquery to be an Index Only Scan; \
plan:\n{oldest_plan}"
);
assert!(
oldest_plan.contains("Heap Fetches: 0"),
"expected zero heap fetches on a vacuumed table; plan:\n{oldest_plan}"
);
}
async fn complete_plans_on_pk_outbox_with_no_seq_scan() {
let pool = shared_seeded_pool().await;
let ids: Vec<uuid::Uuid> = sqlx::query_scalar("SELECT id FROM outbox LIMIT 10")
.fetch_all(&pool)
.await
.unwrap();
let tokens: Vec<Option<uuid::Uuid>> = ids.iter().map(|_| None).collect();
let mut tx = pool.begin().await.unwrap();
let plan: Vec<String> = sqlx::query_scalar(
"EXPLAIN (ANALYZE, BUFFERS, FORMAT TEXT) \
UPDATE outbox o \
SET published_at = now(), attempts = o.attempts + 1, locked_by = NULL, \
claim_token = NULL, updated_at = now() \
FROM UNNEST($1::uuid[], $2::uuid[]) AS f(id, token) \
WHERE o.id = f.id AND o.claim_token = f.token",
)
.bind(&ids)
.bind(&tokens)
.fetch_all(&mut *tx)
.await
.unwrap();
tx.rollback().await.unwrap();
let plan = plan.join("\n");
assert!(
plan.contains("pk_outbox"),
"expected complete's fenced UNNEST join to plan on pk_outbox; plan:\n{plan}"
);
assert!(
plan.contains("Nested Loop"),
"expected the UNNEST join to plan as a nested loop; plan:\n{plan}"
);
assert!(
!plan.contains("Seq Scan on outbox"),
"the fencing join must never fall back to a sequential scan of outbox; plan:\n{plan}"
);
}
async fn extend_lease_plans_on_pk_outbox_with_no_seq_scan() {
let pool = shared_seeded_pool().await;
let ids: Vec<uuid::Uuid> = sqlx::query_scalar("SELECT id FROM outbox LIMIT 10")
.fetch_all(&pool)
.await
.unwrap();
let tokens: Vec<Option<uuid::Uuid>> = ids.iter().map(|_| None).collect();
let mut tx = pool.begin().await.unwrap();
let plan: Vec<String> = sqlx::query_scalar(
"EXPLAIN (ANALYZE, BUFFERS, FORMAT TEXT) \
UPDATE outbox o \
SET available_at = now() + interval '30 seconds', updated_at = now() \
FROM UNNEST($1::uuid[], $2::uuid[]) AS f(id, token) \
WHERE o.id = f.id AND o.claim_token = f.token",
)
.bind(&ids)
.bind(&tokens)
.fetch_all(&mut *tx)
.await
.unwrap();
tx.rollback().await.unwrap();
let plan = plan.join("\n");
assert!(
plan.contains("pk_outbox"),
"expected extend_lease's fenced UNNEST join to plan on pk_outbox; plan:\n{plan}"
);
assert!(
plan.contains("Nested Loop"),
"expected the UNNEST join to plan as a nested loop; plan:\n{plan}"
);
assert!(
!plan.contains("Seq Scan on outbox"),
"the fencing join must never fall back to a sequential scan of outbox; plan:\n{plan}"
);
}
pub(crate) fn trials(rt: &'static tokio::runtime::Runtime) -> Vec<libtest_mimic::Trial> {
vec![
libtest_mimic::Trial::test(
"outbox_plans::list_dead_uses_ix_outbox_dead_cursor_with_no_sort",
move || {
rt.block_on(list_dead_uses_ix_outbox_dead_cursor_with_no_sort());
Ok(())
},
),
libtest_mimic::Trial::test(
"outbox_plans::dead_retention_purge_sub_select_uses_ix_outbox_dead_cursor",
move || {
rt.block_on(dead_retention_purge_sub_select_uses_ix_outbox_dead_cursor());
Ok(())
},
),
libtest_mimic::Trial::test(
"outbox_plans::claim_still_plans_on_ix_outbox_claimable",
move || {
rt.block_on(claim_still_plans_on_ix_outbox_claimable());
Ok(())
},
),
libtest_mimic::Trial::test(
"outbox_plans::stats_pending_and_oldest_pending_are_index_only_scans_with_no_heap_fetches",
move || {
rt.block_on(
stats_pending_and_oldest_pending_are_index_only_scans_with_no_heap_fetches(),
);
Ok(())
},
),
libtest_mimic::Trial::test(
"outbox_plans::complete_plans_on_pk_outbox_with_no_seq_scan",
move || {
rt.block_on(complete_plans_on_pk_outbox_with_no_seq_scan());
Ok(())
},
),
libtest_mimic::Trial::test(
"outbox_plans::extend_lease_plans_on_pk_outbox_with_no_seq_scan",
move || {
rt.block_on(extend_lease_plans_on_pk_outbox_with_no_seq_scan());
Ok(())
},
),
]
}