use super::*;
use crate::ErrorCode;
#[cfg(any(feature = "sync-sender-tcp", feature = "sync-sender-qwp-ws"))]
use tempfile::TempDir;
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_simple() {
let builder = SenderBuilder::from_conf("http::addr=127.0.0.1;").unwrap();
assert_eq!(builder.protocol, Protocol::Http);
assert_specified_eq(&builder.host, "127.0.0.1");
assert_specified_eq(&builder.port, Protocol::Http.default_port());
assert!(!builder.protocol.tls_enabled());
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn https_simple() {
let builder = SenderBuilder::from_conf("https::addr=localhost;").unwrap();
assert_eq!(builder.protocol, Protocol::Https);
assert_specified_eq(&builder.host, "localhost");
assert_specified_eq(&builder.port, Protocol::Https.default_port());
assert!(builder.protocol.tls_enabled());
#[cfg(feature = "tls-webpki-certs")]
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::WebpkiRoots);
#[cfg(not(feature = "tls-webpki-certs"))]
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::OsRoots);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcp_simple() {
let builder = SenderBuilder::from_conf("tcp::addr=127.0.0.1;").unwrap();
assert_eq!(builder.protocol, Protocol::Tcp);
assert_specified_eq(&builder.port, Protocol::Tcp.default_port());
assert_specified_eq(&builder.host, "127.0.0.1");
assert!(!builder.protocol.tls_enabled());
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_simple() {
let builder = SenderBuilder::from_conf("tcps::addr=localhost;").unwrap();
assert_eq!(builder.protocol, Protocol::Tcps);
assert_specified_eq(&builder.host, "localhost");
assert_specified_eq(&builder.port, Protocol::Tcps.default_port());
assert!(builder.protocol.tls_enabled());
#[cfg(feature = "tls-webpki-certs")]
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::WebpkiRoots);
#[cfg(not(feature = "tls-webpki-certs"))]
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::OsRoots);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_simple() {
let builder = SenderBuilder::from_conf("udp::addr=127.0.0.1;").unwrap();
assert_eq!(builder.protocol, Protocol::Udp);
assert_specified_eq(&builder.host, "127.0.0.1");
assert_specified_eq(&builder.port, Protocol::Udp.default_port());
assert!(!builder.protocol.tls_enabled());
let qwp_udp = builder.qwp_udp.as_ref().unwrap();
assert_defaulted_eq(&qwp_udp.max_datagram_size, 1400usize);
assert_defaulted_eq(&qwp_udp.multicast_ttl, 1u32);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_custom_config() {
let builder = SenderBuilder::from_conf(
"udp::addr=239.1.2.3:19002;bind_interface=192.168.1.10;max_datagram_size=1200;multicast_ttl=7;",
)
.unwrap();
assert_eq!(builder.protocol, Protocol::Udp);
assert_specified_eq(&builder.host, "239.1.2.3");
assert_specified_eq(&builder.port, "19002");
assert_specified_eq(&builder.net_interface, Some("192.168.1.10".to_string()));
let qwp_udp = builder.qwp_udp.as_ref().unwrap();
assert_specified_eq(&qwp_udp.max_datagram_size, 1200usize);
assert_specified_eq(&qwp_udp.multicast_ttl, 7u32);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_sender_reports_transport_protocol() {
let sender = SenderBuilder::new(Protocol::Udp, "127.0.0.1", 9007)
.build()
.unwrap();
assert_eq!(sender.protocol(), Protocol::Udp);
}
#[cfg(all(feature = "sync-sender-qwp-ws", feature = "sync-sender-qwp-udp"))]
#[test]
fn qwpws_error_polling_rejects_non_websocket_sender() {
let mut sender = SenderBuilder::new(Protocol::Udp, "127.0.0.1", 9007)
.build()
.unwrap();
let err = sender.poll_qwp_ws_error().unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidApiCall);
assert!(
err.msg()
.contains("poll_qwp_ws_error is only supported for QWP/WebSocket")
);
let err = sender.qwp_ws_errors_dropped().unwrap_err();
assert_eq!(err.code(), ErrorCode::InvalidApiCall);
assert!(
err.msg()
.contains("qwp_ws_errors_dropped is only supported for QWP/WebSocket")
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_store_and_forward_config_parses_java_keys() {
let builder = SenderBuilder::from_conf(
"ws::addr=localhost:9000;\
sf_dir=/tmp/qdb-rust-sf;\
sender_id=primary-1;\
sf_max_segment_bytes=64mb;\
sf_max_total_bytes=4G;\
sf_durability=memory;\
sf_append_deadline_millis=1234;\
poison_min_escalation_window_millis=600;\
auth_timeout_ms=750;",
)
.unwrap();
assert_eq!(builder.protocol, Protocol::Ws);
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.sf_dir, Some(PathBuf::from("/tmp/qdb-rust-sf")));
assert_specified_eq(&qwp_ws.sender_id, "primary-1".to_owned());
assert_specified_eq(&qwp_ws.sf_max_segment_bytes, 64 * 1024 * 1024_u64);
assert_specified_eq(&qwp_ws.sf_max_total_bytes, Some(4 * 1024 * 1024 * 1024_u64));
assert_specified_eq(&qwp_ws.sf_durability, conf::SfDurability::Memory);
assert_specified_eq(&qwp_ws.sf_append_deadline, Duration::from_millis(1234));
assert_specified_eq(
&qwp_ws.poison_min_escalation_window,
Duration::from_millis(600),
);
assert_specified_eq(&qwp_ws.auth_timeout, Duration::from_millis(750));
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_addr_accepts_bracketed_ipv6() {
let builder = SenderBuilder::from_conf("ws::addr=[::1]:9001;").unwrap();
assert_specified_eq(&builder.host, "::1");
assert_specified_eq(&builder.port, "9001");
let defaulted = SenderBuilder::from_conf("ws::addr=[::1];").unwrap();
assert_specified_eq(&defaulted.host, "::1");
assert_specified_eq(&defaulted.port, Protocol::Ws.default_port());
let multi = SenderBuilder::from_conf("ws::addr=[::1]:9000,[2001:db8::1]:9001, localhost:9002;")
.unwrap();
let qwp_ws = multi.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.endpoints,
vec![
conf::QwpWsEndpoint::new("::1".into(), "9000".into()),
conf::QwpWsEndpoint::new("2001:db8::1".into(), "9001".into()),
conf::QwpWsEndpoint::new("localhost".into(), "9002".into()),
],
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_config_accepts_websocket_schemes() {
let plain = SenderBuilder::from_conf("ws::addr=localhost:9000;").unwrap();
assert_eq!(plain.protocol, Protocol::Ws);
let tls = SenderBuilder::from_conf("wss::addr=localhost:9000;").unwrap();
assert_eq!(tls.protocol, Protocol::Wss);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_store_and_forward_defaults_match_java() {
let builder = SenderBuilder::from_conf("ws::addr=localhost:9000;").unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_defaulted_eq(&qwp_ws.sender_id, "default".to_owned());
assert_defaulted_eq(&qwp_ws.sf_max_segment_bytes, 4 * 1024 * 1024_u64);
assert_defaulted_eq(&qwp_ws.sf_max_total_bytes, None);
assert_defaulted_eq(&qwp_ws.sf_durability, conf::SfDurability::Memory);
assert_defaulted_eq(&qwp_ws.sf_sync_interval, None);
assert_eq!(qwp_ws.periodic_sync_interval(), None);
assert_defaulted_eq(&qwp_ws.sf_append_deadline, Duration::from_secs(30));
assert_defaulted_eq(&qwp_ws.auth_timeout, Duration::from_secs(15));
assert_defaulted_eq(&qwp_ws.progress, QwpWsProgress::Background);
assert_defaulted_eq(
&qwp_ws.poison_min_escalation_window,
Duration::from_millis(5000),
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_poison_min_escalation_window_accepts_zero() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;poison_min_escalation_window_millis=0;")
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.poison_min_escalation_window,
Duration::from_millis(0),
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_progress_config_parses_manual_and_background() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;qwp_ws_progress=manual;").unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.progress, QwpWsProgress::Manual);
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;qwp_ws_progress=background;").unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.progress, QwpWsProgress::Background);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_config_rejects_removed_in_flight_keys() {
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;in_flight_window=7;"),
"Unknown config key \"in_flight_window\"",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;max_in_flight=3;"),
"Unknown config key \"max_in_flight\"",
);
}
#[cfg(feature = "sync-sender-http")]
const EGRESS_ONLY_CONFIG_KEYS: &[&str] = &[
"path",
"max_version",
"compression",
"compression_level",
"max_batch_rows",
"client_id",
"target",
"auth",
"failover",
"failover_max_attempts",
"failover_backoff_initial_ms",
"failover_backoff_max_ms",
"failover_max_duration_ms",
"buffer_pool_size",
"on_server_error",
"on_schema_error",
"on_parse_error",
"on_internal_error",
"on_security_error",
"on_write_error",
];
#[cfg(feature = "sync-sender-http")]
#[test]
fn ingress_silently_accepts_every_egress_only_key() {
for key in EGRESS_ONLY_CONFIG_KEYS {
for val in ["1", "primary", "halt", ""] {
let conf = format!("http::addr=127.0.0.1;{key}={val};");
SenderBuilder::from_conf(&conf).unwrap_or_else(|e| {
panic!(
"expected ingress to silently accept egress-only \
key {key}={val:?}, got {}",
e.msg()
)
});
}
}
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn ingress_accepts_full_egress_connect_string_unchanged() {
let conf = "http::addr=127.0.0.1:9000\
;username=u;password=p\
;path=/exec;max_version=1;compression=zstd;compression_level=3\
;max_batch_rows=10000;client_id=svc-a;target=primary\
;failover=on;failover_max_attempts=3\
;on_schema_error=drop;on_parse_error=halt\
;buffer_pool_size=8";
let builder = SenderBuilder::from_conf(conf).unwrap();
assert_eq!(builder.protocol, Protocol::Http);
assert_specified_eq(&builder.host, "127.0.0.1");
assert_specified_eq(&builder.port, "9000");
assert_specified_eq(&builder.username, Some("u".to_string()));
assert_specified_eq(&builder.password, Some("p".to_string()));
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_config_silently_accepts_reserved_on_error_policy_keys() {
for key in [
"on_server_error",
"on_schema_error",
"on_parse_error",
"on_internal_error",
"on_security_error",
"on_write_error",
] {
for val in ["halt", "drop", "auto", "anything", ""] {
let conf = format!("ws::addr=localhost:9000;{key}={val};");
SenderBuilder::from_conf(&conf)
.unwrap_or_else(|e| panic!("expected {key}={val:?} to parse, got {}", e.msg()));
}
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_store_and_forward_size_suffixes_match_java_config_surface() {
for (input, expected) in [
("64k", 64 * 1024_u64),
("64KB", 64 * 1024_u64),
("64m", 64 * 1024 * 1024_u64),
("4g", 4 * 1024 * 1024 * 1024_u64),
("1T", 1024_u64 * 1024 * 1024 * 1024),
] {
let conf = format!("ws::addr=localhost:9000;sf_max_segment_bytes={input};");
let builder = SenderBuilder::from_conf(conf).unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.sf_max_segment_bytes, expected);
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_store_and_forward_config_accepts_and_rejects_java_keys() {
SenderBuilder::from_conf("ws::addr=localhost:9000;request_durable_ack=off;").unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;request_durable_ack=on;").unwrap();
SenderBuilder::from_conf(
"ws::addr=localhost:9000;request_durable_ack=off;durable_ack_keepalive_interval_millis=5000;",
)
.unwrap();
SenderBuilder::from_conf(
"ws::addr=localhost:9000;request_durable_ack=on;durable_ack_keepalive_interval_millis=5000;",
)
.unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;durable_ack_keepalive_interval_millis=5000;")
.unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;durable_ack_keepalive_interval_millis=0;")
.unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;durable_ack_keepalive_interval_millis=-1;")
.unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;drain_orphans=off;").unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;drain_orphans=false;").unwrap();
SenderBuilder::from_conf(
"ws::addr=localhost:9000;drain_orphans=false;max_background_drainers=2;",
)
.unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;max_background_drainers=0;").unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;max_background_drainers=2;").unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;auth_timeout_ms=1;").unwrap();
let builder = SenderBuilder::new(Protocol::Ws, "localhost", 9000)
.auth_timeout(Duration::from_millis(750))
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.auth_timeout, Duration::from_millis(750));
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;sender_id=bad/id;"),
"invalid sender_id [value=bad/id, allowed-chars=[A-Za-z0-9_-]]",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;sf_max_segment_bytes=64mi;"),
"invalid sf_max_segment_bytes [value=64mi]",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;sf_durability=sync;"),
"invalid sf_durability [value=sync, allowed-values=[memory, periodic, flush, append]]",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;qwp_ws_progress=sync;"),
"invalid qwp_ws_progress [value=sync, allowed-values=[background, manual]]",
);
SenderBuilder::from_conf("ws::addr=localhost:9000;sf_append_deadline_millis=1234;").unwrap();
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;sf_append_deadline_millis=0;"),
"\"sf_append_deadline_millis\" must be greater than 0.",
);
for (input, expected) in [(-42, 0), (-1, 0), (0, 0), (5000, 5000), (120000, 120000)] {
let conf = format!("ws::addr=localhost:9000;close_flush_timeout_millis={input};");
let builder = SenderBuilder::from_conf(conf).unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.close_flush_timeout,
std::time::Duration::from_millis(expected as u64),
);
}
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;request_durable_ack=maybe;"),
"invalid request_durable_ack [value=maybe, allowed-values=[on, off]]",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;drain_orphans=maybe;"),
"invalid drain_orphans [value=maybe, allowed-values=[on, off, true, false]]",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;max_background_drainers=-1;"),
"max_background_drainers must be >= 0: -1",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;auth_timeout_ms=0;"),
"auth_timeout_ms must be > 0: 0",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;auth_timeout_ms=-1;"),
"auth_timeout_ms must be > 0: -1",
);
SenderBuilder::from_conf("ws::addr=localhost:9000;drain_orphans=on;").unwrap();
SenderBuilder::from_conf(
"ws::addr=localhost:9000;drain_orphans=true;max_background_drainers=0;",
)
.unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;error_inbox_capacity=64;").unwrap();
SenderBuilder::from_conf("ws::addr=localhost:9000;error_inbox_capacity=16;").unwrap();
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;error_inbox_capacity=15;"),
"error_inbox_capacity must be >= 16: 15",
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_rejects_unknown_config_key_but_tolerates_egress_keys() {
assert_conf_err(
SenderBuilder::from_conf("ws::addr=localhost:9000;totally_bogus_key=1;"),
"Unknown config key \"totally_bogus_key\"",
);
SenderBuilder::from_conf(
"ws::addr=localhost:9000;target=primary;compression=zstd;failover=on;zone=eu-1;max_batch_rows=1000;",
)
.unwrap();
}
#[cfg(all(feature = "sync-sender-qwp-ws", feature = "sync-sender-tcp"))]
#[test]
fn qwpws_store_and_forward_config_is_websocket_only() {
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;sf_dir=/tmp/qdb-rust-sf;"),
"The \"sf_dir\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;close_flush_timeout_millis=5000;"),
"The \"close_flush_timeout_millis\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;auth_timeout_ms=5000;"),
"The \"auth_timeout_ms\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;sf_append_deadline_millis=5000;"),
"The \"sf_append_deadline_millis\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;sf_sync_interval_millis=5000;"),
"The \"sf_sync_interval_millis\" setting is only supported for QWP/WebSocket.",
);
#[cfg(feature = "sync-sender-http")]
assert_conf_err(
SenderBuilder::from_conf("http::addr=localhost:9000;sf_sync_interval_millis=5000;"),
"The \"sf_sync_interval_millis\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;qwp_ws_progress=manual;"),
"The \"qwp_ws_progress\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;request_durable_ack=on;"),
"The \"request_durable_ack\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;request_durable_ack=off;"),
"The \"request_durable_ack\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;drain_orphans=off;"),
"The \"drain_orphans\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf(
"tcp::addr=localhost:9009;durable_ack_keepalive_interval_millis=5000;",
),
"The \"durable_ack_keepalive_interval_millis\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost:9009;max_background_drainers=2;"),
"The \"max_background_drainers\" setting is only supported for QWP/WebSocket.",
);
assert_conf_err(
SenderBuilder::from_conf(
"tcp::addr=localhost:9009;poison_min_escalation_window_millis=5000;",
),
"The \"poison_min_escalation_window_millis\" setting is only supported for QWP/WebSocket.",
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_store_and_forward_reserved_durability_fails_before_connect() {
assert_conf_err(
SenderBuilder::from_conf("ws::addr=127.0.0.1:1;sf_durability=flush;")
.unwrap()
.build(),
"sf_durability=flush is not yet supported (use sf_durability=memory or periodic)",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=127.0.0.1:1;sf_durability=append;")
.unwrap()
.build(),
"sf_durability=append is not yet supported (use sf_durability=memory or periodic)",
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_periodic_durability_config_matches_java() {
let defaulted = SenderBuilder::from_conf(
"ws::addr=localhost:9000;sf_dir=/tmp/qdb-rust-sf;sf_durability=periodic;",
)
.unwrap();
let qwp_ws = defaulted.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.sf_durability, conf::SfDurability::Periodic);
assert_defaulted_eq(&qwp_ws.sf_sync_interval, None);
assert_eq!(
qwp_ws.periodic_sync_interval(),
Some(Duration::from_millis(5000))
);
let explicit = SenderBuilder::from_conf(
"ws::addr=localhost:9000;sf_dir=/tmp/qdb-rust-sf;sf_durability=periodic;\
sf_sync_interval_millis=123;",
)
.unwrap();
let qwp_ws = explicit.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.sf_sync_interval, Some(Duration::from_millis(123)));
assert_eq!(
qwp_ws.periodic_sync_interval(),
Some(Duration::from_millis(123))
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=127.0.0.1:1;sf_sync_interval_millis=5000;")
.unwrap()
.build(),
"sf_sync_interval_millis requires sf_durability=periodic",
);
assert_conf_err(
SenderBuilder::from_conf("ws::addr=127.0.0.1:1;sf_durability=periodic;")
.unwrap()
.build(),
"sf_durability=periodic requires sf_dir",
);
assert_conf_err(
SenderBuilder::from_conf(
"ws::addr=127.0.0.1:1;sf_durability=periodic;sf_sync_interval_millis=0;",
),
"\"sf_sync_interval_millis\" must be greater than 0.",
);
let max_valid = i64::MAX / 1_000_000;
SenderBuilder::from_conf(format!(
"ws::addr=localhost:9000;sf_dir=/tmp/qdb-rust-sf;sf_durability=periodic;\
sf_sync_interval_millis={max_valid};"
))
.unwrap();
assert_conf_err(
SenderBuilder::from_conf(format!(
"ws::addr=localhost:9000;sf_dir=/tmp/qdb-rust-sf;sf_durability=periodic;\
sf_sync_interval_millis={};",
max_valid + 1
)),
format!("\"sf_sync_interval_millis\" must be at most {max_valid}."),
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn invalid_value() {
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost\n;"),
"Config parse error: invalid char '\\n' in value at position 19",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn specified_cant_change() {
let mut builder = SenderBuilder::from_conf("tcp::addr=localhost;").unwrap();
builder = builder.bind_interface("1.1.1.1").unwrap();
assert_conf_err(
builder.bind_interface("1.1.1.2"),
"\"bind_interface\" is already specified",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn missing_addr() {
assert_conf_err(
SenderBuilder::from_conf("tcp::"),
"Missing \"addr\" parameter in config string",
);
}
#[cfg(any(
feature = "sync-sender-tcp",
feature = "sync-sender-http",
feature = "sync-sender-qwp-udp"
))]
#[test]
fn unsupported_service() {
assert_conf_err(
SenderBuilder::from_conf("xaxa::addr=localhost;"),
"Unsupported protocol: xaxa",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn uppercase_scheme_accepted_tcp() {
let builder = SenderBuilder::from_conf("TCP::addr=localhost:9009;").unwrap();
assert_eq!(builder.protocol, Protocol::Tcp);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn uppercase_scheme_accepted_https() {
let builder = SenderBuilder::from_conf("HTTPS::addr=localhost:9000;").unwrap();
assert_eq!(builder.protocol, Protocol::Https);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn uppercase_scheme_accepted_udp() {
let builder = SenderBuilder::from_conf("UDP::addr=localhost:9009;").unwrap();
assert_eq!(builder.protocol, Protocol::Udp);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn uppercase_scheme_accepted_ws() {
let builder = SenderBuilder::from_conf("WS::addr=localhost:9000;").unwrap();
assert_eq!(builder.protocol, Protocol::Ws);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn uppercase_ws_preserves_multi_addr() {
let builder = SenderBuilder::from_conf("WS::addr=h1:9001,h2:9002,h3:9003;").unwrap();
assert_eq!(builder.protocol, Protocol::Ws);
let endpoints: &Vec<conf::QwpWsEndpoint> = &builder.qwp_ws.as_ref().unwrap().endpoints;
assert_eq!(endpoints.len(), 3);
assert_eq!(endpoints[0].host, "h1");
assert_eq!(endpoints[0].port, "9001");
assert_eq!(endpoints[1].host, "h2");
assert_eq!(endpoints[1].port, "9002");
assert_eq!(endpoints[2].host, "h3");
assert_eq!(endpoints[2].port, "9003");
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn mixed_case_ws_preserves_multi_addr() {
let builder = SenderBuilder::from_conf("Ws::addr=h1:9001,h2:9002;").unwrap();
assert_eq!(builder.protocol, Protocol::Ws);
let endpoints: &Vec<conf::QwpWsEndpoint> = &builder.qwp_ws.as_ref().unwrap().endpoints;
assert_eq!(endpoints.len(), 2);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_basic_auth() {
let builder =
SenderBuilder::from_conf("http::addr=localhost;username=user123;password=pass321;")
.unwrap();
let auth = builder.build_auth().unwrap();
match auth.unwrap() {
conf::AuthParams::Basic(conf::BasicAuthParams { username, password }) => {
assert_eq!(username, "user123");
assert_eq!(password, "pass321");
}
_ => {
panic!("Expected AuthParams::Basic");
}
}
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_token_auth() {
let builder = SenderBuilder::from_conf("http::addr=localhost:9000;token=token123;").unwrap();
let auth = builder.build_auth().unwrap();
match auth.unwrap() {
conf::AuthParams::Token(conf::TokenAuthParams { token }) => {
assert_eq!(token, "token123");
}
_ => {
panic!("Expected AuthParams::Token");
}
}
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn incomplete_basic_auth() {
assert_conf_err(
SenderBuilder::from_conf("http::addr=localhost;username=user123;")
.unwrap()
.build(),
"Basic authentication parameter \"username\" is present, but \"password\" is missing.",
);
assert_conf_err(
SenderBuilder::from_conf("http::addr=localhost;password=pass321;")
.unwrap()
.build(),
"Basic authentication parameter \"password\" is present, but \"username\" is missing.",
);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn zero_timeout_forbidden() {
assert_conf_err(
SenderBuilder::from_conf("http::addr=localhost;username=user123;request_timeout=0;"),
"\"request_timeout\" must be greater than 0.",
);
assert_conf_err(
SenderBuilder::new(Protocol::Http, "localhost", 9000)
.request_timeout(Duration::from_millis(0)),
"\"request_timeout\" must be greater than 0.",
);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn misspelled_basic_auth() {
assert_conf_err(
Sender::from_conf("http::addr=localhost;username=user123;pass=pass321;"),
r##"Basic authentication parameter "username" is present, but "password" is missing."##,
);
assert_conf_err(
Sender::from_conf("http::addr=localhost;user=user123;password=pass321;"),
r##"Basic authentication parameter "password" is present, but "username" is missing."##,
);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn inconsistent_http_auth() {
let expected_err_msg = r##"Inconsistent HTTP authentication parameters. Specify either "username" and "password", or just "token"."##;
assert_conf_err(
Sender::from_conf("http::addr=localhost;username=user123;token=token123;"),
expected_err_msg,
);
assert_conf_err(
Sender::from_conf("http::addr=localhost;password=pass321;token=token123;"),
expected_err_msg,
);
}
#[cfg(all(feature = "sync-sender-tcp", feature = "sync-sender-http"))]
#[test]
fn cant_use_basic_auth_with_tcp() {
let builder = SenderBuilder::new(Protocol::Tcp, "localhost", 9000)
.username("user123")
.unwrap()
.password("pass321")
.unwrap();
assert_conf_err(
builder.build_auth(),
"The \"basic_auth\" setting can only be used with the ILP/HTTP protocol.",
);
}
#[cfg(all(feature = "sync-sender-tcp", feature = "sync-sender-http"))]
#[test]
fn cant_use_token_auth_with_tcp() {
let builder = SenderBuilder::new(Protocol::Tcp, "localhost", 9000)
.token("token123")
.unwrap();
assert_conf_err(
builder.build_auth(),
"Token authentication only be used with the ILP/HTTP protocol.",
);
}
#[cfg(all(feature = "sync-sender-tcp", feature = "sync-sender-http"))]
#[test]
fn cant_use_ecdsa_auth_with_http() {
let builder = SenderBuilder::from_conf("http::addr=localhost;")
.unwrap()
.username("key_id123")
.unwrap()
.token("priv_key123")
.unwrap()
.token_x("pub_key1")
.unwrap()
.token_y("pub_key2")
.unwrap();
assert_conf_err(
builder.build_auth(),
"ECDSA authentication is only available with ILP/TCP and not available with ILP/HTTP.",
);
}
#[cfg(all(not(feature = "sync-sender-tcp"), feature = "sync-sender-http"))]
#[test]
fn cant_use_ecdsa_auth_with_http_ex_tcp_support() {
let mk_builder = || {
SenderBuilder::from_conf("http::addr=localhost;")
.unwrap()
.username("key_id123")
.unwrap()
.token("priv_key123")
.unwrap()
};
assert_conf_err(
mk_builder().token_x("pub_key1"),
"cannot specify \"token_x\": ECDSA authentication is only available with ILP/TCP and not available with ILP/HTTP.",
);
assert_conf_err(
mk_builder().token_y("pub_key2"),
"cannot specify \"token_y\": ECDSA authentication is only available with ILP/TCP and not available with ILP/HTTP.",
);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_protocol_version_unsupported() {
for version in ["1", "2", "3"] {
let conf = format!("udp::addr=localhost;protocol_version={version};");
assert_conf_err(
SenderBuilder::from_conf(&conf),
"The \"protocol_version\" setting is not supported for QWP/UDP.",
);
}
for version in [
ProtocolVersion::V1,
ProtocolVersion::V2,
ProtocolVersion::V3,
] {
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).protocol_version(version),
"The \"protocol_version\" setting is not supported for QWP/UDP.",
);
}
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_max_datagram_size_requires_qwp_udp() {
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).max_datagram_size(0),
"\"max_datagram_size\" must be greater than 0.",
);
#[cfg(feature = "sync-sender-http")]
assert_conf_err(
SenderBuilder::new(Protocol::Http, "localhost", 9000).max_datagram_size(1400),
"The \"max_datagram_size\" setting is only supported for QWP/UDP.",
);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_max_datagram_size_accepts_udp_limit_and_rejects_above_it() {
let builder = SenderBuilder::new(Protocol::Udp, "localhost", 9007)
.max_datagram_size(65507)
.unwrap();
let Some(qwp_udp) = builder.qwp_udp.as_ref() else {
panic!("Expected Some(QwpUdpConfig)");
};
assert_specified_eq(&qwp_udp.max_datagram_size, 65507usize);
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).max_datagram_size(65508),
"\"max_datagram_size\" must not exceed 65507 (UDP/IPv4 limit).",
);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_multicast_ttl_requires_qwp_udp() {
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).multicast_ttl(256),
"\"multicast_ttl\" must be between 0 and 255.",
);
#[cfg(feature = "sync-sender-http")]
assert_conf_err(
SenderBuilder::new(Protocol::Http, "localhost", 9000).multicast_ttl(1),
"The \"multicast_ttl\" setting is only supported for QWP/UDP.",
);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_config_string_rejects_invalid_datagram_size_and_multicast_ttl() {
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;max_datagram_size=0;"),
"\"max_datagram_size\" must be greater than 0.",
);
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;max_datagram_size=65508;"),
"\"max_datagram_size\" must not exceed 65507 (UDP/IPv4 limit).",
);
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;multicast_ttl=256;"),
"\"multicast_ttl\" must be between 0 and 255.",
);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_bind_interface_is_supported_via_builder_api() {
let builder = SenderBuilder::new(Protocol::Udp, "239.1.2.3", 9007)
.bind_interface("192.168.1.10")
.unwrap();
assert_eq!(builder.protocol, Protocol::Udp);
assert_specified_eq(&builder.net_interface, Some("192.168.1.10".to_string()));
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_auth_settings_are_rejected_at_config_time() {
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;username=user123;"),
"The \"username\" setting is not supported for QWP/UDP.",
);
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;password=pass321;"),
"The \"password\" setting is not supported for QWP/UDP.",
);
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;token=token123;"),
"The \"token\" setting is not supported for QWP/UDP.",
);
#[cfg(feature = "sync-sender-tcp")]
{
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;token_x=pub_key1;"),
"The \"token_x\" setting is not supported for QWP/UDP.",
);
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;token_y=pub_key2;"),
"The \"token_y\" setting is not supported for QWP/UDP.",
);
}
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).username("user123"),
"The \"username\" setting is not supported for QWP/UDP.",
);
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).password("pass321"),
"The \"password\" setting is not supported for QWP/UDP.",
);
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007).token("token123"),
"The \"token\" setting is not supported for QWP/UDP.",
);
}
#[cfg(feature = "sync-sender-qwp-udp")]
#[test]
fn udp_auth_timeout_is_rejected_at_config_time() {
assert_conf_err(
SenderBuilder::from_conf("udp::addr=localhost;auth_timeout=100;"),
"The \"auth_timeout\" setting is not supported for QWP/UDP.",
);
assert_conf_err(
SenderBuilder::new(Protocol::Udp, "localhost", 9007)
.auth_timeout(Duration::from_millis(100)),
"The \"auth_timeout\" setting is not supported for QWP/UDP.",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn set_auth_specifies_tcp() {
let mut builder = SenderBuilder::new(Protocol::Tcp, "localhost", 9000);
assert_eq!(builder.protocol, Protocol::Tcp);
builder = builder
.username("key_id123")
.unwrap()
.token("priv_key123")
.unwrap()
.token_x("pub_key1")
.unwrap()
.token_y("pub_key2")
.unwrap();
assert_eq!(builder.protocol, Protocol::Tcp);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn set_net_interface_specifies_tcp() {
let builder = SenderBuilder::new(Protocol::Tcp, "localhost", 9000);
assert_eq!(builder.protocol, Protocol::Tcp);
builder.bind_interface("55.88.0.4").unwrap();
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcp_ecdsa_auth() {
let builder = SenderBuilder::from_conf(
"tcp::addr=localhost:9000;username=user123;token=token123;token_x=xtok123;token_y=ytok123;",
)
.unwrap();
let auth = builder.build_auth().unwrap();
match auth.unwrap() {
conf::AuthParams::Ecdsa(conf::EcdsaAuthParams {
key_id,
priv_key,
pub_key_x,
pub_key_y,
}) => {
assert_eq!(key_id, "user123");
assert_eq!(priv_key, "token123");
assert_eq!(pub_key_x, "xtok123");
assert_eq!(pub_key_y, "ytok123");
}
#[cfg(feature = "sync-sender-http")]
_ => {
panic!("Expected AuthParams::Ecdsa");
}
}
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn incomplete_tcp_ecdsa_auth() {
let expected_err_msg = r##"Incomplete ECDSA authentication parameters. Specify either all or none of: "username", "token", "token_x", "token_y"."##;
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost;username=user123;")
.unwrap()
.build(),
expected_err_msg,
);
assert_conf_err(
SenderBuilder::from_conf("tcp::addr=localhost;username=user123;token=token123;")
.unwrap()
.build(),
expected_err_msg,
);
assert_conf_err(
SenderBuilder::from_conf(
"tcp::addr=localhost;username=user123;token=token123;token_x=123;",
)
.unwrap()
.build(),
expected_err_msg,
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn misspelled_tcp_ecdsa_auth() {
assert_conf_err(
Sender::from_conf("tcp::addr=localhost;username=user123;tokenx=123;"),
"Incomplete ECDSA authentication parameters. Specify either all or none of: \"username\", \"token\", \"token_x\", \"token_y\".",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_tls_verify_on() {
let builder = SenderBuilder::from_conf("tcps::addr=localhost;tls_verify=on;").unwrap();
assert!(builder.protocol.tls_enabled());
#[cfg(feature = "tls-webpki-certs")]
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::WebpkiRoots);
#[cfg(not(feature = "tls-webpki-certs"))]
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::OsRoots);
}
#[cfg(feature = "sync-sender-tcp")]
#[cfg(feature = "insecure-skip-verify")]
#[test]
fn tcps_tls_verify_unsafe_off() {
let builder = SenderBuilder::from_conf("tcps::addr=localhost;tls_verify=unsafe_off;").unwrap();
assert!(builder.protocol.tls_enabled());
assert_defaulted_eq(&builder.tls_ca, CertificateAuthority::WebpkiRoots);
assert_specified_eq(&builder.tls_verify, false);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_tls_verify_invalid() {
assert_conf_err(
SenderBuilder::from_conf("tcps::addr=localhost;tls_verify=off;"),
r##"Config parameter "tls_verify" must be either "on" or "unsafe_off".'"##,
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_tls_roots_webpki() {
let builder = SenderBuilder::from_conf("tcps::addr=localhost;tls_ca=webpki_roots;");
#[cfg(feature = "tls-webpki-certs")]
{
let builder = builder.unwrap();
assert!(builder.protocol.tls_enabled());
assert_specified_eq(&builder.tls_ca, CertificateAuthority::WebpkiRoots);
assert_defaulted_eq(&builder.tls_roots, None);
}
#[cfg(not(feature = "tls-webpki-certs"))]
assert_eq!(
"Config parameter \"tls_ca=webpki_roots\" requires the \"tls-webpki-certs\" feature",
builder.unwrap_err().msg()
);
}
#[cfg(feature = "sync-sender-tcp")]
#[cfg(feature = "tls-native-certs")]
#[test]
fn tcps_tls_roots_os() {
let builder = SenderBuilder::from_conf("tcps::addr=localhost;tls_ca=os_roots;").unwrap();
assert!(builder.protocol.tls_enabled());
assert_specified_eq(&builder.tls_ca, CertificateAuthority::OsRoots);
assert_defaulted_eq(&builder.tls_roots, None);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_tls_roots_file() {
use std::io::Write;
let tmp_dir = TempDir::new().unwrap();
let path = tmp_dir.path().join("cacerts.pem");
let mut file = std::fs::File::create(&path).unwrap();
file.write_all(b"dummy").unwrap();
let builder = SenderBuilder::from_conf(format!(
"tcps::addr=localhost;tls_roots={};",
path.to_str().unwrap()
))
.unwrap();
assert_specified_eq(&builder.tls_ca, CertificateAuthority::PemFile);
assert_specified_eq(&builder.tls_roots, path);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_tls_roots_file_missing() {
let err =
SenderBuilder::from_conf("tcps::addr=localhost;tls_roots=/some/invalid/path/cacerts.pem;")
.unwrap_err();
assert_eq!(err.code(), ErrorCode::ConfigError);
assert!(
err.msg()
.contains("Could not open root certificate file from path")
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn tcps_tls_roots_file_with_password() {
use std::io::Write;
let tmp_dir = TempDir::new().unwrap();
let path = tmp_dir.path().join("cacerts.pem");
let mut file = std::fs::File::create(&path).unwrap();
file.write_all(b"dummy").unwrap();
let builder_or_err = SenderBuilder::from_conf(format!(
"tcps::addr=localhost;tls_roots={};tls_roots_password=extremely_secure;",
path.to_str().unwrap()
));
assert_conf_err(
builder_or_err,
"\"tls_roots_password\" is only supported for QWP/WebSocket \
(ws / wss). ILP/TCP and ILP/HTTP transports read unencrypted \
PEM via rustls.",
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpwss_tls_roots_password_accepted() {
use std::io::Write;
let tmp_dir = TempDir::new().unwrap();
let path = tmp_dir.path().join("trust.jks");
let mut file = std::fs::File::create(&path).unwrap();
file.write_all(b"placeholder").unwrap();
let builder = SenderBuilder::from_conf(format!(
"wss::addr=localhost;tls_roots={};tls_roots_password=secret;",
path.to_str().unwrap()
))
.unwrap();
assert_specified_eq(&builder.tls_roots_password, Some("secret".to_string()));
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpwss_tls_roots_password_without_path_rejected() {
let builder_or_err =
SenderBuilder::from_conf("wss::addr=localhost;tls_roots_password=secret;").unwrap();
let err = builder_or_err.build().unwrap_err();
assert!(
err.msg().contains("tls_roots_password") && err.msg().contains("tls_roots"),
"msg: {}",
err.msg()
);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_request_min_throughput() {
let builder =
SenderBuilder::from_conf("http::addr=localhost;request_min_throughput=100;").unwrap();
let Some(http_config) = builder.http else {
panic!("Expected Some(HttpConfig)");
};
assert_specified_eq(&http_config.request_min_throughput, 100u64);
assert_defaulted_eq(&http_config.request_timeout, Duration::from_millis(10000));
assert_defaulted_eq(&http_config.retry_timeout, Duration::from_millis(10000));
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_request_timeout() {
let builder = SenderBuilder::from_conf("http::addr=localhost;request_timeout=100;").unwrap();
let Some(http_config) = builder.http else {
panic!("Expected Some(HttpConfig)");
};
assert_defaulted_eq(&http_config.request_min_throughput, 102400u64);
assert_specified_eq(&http_config.request_timeout, Duration::from_millis(100));
assert_defaulted_eq(&http_config.retry_timeout, Duration::from_millis(10000));
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_retry_timeout() {
let builder = SenderBuilder::from_conf("http::addr=localhost;retry_timeout=100;").unwrap();
let Some(http_config) = builder.http else {
panic!("Expected Some(HttpConfig)");
};
assert_defaulted_eq(&http_config.request_min_throughput, 102400u64);
assert_defaulted_eq(&http_config.request_timeout, Duration::from_millis(10000));
assert_specified_eq(&http_config.retry_timeout, Duration::from_millis(100));
assert_defaulted_eq(&http_config.retry_max_backoff, Duration::from_millis(1000));
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_retry_max_backoff() {
let builder =
SenderBuilder::from_conf("http::addr=localhost;retry_max_backoff_millis=250;").unwrap();
let Some(http_config) = builder.http else {
panic!("Expected Some(HttpConfig)");
};
assert_specified_eq(&http_config.retry_max_backoff, Duration::from_millis(250));
assert_defaulted_eq(&http_config.retry_timeout, Duration::from_millis(10000));
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn http_retry_max_backoff_below_min_rejected() {
let msg = "\"retry_max_backoff_millis\" must be at least 10.";
assert_conf_err(
SenderBuilder::from_conf("http::addr=localhost;retry_max_backoff_millis=0;"),
msg,
);
assert_conf_err(
SenderBuilder::from_conf("http::addr=localhost;retry_max_backoff_millis=3;"),
msg,
);
}
#[cfg(all(feature = "sync-sender-tcp", feature = "sync-sender-http"))]
#[test]
fn retry_max_backoff_rejected_on_non_http() {
assert_conf_err(
SenderBuilder::from_conf("tcps::addr=localhost;retry_max_backoff_millis=250;"),
"retry_max_backoff_millis is supported only in ILP over HTTP.",
);
}
#[cfg(feature = "sync-sender-http")]
#[test]
fn connect_timeout_uses_request_timeout() {
use std::time::Instant;
let request_timeout = Duration::from_millis(10);
let builder = SenderBuilder::new(Protocol::Http, "127.0.0.2", "1111")
.request_timeout(request_timeout)
.unwrap()
.protocol_version(ProtocolVersion::V2)
.unwrap()
.retry_timeout(Duration::from_millis(10))
.unwrap()
.request_min_throughput(0)
.unwrap();
let mut sender = builder.build().unwrap();
let mut buf = sender.new_buffer();
buf.table("x")
.unwrap()
.symbol("x", "x")
.unwrap()
.at_now()
.unwrap();
let start = Instant::now();
sender
.flush(&mut buf)
.expect_err("Request did not time out");
assert!(Instant::now() - start < Duration::from_secs(10));
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn auto_flush_off() {
SenderBuilder::from_conf("tcps::addr=localhost;auto_flush=off;").unwrap();
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn auto_flush_unsupported() {
assert_conf_err(
SenderBuilder::from_conf("tcps::addr=localhost;auto_flush=on;"),
"Invalid auto_flush value 'on'. This client does not support \
auto-flush, so the only accepted value is 'off'",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn auto_flush_rows_unsupported() {
assert_conf_err(
SenderBuilder::from_conf("tcps::addr=localhost;auto_flush_rows=100;"),
"Invalid configuration parameter \"auto_flush_rows\". This client does not support auto-flush",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn auto_flush_bytes_unsupported() {
assert_conf_err(
SenderBuilder::from_conf("tcps::addr=localhost;auto_flush_bytes=100;"),
"Invalid configuration parameter \"auto_flush_bytes\". This client does not support auto-flush",
);
}
#[cfg(feature = "sync-sender-tcp")]
#[test]
fn auto_flush_interval_unsupported() {
assert_conf_err(
SenderBuilder::from_conf("tcps::addr=localhost;auto_flush_interval=500;"),
"Invalid configuration parameter \"auto_flush_interval\". This client does not support auto-flush",
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_defaults_leave_initial_connect_retry_off() {
let builder = SenderBuilder::from_conf("ws::addr=localhost:9000;").unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_defaulted_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Off
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_connect_timeout_defaults_to_unset() {
let builder = SenderBuilder::from_conf("ws::addr=localhost:9000;").unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_defaulted_eq(&qwp_ws.connect_timeout, None);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_connect_timeout_parses() {
let builder = SenderBuilder::from_conf("ws::addr=localhost:9000;connect_timeout=250;").unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(&qwp_ws.connect_timeout, Some(Duration::from_millis(250)));
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_connect_timeout_zero_rejected() {
let err = SenderBuilder::from_conf("ws::addr=localhost:9000;connect_timeout=0;").unwrap_err();
assert_eq!(err.code(), crate::ErrorCode::ConfigError);
assert!(err.msg().contains("connect_timeout"), "msg: {}", err.msg());
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_reconnect_max_duration_implies_initial_connect_retry_sync() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;reconnect_max_duration_millis=120000;")
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_defaulted_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Sync,
);
assert_specified_eq(&qwp_ws.reconnect_max_duration, Duration::from_secs(120));
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_reconnect_initial_backoff_implies_initial_connect_retry_sync() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;reconnect_initial_backoff_millis=250;")
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Sync,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_reconnect_max_backoff_implies_initial_connect_retry_sync() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;reconnect_max_backoff_millis=10000;")
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Sync,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_explicit_initial_connect_retry_off_is_preserved() {
let builder = SenderBuilder::from_conf(
"ws::addr=localhost:9000;reconnect_max_duration_millis=120000;initial_connect_retry=off;",
)
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Off,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_explicit_initial_connect_retry_async_is_preserved() {
let builder = SenderBuilder::from_conf(
"ws::addr=localhost:9000;reconnect_max_duration_millis=120000;initial_connect_retry=async;",
)
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Async,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Async,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_explicit_off_before_reconnect_key_is_preserved() {
let builder = SenderBuilder::from_conf(
"ws::addr=localhost:9000;initial_connect_retry=off;reconnect_max_duration_millis=120000;",
)
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Off,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_multiple_reconnect_keys_promote_once() {
let builder = SenderBuilder::from_conf(
"ws::addr=localhost:9000;\
reconnect_max_duration_millis=120000;\
reconnect_initial_backoff_millis=250;\
reconnect_max_backoff_millis=10000;",
)
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Sync,
);
assert_specified_eq(&qwp_ws.reconnect_max_duration, Duration::from_secs(120));
assert_specified_eq(
&qwp_ws.reconnect_initial_backoff,
Duration::from_millis(250),
);
assert_specified_eq(&qwp_ws.reconnect_max_backoff, Duration::from_secs(10));
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_reconnect_implies_initial_retry_via_builder_api() {
let builder = SenderBuilder::new(Protocol::Ws, "localhost", 9000)
.reconnect_max_duration(Duration::from_secs(120))
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_defaulted_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Sync,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_explicit_sync_after_promotion_stays_explicit() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;reconnect_max_duration_millis=120000;")
.unwrap()
.initial_connect_retry(true)
.unwrap();
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Sync,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Sync,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_explicit_off_after_promotion_overrides_promoted_sync() {
let builder =
SenderBuilder::from_conf("ws::addr=localhost:9000;reconnect_max_duration_millis=120000;")
.unwrap()
.initial_connect_retry(false)
.expect("explicit off must not collide with a build-time default");
let qwp_ws = builder.qwp_ws.as_ref().unwrap();
assert_specified_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Off,
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn qwpws_resolves_initial_connect_retry_off_without_reconnect_keys() {
let qwp_ws = conf::QwpWsConfig::default();
assert_defaulted_eq(
&qwp_ws.initial_connect_retry,
conf::QwpWsInitialConnectMode::Off,
);
assert_eq!(
qwp_ws.resolve_initial_connect_retry(),
conf::QwpWsInitialConnectMode::Off,
);
}
#[test]
fn config_setting_is_specified_reports_variant() {
let mut setting: ConfigSetting<u32> = ConfigSetting::new_default(7);
assert!(!setting.is_specified());
setting.set_specified("test", 42).unwrap();
assert!(setting.is_specified());
}
fn assert_specified_eq<V: PartialEq + Debug, IntoV: Into<V>>(
actual: &ConfigSetting<V>,
expected: IntoV,
) {
let expected = expected.into();
if let ConfigSetting::Specified(actual_value) = actual {
assert_eq!(actual_value, &expected);
} else {
panic!("Expected Specified({expected:?}), but got {actual:?}");
}
}
fn assert_defaulted_eq<V: PartialEq + std::fmt::Debug, IntoV: Into<V>>(
actual: &ConfigSetting<V>,
expected: IntoV,
) {
let expected = expected.into();
if let ConfigSetting::Defaulted(actual_value) = actual {
assert_eq!(actual_value, &expected);
} else {
panic!("Expected Defaulted({expected:?}), but got {actual:?}");
}
}
fn assert_conf_err<T, M: AsRef<str>>(result: Result<T>, expect_msg: M) {
let Err(err) = result else {
panic!("Got Ok, expected ConfigError: {}", expect_msg.as_ref());
};
assert_eq!(err.code(), ErrorCode::ConfigError);
assert_eq!(err.msg(), expect_msg.as_ref());
}