use std::num::NonZeroUsize;
use std::time::Duration;
use crate::config::ClientConfig;
use crate::protocol::request::ConnectionMode;
use crate::session::backoff::BackoffPolicy;
const DEFAULT_KEEPALIVE_SLACK: Duration = Duration::from_millis(3000);
const DEFAULT_OPEN_TIMEOUT: Duration = Duration::from_secs(10);
const DEFAULT_EVENT_CAPACITY: usize = 1024;
const DEFAULT_COMMAND_CAPACITY: usize = 64;
#[derive(Clone, Default, PartialEq, Eq)]
pub(crate) struct Credentials {
pub(crate) user: Option<String>,
pub(crate) password: Option<String>,
pub(crate) adapter_set: Option<String>,
}
impl std::fmt::Debug for Credentials {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("Credentials")
.field("user", &self.user)
.field("password", &self.password.as_ref().map(|_| "<redacted>"))
.field("adapter_set", &self.adapter_set)
.finish()
}
}
#[derive(Debug, Clone)]
pub(crate) struct SessionOptions {
pub(crate) credentials: Credentials,
pub(crate) connection: ConnectionMode,
pub(crate) content_length: Option<u64>,
pub(crate) keepalive_slack: Duration,
pub(crate) open_timeout: Duration,
pub(crate) backoff: BackoffPolicy,
pub(crate) event_capacity: NonZeroUsize,
pub(crate) command_capacity: NonZeroUsize,
}
impl Default for SessionOptions {
fn default() -> Self {
Self {
credentials: Credentials::default(),
connection: ConnectionMode::default(),
content_length: None,
keepalive_slack: DEFAULT_KEEPALIVE_SLACK,
open_timeout: DEFAULT_OPEN_TIMEOUT,
backoff: BackoffPolicy::default(),
event_capacity: NonZeroUsize::new(DEFAULT_EVENT_CAPACITY).unwrap_or(NonZeroUsize::MIN),
command_capacity: NonZeroUsize::new(DEFAULT_COMMAND_CAPACITY)
.unwrap_or(NonZeroUsize::MIN),
}
}
}
impl SessionOptions {
#[must_use]
pub(crate) fn from_client_config(config: ClientConfig) -> Self {
let (adapter_set, credentials, options) = config.into_parts();
let (user, password) = credentials.into_parts();
let retry = options.retry();
let defaults = Self::default();
Self {
credentials: Credentials {
user,
password,
adapter_set: adapter_set.map(|set| set.as_str().to_owned()),
},
connection: ConnectionMode::Streaming {
inactivity_millis: options.inactivity_commitment().map(as_millis),
keepalive_millis: options.keepalive_hint().map(as_millis),
send_sync: if options.send_sync() {
None
} else {
Some(false)
},
},
content_length: options.content_length().map(std::num::NonZeroU64::get),
keepalive_slack: options.keepalive_slack(),
open_timeout: options.open_timeout(),
backoff: BackoffPolicy {
initial: retry.initial_delay(),
max: retry.max_delay(),
max_attempts: retry.max_attempts(),
},
event_capacity: options.session_event_capacity(),
command_capacity: defaults.command_capacity,
}
}
#[must_use = "builders do nothing unless the result is used"]
pub(crate) fn with_credentials(mut self, credentials: Credentials) -> Self {
self.credentials = credentials;
self
}
#[must_use = "builders do nothing unless the result is used"]
pub(crate) fn with_connection(mut self, connection: ConnectionMode) -> Self {
self.connection = connection;
self
}
#[must_use = "builders do nothing unless the result is used"]
pub(crate) fn with_backoff(mut self, backoff: BackoffPolicy) -> Self {
self.backoff = backoff;
self
}
#[must_use]
#[inline]
pub(crate) const fn is_polling(&self) -> bool {
matches!(self.connection, ConnectionMode::Polling { .. })
}
#[must_use]
pub(crate) fn heartbeat_interval(&self) -> Option<Duration> {
match self.connection {
ConnectionMode::Streaming {
inactivity_millis, ..
} => inactivity_millis
.map(Duration::from_millis)
.map(|commitment| commitment / 2),
ConnectionMode::Polling { .. } => None,
}
}
}
fn as_millis(value: Duration) -> u64 {
u64::try_from(value.as_millis()).unwrap_or(u64::MAX)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_credentials_debug_never_shows_the_password() {
let credentials = Credentials {
user: Some("alice".to_owned()),
password: Some("hunter2".to_owned()),
adapter_set: Some("WELCOME".to_owned()),
};
let rendered = format!("{credentials:?}");
assert!(!rendered.contains("hunter2"), "{rendered}");
assert!(rendered.contains("<redacted>"), "{rendered}");
assert!(rendered.contains("alice"), "{rendered}");
}
#[test]
fn test_options_heartbeat_is_half_the_inactivity_commitment() {
let options = SessionOptions::default().with_connection(ConnectionMode::Streaming {
inactivity_millis: Some(8000),
keepalive_millis: None,
send_sync: None,
});
assert_eq!(
options.heartbeat_interval(),
Some(Duration::from_millis(4000))
);
}
#[test]
fn test_options_no_heartbeat_without_a_commitment() {
assert_eq!(SessionOptions::default().heartbeat_interval(), None);
}
fn address() -> crate::config::ServerAddress {
match crate::config::ServerAddress::try_new("wss://push.example.com") {
Ok(address) => address,
Err(error) => unreachable!("the fixture address is valid: {error}"),
}
}
#[test]
fn test_session_options_carry_the_configured_values() {
use std::num::NonZeroU32;
use crate::config::{
AdapterSet, ClientConfig, ConnectionOptions, Credentials as PublicCredentials,
RetryPolicy,
};
let adapters = match AdapterSet::try_new("DEMO") {
Ok(set) => set,
Err(error) => unreachable!("the fixture adapter set is valid: {error}"),
};
let retry = RetryPolicy::default()
.with_initial_delay(Duration::from_millis(100))
.with_max_delay(Duration::from_secs(2))
.with_max_attempts(NonZeroU32::new(3));
let config = ClientConfig::builder(address())
.with_adapter_set(adapters)
.with_credentials(PublicCredentials::new("alice", "hunter2"))
.with_options(
ConnectionOptions::default()
.with_inactivity_commitment(Some(Duration::from_secs(20)))
.with_keepalive_hint(Some(Duration::from_secs(5)))
.with_send_sync(false)
.with_retry(retry),
)
.build();
let config = match config {
Ok(config) => config,
Err(error) => panic!("rejected: {error}"),
};
let session = SessionOptions::from_client_config(config);
assert_eq!(session.credentials.user.as_deref(), Some("alice"));
assert_eq!(session.credentials.adapter_set.as_deref(), Some("DEMO"));
assert_eq!(session.backoff.initial, Duration::from_millis(100));
assert_eq!(session.backoff.max, Duration::from_secs(2));
assert_eq!(session.backoff.max_attempts.map(NonZeroU32::get), Some(3));
match session.connection {
ConnectionMode::Streaming {
inactivity_millis,
keepalive_millis,
send_sync,
} => {
assert_eq!(inactivity_millis, Some(20_000));
assert_eq!(keepalive_millis, Some(5_000));
assert_eq!(send_sync, Some(false));
}
ConnectionMode::Polling { .. } => panic!("expected a streaming connection"),
}
}
#[test]
fn test_session_options_omit_send_sync_when_it_is_enabled() {
let config = match crate::config::ClientConfig::builder(address()).build() {
Ok(config) => config,
Err(error) => panic!("rejected: {error}"),
};
match SessionOptions::from_client_config(config).connection {
ConnectionMode::Streaming { send_sync, .. } => assert_eq!(send_sync, None),
ConnectionMode::Polling { .. } => panic!("expected a streaming connection"),
}
}
#[test]
fn test_options_no_heartbeat_on_a_polling_connection() {
let options = SessionOptions::default().with_connection(ConnectionMode::Polling {
polling_millis: 5000,
idle_millis: Some(10_000),
});
assert_eq!(options.heartbeat_interval(), None);
assert!(options.is_polling());
}
}