use std::io::Write as _;
use std::time::Duration;
use crate::common::OrderCreated;
use reliar_core::Envelope;
use reliar_outbox::{
AcquireRequest, FailedRecord, FailureOutcome, OutboxEnqueue, OutboxStore, PurgeRequest,
WorkerId,
};
use reliar_store_postgres::{PostgresOutboxError, PostgresOutboxSettings, PostgresOutboxStore};
use sqlx::PgPool;
use sqlx::postgres::PgConnectOptions;
use testcontainers::core::{IntoContainerPort, Mount, WaitFor};
use testcontainers::runners::AsyncRunner;
use testcontainers::{GenericImage, ImageExt};
use testcontainers_modules::postgres::Postgres;
const PGDOG_IMAGE: &str = "ghcr.io/pgdogdev/pgdog";
const PGDOG_TAG: &str = "v0.1.46";
fn write_pgdog_config(pg_host: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join(format!("reliar-pgdog-{}", uuid::Uuid::now_v7().simple()));
std::fs::create_dir_all(&dir).expect("create pgdog config dir");
let pgdog_toml = format!(
r#"[general]
host = "0.0.0.0"
port = 6432
pooler_mode = "transaction"
[[databases]]
name = "postgres"
host = "{pg_host}"
port = 5432
database_name = "postgres"
user = "postgres"
"#
);
let users_toml = r#"[[users]]
name = "postgres"
database = "postgres"
password = "postgres"
"#;
let mut f = std::fs::File::create(dir.join("pgdog.toml")).unwrap();
f.write_all(pgdog_toml.as_bytes()).unwrap();
let mut f = std::fs::File::create(dir.join("users.toml")).unwrap();
f.write_all(users_toml.as_bytes()).unwrap();
dir
}
#[allow(
clippy::too_many_lines,
reason = "one end-to-end pooler scenario: URL-options pass-through, the ALTER ROLE fallback, \
and the full enqueue/acquire/complete/fail/purge path — splitting it would scatter \
one ordered narrative across helper functions with no reuse"
)]
async fn pgdog_pooler_passes_options_through_and_behaves_like_direct() {
let network = format!("reliar-pgdog-{}", uuid::Uuid::now_v7().simple());
let pg_name = format!("reliar-pg-{}", uuid::Uuid::now_v7().simple());
let pg = Postgres::default()
.with_tag("18-alpine")
.with_container_name(&pg_name)
.with_network(&network)
.with_label("reliar.test", "true")
.start()
.await
.expect("start postgres");
let pg_direct_port = pg.get_host_port_ipv4(5432).await.expect("postgres port");
let direct_url = format!("postgres://postgres:postgres@127.0.0.1:{pg_direct_port}/postgres");
let direct_pool = PgPool::connect(&direct_url).await.expect("connect direct");
reliar_store_postgres::migrate(
&direct_pool,
reliar_store_postgres::MigrateOptions::default(),
)
.await
.expect("migrate direct");
let config_dir = write_pgdog_config(&pg_name);
let pgdog_before = GenericImage::new(PGDOG_IMAGE, PGDOG_TAG)
.with_exposed_port(6432.tcp())
.with_wait_for(WaitFor::message_on_stderr("PgDog listening on"))
.with_network(&network)
.with_mount(Mount::bind_mount(
config_dir.join("pgdog.toml").to_string_lossy().into_owned(),
"/pgdog/pgdog.toml",
))
.with_mount(Mount::bind_mount(
config_dir.join("users.toml").to_string_lossy().into_owned(),
"/pgdog/users.toml",
))
.with_container_name(format!(
"reliar-pgdog-before-{}",
uuid::Uuid::now_v7().simple()
))
.with_label("reliar.test", "true")
.start()
.await
.expect("start pgdog (before ALTER ROLE)");
let before_port = pgdog_before
.get_host_port_ipv4(6432)
.await
.expect("pgdog (before) port");
let options_url: PgConnectOptions =
format!("postgres://postgres:postgres@127.0.0.1:{before_port}/postgres")
.parse()
.unwrap();
let options_url = options_url.options([("search_path", "reliar,public")]);
let options_pool = PgPool::connect_with(options_url)
.await
.expect("connect through pgdog with URL options");
let options_store = PostgresOutboxStore::new(options_pool);
options_store
.stats()
.await
.expect("URL options pass through, so search_path resolves without ALTER ROLE");
let bare_url: PgConnectOptions =
format!("postgres://postgres:postgres@127.0.0.1:{before_port}/postgres")
.parse()
.unwrap();
let bare_pool = PgPool::connect_with(bare_url)
.await
.expect("connect through pgdog without options");
let bare_store = PostgresOutboxStore::new(bare_pool);
let err = bare_store.stats().await.unwrap_err();
assert!(
matches!(err, PostgresOutboxError::NotMigrated { .. }),
"expected NotMigrated with neither options nor a role default, got {err:?}"
);
sqlx::query("ALTER ROLE postgres SET search_path = reliar, public")
.execute(&direct_pool)
.await
.expect("alter role");
let pgdog_after = GenericImage::new(PGDOG_IMAGE, PGDOG_TAG)
.with_exposed_port(6432.tcp())
.with_wait_for(WaitFor::message_on_stderr("PgDog listening on"))
.with_network(&network)
.with_mount(Mount::bind_mount(
config_dir.join("pgdog.toml").to_string_lossy().into_owned(),
"/pgdog/pgdog.toml",
))
.with_mount(Mount::bind_mount(
config_dir.join("users.toml").to_string_lossy().into_owned(),
"/pgdog/users.toml",
))
.with_container_name(format!(
"reliar-pgdog-after-{}",
uuid::Uuid::now_v7().simple()
))
.with_label("reliar.test", "true")
.start()
.await
.expect("start pgdog (after ALTER ROLE)");
let after_port = pgdog_after
.get_host_port_ipv4(6432)
.await
.expect("pgdog (after) port");
let pool = PgPool::connect(&format!(
"postgres://postgres:postgres@127.0.0.1:{after_port}/postgres"
))
.await
.expect("connect through pgdog (after)");
let store = PostgresOutboxStore::with_settings(pool.clone(), PostgresOutboxSettings::default());
let envelope_a = Envelope::builder(OrderCreated { order_id: 1 }).build();
let envelope_b = Envelope::builder(OrderCreated { order_id: 2 }).build();
for envelope in [envelope_a, envelope_b] {
let mut tx = pool.begin().await.unwrap();
store.enqueue_envelope(&mut tx, envelope).await.unwrap();
tx.commit().await.unwrap();
}
let worker_a = WorkerId::generate();
let worker_b = WorkerId::generate();
let (batch_a, batch_b) = tokio::join!(
store.acquire(AcquireRequest::new(worker_a.clone()).lease(Duration::from_secs(30))),
store.acquire(AcquireRequest::new(worker_b.clone()).lease(Duration::from_secs(30))),
);
let batch_a = batch_a.unwrap();
let batch_b = batch_b.unwrap();
assert_eq!(
batch_a.records.len() + batch_b.records.len(),
2,
"SKIP LOCKED still partitions the seed disjointly and exhaustively through PgDog"
);
let mut all_records: Vec<_> = batch_a.records.into_iter().chain(batch_b.records).collect();
all_records.sort_by_key(|r| r.envelope.id);
let (record_complete, record_fail) = (all_records[0].clone(), all_records[1].clone());
let owner_a = if record_complete.locked_by.as_ref() == Some(&worker_a) {
&worker_a
} else {
&worker_b
};
let affected = store
.complete(owner_a, &[record_complete.record_ref()])
.await
.unwrap();
assert_eq!(affected, 1);
let owner_b = if record_fail.locked_by.as_ref() == Some(&worker_a) {
&worker_a
} else {
&worker_b
};
let affected = store
.fail(
owner_b,
&[FailedRecord::new(
record_fail.record_ref(),
"transient failure through pgdog",
FailureOutcome::Retry {
delay: Duration::from_millis(1),
},
)],
)
.await
.unwrap();
assert_eq!(affected, 1);
sqlx::query(
"UPDATE outbox SET available_at = now() - interval '1 second' WHERE message_id = $1",
)
.bind(record_fail.envelope.id.as_uuid())
.execute(&pool)
.await
.unwrap();
let batch_retry = store
.acquire(AcquireRequest::new(WorkerId::generate()).lease(Duration::from_secs(30)))
.await
.unwrap();
let retried = batch_retry
.records
.iter()
.find(|r| r.envelope.id == record_fail.envelope.id)
.expect("the retried row is claimable again through pgdog");
assert_eq!(retried.attempts, 1);
let retry_worker = retried.locked_by.clone().unwrap();
let extended = store
.extend_lease(
&retry_worker,
&[retried.record_ref()],
Duration::from_secs(120),
)
.await
.unwrap();
assert_eq!(extended, 1);
let released = store
.release(&retry_worker, &[retried.record_ref()])
.await
.unwrap();
assert_eq!(released, 1);
sqlx::query("UPDATE outbox SET published_at = now() - interval '1 hour' WHERE message_id = $1")
.bind(record_complete.envelope.id.as_uuid())
.execute(&pool)
.await
.unwrap();
let report = store
.purge(PurgeRequest::default().published_retention(Some(Duration::ZERO)))
.await
.unwrap();
assert_eq!(
report.published_deleted, 1,
"purge works through the pooler"
);
drop(pgdog_before);
drop(pgdog_after);
let _ = std::fs::remove_dir_all(&config_dir);
}
pub(crate) fn trials(rt: &'static tokio::runtime::Runtime) -> Vec<libtest_mimic::Trial> {
vec![libtest_mimic::Trial::test(
"outbox_pgdog::pgdog_pooler_passes_options_through_and_behaves_like_direct",
move || {
rt.block_on(pgdog_pooler_passes_options_through_and_behaves_like_direct());
Ok(())
},
)]
}