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 = match open_one(&self.config, &self.tls_connector).await {
669 Ok(c) => c,
670 Err(e) => {
671 let mut inner = self.inner.lock().await;
674 inner.total = inner.total.saturating_sub(1);
675 return Err(e);
676 }
677 };
678 let conn = Arc::new(new_conn);
679 self.assign_lease(session_id, conn.clone()).await;
680 return Ok(conn);
681 }
682 unreachable!()
683 }
684
685 async fn assign_lease(&self, session_id: &str, conn: Arc<PooledConnection>) {
686 let mut inner = self.inner.lock().await;
687 inner.leased.insert(session_id.to_string(), conn);
688 }
689
690 async fn reserve_connection_slot(&self) -> Result<(), PgError> {
691 let mut inner = self.inner.lock().await;
692 if inner.total >= self.max_size {
693 return Err(PgError::PoolExhausted(inner.total, self.max_size));
694 }
695 inner.total += 1;
696 Ok(())
697 }
698
699 async fn release_connection_slot(&self) {
700 let mut inner = self.inner.lock().await;
701 inner.total = inner.total.saturating_sub(1);
702 }
703
704 async fn evict_idle(&self) {
708 let now = Instant::now();
709 let to_drop: Vec<Arc<PooledConnection>>;
710 {
711 let mut inner = self.inner.lock().await;
712 let snapshot = std::mem::take(&mut inner.idle);
719 let mut keep: Vec<IdleEntry> = Vec::with_capacity(snapshot.len());
720 let mut drop_list: Vec<Arc<PooledConnection>> = Vec::new();
721 for entry in snapshot.into_iter() {
722 let aged = now.duration_since(entry.since) >= self.idle_timeout;
723 if !aged || keep.len() < self.min_idle {
724 keep.push(entry);
725 } else {
726 drop_list.push(entry.conn);
727 inner.total = inner.total.saturating_sub(1);
728 }
729 }
730 inner.idle = keep;
731 to_drop = drop_list;
732 }
733
734 if !to_drop.is_empty() {
735 tracing::debug!(
736 target: "postgres::pool",
737 count = to_drop.len(),
738 "evicted idle postgres connections"
739 );
740 }
741
742 for conn in to_drop {
747 conn.abort_connection_task();
748 }
749
750 self.evict_secondary_browsers(now).await;
751 }
752
753 async fn evict_secondary_browsers(&self, now: Instant) {
754 let to_drop: Vec<Arc<PooledConnection>> = {
755 let mut map = self.secondary_browsers.lock().await;
756 let drop_keys = map
757 .iter()
758 .filter_map(|(database, entry)| {
759 let aged = now.duration_since(entry.since) >= self.idle_timeout;
760 if aged && Arc::strong_count(&entry.conn) == 1 {
763 Some(database.clone())
764 } else {
765 None
766 }
767 })
768 .collect::<Vec<_>>();
769
770 let mut dropped = Vec::with_capacity(drop_keys.len());
771 for database in drop_keys {
772 if let Some(entry) = map.remove(&database) {
773 dropped.push(entry.conn);
774 }
775 }
776 dropped
777 };
778
779 if to_drop.is_empty() {
780 return;
781 }
782
783 {
784 let mut inner = self.inner.lock().await;
785 for _ in &to_drop {
786 inner.total = inner.total.saturating_sub(1);
787 }
788 }
789
790 tracing::debug!(
791 target: "postgres::pool",
792 count = to_drop.len(),
793 "evicted secondary postgres browser connections"
794 );
795
796 for conn in to_drop {
797 conn.abort_connection_task();
798 }
799 }
800
801 async fn leased_only(&self, session_id: &str) -> Option<Arc<PooledConnection>> {
803 let inner = self.inner.lock().await;
804 inner.leased.get(session_id).cloned()
805 }
806
807 async fn take_lease(&self, session_id: &str) -> Option<Arc<PooledConnection>> {
809 let mut inner = self.inner.lock().await;
810 inner.leased.remove(session_id)
811 }
812}
813
814impl Drop for PgPool {
815 fn drop(&mut self) {
816 self.eviction_cancel.cancel();
819
820 if let Ok(mut inner) = self.inner.try_lock() {
824 let mut conns: Vec<Arc<PooledConnection>> =
825 inner.idle.drain(..).map(|e| e.conn).collect();
826 conns.extend(inner.leased.drain().map(|(_, c)| c));
827 for conn in conns {
828 conn.abort_connection_task();
829 }
830 }
831 if let Ok(mut map) = self.secondary_browsers.try_lock() {
833 for (_, entry) in map.drain() {
834 entry.conn.abort_connection_task();
835 }
836 }
837 }
838}
839
840async fn run_eviction(pool: Weak<PgPool>, cancel: CancellationToken) {
849 let mut ticker = tokio::time::interval(EVICTION_INTERVAL);
850 ticker.tick().await;
853 loop {
854 tokio::select! {
855 _ = cancel.cancelled() => return,
856 _ = ticker.tick() => {
857 let Some(pool) = pool.upgrade() else { return };
858 pool.evict_idle().await;
859 drop(pool);
863 }
864 }
865 }
866}
867
868async fn open_one(cfg: &PgConfig, tls: &TlsConnectorKind) -> Result<PooledConnection, PgError> {
873 let driver_cfg = build_driver_config(cfg)?;
874 match tls {
875 TlsConnectorKind::NoTls => {
876 let (client, connection) = driver_cfg
877 .connect(tokio_postgres::NoTls)
878 .await
879 .map_err(classify_connect_error)?;
880 Ok(spawn_connection(client, connection))
881 }
882 TlsConnectorKind::Rustls(connector) => {
883 let (client, connection) = driver_cfg
886 .connect(connector.clone())
887 .await
888 .map_err(classify_connect_error)?;
889 Ok(spawn_connection(client, connection))
890 }
891 }
892}
893
894fn build_tls_connector(cfg: &PgConfig) -> Result<TlsConnectorKind, PgError> {
899 match cfg.tls {
900 PgTlsMode::Disable => Ok(TlsConnectorKind::NoTls),
901 PgTlsMode::Prefer | PgTlsMode::Require | PgTlsMode::VerifyFull => {
902 let _ = rustls::crypto::ring::default_provider().install_default();
903 let tls_config = build_rustls_config(cfg.tls)?;
904 Ok(TlsConnectorKind::Rustls(MakeRustlsConnect::new(tls_config)))
905 }
906 }
907}
908
909fn build_driver_config(cfg: &PgConfig) -> Result<PgDriverConfig, PgError> {
910 let mut driver = PgDriverConfig::new();
911 driver.host(&cfg.host).port(cfg.port);
912 driver.dbname(&cfg.database).user(&cfg.user);
913
914 let password = match &cfg.auth {
915 PgAuthMethod::Password { password } => password.clone(),
916 PgAuthMethod::Keychain { account } => ssh_commander_keychain::load_password(
917 ssh_commander_keychain::CredentialKind::PostgresPassword,
918 account,
919 )
920 .map_err(|e| PgError::Auth(format!("keychain load failed for {account}: {e}")))?
921 .ok_or_else(|| {
922 PgError::Auth(format!("no keychain entry for postgres account {account}"))
923 })?,
924 };
925 if !password.is_empty() {
926 driver.password(password);
927 }
928
929 if let Some(name) = &cfg.application_name {
930 driver.application_name(name);
931 }
932 if let Some(secs) = cfg.connect_timeout_secs {
933 driver.connect_timeout(Duration::from_secs(secs));
934 }
935
936 driver.ssl_mode(match cfg.tls {
937 PgTlsMode::Disable => PgSslMode::Disable,
938 PgTlsMode::Prefer => PgSslMode::Prefer,
939 PgTlsMode::Require | PgTlsMode::VerifyFull => PgSslMode::Require,
940 });
941 Ok(driver)
942}
943
944fn build_rustls_config(mode: PgTlsMode) -> Result<RustlsClientConfig, PgError> {
945 let mut roots = rustls::RootCertStore::empty();
946 let native = rustls_native_certs::load_native_certs();
947 for cert in native.certs {
948 let _ = roots.add(cert);
949 }
950
951 let cfg = match mode {
952 PgTlsMode::VerifyFull => RustlsClientConfig::builder()
953 .with_root_certificates(roots)
954 .with_no_client_auth(),
955 _ => RustlsClientConfig::builder()
956 .dangerous()
957 .with_custom_certificate_verifier(std::sync::Arc::new(NoCertVerifier))
958 .with_no_client_auth(),
959 };
960 Ok(cfg)
961}
962
963fn classify_connect_error(e: tokio_postgres::Error) -> PgError {
964 if let Some(db_err) = e.as_db_error() {
965 let code = db_err.code().code();
966 if code == "28P01" || code == "28000" {
967 return PgError::Auth(db_err.message().to_string());
968 }
969 }
970 PgError::Connect(e.to_string())
971}
972
973fn spawn_connection<S, T>(
974 client: Client,
975 connection: tokio_postgres::Connection<S, T>,
976) -> PooledConnection
977where
978 S: tokio::io::AsyncRead + tokio::io::AsyncWrite + Unpin + Send + 'static,
979 T: tokio_postgres::tls::TlsStream + Unpin + Send + 'static,
980{
981 let cancel_token = client.cancel_token();
982 let task = tokio::spawn(async move {
983 if let Err(e) = connection.await {
984 tracing::warn!("postgres connection task ended with error: {e}");
985 }
986 });
987 PooledConnection {
988 client,
989 cancel_token,
990 operation_lock: Mutex::new(()),
991 active_cursor: Mutex::new(None),
992 connection_task: StdMutex::new(Some(task)),
993 }
994}
995
996#[derive(Debug)]
997struct NoCertVerifier;
998
999impl rustls::client::danger::ServerCertVerifier for NoCertVerifier {
1000 fn verify_server_cert(
1001 &self,
1002 _end_entity: &rustls::pki_types::CertificateDer<'_>,
1003 _intermediates: &[rustls::pki_types::CertificateDer<'_>],
1004 _server_name: &rustls::pki_types::ServerName<'_>,
1005 _ocsp_response: &[u8],
1006 _now: rustls::pki_types::UnixTime,
1007 ) -> Result<rustls::client::danger::ServerCertVerified, rustls::Error> {
1008 Ok(rustls::client::danger::ServerCertVerified::assertion())
1009 }
1010
1011 fn verify_tls12_signature(
1012 &self,
1013 _message: &[u8],
1014 _cert: &rustls::pki_types::CertificateDer<'_>,
1015 _dss: &rustls::DigitallySignedStruct,
1016 ) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
1017 Ok(rustls::client::danger::HandshakeSignatureValid::assertion())
1018 }
1019
1020 fn verify_tls13_signature(
1021 &self,
1022 _message: &[u8],
1023 _cert: &rustls::pki_types::CertificateDer<'_>,
1024 _dss: &rustls::DigitallySignedStruct,
1025 ) -> Result<rustls::client::danger::HandshakeSignatureValid, rustls::Error> {
1026 Ok(rustls::client::danger::HandshakeSignatureValid::assertion())
1027 }
1028
1029 fn supported_verify_schemes(&self) -> Vec<rustls::SignatureScheme> {
1030 vec![
1031 rustls::SignatureScheme::RSA_PKCS1_SHA256,
1032 rustls::SignatureScheme::RSA_PKCS1_SHA384,
1033 rustls::SignatureScheme::RSA_PKCS1_SHA512,
1034 rustls::SignatureScheme::ECDSA_NISTP256_SHA256,
1035 rustls::SignatureScheme::ECDSA_NISTP384_SHA384,
1036 rustls::SignatureScheme::ED25519,
1037 rustls::SignatureScheme::RSA_PSS_SHA256,
1038 rustls::SignatureScheme::RSA_PSS_SHA384,
1039 rustls::SignatureScheme::RSA_PSS_SHA512,
1040 ]
1041 }
1042}
1043
1044#[cfg(test)]
1045mod tests {
1046 use super::*;
1047
1048 #[test]
1049 fn driver_config_uses_correct_ssl_mode() {
1050 let mut cfg = PgConfig::local("db", "u");
1051 cfg.tls = PgTlsMode::Require;
1052 let driver = build_driver_config(&cfg).expect("driver cfg");
1053 assert!(matches!(driver.get_ssl_mode(), PgSslMode::Require));
1054
1055 cfg.tls = PgTlsMode::Disable;
1056 let driver = build_driver_config(&cfg).expect("driver cfg");
1057 assert!(matches!(driver.get_ssl_mode(), PgSslMode::Disable));
1058 }
1059
1060 #[test]
1061 fn driver_config_omits_password_when_empty() {
1062 let cfg = PgConfig::local("db", "u");
1063 let driver = build_driver_config(&cfg).expect("driver cfg");
1064 assert!(driver.get_password().is_none());
1065 }
1066
1067 #[test]
1072 fn eviction_policy_keeps_min_idle_and_drops_aged() {
1073 fn run_policy(
1077 entries: Vec<(usize, Instant)>,
1078 now: Instant,
1079 idle_timeout: Duration,
1080 min_idle: usize,
1081 ) -> (Vec<usize>, Vec<usize>) {
1082 let mut keep: Vec<usize> = Vec::new();
1083 let mut drop_idx: Vec<usize> = Vec::new();
1084 for (idx, since) in entries {
1085 let aged = now.duration_since(since) >= idle_timeout;
1086 if !aged || keep.len() < min_idle {
1087 keep.push(idx);
1088 } else {
1089 drop_idx.push(idx);
1090 }
1091 }
1092 (keep, drop_idx)
1093 }
1094
1095 let now = Instant::now();
1096 let timeout = Duration::from_secs(300);
1097 let aged = now - timeout - Duration::from_secs(1);
1098 let fresh = now - Duration::from_secs(10);
1099
1100 let (keep, drop_idx) = run_policy(
1103 vec![(0, aged), (1, aged), (2, fresh), (3, aged)],
1104 now,
1105 timeout,
1106 1,
1107 );
1108 assert_eq!(keep, vec![0, 2]);
1109 assert_eq!(drop_idx, vec![1, 3]);
1110
1111 let (keep, drop_idx) = run_policy(vec![(0, aged), (1, fresh), (2, aged)], now, timeout, 0);
1113 assert_eq!(keep, vec![1]);
1114 assert_eq!(drop_idx, vec![0, 2]);
1115 }
1116}