1pub mod config;
36pub mod edit;
37pub mod exec;
38pub mod introspect;
39
40pub use config::{PgAuthMethod, PgConfig, PgTlsMode};
41pub use edit::{InsertColumnInput, InsertedRow, UpdateOutcome};
42pub use exec::{ActiveCursor, ColumnMeta, ExecutionOutcome, PageResult};
43pub use introspect::{
44 ColumnDetail, DbSummary, ObjectType, ObjectTypeKind, Relation, RelationKind, Routine,
45 RoutineKind, SchemaContents, SchemaSummary, Sequence,
46};
47
48use std::collections::HashMap;
49use std::sync::{Arc, Mutex as StdMutex, Weak};
50use std::time::{Duration, Instant};
51
52use rustls::ClientConfig as RustlsClientConfig;
53use tokio::sync::Mutex;
54use tokio::task::JoinHandle;
55use tokio_postgres::config::SslMode as PgSslMode;
56use tokio_postgres::{CancelToken, Client, Config as PgDriverConfig};
57use tokio_postgres_rustls::MakeRustlsConnect;
58use tokio_util::sync::CancellationToken;
59
60const DEFAULT_MAX_POOL_SIZE: usize = 5;
65
66const DEFAULT_IDLE_TIMEOUT: Duration = Duration::from_secs(300);
71
72const EVICTION_INTERVAL: Duration = Duration::from_secs(30);
77
78const DEFAULT_MIN_IDLE_CONNECTIONS: usize = 1;
82
83pub const BROWSER_SESSION_ID: &str = "_browser";
87
88#[derive(Debug, thiserror::Error)]
90pub enum PgError {
91 #[error("postgres connect failed: {0}")]
92 Connect(String),
93 #[error("postgres auth failed: {0}")]
94 Auth(String),
95 #[error("postgres tls setup failed: {0}")]
96 Tls(String),
97 #[error("ssh tunnel error: {0}")]
98 Tunnel(String),
99 #[error("ssh tunnel source not found: {0}")]
102 TunnelSourceMissing(String),
103 #[error("cursor no longer available: {0}")]
108 CursorExpired(String),
109 #[error("invalid postgres input: {0}")]
110 InvalidInput(String),
111 #[error("pool exhausted: {0} of {1} connections leased")]
115 PoolExhausted(usize, usize),
116 #[error("postgres driver error: {0}")]
117 Driver(#[from] tokio_postgres::Error),
118}
119
120pub struct PgPool {
125 config: PgConfig,
126 tls_connector: TlsConnectorKind,
132 inner: Mutex<PoolInner>,
136 max_size: usize,
137 idle_timeout: Duration,
140 min_idle: usize,
143 secondary_browsers: Mutex<HashMap<String, SecondaryBrowserEntry>>,
151 eviction_cancel: CancellationToken,
156}
157
158struct IdleEntry {
163 since: Instant,
164 conn: Arc<PooledConnection>,
165}
166
167struct SecondaryBrowserEntry {
171 since: Instant,
172 conn: Arc<PooledConnection>,
173}
174
175#[derive(Clone)]
180enum TlsConnectorKind {
181 NoTls,
182 Rustls(MakeRustlsConnect),
183}
184
185struct PoolInner {
186 idle: Vec<IdleEntry>,
190 leased: HashMap<String, Arc<PooledConnection>>,
192 total: usize,
195}
196
197struct PooledConnection {
198 client: Client,
199 cancel_token: CancelToken,
200 operation_lock: Mutex<()>,
203 active_cursor: Mutex<Option<ActiveCursor>>,
206 connection_task: StdMutex<Option<JoinHandle<()>>>,
209}
210
211impl PooledConnection {
212 fn abort_connection_task(&self) {
213 if let Ok(mut task) = self.connection_task.lock()
214 && let Some(task) = task.take()
215 {
216 task.abort();
217 }
218 }
219}
220
221impl Drop for PooledConnection {
222 fn drop(&mut self) {
223 self.abort_connection_task();
224 }
225}
226
227impl PgPool {
228 pub async fn connect(cfg: PgConfig) -> Result<Arc<Self>, PgError> {
238 let tls_connector = build_tls_connector(&cfg)?;
241
242 let first = open_one(&cfg, &tls_connector).await?;
245
246 let now = Instant::now();
247 let max_size = cfg
251 .max_pool_size
252 .map(|n| n as usize)
253 .filter(|&n| n > 0)
254 .unwrap_or(DEFAULT_MAX_POOL_SIZE);
255 let idle_timeout = cfg
256 .idle_timeout_secs
257 .map(Duration::from_secs)
258 .unwrap_or(DEFAULT_IDLE_TIMEOUT);
259 let min_idle = cfg
260 .min_idle_connections
261 .map(|n| n as usize)
262 .unwrap_or(DEFAULT_MIN_IDLE_CONNECTIONS);
263
264 let pool = Arc::new(Self {
265 config: cfg,
266 tls_connector,
267 inner: Mutex::new(PoolInner {
268 idle: vec![IdleEntry {
269 since: now,
270 conn: Arc::new(first),
271 }],
272 leased: HashMap::new(),
273 total: 1,
274 }),
275 max_size,
276 idle_timeout,
277 min_idle,
278 secondary_browsers: Mutex::new(HashMap::new()),
279 eviction_cancel: CancellationToken::new(),
280 });
281
282 let weak = Arc::downgrade(&pool);
287 let cancel = pool.eviction_cancel.clone();
288 tokio::spawn(run_eviction(weak, cancel));
289
290 Ok(pool)
291 }
292
293 pub async fn list_databases(&self) -> Result<Vec<DbSummary>, PgError> {
302 let conn = self.lease_for_session(BROWSER_SESSION_ID).await?;
303 let _operation = conn.operation_lock.lock().await;
304 Ok(introspect::list_databases(&conn.client).await?)
305 }
306
307 pub async fn list_schemas(&self) -> Result<Vec<SchemaSummary>, PgError> {
308 self.list_schemas_in(None).await
309 }
310
311 pub async fn list_schemas_in(
318 &self,
319 database: Option<&str>,
320 ) -> Result<Vec<SchemaSummary>, PgError> {
321 let conn = self.browser_connection_for(database).await?;
322 let _operation = conn.operation_lock.lock().await;
323 Ok(introspect::list_schemas(&conn.client).await?)
324 }
325
326 pub async fn list_relations(&self, schema: &str) -> Result<Vec<Relation>, PgError> {
327 self.list_relations_in(schema, None).await
328 }
329
330 pub async fn list_relations_in(
331 &self,
332 schema: &str,
333 database: Option<&str>,
334 ) -> Result<Vec<Relation>, PgError> {
335 let conn = self.browser_connection_for(database).await?;
336 let _operation = conn.operation_lock.lock().await;
337 Ok(introspect::list_relations(&conn.client, schema).await?)
338 }
339
340 pub async fn list_schema_contents_in(
345 &self,
346 schema: &str,
347 database: Option<&str>,
348 ) -> Result<SchemaContents, PgError> {
349 let conn = self.browser_connection_for(database).await?;
350 let _operation = conn.operation_lock.lock().await;
351 Ok(introspect::list_schema_contents(&conn.client, schema).await?)
352 }
353
354 async fn browser_connection_for(
359 &self,
360 database: Option<&str>,
361 ) -> Result<Arc<PooledConnection>, PgError> {
362 let target = database.unwrap_or(self.config.database.as_str());
363 if target == self.config.database {
364 return self.lease_for_session(BROWSER_SESSION_ID).await;
365 }
366 {
368 let mut map = self.secondary_browsers.lock().await;
369 if let Some(entry) = map.get_mut(target) {
370 entry.since = Instant::now();
371 return Ok(entry.conn.clone());
372 }
373 }
374 self.reserve_connection_slot().await?;
375
376 let mut cfg = self.config.clone();
380 cfg.database = target.to_string();
381 let conn = match open_one(&cfg, &self.tls_connector).await {
382 Ok(conn) => conn,
383 Err(e) => {
384 self.release_connection_slot().await;
385 return Err(e);
386 }
387 };
388 let arc = Arc::new(conn);
389 let mut map = self.secondary_browsers.lock().await;
390 if let Some(existing) = map.get_mut(target) {
395 existing.since = Instant::now();
396 let existing = existing.conn.clone();
397 drop(map);
398 self.release_connection_slot().await;
399 return Ok(existing);
400 }
401 map.insert(
402 target.to_string(),
403 SecondaryBrowserEntry {
404 since: Instant::now(),
405 conn: arc.clone(),
406 },
407 );
408 Ok(arc)
409 }
410
411 pub async fn describe_columns(
414 &self,
415 schema: &str,
416 table: &str,
417 ) -> Result<Vec<ColumnDetail>, PgError> {
418 let conn = self.lease_for_session(BROWSER_SESSION_ID).await?;
419 let _operation = conn.operation_lock.lock().await;
420 Ok(introspect::describe_columns(&conn.client, schema, table).await?)
421 }
422
423 pub async fn execute(
428 &self,
429 session_id: &str,
430 sql: &str,
431 page_size: usize,
432 ) -> Result<ExecutionOutcome, PgError> {
433 let conn = self.lease_for_session(session_id).await?;
434 let _operation = conn.operation_lock.lock().await;
435 let previous = conn.active_cursor.lock().await.take();
436 let (outcome, new_cursor) =
437 exec::open_query(&conn.client, sql, page_size, previous).await?;
438 *conn.active_cursor.lock().await = new_cursor;
439 Ok(outcome)
445 }
446
447 pub async fn fetch_page(
448 &self,
449 session_id: &str,
450 cursor_id: &str,
451 count: usize,
452 ) -> Result<PageResult, PgError> {
453 let conn = self
454 .leased_only(session_id)
455 .await
456 .ok_or_else(|| PgError::CursorExpired(format!("no active session {session_id}")))?;
457 let _operation = conn.operation_lock.lock().await;
458 let Some(cursor) = conn.active_cursor.lock().await.clone() else {
459 return Err(PgError::CursorExpired(format!(
460 "session {session_id} has no active cursor"
461 )));
462 };
463 if cursor.cursor_id != cursor_id {
464 return Err(PgError::CursorExpired(format!(
465 "session {session_id} active cursor is {} (looking for {cursor_id})",
466 cursor.cursor_id
467 )));
468 }
469 exec::fetch_page(&conn.client, &cursor, count).await
470 }
471
472 #[allow(clippy::too_many_arguments)]
478 pub async fn update_cell(
479 &self,
480 session_id: &str,
481 schema: &str,
482 table: &str,
483 column: &str,
484 column_type: &str,
485 new_value: Option<&str>,
486 ctid: &str,
487 ) -> Result<UpdateOutcome, PgError> {
488 let conn = self.lease_for_session(session_id).await?;
489 let _operation = conn.operation_lock.lock().await;
490 edit::update_cell(
491 &conn.client,
492 schema,
493 table,
494 column,
495 column_type,
496 new_value,
497 ctid,
498 )
499 .await
500 }
501
502 pub async fn insert_row(
507 &self,
508 session_id: &str,
509 schema: &str,
510 table: &str,
511 inputs: &[InsertColumnInput],
512 return_columns: &[String],
513 ) -> Result<InsertedRow, PgError> {
514 let conn = self.lease_for_session(session_id).await?;
515 let _operation = conn.operation_lock.lock().await;
516 edit::insert_row(&conn.client, schema, table, inputs, return_columns).await
517 }
518
519 pub async fn delete_rows(
524 &self,
525 session_id: &str,
526 schema: &str,
527 table: &str,
528 ctids: &[String],
529 ) -> Result<UpdateOutcome, PgError> {
530 let conn = self.lease_for_session(session_id).await?;
531 let _operation = conn.operation_lock.lock().await;
532 edit::delete_rows(&conn.client, schema, table, ctids).await
533 }
534
535 pub async fn close_query(&self, session_id: &str, cursor_id: &str) -> Result<(), PgError> {
536 let Some(conn) = self.leased_only(session_id).await else {
537 return Ok(()); };
539 let _operation = conn.operation_lock.lock().await;
540 let cursor = {
541 let mut active = conn.active_cursor.lock().await;
542 match active.as_ref() {
543 Some(c) if c.cursor_id == cursor_id => active.take(),
544 _ => None,
545 }
546 };
547 if let Some(cursor) = cursor {
548 exec::close_query(&conn.client, &cursor).await;
549 }
550 Ok(())
551 }
552
553 pub async fn cancel(&self, session_id: &str) -> Result<(), PgError> {
560 let Some(conn) = self.leased_only(session_id).await else {
561 return Ok(());
562 };
563 let token = conn.cancel_token.clone();
564 match &self.tls_connector {
565 TlsConnectorKind::NoTls => token.cancel_query(tokio_postgres::NoTls).await,
566 TlsConnectorKind::Rustls(connector) => token.cancel_query(connector.clone()).await,
567 }
568 .map_err(PgError::Driver)
569 }
570
571 pub async fn release_session(&self, session_id: &str) {
574 let Some(conn) = self.take_lease(session_id).await else {
575 return;
576 };
577 {
580 let _operation = conn.operation_lock.lock().await;
581 let cursor = conn.active_cursor.lock().await.take();
582 if let Some(cursor) = cursor {
583 exec::close_query(&conn.client, &cursor).await;
584 }
585 }
586 let mut inner = self.inner.lock().await;
589 inner.idle.push(IdleEntry {
590 since: Instant::now(),
591 conn,
592 });
593 }
594
595 pub async fn shutdown(&self) {
597 self.eviction_cancel.cancel();
600
601 let mut inner = self.inner.lock().await;
602 let mut conns: Vec<Arc<PooledConnection>> = inner.idle.drain(..).map(|e| e.conn).collect();
603 conns.extend(inner.leased.drain().map(|(_, c)| c));
604 inner.total = 0;
605 drop(inner);
606 let secondaries: Vec<Arc<PooledConnection>> = {
609 let mut map = self.secondary_browsers.lock().await;
610 map.drain().map(|(_, entry)| entry.conn).collect()
611 };
612 let conns_with_secondaries = conns.into_iter().chain(secondaries);
613 let conns: Vec<Arc<PooledConnection>> = conns_with_secondaries.collect();
614 for conn in conns {
615 conn.abort_connection_task();
619 }
620 }
621
622 async fn lease_for_session(&self, session_id: &str) -> Result<Arc<PooledConnection>, PgError> {
629 {
631 let inner = self.inner.lock().await;
632 if let Some(c) = inner.leased.get(session_id) {
633 return Ok(c.clone());
634 }
635 }
636
637 let from_idle = {
641 let mut inner = self.inner.lock().await;
642 inner.idle.pop().map(|e| e.conn)
643 };
644 if let Some(conn) = from_idle {
645 self.assign_lease(session_id, conn.clone()).await;
646 return Ok(conn);
647 }
648
649 let need_new = {
651 let inner = self.inner.lock().await;
652 if inner.total >= self.max_size {
653 return Err(PgError::PoolExhausted(inner.total, self.max_size));
654 }
655 true
656 };
657 if need_new {
658 {
662 let mut inner = self.inner.lock().await;
663 if inner.total >= self.max_size {
664 return Err(PgError::PoolExhausted(inner.total, self.max_size));
665 }
666 inner.total += 1;
667 }
668 let new_conn =
669 match open_one(&self.config, &self.tls_connector).await {
670 Ok(c) => c,
671 Err(e) => {
672 let mut inner = self.inner.lock().await;
675 inner.total = inner.total.saturating_sub(1);
676 return Err(e);
677 }
678 };
679 let conn = Arc::new(new_conn);
680 self.assign_lease(session_id, conn.clone()).await;
681 return Ok(conn);
682 }
683 unreachable!()
684 }
685
686 async fn assign_lease(&self, session_id: &str, conn: Arc<PooledConnection>) {
687 let mut inner = self.inner.lock().await;
688 inner.leased.insert(session_id.to_string(), conn);
689 }
690
691 async fn reserve_connection_slot(&self) -> Result<(), PgError> {
692 let mut inner = self.inner.lock().await;
693 if inner.total >= self.max_size {
694 return Err(PgError::PoolExhausted(inner.total, self.max_size));
695 }
696 inner.total += 1;
697 Ok(())
698 }
699
700 async fn release_connection_slot(&self) {
701 let mut inner = self.inner.lock().await;
702 inner.total = inner.total.saturating_sub(1);
703 }
704
705 async fn evict_idle(&self) {
709 let now = Instant::now();
710 let to_drop: Vec<Arc<PooledConnection>>;
711 {
712 let mut inner = self.inner.lock().await;
713 let snapshot = std::mem::take(&mut inner.idle);
720 let mut keep: Vec<IdleEntry> = Vec::with_capacity(snapshot.len());
721 let mut drop_list: Vec<Arc<PooledConnection>> = Vec::new();
722 for entry in snapshot.into_iter() {
723 let aged = now.duration_since(entry.since) >= self.idle_timeout;
724 if !aged || keep.len() < self.min_idle {
725 keep.push(entry);
726 } else {
727 drop_list.push(entry.conn);
728 inner.total = inner.total.saturating_sub(1);
729 }
730 }
731 inner.idle = keep;
732 to_drop = drop_list;
733 }
734
735 if !to_drop.is_empty() {
736 tracing::debug!(
737 target: "postgres::pool",
738 count = to_drop.len(),
739 "evicted idle postgres connections"
740 );
741 }
742
743 for conn in to_drop {
748 conn.abort_connection_task();
749 }
750
751 self.evict_secondary_browsers(now).await;
752 }
753
754 async fn evict_secondary_browsers(&self, now: Instant) {
755 let to_drop: Vec<Arc<PooledConnection>> = {
756 let mut map = self.secondary_browsers.lock().await;
757 let drop_keys = map
758 .iter()
759 .filter_map(|(database, entry)| {
760 let aged = now.duration_since(entry.since) >= self.idle_timeout;
761 if aged && Arc::strong_count(&entry.conn) == 1 {
764 Some(database.clone())
765 } else {
766 None
767 }
768 })
769 .collect::<Vec<_>>();
770
771 let mut dropped = Vec::with_capacity(drop_keys.len());
772 for database in drop_keys {
773 if let Some(entry) = map.remove(&database) {
774 dropped.push(entry.conn);
775 }
776 }
777 dropped
778 };
779
780 if to_drop.is_empty() {
781 return;
782 }
783
784 {
785 let mut inner = self.inner.lock().await;
786 for _ in &to_drop {
787 inner.total = inner.total.saturating_sub(1);
788 }
789 }
790
791 tracing::debug!(
792 target: "postgres::pool",
793 count = to_drop.len(),
794 "evicted secondary postgres browser connections"
795 );
796
797 for conn in to_drop {
798 conn.abort_connection_task();
799 }
800 }
801
802 async fn leased_only(&self, session_id: &str) -> Option<Arc<PooledConnection>> {
804 let inner = self.inner.lock().await;
805 inner.leased.get(session_id).cloned()
806 }
807
808 async fn take_lease(&self, session_id: &str) -> Option<Arc<PooledConnection>> {
810 let mut inner = self.inner.lock().await;
811 inner.leased.remove(session_id)
812 }
813}
814
815impl Drop for PgPool {
816 fn drop(&mut self) {
817 self.eviction_cancel.cancel();
820
821 if let Ok(mut inner) = self.inner.try_lock() {
825 let mut conns: Vec<Arc<PooledConnection>> =
826 inner.idle.drain(..).map(|e| e.conn).collect();
827 conns.extend(inner.leased.drain().map(|(_, c)| c));
828 for conn in conns {
829 conn.abort_connection_task();
830 }
831 }
832 if let Ok(mut map) = self.secondary_browsers.try_lock() {
834 for (_, entry) in map.drain() {
835 entry.conn.abort_connection_task();
836 }
837 }
838 }
839}
840
841async fn run_eviction(pool: Weak<PgPool>, cancel: CancellationToken) {
850 let mut ticker = tokio::time::interval(EVICTION_INTERVAL);
851 ticker.tick().await;
854 loop {
855 tokio::select! {
856 _ = cancel.cancelled() => return,
857 _ = ticker.tick() => {
858 let Some(pool) = pool.upgrade() else { return };
859 pool.evict_idle().await;
860 drop(pool);
864 }
865 }
866 }
867}
868
869async fn open_one(
874 cfg: &PgConfig,
875 tls: &TlsConnectorKind,
876) -> Result<PooledConnection, PgError> {
877 let driver_cfg = build_driver_config(cfg)?;
878 match tls {
879 TlsConnectorKind::NoTls => {
880 let (client, connection) = driver_cfg
881 .connect(tokio_postgres::NoTls)
882 .await
883 .map_err(classify_connect_error)?;
884 Ok(spawn_connection(client, connection))
885 }
886 TlsConnectorKind::Rustls(connector) => {
887 let (client, connection) = driver_cfg
890 .connect(connector.clone())
891 .await
892 .map_err(classify_connect_error)?;
893 Ok(spawn_connection(client, connection))
894 }
895 }
896}
897
898fn build_tls_connector(cfg: &PgConfig) -> Result<TlsConnectorKind, PgError> {
903 match cfg.tls {
904 PgTlsMode::Disable => Ok(TlsConnectorKind::NoTls),
905 PgTlsMode::Prefer | PgTlsMode::Require | PgTlsMode::VerifyFull => {
906 let _ = rustls::crypto::ring::default_provider().install_default();
907 let tls_config = build_rustls_config(cfg.tls)?;
908 Ok(TlsConnectorKind::Rustls(MakeRustlsConnect::new(tls_config)))
909 }
910 }
911}
912
913fn build_driver_config(cfg: &PgConfig) -> Result<PgDriverConfig, PgError> {
914 let mut driver = PgDriverConfig::new();
915 driver.host(&cfg.host).port(cfg.port);
916 driver.dbname(&cfg.database).user(&cfg.user);
917
918 let password = match &cfg.auth {
919 PgAuthMethod::Password { password } => password.clone(),
920 PgAuthMethod::Keychain { account } => ssh_commander_keychain::load_password(
921 ssh_commander_keychain::CredentialKind::PostgresPassword,
922 account,
923 )
924 .map_err(|e| PgError::Auth(format!("keychain load failed for {account}: {e}")))?
925 .ok_or_else(|| {
926 PgError::Auth(format!("no keychain entry for postgres account {account}"))
927 })?,
928 };
929 if !password.is_empty() {
930 driver.password(password);
931 }
932
933 if let Some(name) = &cfg.application_name {
934 driver.application_name(name);
935 }
936 if let Some(secs) = cfg.connect_timeout_secs {
937 driver.connect_timeout(Duration::from_secs(secs));
938 }
939
940 driver.ssl_mode(match cfg.tls {
941 PgTlsMode::Disable => PgSslMode::Disable,
942 PgTlsMode::Prefer => PgSslMode::Prefer,
943 PgTlsMode::Require | PgTlsMode::VerifyFull => PgSslMode::Require,
944 });
945 Ok(driver)
946}
947
948fn build_rustls_config(mode: PgTlsMode) -> Result<RustlsClientConfig, PgError> {
949 let mut roots = rustls::RootCertStore::empty();
950 let native = rustls_native_certs::load_native_certs();
951 for cert in native.certs {
952 let _ = roots.add(cert);
953 }
954
955 let cfg = match mode {
956 PgTlsMode::VerifyFull => RustlsClientConfig::builder()
957 .with_root_certificates(roots)
958 .with_no_client_auth(),
959 _ => RustlsClientConfig::builder()
960 .dangerous()
961 .with_custom_certificate_verifier(std::sync::Arc::new(NoCertVerifier))
962 .with_no_client_auth(),
963 };
964 Ok(cfg)
965}
966
967fn classify_connect_error(e: tokio_postgres::Error) -> PgError {
968 if let Some(db_err) = e.as_db_error() {
969 let code = db_err.code().code();
970 if code == "28P01" || code == "28000" {
971 return PgError::Auth(db_err.message().to_string());
972 }
973 }
974 PgError::Connect(e.to_string())
975}
976
977fn spawn_connection<S, T>(
978 client: Client,
979 connection: tokio_postgres::Connection<S, T>,
980) -> PooledConnection
981where
982 S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static,
983 T: tokio_postgres::tls::TlsStream + Unpin + Send + 'static,
984{
985 let cancel_token = client.cancel_token();
986 let task = tokio::spawn(async move {
987 if let Err(e) = connection.await {
988 tracing::warn!("postgres connection task ended with error: {e}");
989 }
990 });
991 PooledConnection {
992 client,
993 cancel_token,
994 operation_lock: Mutex::new(()),
995 active_cursor: Mutex::new(None),
996 connection_task: StdMutex::new(Some(task)),
997 }
998}
999
1000#[derive(Debug)]
1001struct NoCertVerifier;
1002
1003impl rustls::client::danger::ServerCertVerifier for NoCertVerifier {
1004 fn verify_server_cert(
1005 &self,
1006 _end_entity: &rustls::pki_types::CertificateDer<'_>,
1007 _intermediates: &[rustls::pki_types::CertificateDer<'_>],
1008 _server_name: &rustls::pki_types::ServerName<'_>,
1009 _ocsp_response: &[u8],
1010 _now: rustls::pki_types::UnixTime,
1011 ) -> Result<rustls::client::danger::ServerCertVerified, rustls::Error> {
1012 Ok(rustls::client::danger::ServerCertVerified::assertion())
1013 }
1014
1015 fn verify_tls12_signature(
1016 &self,
1017 _message: &[u8],
1018 _cert: &rustls::pki_types::CertificateDer<'_>,
1019 _dss: &rustls::DigitallySignedStruct,
1020 ) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
1021 Ok(rustls::client::danger::HandshakeSignatureValid::assertion())
1022 }
1023
1024 fn verify_tls13_signature(
1025 &self,
1026 _message: &[u8],
1027 _cert: &rustls::pki_types::CertificateDer<'_>,
1028 _dss: &rustls::DigitallySignedStruct,
1029 ) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
1030 Ok(rustls::client::danger::HandshakeSignatureValid::assertion())
1031 }
1032
1033 fn supported_verify_schemes(&self) -> Vec<rustls::SignatureScheme> {
1034 vec![
1035 rustls::SignatureScheme::RSA_PKCS1_SHA256,
1036 rustls::SignatureScheme::RSA_PKCS1_SHA384,
1037 rustls::SignatureScheme::RSA_PKCS1_SHA512,
1038 rustls::SignatureScheme::ECDSA_NISTP256_SHA256,
1039 rustls::SignatureScheme::ECDSA_NISTP384_SHA384,
1040 rustls::SignatureScheme::ED25519,
1041 rustls::SignatureScheme::RSA_PSS_SHA256,
1042 rustls::SignatureScheme::RSA_PSS_SHA384,
1043 rustls::SignatureScheme::RSA_PSS_SHA512,
1044 ]
1045 }
1046}
1047
1048#[cfg(test)]
1049mod tests {
1050 use super::*;
1051
1052 #[test]
1053 fn driver_config_uses_correct_ssl_mode() {
1054 let mut cfg = PgConfig::local("db", "u");
1055 cfg.tls = PgTlsMode::Require;
1056 let driver = build_driver_config(&cfg).expect("driver cfg");
1057 assert!(matches!(driver.get_ssl_mode(), PgSslMode::Require));
1058
1059 cfg.tls = PgTlsMode::Disable;
1060 let driver = build_driver_config(&cfg).expect("driver cfg");
1061 assert!(matches!(driver.get_ssl_mode(), PgSslMode::Disable));
1062 }
1063
1064 #[test]
1065 fn driver_config_omits_password_when_empty() {
1066 let cfg = PgConfig::local("db", "u");
1067 let driver = build_driver_config(&cfg).expect("driver cfg");
1068 assert!(driver.get_password().is_none());
1069 }
1070
1071 #[test]
1076 fn eviction_policy_keeps_min_idle_and_drops_aged() {
1077 fn run_policy(
1081 entries: Vec<(usize, Instant)>,
1082 now: Instant,
1083 idle_timeout: Duration,
1084 min_idle: usize,
1085 ) -> (Vec<usize>, Vec<usize>) {
1086 let mut keep: Vec<usize> = Vec::new();
1087 let mut drop_idx: Vec<usize> = Vec::new();
1088 for (idx, since) in entries {
1089 let aged = now.duration_since(since) >= idle_timeout;
1090 if !aged || keep.len() < min_idle {
1091 keep.push(idx);
1092 } else {
1093 drop_idx.push(idx);
1094 }
1095 }
1096 (keep, drop_idx)
1097 }
1098
1099 let now = Instant::now();
1100 let timeout = Duration::from_secs(300);
1101 let aged = now - timeout - Duration::from_secs(1);
1102 let fresh = now - Duration::from_secs(10);
1103
1104 let (keep, drop_idx) = run_policy(
1107 vec![(0, aged), (1, aged), (2, fresh), (3, aged)],
1108 now,
1109 timeout,
1110 1,
1111 );
1112 assert_eq!(keep, vec![0, 2]);
1113 assert_eq!(drop_idx, vec![1, 3]);
1114
1115 let (keep, drop_idx) = run_policy(vec![(0, aged), (1, fresh), (2, aged)], now, timeout, 0);
1117 assert_eq!(keep, vec![1]);
1118 assert_eq!(drop_idx, vec![0, 2]);
1119 }
1120}