1use std::time::{Duration, SystemTime, UNIX_EPOCH};
34
35use faucet_core::FaucetError;
36use pgwire_replication::{Lsn, ReplicationClient, ReplicationConfig, TlsConfig};
37use sqlx::postgres::PgConnectOptions;
38
39pub use pgwire_replication::ReplicationEvent;
42use sqlx::{Executor, PgConnection};
43use tracing::debug;
44
45pub const POSTGRES_EPOCH_MICROS: i64 = 946_684_800_000_000;
48
49pub struct Client {
58 _private: (),
59}
60
61pub struct Duplex {
64 inner: ReplicationClient,
65}
66
67#[derive(Clone, Debug)]
73pub struct ReplicationParams<'a> {
74 pub connection_url: &'a str,
76 pub slot_name: &'a str,
78 pub publication_name: &'a str,
80 pub proto_version: u32,
82 pub create_slot_if_missing: bool,
84 pub start_lsn: Option<u64>,
87 pub status_update_interval: Duration,
90 pub tcp_keepalive: Duration,
93 pub slot_type: crate::config::SlotType,
95 pub tls: &'a crate::config::CdcTls,
97}
98
99#[cfg_attr(test, derive(Debug))]
102struct PgCoords {
103 host: String,
104 port: u16,
105 user: String,
106 password: String,
107 dbname: String,
108}
109
110fn parse_url(url: &str) -> Result<PgCoords, FaucetError> {
111 let parsed = url::Url::parse(url)
114 .map_err(|e| FaucetError::Config(format!("postgres-cdc: invalid connection URL: {e}")))?;
115
116 let host = parsed
122 .host_str()
123 .filter(|h| !h.is_empty())
124 .ok_or_else(|| {
125 FaucetError::Config(
126 "postgres-cdc: connection URL is missing a host (expected \
127 postgres://user@host[:port]/dbname)"
128 .to_owned(),
129 )
130 })?
131 .to_owned();
132 let port = parsed.port().unwrap_or(5432);
133 let user = parsed.username().to_owned();
134 if user.is_empty() {
135 return Err(FaucetError::Config(
136 "postgres-cdc: connection URL is missing a user (expected \
137 postgres://user@host[:port]/dbname)"
138 .to_owned(),
139 ));
140 }
141 let password = parsed.password().unwrap_or("").to_owned();
142 let dbname = parsed.path().trim_start_matches('/').to_owned();
143 let dbname = if dbname.is_empty() {
144 "postgres".to_owned()
145 } else {
146 dbname
147 };
148
149 Ok(PgCoords {
150 host,
151 port,
152 user,
153 password,
154 dbname,
155 })
156}
157
158pub async fn connect(params: &ReplicationParams<'_>) -> Result<Client, FaucetError> {
166 let _ = parse_url(params.connection_url)?;
169 Ok(Client { _private: () })
170}
171
172pub async fn ensure_slot(
179 _client: &Client,
180 connection_url: &str,
181 slot_name: &str,
182 create_if_missing: bool,
183 slot_type: crate::config::SlotType,
184 tls: &crate::config::CdcTls,
185) -> Result<(), FaucetError> {
186 use crate::config::SlotType;
187 let opts: PgConnectOptions = connection_url
189 .parse()
190 .map_err(|e| FaucetError::Config(format!("postgres-cdc: invalid connection URL: {e}")))?;
191 let opts = apply_cdc_tls(opts, tls);
192
193 use sqlx::ConnectOptions as _;
194 let mut conn: PgConnection = opts
195 .connect()
196 .await
197 .map_err(|e| pg_err("postgres-cdc ensure_slot connect", e))?;
198
199 let row: Option<(String,)> =
201 sqlx::query_as("SELECT slot_name::text FROM pg_replication_slots WHERE slot_name = $1")
202 .bind(slot_name)
203 .fetch_optional(&mut conn)
204 .await
205 .map_err(|e| pg_err("postgres-cdc slot lookup", e))?;
206
207 if row.is_some() {
208 debug!("postgres-cdc: replication slot '{slot_name}' already exists");
209 return Ok(());
210 }
211
212 if !create_if_missing {
213 return Err(FaucetError::Source(format!(
214 "postgres-cdc: replication slot '{slot_name}' does not exist \
215 and create_slot_if_missing = false"
216 )));
217 }
218
219 let temporary = matches!(slot_type, SlotType::Temporary);
226 let sql = format!(
227 "SELECT pg_create_logical_replication_slot({}, 'pgoutput', {})",
228 quote_literal(slot_name),
229 temporary
230 );
231 conn.execute(sql.as_str())
232 .await
233 .map_err(|e| pg_err("postgres-cdc create slot", e))?;
234
235 if temporary {
236 debug!("postgres-cdc: created temporary replication slot '{slot_name}'");
237 } else {
238 tracing::warn!(
241 "postgres-cdc: created PERMANENT replication slot '{slot_name}' — it will retain \
242 WAL on the server until consumed or explicitly dropped (drop_slot). Use \
243 slot_type=temporary for ephemeral runs."
244 );
245 }
246 Ok(())
247}
248
249pub async fn drop_slot(
253 connection_url: &str,
254 slot_name: &str,
255 tls: &crate::config::CdcTls,
256) -> Result<(), FaucetError> {
257 let opts: PgConnectOptions = connection_url
258 .parse()
259 .map_err(|e| FaucetError::Config(format!("postgres-cdc: invalid connection URL: {e}")))?;
260 let opts = apply_cdc_tls(opts, tls);
261 use sqlx::ConnectOptions as _;
262 let mut conn: PgConnection = opts
263 .connect()
264 .await
265 .map_err(|e| pg_err("postgres-cdc drop_slot connect", e))?;
266
267 let exists: Option<(String,)> =
268 sqlx::query_as("SELECT slot_name::text FROM pg_replication_slots WHERE slot_name = $1")
269 .bind(slot_name)
270 .fetch_optional(&mut conn)
271 .await
272 .map_err(|e| pg_err("postgres-cdc slot lookup", e))?;
273 if exists.is_none() {
274 debug!("postgres-cdc: replication slot '{slot_name}' already absent; drop is a no-op");
275 return Ok(());
276 }
277
278 sqlx::query("SELECT pg_drop_replication_slot($1)")
279 .bind(slot_name)
280 .execute(&mut conn)
281 .await
282 .map_err(|e| pg_err("postgres-cdc drop slot", e))?;
283 debug!("postgres-cdc: dropped replication slot '{slot_name}'");
284 Ok(())
285}
286
287fn apply_cdc_tls(opts: PgConnectOptions, tls: &crate::config::CdcTls) -> PgConnectOptions {
296 use crate::config::CdcTls;
297 use sqlx::postgres::PgSslMode;
298 match tls {
299 CdcTls::Disable => opts.ssl_mode(PgSslMode::Disable),
300 CdcTls::Require => opts.ssl_mode(PgSslMode::Require),
301 CdcTls::VerifyCa { ca_path } => {
302 let o = opts.ssl_mode(PgSslMode::VerifyCa);
303 match ca_path {
304 Some(p) => o.ssl_root_cert(p),
305 None => o,
306 }
307 }
308 CdcTls::VerifyFull { ca_path } => {
309 let o = opts.ssl_mode(PgSslMode::VerifyFull);
310 match ca_path {
311 Some(p) => o.ssl_root_cert(p),
312 None => o,
313 }
314 }
315 }
316}
317
318fn tls_config(tls: &crate::config::CdcTls) -> TlsConfig {
321 use crate::config::CdcTls;
322 use std::path::PathBuf;
323 match tls {
324 CdcTls::Disable => TlsConfig::disabled(),
325 CdcTls::Require => TlsConfig::require(),
326 CdcTls::VerifyCa { ca_path } => TlsConfig::verify_ca(ca_path.clone().map(PathBuf::from)),
327 CdcTls::VerifyFull { ca_path } => {
328 TlsConfig::verify_full(ca_path.clone().map(PathBuf::from))
329 }
330 }
331}
332
333pub async fn advance_slot(
349 connection_url: &str,
350 slot_name: &str,
351 lsn: u64,
352 tls: &crate::config::CdcTls,
353) -> Result<(), FaucetError> {
354 if lsn == 0 {
355 return Ok(());
356 }
357 let opts: PgConnectOptions = connection_url
358 .parse()
359 .map_err(|e| FaucetError::Config(format!("postgres-cdc: invalid connection URL: {e}")))?;
360 let opts = apply_cdc_tls(opts, tls);
361
362 use sqlx::ConnectOptions as _;
363 let mut conn: PgConnection = opts
364 .connect()
365 .await
366 .map_err(|e| pg_err("postgres-cdc advance_slot connect", e))?;
367
368 sqlx::query("SELECT pg_replication_slot_advance($1, $2::pg_lsn)")
372 .bind(slot_name)
373 .bind(crate::state::format_lsn(lsn))
374 .execute(&mut conn)
375 .await
376 .map_err(|e| pg_err("postgres-cdc advance_slot", e))?;
377
378 debug!("postgres-cdc: advanced slot '{slot_name}' confirmed_flush_lsn to {lsn:#x}");
379 Ok(())
380}
381
382pub async fn ensure_slot_and_current_lsn(
399 connection_url: &str,
400 slot_name: &str,
401 create_if_missing: bool,
402 slot_type: crate::config::SlotType,
403 tls: &crate::config::CdcTls,
404) -> Result<u64, FaucetError> {
405 let client = Client { _private: () };
408 ensure_slot(
409 &client,
410 connection_url,
411 slot_name,
412 create_if_missing,
413 slot_type,
414 tls,
415 )
416 .await?;
417
418 let opts: PgConnectOptions = connection_url
419 .parse()
420 .map_err(|e| FaucetError::Config(format!("postgres-cdc: invalid connection URL: {e}")))?;
421 let opts = apply_cdc_tls(opts, tls);
422 use sqlx::ConnectOptions as _;
423 let mut conn: PgConnection = opts
424 .connect()
425 .await
426 .map_err(|e| pg_err("postgres-cdc capture_position connect", e))?;
427 let slot_lsn: Option<(Option<String>, Option<String>)> = sqlx::query_as(
429 "SELECT confirmed_flush_lsn::text, restart_lsn::text \
430 FROM pg_replication_slots WHERE slot_name = $1",
431 )
432 .bind(slot_name)
433 .fetch_optional(&mut conn)
434 .await
435 .map_err(|e| pg_err("postgres-cdc slot lsn lookup", e))?;
436 let anchor = match slot_lsn {
437 Some((confirmed, restart)) => confirmed.or(restart),
438 None => None,
439 };
440 let lsn_text = match anchor {
441 Some(t) => t,
442 None => {
443 let (cur,): (String,) = sqlx::query_as("SELECT pg_current_wal_lsn()::text")
446 .fetch_one(&mut conn)
447 .await
448 .map_err(|e| {
449 FaucetError::Source(format!("postgres-cdc pg_current_wal_lsn: {e}"))
450 })?;
451 cur
452 }
453 };
454 crate::state::parse_lsn(&lsn_text)
455}
456
457pub async fn slot_positions(
460 connection_url: &str,
461 slot_name: &str,
462 tls: &crate::config::CdcTls,
463) -> Result<Option<(u64, Option<u64>)>, FaucetError> {
464 let opts: PgConnectOptions = connection_url
465 .parse()
466 .map_err(|e| FaucetError::Config(format!("postgres-cdc: invalid connection URL: {e}")))?;
467 let opts = apply_cdc_tls(opts, tls);
468 use sqlx::ConnectOptions as _;
469 let mut conn: PgConnection = opts
470 .connect()
471 .await
472 .map_err(|e| pg_err("postgres-cdc lag connect", e))?;
473 let row: Option<(String, Option<String>)> = sqlx::query_as(
474 "SELECT pg_current_wal_lsn()::text, confirmed_flush_lsn::text \
475 FROM pg_replication_slots WHERE slot_name = $1",
476 )
477 .bind(slot_name)
478 .fetch_optional(&mut conn)
479 .await
480 .map_err(|e| pg_err("postgres-cdc lag query", e))?;
481 let Some((current, confirmed)) = row else {
482 return Ok(None);
483 };
484 let confirmed = confirmed
485 .as_deref()
486 .map(crate::state::parse_lsn)
487 .transpose()?;
488 Ok(Some((crate::state::parse_lsn(¤t)?, confirmed)))
489}
490
491pub async fn start_replication(
497 _client: &Client,
498 params: &ReplicationParams<'_>,
499) -> Result<Duplex, FaucetError> {
500 if params.proto_version != 1 {
501 return Err(FaucetError::Config(format!(
502 "postgres-cdc: pgwire-replication 0.3.2 supports proto_version = 1 only; \
503 got {}",
504 params.proto_version
505 )));
506 }
507
508 let coords = parse_url(params.connection_url)?;
509
510 let start_lsn = Lsn::from_u64(params.start_lsn.unwrap_or(0));
511
512 let cfg = ReplicationConfig {
513 host: coords.host,
514 port: coords.port,
515 user: coords.user,
516 password: coords.password,
517 database: coords.dbname,
518 tls: tls_config(params.tls),
519 slot: params.slot_name.to_owned(),
520 publication: params.publication_name.to_owned(),
521 start_lsn,
522 stop_at_lsn: None,
523 status_interval: params.status_update_interval,
526 idle_wakeup_interval: params.status_update_interval,
528 buffer_events: 8192,
529 };
530
531 let inner = ReplicationClient::connect(cfg)
532 .await
533 .map_err(|e| FaucetError::Source(format!("postgres-cdc start_replication: {e}")))?;
534
535 Ok(Duplex { inner })
536}
537
538pub async fn send_status_update(
548 duplex: &mut Duplex,
549 confirmed_lsn: u64,
550 _reply_requested: bool,
551) -> Result<(), FaucetError> {
552 duplex
553 .inner
554 .update_applied_lsn(Lsn::from_u64(confirmed_lsn));
555 Ok(())
556}
557
558pub async fn recv(duplex: &mut Duplex) -> Result<Option<ReplicationEvent>, FaucetError> {
581 loop {
582 match duplex
583 .inner
584 .recv()
585 .await
586 .map_err(|e| FaucetError::Source(format!("postgres-cdc recv: {e}")))?
587 {
588 None => return Ok(None),
589
590 Some(ReplicationEvent::StoppedAt { .. }) => {
591 return Ok(None);
592 }
593
594 Some(ReplicationEvent::KeepAlive { .. }) => {
595 }
598
599 Some(ev) => {
600 return Ok(Some(ev));
605 }
606 }
607 }
608}
609
610pub fn postgres_clock_now() -> i64 {
616 let now = SystemTime::now()
617 .duration_since(UNIX_EPOCH)
618 .unwrap_or_default();
619 let unix_micros = (now.as_secs() as i64) * 1_000_000 + (now.subsec_micros() as i64);
620 unix_micros - POSTGRES_EPOCH_MICROS
621}
622
623pub fn postgres_clock_to_unix_ms(ts: i64) -> i64 {
626 (POSTGRES_EPOCH_MICROS.saturating_add(ts)) / 1_000
627}
628
629#[allow(dead_code)]
635fn quote_slot(s: &str) -> String {
636 format!("\"{}\"", s.replace('"', "\"\""))
637}
638
639fn escape_simple(s: &str) -> String {
642 s.replace('\'', "''")
643}
644
645fn quote_literal(s: &str) -> String {
647 format!("'{}'", escape_simple(s))
648}
649
650pub fn is_slot_active_error(err: &FaucetError) -> bool {
661 sqlstate_of(err) == Some(SQLSTATE_OBJECT_IN_USE)
662}
663
664pub const SQLSTATE_OBJECT_IN_USE: &str = "55006";
667
668#[derive(Debug, thiserror::Error)]
678#[error("{context}: {message}")]
679pub struct PostgresError {
680 pub sqlstate: Option<String>,
682 context: String,
683 message: String,
684}
685
686impl PostgresError {
687 pub fn new(
689 sqlstate: Option<String>,
690 context: impl Into<String>,
691 message: impl Into<String>,
692 ) -> Self {
693 Self {
694 sqlstate,
695 context: context.into(),
696 message: message.into(),
697 }
698 }
699}
700
701pub(crate) fn pg_err(context: &str, e: sqlx::Error) -> FaucetError {
704 let sqlstate = e
705 .as_database_error()
706 .and_then(|d| d.code())
707 .map(|c| c.to_string());
708 FaucetError::Custom(Box::new(PostgresError {
709 sqlstate,
710 context: context.to_string(),
711 message: e.to_string(),
712 }))
713}
714
715fn sqlstate_of(err: &FaucetError) -> Option<&str> {
717 match err {
718 FaucetError::Custom(inner) => inner
719 .downcast_ref::<PostgresError>()
720 .and_then(|pg| pg.sqlstate.as_deref()),
721 _ => None,
722 }
723}
724
725fn slot_acquire_backoff(attempt: u32) -> Duration {
728 let factor = 1u64.checked_shl(attempt).unwrap_or(u64::MAX);
729 let ms = 250u64.saturating_mul(factor).min(4000);
730 Duration::from_millis(ms)
731}
732
733pub async fn retry_on_slot_active<F, Fut, T>(max_retries: u32, op: F) -> Result<T, FaucetError>
738where
739 F: Fn() -> Fut,
740 Fut: std::future::Future<Output = Result<T, FaucetError>>,
741{
742 let mut attempt = 0u32;
743 loop {
744 match op().await {
745 Ok(value) => return Ok(value),
746 Err(e) if attempt < max_retries && is_slot_active_error(&e) => {
747 let backoff = slot_acquire_backoff(attempt);
748 tracing::warn!(
749 attempt = attempt + 1,
750 max_retries,
751 backoff_ms = backoff.as_millis() as u64,
752 error = %e,
753 "postgres-cdc: replication slot still active; retrying after backoff"
754 );
755 tokio::time::sleep(backoff).await;
756 attempt += 1;
757 }
758 Err(e) => return Err(e),
759 }
760 }
761}
762
763#[cfg(test)]
764mod tests {
765 use super::*;
766 use crate::config::CdcTls;
767 use chrono::{TimeZone, Utc};
768 use pgwire_replication::SslMode;
769
770 #[test]
777 fn pg_err_preserves_context_and_message_without_a_sqlstate() {
778 let err = pg_err("postgres-cdc slot lookup", sqlx::Error::PoolTimedOut);
779 let FaucetError::Custom(inner) = &err else {
780 panic!("pg_err must produce FaucetError::Custom, got {err:?}");
781 };
782 let pg = inner
783 .downcast_ref::<PostgresError>()
784 .expect("the boxed error must still be a PostgresError");
785 assert_eq!(pg.context, "postgres-cdc slot lookup");
786 assert_eq!(pg.message, sqlx::Error::PoolTimedOut.to_string());
787 assert_eq!(
788 pg.sqlstate, None,
789 "a pool timeout is not a database-reported failure, so it has no SQLSTATE"
790 );
791 let rendered = err.to_string();
793 assert!(
794 rendered.contains("postgres-cdc slot lookup"),
795 "context must survive into the display form: {rendered}"
796 );
797 }
798
799 #[test]
804 fn sqlstate_of_reads_only_the_typed_carrier() {
805 let carried = PostgresError::new(
806 Some(SQLSTATE_OBJECT_IN_USE.to_string()),
807 "acquire slot",
808 "replication slot \"s\" is active for PID 42",
809 );
810 let err = FaucetError::Custom(Box::new(carried));
811 assert_eq!(sqlstate_of(&err), Some(SQLSTATE_OBJECT_IN_USE));
812 assert!(
813 is_slot_active_error(&err),
814 "the typed 55006 carrier is what slot-active detection keys off"
815 );
816
817 let codeless = FaucetError::Custom(Box::new(PostgresError::new(
819 None,
820 "acquire slot",
821 "replication slot is active for PID 42",
822 )));
823 assert_eq!(sqlstate_of(&codeless), None);
824 assert!(!is_slot_active_error(&codeless));
825
826 let lookalike = FaucetError::Source(
829 "replication slot \"s\" is active for PID 42 (SQLSTATE 55006)".into(),
830 );
831 assert_eq!(sqlstate_of(&lookalike), None);
832 assert!(
833 !is_slot_active_error(&lookalike),
834 "prose must never be mistaken for a SQLSTATE"
835 );
836
837 let other_code = FaucetError::Custom(Box::new(PostgresError::new(
839 Some("42P01".to_string()),
840 "acquire slot",
841 "relation does not exist",
842 )));
843 assert_eq!(sqlstate_of(&other_code), Some("42P01"));
844 assert!(!is_slot_active_error(&other_code));
845 }
846
847 #[test]
848 fn apply_cdc_tls_sets_ssl_mode_on_control_plane_options() {
849 let base: PgConnectOptions = "postgres://u:p@h:5432/db".parse().unwrap();
852 let dbg = |tls: &CdcTls| format!("{:?}", apply_cdc_tls(base.clone(), tls));
853 assert!(
854 dbg(&CdcTls::Disable).contains("Disable"),
855 "{}",
856 dbg(&CdcTls::Disable)
857 );
858 assert!(dbg(&CdcTls::Require).contains("Require"));
859 assert!(dbg(&CdcTls::VerifyCa { ca_path: None }).contains("VerifyCa"));
860 assert!(
861 dbg(&CdcTls::VerifyFull {
862 ca_path: Some("/ca.pem".into())
863 })
864 .contains("VerifyFull")
865 );
866 }
867
868 #[test]
869 fn tls_config_maps_each_mode() {
870 assert_eq!(tls_config(&CdcTls::Disable).mode, SslMode::Disable);
871 assert_eq!(tls_config(&CdcTls::Require).mode, SslMode::Require);
872 assert_eq!(
873 tls_config(&CdcTls::VerifyCa { ca_path: None }).mode,
874 SslMode::VerifyCa
875 );
876 assert_eq!(
877 tls_config(&CdcTls::VerifyFull {
878 ca_path: Some("/ca.pem".into())
879 })
880 .mode,
881 SslMode::VerifyFull
882 );
883 }
884
885 fn postgres_clock_to_datetime(ts: i64) -> chrono::DateTime<Utc> {
888 Utc.timestamp_micros(POSTGRES_EPOCH_MICROS.saturating_add(ts))
889 .single()
890 .unwrap_or_else(Utc::now)
891 }
892
893 #[test]
894 fn postgres_clock_round_trip() {
895 let dt = Utc.with_ymd_and_hms(2026, 5, 17, 12, 0, 0).unwrap();
896 let pg_ts = dt.timestamp_micros() - POSTGRES_EPOCH_MICROS;
897 let back = postgres_clock_to_datetime(pg_ts);
898 assert_eq!(back, dt);
899 }
900
901 #[test]
902 fn unix_ms_conversion() {
903 let dt = Utc.with_ymd_and_hms(2026, 5, 17, 12, 0, 0).unwrap();
904 let pg_ts = dt.timestamp_micros() - POSTGRES_EPOCH_MICROS;
905 assert_eq!(postgres_clock_to_unix_ms(pg_ts), 1_779_019_200_000);
906 }
907
908 #[test]
909 fn quote_slot_simple() {
910 assert_eq!(quote_slot("faucet_slot"), "\"faucet_slot\"");
911 }
912
913 #[test]
914 fn escape_simple_doubles_quotes() {
915 assert_eq!(escape_simple("foo'bar"), "foo''bar");
916 }
917
918 #[test]
919 fn parse_url_extracts_all_components() {
920 let c = parse_url("postgres://alice:secret@db.example.com:5544/analytics").unwrap();
921 assert_eq!(c.host, "db.example.com");
922 assert_eq!(c.port, 5544);
923 assert_eq!(c.user, "alice");
924 assert_eq!(c.password, "secret");
925 assert_eq!(c.dbname, "analytics");
926 }
927
928 #[test]
929 fn parse_url_defaults_port_and_dbname() {
930 let c = parse_url("postgres://alice@db.example.com").unwrap();
931 assert_eq!(c.port, 5432);
932 assert_eq!(c.dbname, "postgres");
933 assert_eq!(c.password, "");
934 }
935
936 #[test]
937 fn parse_url_rejects_missing_host() {
938 let err = parse_url("postgres:///analytics").unwrap_err();
940 assert!(format!("{err}").contains("missing a host"), "{err}");
941 }
942
943 #[test]
944 fn parse_url_rejects_missing_user() {
945 let err = parse_url("postgres://db.example.com/analytics").unwrap_err();
946 assert!(format!("{err}").contains("missing a user"), "{err}");
947 }
948
949 #[test]
950 fn slot_active_is_classified_by_sqlstate_not_message_text() {
951 assert!(is_slot_active_error(&slot_in_use()));
956
957 assert!(!is_slot_active_error(&FaucetError::Source(
959 "replication slot \"orders_is_active\" is active for PID 1".into()
960 )));
961
962 assert!(!is_slot_active_error(&FaucetError::Custom(Box::new(
964 PostgresError::new(
965 Some("42P01".to_string()),
966 "postgres-cdc slot lookup",
967 "relation does not exist",
968 )
969 ))));
970
971 assert!(!is_slot_active_error(&FaucetError::Config(
973 "bad url".into()
974 )));
975 }
976
977 #[test]
978 fn slot_acquire_backoff_grows_and_is_capped() {
979 assert_eq!(slot_acquire_backoff(0), Duration::from_millis(250));
980 assert_eq!(slot_acquire_backoff(1), Duration::from_millis(500));
981 assert_eq!(slot_acquire_backoff(2), Duration::from_millis(1000));
982 assert_eq!(slot_acquire_backoff(20), Duration::from_millis(4000));
984 assert_eq!(slot_acquire_backoff(64), Duration::from_millis(4000));
985 }
986
987 fn slot_in_use() -> FaucetError {
989 FaucetError::Custom(Box::new(PostgresError::new(
990 Some(SQLSTATE_OBJECT_IN_USE.to_string()),
991 "postgres-cdc create slot",
992 "replication slot \"s\" is active for PID 1",
993 )))
994 }
995
996 #[tokio::test]
997 async fn retry_on_slot_active_retries_then_succeeds() {
998 use std::sync::atomic::{AtomicU32, Ordering};
999 let calls = AtomicU32::new(0);
1000 let result = retry_on_slot_active(5, || {
1001 let n = calls.fetch_add(1, Ordering::SeqCst);
1002 async move {
1003 if n < 2 {
1004 Err(slot_in_use())
1005 } else {
1006 Ok::<u32, FaucetError>(42)
1007 }
1008 }
1009 })
1010 .await;
1011 assert_eq!(result.unwrap(), 42);
1012 assert_eq!(calls.load(Ordering::SeqCst), 3, "2 failures + 1 success");
1013 }
1014
1015 #[tokio::test]
1016 async fn retry_on_slot_active_gives_up_after_max_retries() {
1017 use std::sync::atomic::{AtomicU32, Ordering};
1018 let calls = AtomicU32::new(0);
1019 let result: Result<(), _> = retry_on_slot_active(2, || {
1020 calls.fetch_add(1, Ordering::SeqCst);
1021 async { Err(slot_in_use()) }
1022 })
1023 .await;
1024 assert!(result.is_err());
1025 assert_eq!(
1026 calls.load(Ordering::SeqCst),
1027 3,
1028 "initial attempt + 2 retries"
1029 );
1030 }
1031
1032 #[tokio::test]
1033 async fn retry_on_slot_active_does_not_retry_unrelated_errors() {
1034 use std::sync::atomic::{AtomicU32, Ordering};
1035 let calls = AtomicU32::new(0);
1036 let result: Result<(), _> = retry_on_slot_active(5, || {
1037 calls.fetch_add(1, Ordering::SeqCst);
1038 async { Err(FaucetError::Source("connection refused".into())) }
1039 })
1040 .await;
1041 assert!(result.is_err());
1042 assert_eq!(
1043 calls.load(Ordering::SeqCst),
1044 1,
1045 "a non-slot-active error must not be retried"
1046 );
1047 }
1048}