1use std::path::PathBuf;
68use std::str::FromStr;
69
70use crate::egress::auth::AuthMode;
71use crate::error::{Result, fmt};
72use crate::ingress::CertificateAuthority;
73
74pub const DEFAULT_PATH: &str = "/read/v1";
76
77pub const HIGHEST_KNOWN_VERSION: u8 = crate::egress::wire::PROTOCOL_VERSION;
89
90const DEFAULT_PLAIN_PORT: &str = "9000";
92const DEFAULT_TLS_PORT: &str = "9000";
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq)]
112#[non_exhaustive]
113pub enum Compression {
114 Raw,
118 Zstd,
122 Auto,
129}
130
131impl Compression {
132 pub fn accept_encoding(self, level: u8) -> String {
138 match self {
139 Compression::Raw => "raw".to_string(),
140 Compression::Zstd => format!("zstd;level={}", level),
141 Compression::Auto => format!("zstd;level={},raw", level),
142 }
143 }
144
145 pub fn header_token(self) -> &'static str {
150 match self {
151 Compression::Raw => "raw",
152 Compression::Zstd => "zstd",
153 Compression::Auto => "zstd,raw",
154 }
155 }
156}
157
158#[derive(Debug, Clone, Copy, PartialEq, Eq)]
161#[non_exhaustive]
162pub enum Target {
163 Any,
165 Primary,
170 Replica,
174}
175
176#[derive(Debug, Clone, PartialEq, Eq, Hash)]
191#[non_exhaustive]
192pub struct Endpoint {
193 pub host: String,
199 pub port: u16,
203}
204
205impl Endpoint {
206 pub fn new<S: Into<String>>(host: S, port: u16) -> Self {
215 Endpoint {
216 host: host.into(),
217 port,
218 }
219 }
220}
221
222impl std::fmt::Display for Endpoint {
228 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
229 if self.host.contains(':') {
230 write!(f, "[{}]:{}", self.host, self.port)
231 } else {
232 write!(f, "{}:{}", self.host, self.port)
233 }
234 }
235}
236
237pub const DEFAULT_FAILOVER_ENABLED: bool = true;
247
248pub const DEFAULT_FAILOVER_MAX_ATTEMPTS: u32 = 8;
252
253pub const DEFAULT_FAILOVER_BACKOFF_INITIAL_MS: u64 = 50;
258
259pub const DEFAULT_FAILOVER_BACKOFF_MAX_MS: u64 = 1_000;
264
265pub const MAX_FAILOVER_MAX_ATTEMPTS: u32 = 1024;
277
278pub const MAX_ADDRS: usize = 1024;
288
289pub const MAX_FAILOVER_BACKOFF_MAX_MS: u64 = 60 * 60 * 1_000;
295
296pub const DEFAULT_AUTH_TIMEOUT_MS: u64 = 15_000;
302
303pub const MAX_AUTH_TIMEOUT_MS: u64 = 60 * 60 * 1_000;
307
308pub const DEFAULT_SERVER_INFO_TIMEOUT_MS: u64 = 5_000;
319
320pub const MAX_SERVER_INFO_TIMEOUT_MS: u64 = 60 * 60 * 1_000;
323
324pub const DEFAULT_FAILOVER_MAX_DURATION_MS: u64 = 30_000;
327
328pub const MAX_FAILOVER_MAX_DURATION_MS: u64 = 60 * 60 * 1_000;
332
333pub const MAX_CONNECT_TIMEOUT_MS: u64 = 60 * 60 * 1_000;
338
339pub const DEFAULT_COMPRESSION_LEVEL: u8 = 1;
356
357pub const MIN_COMPRESSION_LEVEL: u8 = 1;
361
362pub const MAX_COMPRESSION_LEVEL: u8 = 22;
367
368#[derive(Debug, Clone, Copy, PartialEq, Eq)]
370#[non_exhaustive]
371pub enum TlsVerify {
372 On,
373 UnsafeOff,
376}
377
378#[derive(Debug, Clone)]
409#[non_exhaustive]
410pub struct ReaderConfig {
411 pub(crate) addrs: Vec<Endpoint>,
419 pub tls: bool,
420 pub path: String,
421 pub max_version: u8,
422 pub compression: Compression,
423 pub compression_level: u8,
429 pub max_batch_rows: u64,
430 pub client_id: Option<String>,
431 pub target: Target,
432 pub failover: bool,
445 pub failover_max_attempts: u32,
450 pub failover_backoff_initial_ms: u64,
454 pub failover_backoff_max_ms: u64,
457 pub failover_max_duration_ms: u64,
465 pub auth_timeout_ms: u64,
472 pub server_info_timeout_ms: u64,
481 pub connect_timeout_ms: u64,
492 pub zone: Option<String>,
501 pub auth: AuthMode,
502 pub tls_verify: TlsVerify,
503 pub tls_ca: CertificateAuthority,
504 pub tls_roots: Option<PathBuf>,
505 pub tls_roots_password: Option<String>,
515}
516
517pub(crate) const INGRESS_ONLY_CONFIG_KEYS: &[&str] = &[
528 "token_x",
530 "token_y",
531 "bind_interface",
533 "max_datagram_size",
535 "multicast_ttl",
536 "auto_flush",
538 "auto_flush_rows",
539 "auto_flush_bytes",
540 "auto_flush_interval",
541 "init_buf_size",
543 "max_buf_size",
544 "max_name_len",
545 "protocol_version",
546 "request_min_throughput",
548 "request_timeout",
549 "retry_timeout",
550 "retry_max_backoff_millis",
551 "auth_timeout",
554 "qwp_ws_progress",
556 "max_frame_rejections",
557 "poison_min_escalation_window_millis",
558 "sf_dir",
560 "sender_id",
561 "sf_max_segment_bytes",
562 "sf_max_total_bytes",
563 "sf_durability",
564 "sf_sync_interval_millis",
565 "sf_append_deadline_millis",
566 "reconnect_max_duration_millis",
568 "reconnect_initial_backoff_millis",
569 "reconnect_max_backoff_millis",
570 "initial_connect_retry",
571 "close_flush_timeout_millis",
573 "request_durable_ack",
574 "durable_ack_keepalive_interval_millis",
575 "drain_orphans",
576 "max_background_drainers",
577 "error_inbox_capacity",
578 "sender_pool_min",
584 "sender_pool_max",
585 "query_pool_min",
586 "query_pool_max",
587 "acquire_timeout_ms",
588 "idle_timeout_ms",
589 "lazy_connect",
590 "pool_reap",
591];
592
593impl ReaderConfig {
594 pub fn from_conf<T: AsRef<str>>(conf: T) -> Result<Self> {
596 let conf_str = conf.as_ref();
597
598 let addr_scan = crate::ingress::scan_qwp_ws_addr_params(conf_str)
609 .map_err(|e| fmt!(ConfigError, "{}", e.msg()))?;
610 let conf_to_parse = addr_scan
611 .as_ref()
612 .map(|s| s.sanitized_conf.as_str())
613 .unwrap_or(conf_str);
614
615 let conf = questdb_confstr::parse_conf_str(conf_to_parse)
616 .map_err(|e| fmt!(ConfigError, "Config parse error: {}", e))?;
617 let scheme = conf.service();
618 let tls = match scheme {
619 "ws" => false,
620 "wss" => true,
621 other => {
622 return Err(fmt!(
623 ConfigError,
624 "Unknown scheme \"{}\" — expected \"ws\" or \"wss\"",
625 other
626 ));
627 }
628 };
629 let params = conf.params();
630
631 let addr_values: Vec<&str> = match &addr_scan {
636 Some(s) if !s.addr_values.is_empty() => {
637 s.addr_values.iter().map(String::as_str).collect()
638 }
639 _ => {
640 let addr = params.get("addr").ok_or_else(|| {
641 fmt!(ConfigError, "Missing \"addr\" parameter in config string")
642 })?;
643 vec![addr.as_str()]
644 }
645 };
646
647 let default_port = if tls {
648 DEFAULT_TLS_PORT
649 } else {
650 DEFAULT_PLAIN_PORT
651 };
652 let mut addrs: Vec<Endpoint> = Vec::new();
653 let mut i: usize = 0;
654 for addr in addr_values {
655 for entry in addr.split(',').map(str::trim) {
656 if entry.is_empty() {
657 return Err(fmt!(ConfigError, "Empty entry {} in \"addr\" list", i));
658 }
659 let (host, port_str) = if let Some(rest) = entry.strip_prefix('[') {
667 let close = rest.find(']').ok_or_else(|| {
668 fmt!(
669 ConfigError,
670 "Bracketed addr entry {} missing closing ']': {:?}",
671 i,
672 entry
673 )
674 })?;
675 let host = rest[..close].to_string();
676 let after = &rest[close + 1..];
677 let port_str = if after.is_empty() {
678 default_port.to_string()
679 } else if let Some(p) = after.strip_prefix(':') {
680 p.to_string()
681 } else {
682 return Err(fmt!(
683 ConfigError,
684 "Unexpected characters after ']' in addr entry {}: {:?}",
685 i,
686 entry
687 ));
688 };
689 (host, port_str)
690 } else {
691 if entry.bytes().filter(|&b| b == b':').count() > 1 {
697 return Err(fmt!(
698 ConfigError,
699 "addr entry {} contains multiple ':' — IPv6 literals \
700 must be bracketed (e.g. [::1]:9000): {:?}",
701 i,
702 entry
703 ));
704 }
705 match entry.rsplit_once(':') {
706 Some((h, p)) => (h.to_string(), p.to_string()),
707 None => (entry.to_string(), default_port.to_string()),
708 }
709 };
710 if host.is_empty() {
711 return Err(fmt!(
712 ConfigError,
713 "Empty host in \"addr\" entry {}: {:?}",
714 i,
715 entry
716 ));
717 }
718 let port: u16 = port_str.parse().map_err(|_| {
719 fmt!(
720 ConfigError,
721 "Invalid port in \"addr\" entry {}: {:?}",
722 i,
723 entry
724 )
725 })?;
726 if port == 0 {
734 return Err(fmt!(
735 ConfigError,
736 "Port 0 is not a valid connect target in \"addr\" entry {}: {:?}",
737 i,
738 entry
739 ));
740 }
741 addrs.push(Endpoint { host, port });
742 i += 1;
743 }
744 }
745 if addrs.is_empty() {
746 return Err(fmt!(ConfigError, "\"addr\" parameter is empty"));
747 }
748 if addrs.len() > MAX_ADDRS {
749 return Err(fmt!(
750 ConfigError,
751 "\"addr\" list length {} exceeds the hard cap of {}",
752 addrs.len(),
753 MAX_ADDRS
754 ));
755 }
756
757 let mut path: String = DEFAULT_PATH.to_string();
759 let mut max_version: u8 = HIGHEST_KNOWN_VERSION;
760 let mut compression = Compression::Raw;
761 let mut compression_level: u8 = DEFAULT_COMPRESSION_LEVEL;
762 let mut max_batch_rows: u64 = 0;
763 let mut client_id: Option<String> = None;
764 let mut target = Target::Any;
765 let mut failover = DEFAULT_FAILOVER_ENABLED;
766 let mut failover_max_attempts: u32 = DEFAULT_FAILOVER_MAX_ATTEMPTS;
767 let mut failover_backoff_initial_ms: u64 = DEFAULT_FAILOVER_BACKOFF_INITIAL_MS;
768 let mut failover_backoff_max_ms: u64 = DEFAULT_FAILOVER_BACKOFF_MAX_MS;
769 let mut failover_max_duration_ms: u64 = DEFAULT_FAILOVER_MAX_DURATION_MS;
770 let mut auth_timeout_ms: u64 = DEFAULT_AUTH_TIMEOUT_MS;
771 let server_info_timeout_ms: u64 = DEFAULT_SERVER_INFO_TIMEOUT_MS;
772 let mut connect_timeout_ms: u64 = 0;
773 let mut zone: Option<String> = None;
774 let mut tls_verify = TlsVerify::On;
775 let mut tls_ca = default_tls_ca();
776 let mut tls_ca_explicit = false;
777 let mut tls_roots: Option<PathBuf> = None;
778 let mut tls_roots_password: Option<String> = None;
779
780 let mut username: Option<String> = None;
781 let mut password: Option<String> = None;
782 let mut token: Option<String> = None;
783 let mut auth_verbatim: Option<String> = None;
784
785 for (key, val) in params.iter() {
786 let key = key.as_str();
787 let val = val.as_str();
788 match key {
789 "addr" => {} "path" => {
791 if !val.starts_with('/') {
792 return Err(fmt!(
793 ConfigError,
794 "\"path\" must start with '/' (got {:?})",
795 val
796 ));
797 }
798 path = val.to_string();
799 }
800 "max_version" => {
801 let v: u8 = parse_value("max_version", val)?;
802 if !(1..=HIGHEST_KNOWN_VERSION).contains(&v) {
803 return Err(fmt!(
804 ConfigError,
805 "\"max_version\" must be in 1..={} (got {})",
806 HIGHEST_KNOWN_VERSION,
807 v
808 ));
809 }
810 max_version = v;
811 }
812 "compression" => {
813 compression = match val {
814 "raw" => Compression::Raw,
815 "zstd" => Compression::Zstd,
816 "auto" => Compression::Auto,
817 other => {
818 return Err(fmt!(
819 ConfigError,
820 "\"compression\" must be one of raw|zstd|auto (got {:?})",
821 other
822 ));
823 }
824 };
825 }
826 "compression_level" => {
827 let v: u8 = parse_value("compression_level", val)?;
828 if !(MIN_COMPRESSION_LEVEL..=MAX_COMPRESSION_LEVEL).contains(&v) {
829 return Err(fmt!(
830 ConfigError,
831 "\"compression_level\" must be in {}..={} (got {})",
832 MIN_COMPRESSION_LEVEL,
833 MAX_COMPRESSION_LEVEL,
834 v
835 ));
836 }
837 compression_level = v;
838 }
839 "max_batch_rows" => {
840 max_batch_rows = parse_value("max_batch_rows", val)?;
841 }
842 "client_id" => {
843 reject_crlf("client_id", val)?;
844 client_id = Some(val.to_string());
845 }
846 "target" => {
847 target = match val {
848 "any" => Target::Any,
849 "primary" => Target::Primary,
850 "replica" => Target::Replica,
851 other => {
852 return Err(fmt!(
853 ConfigError,
854 "\"target\" must be one of any|primary|replica (got {:?})",
855 other
856 ));
857 }
858 };
859 }
860 "username" => username = Some(val.to_string()),
861 "password" => password = Some(val.to_string()),
862 "token" => token = Some(val.to_string()),
863 "auth" => auth_verbatim = Some(val.to_string()),
864 "tls_verify" => {
865 tls_verify = match val {
866 "on" => TlsVerify::On,
867 "unsafe_off" => TlsVerify::UnsafeOff,
868 other => {
869 return Err(fmt!(
870 ConfigError,
871 "\"tls_verify\" must be \"on\" or \"unsafe_off\" (got {:?})",
872 other
873 ));
874 }
875 };
876 }
877 "tls_ca" => {
878 tls_ca = parse_tls_ca(val)?;
879 tls_ca_explicit = true;
880 }
881 "tls_roots" => {
882 let path = PathBuf::from_str(val).map_err(|e| {
883 fmt!(
884 ConfigError,
885 "Invalid path for \"tls_roots\" ({:?}): {}",
886 val,
887 e
888 )
889 })?;
890 tls_roots = Some(path);
891 }
892 "tls_roots_password" => {
893 tls_roots_password = Some(val.to_string());
894 }
895
896 "failover" => {
897 failover = parse_bool("failover", val)?;
898 }
899 "failover_max_attempts" => {
900 failover_max_attempts = parse_value("failover_max_attempts", val)?;
901 }
902 "failover_backoff_initial_ms" => {
903 failover_backoff_initial_ms = parse_value("failover_backoff_initial_ms", val)?;
904 }
905 "failover_backoff_max_ms" => {
906 failover_backoff_max_ms = parse_value("failover_backoff_max_ms", val)?;
907 }
908 "failover_max_duration_ms" => {
909 failover_max_duration_ms = parse_value("failover_max_duration_ms", val)?;
910 }
911 "auth_timeout_ms" => {
912 auth_timeout_ms = parse_value("auth_timeout_ms", val)?;
913 }
914 "connect_timeout" => {
915 connect_timeout_ms = parse_value("connect_timeout", val)?;
916 }
917 "zone" => {
918 reject_crlf("zone", val)?;
923 let trimmed = val.trim();
924 zone = if trimmed.is_empty() {
925 None
926 } else {
927 Some(trimmed.to_string())
928 };
929 }
930
931 "on_server_error" | "on_schema_error" | "on_parse_error" | "on_internal_error"
938 | "on_security_error" | "on_write_error" => {}
939
940 "buffer_pool_size" => {}
948
949 other if INGRESS_ONLY_CONFIG_KEYS.contains(&other) => {}
956
957 other => {
958 return Err(fmt!(ConfigError, "Unknown config key \"{}\"", other));
959 }
960 }
961 }
962
963 #[cfg(not(feature = "sync-reader-zstd"))]
965 {
966 if !matches!(compression, Compression::Raw) {
967 let user_token = match compression {
968 Compression::Raw => "raw",
969 Compression::Zstd => "zstd",
970 Compression::Auto => "auto",
971 };
972 return Err(fmt!(
973 ConfigError,
974 "\"compression={}\" requires the `sync-reader-zstd` crate feature; \
975 either enable it or use \"raw\"",
976 user_token
977 ));
978 }
979 }
980
981 if !tls && (tls_roots.is_some() || tls_ca_explicit || tls_roots_password.is_some()) {
990 return Err(fmt!(
991 ConfigError,
992 "TLS-related keys require the \"wss\" scheme"
993 ));
994 }
995
996 if tls_roots_password.is_some() && tls_roots.is_none() {
1002 return Err(fmt!(
1003 ConfigError,
1004 "\"tls_roots_password\" requires \"tls_roots\" \
1005 (the password unlocks the keystore at that path)"
1006 ));
1007 }
1008
1009 if tls_roots.is_some() {
1013 if tls_ca_explicit && tls_ca != CertificateAuthority::PemFile {
1014 return Err(fmt!(
1015 ConfigError,
1016 "\"tls_roots\" requires \"tls_ca=pem_file\" (or omit \"tls_ca\")"
1017 ));
1018 }
1019 tls_ca = CertificateAuthority::PemFile;
1020 }
1021
1022 let auth = AuthMode::from_parts(
1023 username.as_deref(),
1024 password.as_deref(),
1025 token.as_deref(),
1026 auth_verbatim.as_deref(),
1027 )?;
1028
1029 let cfg = ReaderConfig {
1030 addrs,
1031 tls,
1032 path,
1033 max_version,
1034 compression,
1035 compression_level,
1036 max_batch_rows,
1037 client_id,
1038 target,
1039 failover,
1040 failover_max_attempts,
1041 failover_backoff_initial_ms,
1042 failover_backoff_max_ms,
1043 failover_max_duration_ms,
1044 auth_timeout_ms,
1045 server_info_timeout_ms,
1046 connect_timeout_ms,
1047 zone,
1048 auth,
1049 tls_verify,
1050 tls_ca,
1051 tls_roots,
1052 tls_roots_password,
1053 };
1054 cfg.validate()?;
1055 Ok(cfg)
1056 }
1057
1058 pub fn validate(&self) -> Result<()> {
1068 if self.addrs.is_empty() {
1069 return Err(fmt!(ConfigError, "\"addr\" parameter is empty"));
1070 }
1071 if self.addrs.len() > MAX_ADDRS {
1072 return Err(fmt!(
1073 ConfigError,
1074 "\"addr\" list length {} exceeds the hard cap of {}",
1075 self.addrs.len(),
1076 MAX_ADDRS
1077 ));
1078 }
1079 if !(1..=HIGHEST_KNOWN_VERSION).contains(&self.max_version) {
1080 return Err(fmt!(
1081 ConfigError,
1082 "\"max_version\" must be in 1..={} (got {})",
1083 HIGHEST_KNOWN_VERSION,
1084 self.max_version
1085 ));
1086 }
1087 if !(MIN_COMPRESSION_LEVEL..=MAX_COMPRESSION_LEVEL).contains(&self.compression_level) {
1088 return Err(fmt!(
1089 ConfigError,
1090 "\"compression_level\" must be in {}..={} (got {})",
1091 MIN_COMPRESSION_LEVEL,
1092 MAX_COMPRESSION_LEVEL,
1093 self.compression_level
1094 ));
1095 }
1096 if self.failover_max_attempts == 0 {
1097 return Err(fmt!(
1098 ConfigError,
1099 "\"failover_max_attempts\" must be >= 1 (use \"failover=off\" to disable failover entirely)"
1100 ));
1101 }
1102 if self.failover_max_attempts > MAX_FAILOVER_MAX_ATTEMPTS {
1103 return Err(fmt!(
1104 ConfigError,
1105 "\"failover_max_attempts\" {} exceeds the hard cap of {}",
1106 self.failover_max_attempts,
1107 MAX_FAILOVER_MAX_ATTEMPTS
1108 ));
1109 }
1110 if self.failover_backoff_max_ms < self.failover_backoff_initial_ms {
1111 return Err(fmt!(
1112 ConfigError,
1113 "\"failover_backoff_max_ms\" ({}) must be >= \"failover_backoff_initial_ms\" ({})",
1114 self.failover_backoff_max_ms,
1115 self.failover_backoff_initial_ms
1116 ));
1117 }
1118 if self.failover_backoff_max_ms > MAX_FAILOVER_BACKOFF_MAX_MS {
1119 return Err(fmt!(
1120 ConfigError,
1121 "\"failover_backoff_max_ms\" {} exceeds the hard cap of {} (1 hour)",
1122 self.failover_backoff_max_ms,
1123 MAX_FAILOVER_BACKOFF_MAX_MS
1124 ));
1125 }
1126 if self.failover_max_duration_ms > MAX_FAILOVER_MAX_DURATION_MS {
1131 return Err(fmt!(
1132 ConfigError,
1133 "\"failover_max_duration_ms\" {} exceeds the hard cap of {} (1 hour)",
1134 self.failover_max_duration_ms,
1135 MAX_FAILOVER_MAX_DURATION_MS
1136 ));
1137 }
1138 if self.auth_timeout_ms == 0 {
1139 return Err(fmt!(
1140 ConfigError,
1141 "\"auth_timeout_ms\" must be > 0 (no sentinel for \"unbounded\" — \
1142 set a value high enough for your slowest peer's upgrade response)"
1143 ));
1144 }
1145 if self.auth_timeout_ms > MAX_AUTH_TIMEOUT_MS {
1146 return Err(fmt!(
1147 ConfigError,
1148 "\"auth_timeout_ms\" {} exceeds the hard cap of {} (1 hour)",
1149 self.auth_timeout_ms,
1150 MAX_AUTH_TIMEOUT_MS
1151 ));
1152 }
1153 if self.server_info_timeout_ms == 0 {
1154 return Err(fmt!(ConfigError, "\"server_info_timeout_ms\" must be > 0"));
1155 }
1156 if self.server_info_timeout_ms > MAX_SERVER_INFO_TIMEOUT_MS {
1157 return Err(fmt!(
1158 ConfigError,
1159 "\"server_info_timeout_ms\" {} exceeds the hard cap of {} (1 hour)",
1160 self.server_info_timeout_ms,
1161 MAX_SERVER_INFO_TIMEOUT_MS
1162 ));
1163 }
1164 if self.connect_timeout_ms > MAX_CONNECT_TIMEOUT_MS {
1168 return Err(fmt!(
1169 ConfigError,
1170 "\"connect_timeout\" {} exceeds the hard cap of {} (1 hour)",
1171 self.connect_timeout_ms,
1172 MAX_CONNECT_TIMEOUT_MS
1173 ));
1174 }
1175 if let Some(id) = &self.client_id {
1183 reject_crlf("client_id", id)?;
1184 }
1185 if let Some(z) = &self.zone {
1186 reject_crlf("zone", z)?;
1187 }
1188 self.auth.validate()?;
1189 #[cfg(not(feature = "insecure-skip-verify"))]
1196 if matches!(self.tls_verify, TlsVerify::UnsafeOff) {
1197 return Err(fmt!(
1198 ConfigError,
1199 "\"tls_verify=unsafe_off\" requires the \"insecure-skip-verify\" crate feature"
1200 ));
1201 }
1202 Ok(())
1203 }
1204
1205 pub(crate) fn failover_reconnect_rounds(&self) -> u32 {
1206 self.failover_max_attempts.saturating_sub(1)
1207 }
1208
1209 pub fn addrs(&self) -> &[Endpoint] {
1213 &self.addrs
1214 }
1215
1216 pub fn url_for(&self, idx: usize) -> String {
1219 let ep = &self.addrs[idx];
1220 let scheme = if self.tls { "wss" } else { "ws" };
1221 format!("{}://{}{}", scheme, ep, self.path)
1224 }
1225
1226 pub fn url(&self) -> String {
1228 self.url_for(0)
1229 }
1230
1231 pub fn upgrade_headers(&self) -> Vec<(&'static str, String)> {
1235 let mut headers = Vec::with_capacity(8);
1236 headers.push(("X-QWP-Max-Version", self.max_version.to_string()));
1237 if let Some(id) = &self.client_id {
1238 headers.push(("X-QWP-Client-Id", id.clone()));
1239 }
1240 headers.push((
1245 "X-QWP-Accept-Encoding",
1246 self.compression.accept_encoding(self.compression_level),
1247 ));
1248 if self.max_batch_rows > 0 {
1249 headers.push(("X-QWP-Max-Batch-Rows", self.max_batch_rows.to_string()));
1250 }
1251 if let Some(v) = self.auth.header_value() {
1252 headers.push(("Authorization", v));
1253 }
1254 headers
1255 }
1256}
1257
1258fn default_tls_ca() -> CertificateAuthority {
1263 #[cfg(feature = "tls-webpki-certs")]
1264 {
1265 CertificateAuthority::WebpkiRoots
1266 }
1267 #[cfg(all(not(feature = "tls-webpki-certs"), feature = "tls-native-certs"))]
1268 {
1269 CertificateAuthority::OsRoots
1270 }
1271 #[cfg(not(any(feature = "tls-webpki-certs", feature = "tls-native-certs")))]
1272 {
1273 CertificateAuthority::PemFile
1274 }
1275}
1276
1277fn parse_tls_ca(val: &str) -> Result<CertificateAuthority> {
1278 Ok(match val {
1279 #[cfg(feature = "tls-webpki-certs")]
1280 "webpki_roots" => CertificateAuthority::WebpkiRoots,
1281 #[cfg(not(feature = "tls-webpki-certs"))]
1282 "webpki_roots" => {
1283 return Err(fmt!(
1284 ConfigError,
1285 "\"tls_ca=webpki_roots\" requires the \"tls-webpki-certs\" feature"
1286 ));
1287 }
1288 #[cfg(feature = "tls-native-certs")]
1289 "os_roots" => CertificateAuthority::OsRoots,
1290 #[cfg(not(feature = "tls-native-certs"))]
1291 "os_roots" => {
1292 return Err(fmt!(
1293 ConfigError,
1294 "\"tls_ca=os_roots\" requires the \"tls-native-certs\" feature"
1295 ));
1296 }
1297 #[cfg(all(feature = "tls-webpki-certs", feature = "tls-native-certs"))]
1298 "webpki_and_os_roots" => CertificateAuthority::WebpkiAndOsRoots,
1299 #[cfg(not(all(feature = "tls-webpki-certs", feature = "tls-native-certs")))]
1300 "webpki_and_os_roots" => {
1301 return Err(fmt!(
1302 ConfigError,
1303 "\"tls_ca=webpki_and_os_roots\" requires both the \"tls-webpki-certs\" and \"tls-native-certs\" features"
1304 ));
1305 }
1306 "pem_file" => CertificateAuthority::PemFile,
1307 other => {
1308 return Err(fmt!(
1309 ConfigError,
1310 "\"tls_ca\" must be one of webpki_roots|os_roots|webpki_and_os_roots|pem_file (got {:?})",
1311 other
1312 ));
1313 }
1314 })
1315}
1316
1317fn parse_value<T>(name: &str, raw: &str) -> Result<T>
1318where
1319 T: FromStr,
1320{
1321 raw.parse::<T>()
1322 .map_err(|_| fmt!(ConfigError, "Could not parse \"{}\" value: {:?}", name, raw))
1323}
1324
1325fn parse_bool(name: &str, raw: &str) -> Result<bool> {
1326 match raw {
1327 "true" | "on" | "yes" | "1" => Ok(true),
1328 "false" | "off" | "no" | "0" => Ok(false),
1329 _ => Err(fmt!(
1330 ConfigError,
1331 "\"{}\" must be a boolean (got {:?})",
1332 name,
1333 raw
1334 )),
1335 }
1336}
1337
1338fn reject_crlf(name: &str, val: &str) -> Result<()> {
1346 if val.contains('\n') || val.contains('\r') {
1347 return Err(fmt!(ConfigError, "\"{}\" must not contain CR or LF", name));
1348 }
1349 Ok(())
1350}
1351
1352#[cfg(test)]
1353mod tests {
1354 use super::*;
1355 use crate::error::ErrorCode;
1356
1357 #[test]
1358 fn minimal_plain_conf() {
1359 let c = ReaderConfig::from_conf("ws::addr=localhost:9000").unwrap();
1360 assert_eq!(c.addrs.len(), 1);
1361 assert_eq!(c.addrs[0], Endpoint::new("localhost", 9000));
1362 assert!(!c.tls);
1363 assert_eq!(c.path, DEFAULT_PATH);
1364 assert_eq!(c.max_version, HIGHEST_KNOWN_VERSION);
1365 assert_eq!(c.compression, Compression::Raw);
1366 assert_eq!(c.url(), "ws://localhost:9000/read/v1");
1367 }
1368
1369 #[test]
1370 fn tls_scheme_changes_url() {
1371 let c = ReaderConfig::from_conf("wss::addr=h:8443").unwrap();
1372 assert!(c.tls);
1373 assert_eq!(c.url(), "wss://h:8443/read/v1");
1374 }
1375
1376 #[test]
1377 fn ws_scheme_is_plain() {
1378 let c = ReaderConfig::from_conf("ws::addr=localhost:9000").unwrap();
1379 assert!(!c.tls);
1380 assert_eq!(c.url(), "ws://localhost:9000/read/v1");
1381 }
1382
1383 #[test]
1384 fn wss_scheme_is_tls() {
1385 let c = ReaderConfig::from_conf("wss::addr=h:8443").unwrap();
1386 assert!(c.tls);
1387 assert_eq!(c.url(), "wss://h:8443/read/v1");
1388 }
1389
1390 #[test]
1391 fn unknown_scheme_rejected() {
1392 let err = ReaderConfig::from_conf("http::addr=h:1").unwrap_err();
1393 assert_eq!(err.code(), ErrorCode::ConfigError);
1394 }
1395
1396 #[test]
1397 fn missing_addr_rejected() {
1398 let err = ReaderConfig::from_conf("ws::path=/read/v1").unwrap_err();
1399 assert_eq!(err.code(), ErrorCode::ConfigError);
1400 }
1401
1402 #[test]
1403 fn unknown_key_rejected() {
1404 let err = ReaderConfig::from_conf("ws::addr=h:1;mystery=x").unwrap_err();
1405 assert_eq!(err.code(), ErrorCode::ConfigError);
1406 }
1407
1408 #[test]
1409 fn basic_auth_in_conf() {
1410 let c = ReaderConfig::from_conf("ws::addr=h:1;username=admin;password=quest").unwrap();
1411 assert_eq!(
1412 c.auth.header_value(),
1413 Some("Basic YWRtaW46cXVlc3Q=".to_string())
1414 );
1415 }
1416
1417 #[test]
1418 fn bearer_in_conf() {
1419 let c = ReaderConfig::from_conf("ws::addr=h:1;token=tok").unwrap();
1420 assert_eq!(c.auth.header_value(), Some("Bearer tok".to_string()));
1421 }
1422
1423 #[test]
1424 fn auth_modes_mutually_exclusive() {
1425 let err =
1426 ReaderConfig::from_conf("ws::addr=h:1;username=u;password=p;token=t").unwrap_err();
1427 assert_eq!(err.code(), ErrorCode::ConfigError);
1428 }
1429
1430 #[cfg(not(feature = "sync-reader-zstd"))]
1431 #[test]
1432 fn compression_zstd_rejected_without_feature() {
1433 let err = ReaderConfig::from_conf("ws::addr=h:1;compression=zstd").unwrap_err();
1434 assert_eq!(err.code(), ErrorCode::ConfigError);
1435 let err = ReaderConfig::from_conf("ws::addr=h:1;compression=auto").unwrap_err();
1436 assert_eq!(err.code(), ErrorCode::ConfigError);
1437 }
1438
1439 #[cfg(feature = "sync-reader-zstd")]
1440 #[test]
1441 fn compression_zstd_accepted_with_feature() {
1442 let c = ReaderConfig::from_conf("ws::addr=h:1;compression=zstd").unwrap();
1443 assert_eq!(c.compression, Compression::Zstd);
1444 let c = ReaderConfig::from_conf("ws::addr=h:1;compression=auto").unwrap();
1445 assert_eq!(c.compression, Compression::Auto);
1446 }
1447
1448 #[test]
1449 fn invalid_compression_value() {
1450 let err = ReaderConfig::from_conf("ws::addr=h:1;compression=xyz").unwrap_err();
1451 assert_eq!(err.code(), ErrorCode::ConfigError);
1452 }
1453
1454 #[test]
1455 fn compression_level_default_is_one() {
1456 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1457 assert_eq!(c.compression_level, DEFAULT_COMPRESSION_LEVEL);
1458 assert_eq!(c.compression_level, 1);
1459 }
1460
1461 #[cfg(feature = "sync-reader-zstd")]
1462 #[test]
1463 fn compression_level_parses_and_is_emitted() {
1464 let c =
1465 ReaderConfig::from_conf("ws::addr=h:1;compression=zstd;compression_level=9").unwrap();
1466 assert_eq!(c.compression_level, 9);
1467 let headers = c.upgrade_headers();
1468 let accept = headers
1469 .iter()
1470 .find(|(n, _)| *n == "X-QWP-Accept-Encoding")
1471 .expect("accept-encoding header present");
1472 assert_eq!(accept.1, "zstd;level=9");
1473 }
1474
1475 #[cfg(feature = "sync-reader-zstd")]
1476 #[test]
1477 fn compression_level_emitted_for_auto() {
1478 let c =
1479 ReaderConfig::from_conf("ws::addr=h:1;compression=auto;compression_level=7").unwrap();
1480 let headers = c.upgrade_headers();
1481 let accept = headers
1482 .iter()
1483 .find(|(n, _)| *n == "X-QWP-Accept-Encoding")
1484 .expect("accept-encoding header present");
1485 assert_eq!(accept.1, "zstd;level=7,raw");
1487 }
1488
1489 #[test]
1490 fn compression_level_ignored_for_raw() {
1491 let c = ReaderConfig::from_conf("ws::addr=h:1;compression_level=15").unwrap();
1495 let headers = c.upgrade_headers();
1496 let accept = headers
1497 .iter()
1498 .find(|(n, _)| *n == "X-QWP-Accept-Encoding")
1499 .expect("accept-encoding header present");
1500 assert_eq!(accept.1, "raw");
1501 }
1502
1503 #[test]
1504 fn compression_level_out_of_range_rejected() {
1505 for bad in ["0", "23", "100"] {
1506 let err = ReaderConfig::from_conf(format!("ws::addr=h:1;compression_level={}", bad))
1507 .unwrap_err();
1508 assert_eq!(
1509 err.code(),
1510 ErrorCode::ConfigError,
1511 "compression_level={} must be rejected",
1512 bad
1513 );
1514 }
1515 }
1516
1517 #[test]
1518 fn compression_level_accepts_full_range() {
1519 for ok in [
1520 MIN_COMPRESSION_LEVEL,
1521 DEFAULT_COMPRESSION_LEVEL,
1522 MAX_COMPRESSION_LEVEL,
1523 ] {
1524 let c = ReaderConfig::from_conf(format!("ws::addr=h:1;compression_level={}", ok))
1525 .expect("level in-range");
1526 assert_eq!(c.compression_level, ok);
1527 }
1528 }
1529
1530 #[test]
1531 fn target_parses() {
1532 let c = ReaderConfig::from_conf("ws::addr=h:1;target=primary").unwrap();
1533 assert_eq!(c.target, Target::Primary);
1534 }
1535
1536 #[test]
1537 fn multi_addr_parses() {
1538 let c = ReaderConfig::from_conf("ws::addr=h1:9000,h2:9001,h3,h4:9999;").unwrap();
1539 assert_eq!(c.addrs.len(), 4);
1540 assert_eq!(c.addrs[0], Endpoint::new("h1", 9000));
1541 assert_eq!(c.addrs[1], Endpoint::new("h2", 9001));
1542 assert_eq!(c.addrs[2], Endpoint::new("h3", 9000)); assert_eq!(c.addrs[3], Endpoint::new("h4", 9999));
1544 }
1545
1546 #[test]
1547 fn empty_addr_entry_rejected() {
1548 let err = ReaderConfig::from_conf("ws::addr=h1:9000,,h2:9001;").unwrap_err();
1549 assert_eq!(err.code(), ErrorCode::ConfigError);
1550 }
1551
1552 #[test]
1553 fn target_invalid_rejected() {
1554 let err = ReaderConfig::from_conf("ws::addr=h:1;target=leader").unwrap_err();
1555 assert_eq!(err.code(), ErrorCode::ConfigError);
1556 }
1557
1558 #[test]
1559 fn upgrade_headers_default() {
1560 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1561 let h = c.upgrade_headers();
1562 assert_eq!(h.len(), 2);
1564 assert_eq!(h[0], ("X-QWP-Max-Version", "1".to_string()));
1565 assert_eq!(h[1], ("X-QWP-Accept-Encoding", "raw".to_string()));
1566 }
1567
1568 #[test]
1569 fn upgrade_headers_full_set() {
1570 let c = ReaderConfig::from_conf(
1571 "ws::addr=h:1;client_id=app1;max_batch_rows=1000;username=u;password=p",
1572 )
1573 .unwrap();
1574 let h = c.upgrade_headers();
1575 let names: Vec<_> = h.iter().map(|(n, _)| *n).collect();
1576 assert!(names.contains(&"X-QWP-Max-Version"));
1577 assert!(names.contains(&"X-QWP-Client-Id"));
1578 assert!(names.contains(&"X-QWP-Accept-Encoding"));
1579 assert!(names.contains(&"X-QWP-Max-Batch-Rows"));
1580 assert!(names.contains(&"Authorization"));
1581 assert!(!names.contains(&"X-QWP-Request-Durable-Ack"));
1582
1583 let c = ReaderConfig::from_conf("ws::addr=h:1;max_batch_rows=0").unwrap();
1585 let h = c.upgrade_headers();
1586 assert!(h.iter().all(|(n, _)| *n != "X-QWP-Max-Batch-Rows"));
1587 }
1588
1589 #[test]
1590 fn path_must_start_with_slash() {
1591 let err = ReaderConfig::from_conf("ws::addr=h:1;path=read/v1").unwrap_err();
1592 assert_eq!(err.code(), ErrorCode::ConfigError);
1593 }
1594
1595 #[test]
1596 fn default_port_when_omitted() {
1597 let c = ReaderConfig::from_conf("ws::addr=localhost").unwrap();
1598 assert_eq!(c.addrs[0].port, 9000);
1599 }
1600
1601 #[test]
1602 fn invalid_port_rejected() {
1603 let err = ReaderConfig::from_conf("ws::addr=h:notaport").unwrap_err();
1604 assert_eq!(err.code(), ErrorCode::ConfigError);
1605 }
1606
1607 #[test]
1608 fn port_zero_rejected() {
1609 let err = ReaderConfig::from_conf("ws::addr=h:0").unwrap_err();
1614 assert_eq!(err.code(), ErrorCode::ConfigError);
1615 assert!(
1616 err.msg().contains("Port 0"),
1617 "diagnostic must name the offending value; got: {}",
1618 err.msg()
1619 );
1620 let err = ReaderConfig::from_conf("ws::addr=a:9000,b:0").unwrap_err();
1623 assert_eq!(err.code(), ErrorCode::ConfigError);
1624 let err = ReaderConfig::from_conf("ws::addr=[::1]:0").unwrap_err();
1627 assert_eq!(err.code(), ErrorCode::ConfigError);
1628 }
1629
1630 #[test]
1631 fn tls_keys_with_plain_scheme_rejected() {
1632 let err = ReaderConfig::from_conf("ws::addr=h:1;tls_roots=/tmp/x").unwrap_err();
1633 assert_eq!(err.code(), ErrorCode::ConfigError);
1634 let err = ReaderConfig::from_conf("ws::addr=h:1;tls_ca=pem_file").unwrap_err();
1635 assert_eq!(err.code(), ErrorCode::ConfigError);
1636 }
1637
1638 #[cfg(not(feature = "insecure-skip-verify"))]
1647 #[test]
1648 fn validate_rejects_unsafe_off_when_feature_disabled() {
1649 let err = ReaderConfig::from_conf("wss::addr=h:1;tls_verify=unsafe_off").unwrap_err();
1651 assert_eq!(err.code(), ErrorCode::ConfigError);
1652 assert!(
1653 err.msg().contains("insecure-skip-verify"),
1654 "msg: {}",
1655 err.msg()
1656 );
1657
1658 let mut cfg = ReaderConfig::from_conf("wss::addr=h:1").unwrap();
1662 assert_eq!(cfg.tls_verify, TlsVerify::On);
1663 cfg.tls_verify = TlsVerify::UnsafeOff;
1664 let err = cfg.validate().unwrap_err();
1665 assert_eq!(err.code(), ErrorCode::ConfigError);
1666 assert!(
1667 err.msg().contains("insecure-skip-verify"),
1668 "msg: {}",
1669 err.msg()
1670 );
1671 }
1672
1673 #[cfg(feature = "insecure-skip-verify")]
1674 #[test]
1675 fn validate_accepts_unsafe_off_when_feature_enabled() {
1676 let cfg = ReaderConfig::from_conf("wss::addr=h:1;tls_verify=unsafe_off").unwrap();
1679 assert_eq!(cfg.tls_verify, TlsVerify::UnsafeOff);
1680 cfg.validate()
1681 .expect("unsafe_off must pass validate when feature is on");
1682 }
1683
1684 #[test]
1685 fn tls_roots_password_without_tls_roots_rejected() {
1686 let err = ReaderConfig::from_conf("wss::addr=h:1;tls_roots_password=secret").unwrap_err();
1691 assert_eq!(err.code(), ErrorCode::ConfigError);
1692 assert!(
1693 err.msg().contains("tls_roots_password") && err.msg().contains("tls_roots"),
1694 "msg: {}",
1695 err.msg()
1696 );
1697 }
1698
1699 #[test]
1700 fn tls_roots_password_without_tls_scheme_rejected() {
1701 let err =
1705 ReaderConfig::from_conf("ws::addr=h:1;tls_roots=/tmp/r;tls_roots_password=secret")
1706 .unwrap_err();
1707 assert_eq!(err.code(), ErrorCode::ConfigError);
1708 }
1709
1710 #[test]
1711 fn tls_roots_password_with_tls_roots_accepted() {
1712 let c = ReaderConfig::from_conf(
1713 "wss::addr=h:1;tls_roots=/path/to/store.jks;tls_roots_password=secret",
1714 )
1715 .unwrap();
1716 assert_eq!(c.tls_ca, CertificateAuthority::PemFile);
1717 assert_eq!(c.tls_roots_password.as_deref(), Some("secret"));
1718 }
1719
1720 #[test]
1721 fn tls_roots_implies_pem_file_ca() {
1722 let c = ReaderConfig::from_conf("wss::addr=h:1;tls_roots=/path/to/roots.pem").unwrap();
1723 assert_eq!(c.tls_ca, CertificateAuthority::PemFile);
1724 assert_eq!(
1725 c.tls_roots.as_deref(),
1726 Some(std::path::Path::new("/path/to/roots.pem"))
1727 );
1728 }
1729
1730 #[test]
1731 fn tls_roots_with_conflicting_ca_rejected() {
1732 #[cfg(feature = "tls-webpki-certs")]
1733 {
1734 let err = ReaderConfig::from_conf("wss::addr=h:1;tls_ca=webpki_roots;tls_roots=/tmp/x")
1735 .unwrap_err();
1736 assert_eq!(err.code(), ErrorCode::ConfigError);
1737 }
1738 }
1739
1740 #[test]
1741 fn tls_ca_pem_file_explicit() {
1742 let c =
1743 ReaderConfig::from_conf("wss::addr=h:1;tls_ca=pem_file;tls_roots=/tmp/r.pem").unwrap();
1744 assert_eq!(c.tls_ca, CertificateAuthority::PemFile);
1745 }
1746
1747 #[test]
1748 fn tls_ca_invalid_value_rejected() {
1749 let err = ReaderConfig::from_conf("wss::addr=h:1;tls_ca=mystery").unwrap_err();
1750 assert_eq!(err.code(), ErrorCode::ConfigError);
1751 }
1752
1753 #[cfg(feature = "tls-webpki-certs")]
1754 #[test]
1755 fn tls_ca_webpki_roots_default() {
1756 let c = ReaderConfig::from_conf("wss::addr=h:1").unwrap();
1757 assert_eq!(c.tls_ca, CertificateAuthority::WebpkiRoots);
1758 assert_eq!(c.tls_roots, None);
1759 }
1760
1761 #[test]
1762 fn durable_ack_key_rejected() {
1763 let err = ReaderConfig::from_conf("ws::addr=h:1;durable_ack=true").unwrap_err();
1770 assert!(
1771 err.msg().to_lowercase().contains("durable_ack")
1772 || err.msg().to_lowercase().contains("unknown")
1773 );
1774 }
1775
1776 #[test]
1777 fn failover_defaults() {
1778 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1779 assert!(c.failover);
1780 assert_eq!(c.failover_max_attempts, DEFAULT_FAILOVER_MAX_ATTEMPTS);
1781 assert_eq!(
1782 c.failover_backoff_initial_ms,
1783 DEFAULT_FAILOVER_BACKOFF_INITIAL_MS
1784 );
1785 assert_eq!(c.failover_backoff_max_ms, DEFAULT_FAILOVER_BACKOFF_MAX_MS);
1786 }
1787
1788 #[test]
1789 fn failover_keys_parsed() {
1790 let c = ReaderConfig::from_conf(
1791 "ws::addr=h:1;failover=off;failover_max_attempts=3;failover_backoff_initial_ms=100;failover_backoff_max_ms=2000",
1792 )
1793 .unwrap();
1794 assert!(!c.failover);
1795 assert_eq!(c.failover_max_attempts, 3);
1796 assert_eq!(c.failover_backoff_initial_ms, 100);
1797 assert_eq!(c.failover_backoff_max_ms, 2000);
1798 }
1799
1800 #[test]
1801 fn failover_backoff_initial_zero_disables_sleep() {
1802 let c = ReaderConfig::from_conf("ws::addr=h:1;failover_backoff_initial_ms=0").unwrap();
1803 assert_eq!(c.failover_backoff_initial_ms, 0);
1804 assert_eq!(c.failover_backoff_max_ms, DEFAULT_FAILOVER_BACKOFF_MAX_MS);
1805 }
1806
1807 #[test]
1808 fn failover_backoff_max_below_initial_rejected() {
1809 let err = ReaderConfig::from_conf(
1810 "ws::addr=h:1;failover_backoff_initial_ms=500;failover_backoff_max_ms=100",
1811 )
1812 .unwrap_err();
1813 assert_eq!(err.code(), ErrorCode::ConfigError);
1814 }
1815
1816 #[test]
1817 fn failover_invalid_attempts_rejected() {
1818 let err = ReaderConfig::from_conf("ws::addr=h:1;failover_max_attempts=abc").unwrap_err();
1819 assert_eq!(err.code(), ErrorCode::ConfigError);
1820 }
1821
1822 #[test]
1823 fn failover_max_attempts_above_cap_rejected() {
1824 let conf = format!(
1825 "ws::addr=h:1;failover_max_attempts={}",
1826 MAX_FAILOVER_MAX_ATTEMPTS + 1
1827 );
1828 let err = ReaderConfig::from_conf(&conf).unwrap_err();
1829 assert_eq!(err.code(), ErrorCode::ConfigError);
1830 assert!(err.msg().contains("exceeds the hard cap"));
1831 }
1832
1833 #[test]
1834 fn failover_max_attempts_at_cap_accepted() {
1835 let conf = format!(
1836 "ws::addr=h:1;failover_max_attempts={}",
1837 MAX_FAILOVER_MAX_ATTEMPTS
1838 );
1839 let c = ReaderConfig::from_conf(&conf).unwrap();
1840 assert_eq!(c.failover_max_attempts, MAX_FAILOVER_MAX_ATTEMPTS);
1841 }
1842
1843 #[test]
1844 fn failover_backoff_max_above_cap_rejected() {
1845 let conf = format!(
1850 "ws::addr=h:1;failover_backoff_initial_ms=1;failover_backoff_max_ms={}",
1851 MAX_FAILOVER_BACKOFF_MAX_MS + 1
1852 );
1853 let err = ReaderConfig::from_conf(&conf).unwrap_err();
1854 assert_eq!(err.code(), ErrorCode::ConfigError);
1855 assert!(
1856 err.msg().contains("exceeds the hard cap"),
1857 "msg: {}",
1858 err.msg()
1859 );
1860 }
1861
1862 #[test]
1863 fn failover_backoff_max_at_cap_accepted() {
1864 let conf = format!(
1865 "ws::addr=h:1;failover_backoff_initial_ms=1;failover_backoff_max_ms={}",
1866 MAX_FAILOVER_BACKOFF_MAX_MS
1867 );
1868 let c = ReaderConfig::from_conf(&conf).unwrap();
1869 assert_eq!(c.failover_backoff_max_ms, MAX_FAILOVER_BACKOFF_MAX_MS);
1870 }
1871
1872 #[test]
1875 fn zone_unset_is_none_by_default() {
1876 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1877 assert_eq!(c.zone, None);
1878 }
1879
1880 #[test]
1881 fn zone_parses() {
1882 let c = ReaderConfig::from_conf("ws::addr=h:1;zone=eu-west-1a").unwrap();
1883 assert_eq!(c.zone.as_deref(), Some("eu-west-1a"));
1884 }
1885
1886 #[test]
1887 fn zone_empty_or_whitespace_normalises_to_none() {
1888 let c = ReaderConfig::from_conf("ws::addr=h:1;zone=").unwrap();
1889 assert_eq!(c.zone, None, "empty value collapses to unset");
1890 let c = ReaderConfig::from_conf("ws::addr=h:1;zone= ").unwrap();
1891 assert_eq!(c.zone, None, "whitespace-only collapses to unset");
1892 }
1893
1894 #[test]
1895 fn zone_trims_value() {
1896 let c = ReaderConfig::from_conf("ws::addr=h:1;zone= eu-west-1a ").unwrap();
1897 assert_eq!(c.zone.as_deref(), Some("eu-west-1a"));
1898 }
1899
1900 #[test]
1901 fn zone_rejects_cr_lf() {
1902 let err = ReaderConfig::from_conf("ws::addr=h:1;zone=eu\nwest").unwrap_err();
1905 assert_eq!(err.code(), ErrorCode::ConfigError);
1906 let err = ReaderConfig::from_conf("ws::addr=h:1;zone=eu\rwest").unwrap_err();
1907 assert_eq!(err.code(), ErrorCode::ConfigError);
1908 }
1909
1910 #[test]
1911 fn auth_timeout_defaults_to_15s() {
1912 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1913 assert_eq!(c.auth_timeout_ms, DEFAULT_AUTH_TIMEOUT_MS);
1914 assert_eq!(DEFAULT_AUTH_TIMEOUT_MS, 15_000);
1915 }
1916
1917 #[test]
1918 fn auth_timeout_parses() {
1919 let c = ReaderConfig::from_conf("ws::addr=h:1;auth_timeout_ms=3000").unwrap();
1920 assert_eq!(c.auth_timeout_ms, 3_000);
1921 }
1922
1923 #[test]
1924 fn auth_timeout_zero_rejected() {
1925 let err = ReaderConfig::from_conf("ws::addr=h:1;auth_timeout_ms=0").unwrap_err();
1929 assert_eq!(err.code(), ErrorCode::ConfigError);
1930 assert!(err.msg().contains("auth_timeout_ms"), "msg: {}", err.msg());
1931 }
1932
1933 #[test]
1934 fn auth_timeout_above_cap_rejected() {
1935 let conf = format!("ws::addr=h:1;auth_timeout_ms={}", MAX_AUTH_TIMEOUT_MS + 1);
1936 let err = ReaderConfig::from_conf(&conf).unwrap_err();
1937 assert_eq!(err.code(), ErrorCode::ConfigError);
1938 assert!(err.msg().contains("exceeds the hard cap"));
1939 }
1940
1941 #[test]
1942 fn auth_timeout_at_cap_accepted() {
1943 let conf = format!("ws::addr=h:1;auth_timeout_ms={}", MAX_AUTH_TIMEOUT_MS);
1944 let c = ReaderConfig::from_conf(&conf).unwrap();
1945 assert_eq!(c.auth_timeout_ms, MAX_AUTH_TIMEOUT_MS);
1946 }
1947
1948 #[test]
1949 fn connect_timeout_defaults_to_os_default() {
1950 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1951 assert_eq!(
1952 c.connect_timeout_ms, 0,
1953 "default is the OS-default dial (0)"
1954 );
1955 }
1956
1957 #[test]
1958 fn connect_timeout_parses_from_connect_string() {
1959 let c = ReaderConfig::from_conf("ws::addr=h:1;connect_timeout=250").unwrap();
1960 assert_eq!(c.connect_timeout_ms, 250);
1961 }
1962
1963 #[test]
1964 fn connect_timeout_zero_is_os_default() {
1965 let c = ReaderConfig::from_conf("ws::addr=h:1;connect_timeout=0").unwrap();
1968 assert_eq!(c.connect_timeout_ms, 0);
1969 }
1970
1971 #[test]
1972 fn connect_timeout_at_cap_accepted() {
1973 let conf = format!("ws::addr=h:1;connect_timeout={}", MAX_CONNECT_TIMEOUT_MS);
1974 let c = ReaderConfig::from_conf(&conf).unwrap();
1975 assert_eq!(c.connect_timeout_ms, MAX_CONNECT_TIMEOUT_MS);
1976 }
1977
1978 #[test]
1979 fn connect_timeout_above_cap_rejected() {
1980 let conf = format!(
1981 "ws::addr=h:1;connect_timeout={}",
1982 MAX_CONNECT_TIMEOUT_MS + 1
1983 );
1984 let err = ReaderConfig::from_conf(&conf).unwrap_err();
1985 assert_eq!(err.code(), ErrorCode::ConfigError);
1986 assert!(err.msg().contains("connect_timeout"), "msg: {}", err.msg());
1987 }
1988
1989 #[test]
1990 fn failover_max_duration_defaults_to_30s() {
1991 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
1992 assert_eq!(c.failover_max_duration_ms, DEFAULT_FAILOVER_MAX_DURATION_MS);
1993 assert_eq!(DEFAULT_FAILOVER_MAX_DURATION_MS, 30_000);
1994 }
1995
1996 #[test]
1997 fn failover_max_duration_parses() {
1998 let c = ReaderConfig::from_conf("ws::addr=h:1;failover_max_duration_ms=60000").unwrap();
1999 assert_eq!(c.failover_max_duration_ms, 60_000);
2000 }
2001
2002 #[test]
2003 fn failover_max_duration_zero_is_unbounded() {
2004 let c = ReaderConfig::from_conf("ws::addr=h:1;failover_max_duration_ms=0").unwrap();
2007 assert_eq!(c.failover_max_duration_ms, 0);
2008 }
2009
2010 #[test]
2011 fn failover_max_duration_above_cap_rejected() {
2012 let conf = format!(
2013 "ws::addr=h:1;failover_max_duration_ms={}",
2014 MAX_FAILOVER_MAX_DURATION_MS + 1
2015 );
2016 let err = ReaderConfig::from_conf(&conf).unwrap_err();
2017 assert_eq!(err.code(), ErrorCode::ConfigError);
2018 assert!(err.msg().contains("exceeds the hard cap"));
2019 }
2020
2021 #[test]
2024 fn server_info_timeout_defaults_to_5s() {
2025 let c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
2026 assert_eq!(c.server_info_timeout_ms, DEFAULT_SERVER_INFO_TIMEOUT_MS);
2027 assert_eq!(DEFAULT_SERVER_INFO_TIMEOUT_MS, 5_000);
2028 }
2029
2030 #[test]
2031 fn server_info_timeout_is_not_parsed_from_connect_string() {
2032 let err = ReaderConfig::from_conf("ws::addr=h:1;server_info_timeout_ms=1000").unwrap_err();
2037 assert_eq!(err.code(), ErrorCode::ConfigError);
2038 assert!(
2039 err.msg().contains("Unknown config key"),
2040 "msg: {}",
2041 err.msg()
2042 );
2043 }
2044
2045 const RESERVED_ON_ERROR_KEYS: &[&str] = &[
2055 "on_server_error",
2056 "on_schema_error",
2057 "on_parse_error",
2058 "on_internal_error",
2059 "on_security_error",
2060 "on_write_error",
2061 ];
2062
2063 #[test]
2064 fn reserved_on_error_policy_keys_all_together_are_accepted_silently() {
2065 let conf = "ws::addr=h:1\
2066 ;on_server_error=halt\
2067 ;on_schema_error=drop\
2068 ;on_parse_error=halt\
2069 ;on_internal_error=halt\
2070 ;on_security_error=halt\
2071 ;on_write_error=drop";
2072 let c = ReaderConfig::from_conf(conf).unwrap();
2073 assert_eq!(c.addrs.len(), 1);
2074 assert_eq!(c.addrs[0].host, "h");
2075 assert_eq!(c.addrs[0].port, 1);
2076 }
2077
2078 #[test]
2079 fn reserved_on_error_policy_keys_each_accepted_individually() {
2080 for key in RESERVED_ON_ERROR_KEYS {
2081 let conf = format!("ws::addr=h:1;{key}=halt");
2082 ReaderConfig::from_conf(&conf)
2083 .unwrap_or_else(|e| panic!("expected {key:?} to parse, got {}", e.msg()));
2084 }
2085 }
2086
2087 #[test]
2088 fn reserved_on_error_policy_keys_accept_any_value_without_validation() {
2089 for key in RESERVED_ON_ERROR_KEYS {
2096 for val in ["halt", "drop", "auto", "anything", ""] {
2097 let conf = format!("ws::addr=h:1;{key}={val}");
2098 ReaderConfig::from_conf(&conf)
2099 .unwrap_or_else(|e| panic!("expected {key}={val:?} to parse, got {}", e.msg()));
2100 }
2101 }
2102 }
2103
2104 #[test]
2105 fn reserved_on_error_policy_keys_do_not_swallow_other_settings() {
2106 let conf = "ws::addr=h:1;on_schema_error=drop;target=primary;zone=eu-1";
2110 let c = ReaderConfig::from_conf(conf).unwrap();
2111 assert_eq!(c.target, Target::Primary);
2112 assert_eq!(c.zone.as_deref(), Some("eu-1"));
2113 }
2114
2115 #[test]
2116 fn reserved_on_error_policy_keys_typo_still_rejected() {
2117 for typo in [
2120 "on_server_err",
2121 "on_schema_errors",
2122 "on_parse",
2123 "On_Write_Error",
2124 ] {
2125 let conf = format!("ws::addr=h:1;{typo}=halt");
2126 let err = ReaderConfig::from_conf(&conf)
2127 .err()
2128 .unwrap_or_else(|| panic!("expected {typo:?} to be rejected"));
2129 assert_eq!(err.code(), ErrorCode::ConfigError);
2130 assert!(
2131 err.msg().contains("Unknown config key"),
2132 "typo {typo:?}: msg: {}",
2133 err.msg()
2134 );
2135 }
2136 }
2137
2138 #[test]
2146 fn reserved_buffer_pool_size_accepts_any_value_without_validation() {
2147 for val in ["1", "4", "1024", "0", "-1", "not-a-number", ""] {
2152 let conf = format!("ws::addr=h:1;buffer_pool_size={val}");
2153 ReaderConfig::from_conf(&conf).unwrap_or_else(|e| {
2154 panic!(
2155 "expected buffer_pool_size={val:?} to parse, got {}",
2156 e.msg()
2157 )
2158 });
2159 }
2160 }
2161
2162 #[test]
2163 fn reserved_buffer_pool_size_does_not_swallow_other_settings() {
2164 let conf = "ws::addr=h:1;buffer_pool_size=8;target=replica;zone=us-2";
2165 let c = ReaderConfig::from_conf(conf).unwrap();
2166 assert_eq!(c.target, Target::Replica);
2167 assert_eq!(c.zone.as_deref(), Some("us-2"));
2168 }
2169
2170 #[test]
2173 fn egress_silently_accepts_every_ingress_only_key() {
2174 for key in INGRESS_ONLY_CONFIG_KEYS {
2179 for val in ["1", "off", "anything", ""] {
2180 let conf = format!("ws::addr=h:1;{key}={val}");
2181 ReaderConfig::from_conf(&conf).unwrap_or_else(|e| {
2182 panic!(
2183 "expected egress to silently accept ingress-only \
2184 key {key}={val:?}, got {}",
2185 e.msg()
2186 )
2187 });
2188 }
2189 }
2190 }
2191
2192 #[test]
2193 fn egress_rejects_removed_in_flight_keys() {
2194 for key in ["in_flight_window", "max_in_flight"] {
2200 let conf = format!("ws::addr=h:1;{key}=8");
2201 let err = ReaderConfig::from_conf(&conf).unwrap_err();
2202 assert_eq!(err.code(), ErrorCode::ConfigError, "key: {key}");
2203 assert!(
2204 err.msg().contains(&format!("Unknown config key \"{key}\"")),
2205 "key: {key}, msg: {}",
2206 err.msg()
2207 );
2208 }
2209 }
2210
2211 #[test]
2212 fn egress_accepts_full_ingress_connect_string_unchanged() {
2213 let conf = "ws::addr=h:9000\
2218 ;username=u;password=p\
2219 ;init_buf_size=65536;max_buf_size=1048576;max_name_len=127\
2220 ;auto_flush=off;auto_flush_rows=1000\
2221 ;protocol_version=2\
2222 ;tls_verify=on\
2223 ;target=primary;zone=eu-west-1a";
2224 let c = ReaderConfig::from_conf(conf).unwrap();
2225 assert_eq!(c.addrs.len(), 1);
2226 assert_eq!(c.addrs[0], Endpoint::new("h", 9000));
2227 assert_eq!(c.target, Target::Primary);
2228 assert_eq!(c.zone.as_deref(), Some("eu-west-1a"));
2229 assert!(matches!(c.auth, AuthMode::Basic { .. }));
2231 }
2232
2233 #[test]
2234 fn reserved_buffer_pool_size_typo_still_rejected() {
2235 for typo in [
2236 "buffer_pool",
2237 "buffer_pool_sizes",
2238 "Buffer_Pool_Size",
2239 "buffer_size",
2240 ] {
2241 let conf = format!("ws::addr=h:1;{typo}=4");
2242 let err = ReaderConfig::from_conf(&conf)
2243 .err()
2244 .unwrap_or_else(|| panic!("expected {typo:?} to be rejected"));
2245 assert_eq!(err.code(), ErrorCode::ConfigError);
2246 assert!(
2247 err.msg().contains("Unknown config key"),
2248 "typo {typo:?}: msg: {}",
2249 err.msg()
2250 );
2251 }
2252 }
2253
2254 #[test]
2255 fn server_info_timeout_zero_rejected_by_validate() {
2256 let mut c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
2259 c.server_info_timeout_ms = 0;
2260 let err = c.validate().unwrap_err();
2261 assert_eq!(err.code(), ErrorCode::ConfigError);
2262 assert!(err.msg().contains("server_info_timeout_ms"));
2263 }
2264
2265 #[test]
2266 fn server_info_timeout_above_cap_rejected_by_validate() {
2267 let mut c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
2268 c.server_info_timeout_ms = MAX_SERVER_INFO_TIMEOUT_MS + 1;
2269 let err = c.validate().unwrap_err();
2270 assert_eq!(err.code(), ErrorCode::ConfigError);
2271 assert!(err.msg().contains("exceeds the hard cap"));
2272 }
2273
2274 #[test]
2275 fn server_info_timeout_at_cap_accepted() {
2276 let mut c = ReaderConfig::from_conf("ws::addr=h:1").unwrap();
2277 c.server_info_timeout_ms = MAX_SERVER_INFO_TIMEOUT_MS;
2278 c.validate().unwrap();
2279 }
2280
2281 #[test]
2282 fn addrs_above_cap_rejected() {
2283 let mut addr = String::from("ws::addr=");
2288 for i in 0..(MAX_ADDRS + 1) {
2289 if i > 0 {
2290 addr.push(',');
2291 }
2292 addr.push_str(&format!("h{}:9000", i));
2293 }
2294 let err = ReaderConfig::from_conf(&addr).unwrap_err();
2295 assert_eq!(err.code(), ErrorCode::ConfigError);
2296 assert!(
2297 err.msg().contains("exceeds the hard cap"),
2298 "msg: {}",
2299 err.msg()
2300 );
2301 }
2302
2303 #[test]
2304 fn failover_max_attempts_zero_rejected() {
2305 let err = ReaderConfig::from_conf("ws::addr=h:1;failover_max_attempts=0").unwrap_err();
2308 assert_eq!(err.code(), ErrorCode::ConfigError);
2309 assert!(
2310 err.msg().contains("failover_max_attempts"),
2311 "msg: {}",
2312 err.msg()
2313 );
2314 }
2315
2316 #[test]
2317 fn failover_max_attempts_counts_initial_execute_attempt() {
2318 let c = ReaderConfig::from_conf("ws::addr=h:1;failover_max_attempts=1").unwrap();
2319 assert_eq!(c.failover_reconnect_rounds(), 0);
2320
2321 let c = ReaderConfig::from_conf("ws::addr=h:1;failover_max_attempts=8").unwrap();
2322 assert_eq!(c.failover_reconnect_rounds(), 7);
2323 }
2324
2325 #[test]
2326 fn endpoint_display_common_cases() {
2327 assert_eq!(
2332 Endpoint::new("localhost", 9000).to_string(),
2333 "localhost:9000"
2334 );
2335 assert_eq!(Endpoint::new("db-a", 9000).to_string(), "db-a:9000");
2336 assert_eq!(
2337 Endpoint::new("127.0.0.1", 9000).to_string(),
2338 "127.0.0.1:9000"
2339 );
2340 let ep = Endpoint::new("example.com", 1234);
2346 let conf = format!("ws::addr={}", ep);
2347 let parsed = ReaderConfig::from_conf(&conf).expect("parse round-trip");
2348 assert_eq!(parsed.addrs(), &[ep]);
2349 }
2350
2351 #[test]
2352 fn endpoint_display_ipv6_brackets() {
2353 assert_eq!(Endpoint::new("::1", 9000).to_string(), "[::1]:9000");
2357 assert_eq!(
2358 Endpoint::new("2001:db8::1", 443).to_string(),
2359 "[2001:db8::1]:443"
2360 );
2361 }
2362
2363 #[test]
2364 fn ipv6_addr_parses_with_explicit_port() {
2365 let c = ReaderConfig::from_conf("ws::addr=[::1]:9000").unwrap();
2366 assert_eq!(c.addrs.len(), 1);
2367 assert_eq!(c.addrs[0], Endpoint::new("::1", 9000));
2369 assert_eq!(c.url_for(0), "ws://[::1]:9000/read/v1");
2370 }
2371
2372 #[test]
2373 fn ipv6_addr_default_port() {
2374 let c = ReaderConfig::from_conf("ws::addr=[2001:db8::1]").unwrap();
2375 assert_eq!(c.addrs[0], Endpoint::new("2001:db8::1", 9000));
2376 assert_eq!(c.url_for(0), "ws://[2001:db8::1]:9000/read/v1");
2377 }
2378
2379 #[test]
2380 fn ipv6_addr_in_multi_addr_list() {
2381 let c = ReaderConfig::from_conf("ws::addr=[::1]:9000,h2:9001,[2001:db8::5]").unwrap();
2382 assert_eq!(c.addrs.len(), 3);
2383 assert_eq!(c.addrs[0], Endpoint::new("::1", 9000));
2384 assert_eq!(c.addrs[1], Endpoint::new("h2", 9001));
2385 assert_eq!(c.addrs[2], Endpoint::new("2001:db8::5", 9000));
2386 }
2387
2388 #[test]
2389 fn ipv6_addr_missing_close_bracket_rejected() {
2390 let err = ReaderConfig::from_conf("ws::addr=[::1:9000").unwrap_err();
2391 assert_eq!(err.code(), ErrorCode::ConfigError);
2392 }
2393
2394 #[test]
2395 fn ipv6_addr_garbage_after_bracket_rejected() {
2396 let err = ReaderConfig::from_conf("ws::addr=[::1]junk").unwrap_err();
2397 assert_eq!(err.code(), ErrorCode::ConfigError);
2398 }
2399
2400 #[test]
2401 fn unbracketed_ipv6_rejected() {
2402 for bad in [
2403 "ws::addr=::1",
2404 "ws::addr=::1:9000",
2405 "ws::addr=2001:db8::1",
2406 "ws::addr=fe80::1%eth0",
2407 "ws::addr=h1:9000,::1:9001",
2408 ] {
2409 let err = ReaderConfig::from_conf(bad).unwrap_err();
2410 assert_eq!(
2411 err.code(),
2412 ErrorCode::ConfigError,
2413 "expected reject for {bad:?}"
2414 );
2415 let msg = err.msg();
2416 assert!(
2417 msg.contains("multiple ':'") || msg.contains("bracketed"),
2418 "expected diagnostic to mention bracketing, got {msg:?}"
2419 );
2420 }
2421 }
2422
2423 #[test]
2424 fn single_colon_host_port_still_accepted() {
2425 let c = ReaderConfig::from_conf("ws::addr=h1:9000").unwrap();
2426 assert_eq!(c.addrs[0], Endpoint::new("h1", 9000));
2427 }
2428
2429 #[test]
2430 fn url_for_uses_endpoint_display() {
2431 let c = ReaderConfig::from_conf("ws::addr=db-a:9000;path=/exec").unwrap();
2435 assert_eq!(c.url_for(0), "ws://db-a:9000/exec");
2436 }
2437
2438 #[test]
2439 fn validate_accepts_parsed_default_config() {
2440 let c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2441 c.validate().expect("a freshly-parsed config must validate");
2442 }
2443
2444 #[test]
2445 fn validate_rejects_post_parse_backoff_overflow() {
2446 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2447 c.failover_backoff_max_ms = u64::MAX;
2448 let err = c.validate().unwrap_err();
2449 assert_eq!(err.code(), ErrorCode::ConfigError);
2450 assert!(
2451 err.msg().contains("failover_backoff_max_ms"),
2452 "got: {}",
2453 err.msg()
2454 );
2455 }
2456
2457 #[test]
2458 fn validate_rejects_post_parse_max_attempts_overflow() {
2459 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2460 c.failover_max_attempts = MAX_FAILOVER_MAX_ATTEMPTS + 1;
2461 let err = c.validate().unwrap_err();
2462 assert_eq!(err.code(), ErrorCode::ConfigError);
2463 }
2464
2465 #[test]
2466 fn validate_rejects_post_parse_max_attempts_zero() {
2467 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2468 c.failover_max_attempts = 0;
2469 let err = c.validate().unwrap_err();
2470 assert_eq!(err.code(), ErrorCode::ConfigError);
2471 }
2472
2473 #[test]
2474 fn validate_accepts_post_parse_backoff_zero_initial() {
2475 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2476 c.failover_backoff_initial_ms = 0;
2477 c.validate().unwrap();
2478 }
2479
2480 #[test]
2481 fn validate_rejects_post_parse_backoff_inversion() {
2482 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2483 c.failover_backoff_initial_ms = 1000;
2484 c.failover_backoff_max_ms = 50;
2485 let err = c.validate().unwrap_err();
2486 assert_eq!(err.code(), ErrorCode::ConfigError);
2487 }
2488
2489 #[test]
2490 fn validate_rejects_post_parse_max_version_out_of_range() {
2491 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2492 c.max_version = 0;
2493 let err = c.validate().unwrap_err();
2494 assert_eq!(err.code(), ErrorCode::ConfigError);
2495 c.max_version = HIGHEST_KNOWN_VERSION + 1;
2496 let err = c.validate().unwrap_err();
2497 assert_eq!(err.code(), ErrorCode::ConfigError);
2498 }
2499
2500 #[test]
2509 fn validate_rejects_post_parse_client_id_with_crlf() {
2510 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2513 c.client_id = Some("foo\r\nAuthorization: Bearer attacker".into());
2514 let err = c.validate().unwrap_err();
2515 assert_eq!(err.code(), ErrorCode::ConfigError);
2516 assert!(
2517 err.msg().contains("client_id"),
2518 "error message must name the offending field; got: {}",
2519 err.msg()
2520 );
2521 c.client_id = Some("foo\nbar".into());
2523 assert_eq!(c.validate().unwrap_err().code(), ErrorCode::ConfigError);
2524 c.client_id = Some("foo\rbar".into());
2525 assert_eq!(c.validate().unwrap_err().code(), ErrorCode::ConfigError);
2526 }
2527
2528 #[test]
2529 fn validate_rejects_post_parse_zone_with_crlf() {
2530 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2531 c.zone = Some("eu-west-1a\r\nX-Injected: 1".into());
2532 let err = c.validate().unwrap_err();
2533 assert_eq!(err.code(), ErrorCode::ConfigError);
2534 assert!(err.msg().contains("zone"));
2535 }
2536
2537 #[test]
2538 fn validate_rejects_post_parse_verbatim_auth_with_control_bytes() {
2539 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2545 c.auth = AuthMode::Verbatim {
2546 value: "Bearer xx\r\nX-Injected: 1".into(),
2547 };
2548 let err = c.validate().unwrap_err();
2549 assert_eq!(err.code(), ErrorCode::AuthError);
2550 c.auth = AuthMode::Verbatim {
2552 value: "Bearer\nyy".into(),
2553 };
2554 assert_eq!(c.validate().unwrap_err().code(), ErrorCode::AuthError);
2555 }
2556
2557 #[test]
2558 fn validate_rejects_post_parse_bearer_token_with_control_bytes() {
2559 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2560 c.auth = AuthMode::Bearer {
2561 token: "abc\r\ndef".into(),
2562 };
2563 let err = c.validate().unwrap_err();
2564 assert_eq!(err.code(), ErrorCode::AuthError);
2565 }
2566
2567 #[test]
2568 fn validate_rejects_post_parse_basic_auth_with_control_bytes() {
2569 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2570 c.auth = AuthMode::Basic {
2571 username: "user\nfoo".into(),
2572 password: "pw".into(),
2573 };
2574 assert_eq!(c.validate().unwrap_err().code(), ErrorCode::AuthError);
2575 c.auth = AuthMode::Basic {
2576 username: "user".into(),
2577 password: "pw\r\nX-Injected: 1".into(),
2578 };
2579 assert_eq!(c.validate().unwrap_err().code(), ErrorCode::AuthError);
2580 }
2581
2582 #[test]
2583 fn validate_rejects_post_parse_basic_username_with_colon() {
2584 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2589 c.auth = AuthMode::Basic {
2590 username: "admin:override".into(),
2591 password: "real".into(),
2592 };
2593 let err = c.validate().unwrap_err();
2594 assert_eq!(err.code(), ErrorCode::AuthError);
2595 }
2596
2597 #[test]
2598 fn validate_accepts_post_parse_clean_string_fields() {
2599 let mut c = ReaderConfig::from_conf("ws::addr=h:9000").unwrap();
2603 c.client_id = Some("benign-id".into());
2604 c.zone = Some("eu-west-1a".into());
2605 c.auth = AuthMode::Bearer {
2606 token: "benign.token.value".into(),
2607 };
2608 c.validate().expect("clean string fields must validate");
2609 }
2610
2611 #[test]
2612 fn addr_comma_list_collects_all_endpoints() {
2613 let c = ReaderConfig::from_conf("ws::addr=h1:9000,h2:9001,h3:9002").unwrap();
2614 assert_eq!(
2615 c.addrs,
2616 vec![
2617 Endpoint::new("h1", 9000),
2618 Endpoint::new("h2", 9001),
2619 Endpoint::new("h3", 9002),
2620 ]
2621 );
2622 }
2623
2624 #[test]
2625 fn addr_repeated_key_collects_all_endpoints() {
2626 let c = ReaderConfig::from_conf("ws::addr=h1:9000;addr=h2:9001;addr=h3:9002;").unwrap();
2629 assert_eq!(
2630 c.addrs,
2631 vec![
2632 Endpoint::new("h1", 9000),
2633 Endpoint::new("h2", 9001),
2634 Endpoint::new("h3", 9002),
2635 ]
2636 );
2637 }
2638
2639 #[test]
2640 fn addr_mixed_comma_and_repeated_key_collects_all_endpoints() {
2641 let c = ReaderConfig::from_conf("ws::addr=h1:9000,h2:9001;addr=h3:9002,h4:9003;").unwrap();
2644 assert_eq!(
2645 c.addrs,
2646 vec![
2647 Endpoint::new("h1", 9000),
2648 Endpoint::new("h2", 9001),
2649 Endpoint::new("h3", 9002),
2650 Endpoint::new("h4", 9003),
2651 ]
2652 );
2653 }
2654
2655 #[test]
2656 fn addr_repeated_key_rejects_empty_entry() {
2657 let err = ReaderConfig::from_conf("ws::addr=h1:9000;addr=,;addr=h2:9001;").unwrap_err();
2660 assert_eq!(err.code(), ErrorCode::ConfigError);
2661 assert!(
2662 err.msg().contains("Empty entry"),
2663 "unexpected msg: {}",
2664 err.msg()
2665 );
2666 }
2667
2668 #[test]
2669 fn addr_repeated_key_propagates_invalid_port() {
2670 let err = ReaderConfig::from_conf("ws::addr=h1:9000;addr=h2:notaport;").unwrap_err();
2673 assert_eq!(err.code(), ErrorCode::ConfigError);
2674 assert!(
2675 err.msg().contains("Invalid port in \"addr\" entry 1"),
2676 "unexpected msg: {}",
2677 err.msg()
2678 );
2679 }
2680}