#![allow(dead_code)]
pub(crate) mod backoff;
pub(crate) mod liveness;
pub(crate) mod options;
use std::collections::{BTreeMap, HashMap};
use std::num::NonZeroU32;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Duration;
use tokio::sync::{mpsc, watch};
use tokio::time::Instant;
use crate::config::ClientConfig;
pub(crate) use crate::protocol::ProtocolError;
use crate::protocol::request::{
BindSession, ConnectionMode, CreateSession, DestroyCause, DestroySession, ForceRebind,
Heartbeat, MaxFrequencyLimit, MessageSend, ReconfigureSubscription, RequestId, Subscribe,
TlcpRequest, TransportKind, Unsubscribe,
};
use crate::protocol::response::{Notification, parse_line};
use crate::session::backoff::Backoff;
use crate::session::liveness::{Liveness, LivenessAction};
use crate::session::options::SessionOptions;
use crate::subscription::manager::SubscriptionManager;
use crate::transport::http::{HttpMode, HttpTransport};
use crate::transport::ws::WsTransport;
use crate::transport::{AnyTransport, TransportError};
pub(crate) use crate::protocol::request::{
DecimalNumber, RequestedBufferSize, RequestedMaxFrequency, SequenceName, Snapshot,
SubscriptionMode,
};
pub(crate) use crate::protocol::response::{
Bandwidth, FilteringMode, MaxFrequency, MessageSequence,
};
pub(crate) use crate::subscription::manager::{
CommandFields, SubscriptionError, SubscriptionEvent,
};
pub(crate) use crate::transport::{EncodedRequest, StreamOpen, Transport, TransportProperties};
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub(crate) enum SessionError {
#[error("invalid session configuration: {reason}")]
Configuration {
reason: String,
},
#[error(transparent)]
Protocol(#[from] ProtocolError),
#[error("the session driver has stopped")]
Stopped,
#[error("the {what} identifier space is exhausted")]
Exhausted {
what: &'static str,
},
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub(crate) struct SessionId(String);
impl SessionId {
#[must_use]
#[inline]
pub(crate) fn new(id: impl Into<String>) -> Self {
Self(id.into())
}
#[must_use]
#[inline]
pub(crate) fn as_str(&self) -> &str {
&self.0
}
}
impl std::fmt::Display for SessionId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(&self.0)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub(crate) struct SubscriptionKey(u64);
impl SubscriptionKey {
#[must_use]
#[inline]
pub(crate) const fn get(self) -> u64 {
self.0
}
#[cfg(feature = "test-util")]
#[must_use]
pub(crate) const fn from_raw(id: u64) -> Self {
Self(id)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct SubscriptionSpec {
pub(crate) group: String,
pub(crate) schema: String,
pub(crate) mode: SubscriptionMode,
pub(crate) data_adapter: Option<String>,
pub(crate) selector: Option<String>,
pub(crate) requested_buffer_size: Option<RequestedBufferSize>,
pub(crate) requested_max_frequency: Option<RequestedMaxFrequency>,
pub(crate) snapshot: Option<Snapshot>,
pub(crate) declared_items: Vec<String>,
pub(crate) declared_fields: Vec<String>,
}
impl SubscriptionSpec {
#[must_use]
pub(crate) fn new(
group: impl Into<String>,
schema: impl Into<String>,
mode: SubscriptionMode,
) -> Self {
Self {
group: group.into(),
schema: schema.into(),
mode,
data_adapter: None,
selector: None,
requested_buffer_size: None,
requested_max_frequency: None,
snapshot: None,
declared_items: Vec::new(),
declared_fields: Vec::new(),
}
}
#[must_use]
#[inline]
fn wants_snapshot(&self) -> bool {
matches!(self.snapshot, Some(Snapshot::On | Snapshot::Length(_)))
}
#[must_use]
fn new_manager(&self) -> SubscriptionManager {
SubscriptionManager::new(
self.mode,
self.wants_snapshot(),
self.declared_items.clone(),
self.declared_fields.clone(),
)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct OutgoingMessage {
pub(crate) message: String,
pub(crate) sequence: Option<SequenceName>,
pub(crate) msg_prog: Option<NonZeroU32>,
pub(crate) max_wait_millis: Option<u64>,
pub(crate) outcome: Option<bool>,
}
impl OutgoingMessage {
#[must_use]
pub(crate) fn fire_and_forget(message: impl Into<String>) -> Self {
Self {
message: message.into(),
sequence: None,
msg_prog: None,
max_wait_millis: None,
outcome: Some(false),
}
}
#[must_use]
pub(crate) fn numbered(message: impl Into<String>, msg_prog: NonZeroU32) -> Self {
Self {
message: message.into(),
sequence: None,
msg_prog: Some(msg_prog),
max_wait_millis: None,
outcome: None,
}
}
#[must_use = "builders do nothing unless the result is used"]
pub(crate) fn in_sequence(mut self, sequence: SequenceName) -> Self {
self.sequence = Some(sequence);
self
}
#[must_use]
fn identity(&self) -> (MessageSequence, Option<u64>) {
let sequence = match &self.sequence {
Some(name) => MessageSequence::Named(name.as_str().to_owned()),
None => MessageSequence::Unspecified,
};
(sequence, self.msg_prog.map(|prog| u64::from(prog.get())))
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ServerCause {
pub(crate) code: i64,
pub(crate) message: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum BindKind {
Created,
Recreated {
previous: Option<SessionId>,
},
Rebound,
Recovering {
requested_progressive: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct RecoveryOutcome {
pub(crate) requested: u64,
pub(crate) resumed_at: u64,
pub(crate) kind: RecoveryKind,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum RecoveryKind {
Exact,
Duplicated {
count: u64,
},
Gap {
missing: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum UnbindReason {
ForcedByClient {
expected_delay: Duration,
},
ContentLengthReached {
expected_delay: Duration,
},
PollCycleExpired {
expected_delay: Duration,
},
Looped {
expected_delay: Duration,
},
ConnectionFailed {
detail: String,
},
KeepaliveExpired {
budget: Duration,
},
Rejected {
cause: ServerCause,
},
ServerRefresh {
cause: ServerCause,
},
}
impl UnbindReason {
#[must_use]
const fn is_clean_loop(&self) -> bool {
matches!(
self,
Self::ForcedByClient { .. }
| Self::ContentLengthReached { .. }
| Self::PollCycleExpired { .. }
| Self::Looped { .. }
)
}
#[must_use]
const fn expected_delay(&self) -> Option<Duration> {
match self {
Self::ForcedByClient { expected_delay }
| Self::ContentLengthReached { expected_delay }
| Self::PollCycleExpired { expected_delay }
| Self::Looped { expected_delay } => Some(*expected_delay),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum SessionClosed {
ByClient {
destroy_confirmed: bool,
cause: Option<ServerCause>,
},
ByServer {
cause: ServerCause,
},
RetriesExhausted {
attempts: u32,
last: Option<UnbindReason>,
},
Internal {
reason: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct BoundInfo {
pub(crate) session_id: SessionId,
pub(crate) kind: BindKind,
pub(crate) keep_alive: Duration,
pub(crate) request_limit_bytes: u64,
pub(crate) control_link: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ResubscribedEntry {
pub(crate) key: SubscriptionKey,
pub(crate) subscription_id: NonZeroU32,
pub(crate) previously_active: bool,
pub(crate) snapshot_requested: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum SubscriptionOperation {
Subscribe,
Unsubscribe,
Reconfigure,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ControlTarget {
Session,
Subscription {
key: SubscriptionKey,
operation: SubscriptionOperation,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ControlOutcome {
Accepted,
Rejected {
cause: ServerCause,
},
NotSent {
reason: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum MessageResult {
Done {
response: String,
},
Failed {
cause: ServerCause,
},
Refused {
cause: ServerCause,
},
NotSent {
reason: String,
},
}
pub(crate) type SubscriptionOutcome = Result<SubscriptionEvent, SubscriptionError>;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum ServerAnnouncement {
Name(String),
ClientIp(String),
Clock {
elapsed_seconds: u64,
},
Bandwidth(Bandwidth),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) enum SessionEvent {
Bound(Box<BoundInfo>),
Recovered(RecoveryOutcome),
Resubscribed(Vec<ResubscribedEntry>),
Unbound {
reason: UnbindReason,
retry_in: Option<Duration>,
},
Message {
progressive: Option<u64>,
sequence: MessageSequence,
prog: Option<u64>,
result: MessageResult,
},
Subscription {
progressive: u64,
key: SubscriptionKey,
outcome: Box<SubscriptionOutcome>,
},
ControlResponse {
request_id: Option<RequestId>,
target: ControlTarget,
outcome: ControlOutcome,
},
ServerInfo(ServerAnnouncement),
Unparsed {
line: String,
error: ProtocolError,
},
Closed(SessionClosed),
}
#[derive(Debug)]
pub(crate) enum SessionCommand {
Subscribe {
key: SubscriptionKey,
spec: Box<SubscriptionSpec>,
},
Unsubscribe {
key: SubscriptionKey,
},
Reconfigure {
key: SubscriptionKey,
max_frequency: MaxFrequencyLimit,
},
SendMessage {
message: Box<OutgoingMessage>,
},
ForceRebind {
close_socket: Option<bool>,
},
Destroy {
cause: Option<DestroyCause>,
},
Shutdown,
}
#[derive(Debug, Clone)]
pub(crate) struct SessionHandle {
commands: mpsc::Sender<SessionCommand>,
next_key: Arc<AtomicU64>,
stop: Arc<watch::Sender<bool>>,
}
impl SessionHandle {
pub(crate) fn allocate_key(&self) -> Result<SubscriptionKey, SessionError> {
self.next_key
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |key| {
key.checked_add(1)
})
.map(SubscriptionKey)
.map_err(|_| SessionError::Exhausted {
what: "subscription key",
})
}
pub(crate) async fn subscribe_with_key(
&self,
key: SubscriptionKey,
spec: SubscriptionSpec,
) -> Result<(), SessionError> {
self.send(SessionCommand::Subscribe {
key,
spec: Box::new(spec),
})
.await
}
pub(crate) async fn subscribe(
&self,
spec: SubscriptionSpec,
) -> Result<SubscriptionKey, SessionError> {
let key = self.allocate_key()?;
self.subscribe_with_key(key, spec).await?;
Ok(key)
}
pub(crate) fn stop(&self) {
let _ = self.stop.send(true);
}
pub(crate) fn stop_signal(&self) -> watch::Receiver<bool> {
self.stop.subscribe()
}
pub(crate) async fn unsubscribe(&self, key: SubscriptionKey) -> Result<(), SessionError> {
self.send(SessionCommand::Unsubscribe { key }).await
}
pub(crate) async fn reconfigure(
&self,
key: SubscriptionKey,
max_frequency: MaxFrequencyLimit,
) -> Result<(), SessionError> {
self.send(SessionCommand::Reconfigure { key, max_frequency })
.await
}
pub(crate) async fn send_message(&self, message: OutgoingMessage) -> Result<(), SessionError> {
self.send(SessionCommand::SendMessage {
message: Box::new(message),
})
.await
}
pub(crate) async fn force_rebind(
&self,
close_socket: Option<bool>,
) -> Result<(), SessionError> {
self.send(SessionCommand::ForceRebind { close_socket })
.await
}
pub(crate) async fn destroy(&self, cause: Option<DestroyCause>) -> Result<(), SessionError> {
self.send(SessionCommand::Destroy { cause }).await
}
pub(crate) async fn shutdown(&self) -> Result<(), SessionError> {
self.send(SessionCommand::Shutdown).await
}
async fn send(&self, command: SessionCommand) -> Result<(), SessionError> {
self.commands
.send(command)
.await
.map_err(|_| SessionError::Stopped)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Recovery {
RecreateSession,
RetryLater,
Fatal,
}
#[must_use]
pub(crate) fn classify(code: i64) -> Recovery {
match code {
4 => Recovery::RecreateSession,
20 => Recovery::RecreateSession,
48 => Recovery::RecreateSession,
5 | 6 | 10 => Recovery::RetryLater,
7..=9 => Recovery::RetryLater,
33 | 34 => Recovery::RetryLater,
21 => Recovery::Fatal,
_ => Recovery::Fatal,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SubscriptionState {
Pending,
Active,
Removing,
}
#[derive(Debug, Clone)]
struct RegistryEntry {
spec: SubscriptionSpec,
wire: Option<NonZeroU32>,
state: SubscriptionState,
was_active: bool,
manager: SubscriptionManager,
}
#[derive(Debug, Default)]
struct Registry {
entries: BTreeMap<SubscriptionKey, RegistryEntry>,
by_sub_id: BTreeMap<u32, SubscriptionKey>,
next_sub_id: u32,
}
impl Registry {
fn new() -> Self {
Self {
entries: BTreeMap::new(),
by_sub_id: BTreeMap::new(),
next_sub_id: 1,
}
}
fn reset_for_new_session(&mut self) {
self.entries
.retain(|_, entry| entry.state != SubscriptionState::Removing);
for entry in self.entries.values_mut() {
entry.wire = None;
entry.state = SubscriptionState::Pending;
entry.manager = entry.spec.new_manager();
}
self.by_sub_id.clear();
self.next_sub_id = 1;
}
fn allocate(&mut self) -> Option<NonZeroU32> {
let id = NonZeroU32::new(self.next_sub_id)?;
self.next_sub_id = self.next_sub_id.checked_add(1)?;
Some(id)
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum PendingKind {
Subscribe(SubscriptionKey),
Unsubscribe(SubscriptionKey),
Reconfigure {
key: SubscriptionKey,
max_frequency: MaxFrequencyLimit,
},
SendMessage {
sequence: MessageSequence,
prog: Option<u64>,
},
ForceRebind,
Destroy,
}
impl PendingKind {
#[must_use]
fn target(&self) -> ControlTarget {
match self {
Self::Subscribe(key) => ControlTarget::Subscription {
key: *key,
operation: SubscriptionOperation::Subscribe,
},
Self::Unsubscribe(key) => ControlTarget::Subscription {
key: *key,
operation: SubscriptionOperation::Unsubscribe,
},
Self::Reconfigure { key, .. } => ControlTarget::Subscription {
key: *key,
operation: SubscriptionOperation::Reconfigure,
},
Self::SendMessage { .. } | Self::ForceRebind | Self::Destroy => ControlTarget::Session,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Pending {
kind: PendingKind,
generation: u64,
}
#[must_use]
fn as_requested_frequency(limit: MaxFrequencyLimit) -> RequestedMaxFrequency {
match limit {
MaxFrequencyLimit::Unlimited => RequestedMaxFrequency::Unlimited,
MaxFrequencyLimit::Limited(number) => RequestedMaxFrequency::Limited(number),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum OpenIntent {
Create,
Rebind,
Recover {
requested: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum StreamState {
Establishing,
Bound,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct Reopen {
intent: OpenIntent,
delay: Duration,
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum Flow {
Continue,
Unbind(UnbindReason),
Closed(SessionClosed),
}
enum Woke {
Line(Option<Result<String, TransportError>>),
Command(Option<SessionCommand>),
Timer,
Stop,
}
enum Opened {
Established,
Failed {
detail: String,
},
Stopped,
}
pub(crate) fn connect_configured(
config: ClientConfig,
) -> Result<
(
SessionDriver<AnyTransport>,
SessionHandle,
mpsc::Receiver<SessionEvent>,
),
crate::error::Error,
> {
let choice = config.transport();
let address = config.address().as_str().to_owned();
let transport = match choice {
crate::config::Transport::WebSocket => {
AnyTransport::WebSocket(WsTransport::try_new(address)?)
}
crate::config::Transport::HttpStreaming => {
AnyTransport::Http(HttpTransport::try_new(address, HttpMode::Streaming)?)
}
crate::config::Transport::HttpPolling => {
AnyTransport::Http(HttpTransport::try_new(address, HttpMode::Polling)?)
}
};
let mut options = SessionOptions::from_client_config(config);
if matches!(choice, crate::config::Transport::HttpPolling) {
let idle_millis = match options.connection {
ConnectionMode::Streaming {
keepalive_millis, ..
} => keepalive_millis,
ConnectionMode::Polling { idle_millis, .. } => idle_millis,
}
.unwrap_or(DEFAULT_IDLE_MILLIS);
options = options.with_connection(ConnectionMode::Polling {
polling_millis: DEFAULT_POLLING_MILLIS,
idle_millis: Some(idle_millis),
});
}
connect(transport, options).map_err(|error| crate::error::Error::Internal {
reason: error.to_string(),
})
}
const DEFAULT_POLLING_MILLIS: u64 = 5_000;
const DEFAULT_IDLE_MILLIS: u64 = 15_000;
pub(crate) fn connect<T: Transport>(
transport: T,
options: SessionOptions,
) -> Result<
(
SessionDriver<T>,
SessionHandle,
mpsc::Receiver<SessionEvent>,
),
SessionError,
> {
let properties = transport.properties();
if properties.is_polling != options.is_polling() {
return Err(SessionError::Configuration {
reason: format!(
"transport declares is_polling={} but the connection mode asks for polling={}",
properties.is_polling,
options.is_polling()
),
});
}
let (command_tx, command_rx) = mpsc::channel(options.command_capacity.get());
let (event_tx, event_rx) = mpsc::channel(options.event_capacity.get());
let (stop_tx, stop_rx) = watch::channel(false);
let handle = SessionHandle {
commands: command_tx,
next_key: Arc::new(AtomicU64::new(1)),
stop: Arc::new(stop_tx),
};
let driver = SessionDriver::new(
transport, properties, options, command_rx, event_tx, stop_rx,
);
Ok((driver, handle, event_rx))
}
async fn stopped(stop: &mut watch::Receiver<bool>) {
loop {
if *stop.borrow_and_update() {
return;
}
if stop.changed().await.is_err() {
return;
}
}
}
const MAX_LINES_PER_TURN: u32 = 32;
const MAX_HEARTBEAT_FAILURES: u32 = 3;
#[derive(Debug)]
pub(crate) struct SessionDriver<T: Transport> {
transport: T,
properties: TransportProperties,
options: SessionOptions,
commands: mpsc::Receiver<SessionCommand>,
events: mpsc::Sender<SessionEvent>,
stop: watch::Receiver<bool>,
stopping: bool,
stream_state: StreamState,
intent: OpenIntent,
session: Option<SessionId>,
previous_session: Option<SessionId>,
progressive: u64,
skip_remaining: u64,
prog_honoured: bool,
registry: Registry,
pending: HashMap<String, Pending>,
generation: u64,
next_request_id: u64,
force_rebind_outstanding: bool,
destroy_requested: bool,
refresh_streak: u32,
heartbeat_failures: u32,
negotiated_keep_alive: Option<Duration>,
liveness: Liveness,
backoff: Backoff,
}
impl<T: Transport> SessionDriver<T> {
fn new(
transport: T,
properties: TransportProperties,
options: SessionOptions,
commands: mpsc::Receiver<SessionCommand>,
events: mpsc::Sender<SessionEvent>,
stop: watch::Receiver<bool>,
) -> Self {
let liveness = Liveness::new(
options.keepalive_slack,
options.heartbeat_interval(),
Instant::now(),
);
let backoff = Backoff::new(options.backoff);
Self {
transport,
properties,
options,
commands,
events,
stop,
stopping: false,
stream_state: StreamState::Establishing,
intent: OpenIntent::Create,
session: None,
previous_session: None,
progressive: 0,
skip_remaining: 0,
prog_honoured: false,
registry: Registry::new(),
pending: HashMap::new(),
generation: 1,
next_request_id: 1,
force_rebind_outstanding: false,
destroy_requested: false,
refresh_streak: 0,
heartbeat_failures: 0,
negotiated_keep_alive: None,
liveness,
backoff,
}
}
#[must_use]
fn establishment_budget(&self) -> Duration {
let ConnectionMode::Polling { idle_millis, .. } = self.options.connection else {
return self.options.open_timeout;
};
let requested = Duration::from_millis(idle_millis.unwrap_or(0));
let idle = self
.negotiated_keep_alive
.unwrap_or(requested)
.max(requested);
self.options
.open_timeout
.checked_add(idle)
.unwrap_or(Duration::MAX)
}
#[must_use]
#[inline]
pub(crate) const fn progressive(&self) -> u64 {
self.progressive
}
#[must_use]
#[inline]
pub(crate) const fn session_id(&self) -> Option<&SessionId> {
self.session.as_ref()
}
pub(crate) async fn run(mut self) -> SessionClosed {
let closed = self.drive().await;
if let Err(error) = self.transport.close().await {
tracing::warn!(%error, "transport did not close cleanly");
}
self.emit(SessionEvent::Closed(closed.clone())).await;
tracing::info!(?closed, "session ended");
closed
}
async fn drive(&mut self) -> SessionClosed {
let mut next = Reopen {
intent: OpenIntent::Create,
delay: Duration::ZERO,
};
loop {
if let Some(closed) = self.wait_before_reopen(next.delay).await {
return closed;
}
let request = match self.build_open(&next.intent) {
Ok(request) => request,
Err(error) => {
return SessionClosed::Internal {
reason: error.to_string(),
};
}
};
self.intent = next.intent.clone();
self.stream_state = StreamState::Establishing;
self.skip_remaining = 0;
self.prog_honoured = false;
let flow = match self.open_stream(request).await {
Opened::Established => {
self.liveness
.on_stream_opened(self.establishment_budget(), Instant::now());
self.pump().await
}
Opened::Failed { detail } => {
Flow::Unbind(UnbindReason::ConnectionFailed { detail })
}
Opened::Stopped => {
return SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
};
}
};
match flow {
Flow::Closed(closed) => return closed,
Flow::Continue | Flow::Unbind(_) => {}
}
let reason = match flow {
Flow::Unbind(reason) => reason,
_ => UnbindReason::ConnectionFailed {
detail: "stream ended without a reason".to_owned(),
},
};
self.liveness.on_stream_closed();
self.stream_state = StreamState::Establishing;
if !reason.is_clean_loop()
&& let Err(error) = self.transport.close().await
{
tracing::debug!(%error, "closing the failed stream connection");
}
match self.plan_after(&reason) {
Some(reopen) => {
self.emit(SessionEvent::Unbound {
reason: reason.clone(),
retry_in: Some(reopen.delay),
})
.await;
if self.stopping {
return SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
};
}
tracing::debug!(?reason, delay = ?reopen.delay, intent = ?reopen.intent, "reopening");
next = reopen;
}
None => {
self.emit(SessionEvent::Unbound {
reason: reason.clone(),
retry_in: None,
})
.await;
return SessionClosed::RetriesExhausted {
attempts: self.backoff.attempts(),
last: Some(reason),
};
}
}
}
}
async fn open_stream(&mut self, request: StreamOpen) -> Opened {
let budget = self.establishment_budget();
let transport = &mut self.transport;
let stop = &mut self.stop;
tokio::select! {
biased;
result = tokio::time::timeout(budget, transport.open_stream(request)) => match result {
Ok(Ok(())) => Opened::Established,
Ok(Err(error)) => Opened::Failed { detail: error.to_string() },
Err(_elapsed) => {
tracing::warn!(?budget, "the stream connection could not be established in time");
Opened::Failed {
detail: format!("the stream connection was not established within {budget:?}"),
}
}
},
() = stopped(stop) => Opened::Stopped,
}
}
async fn wait_before_reopen(&mut self, delay: Duration) -> Option<SessionClosed> {
if delay.is_zero() {
return None;
}
let deadline = Instant::now()
.checked_add(delay)
.unwrap_or_else(Instant::now);
loop {
let woke = tokio::select! {
biased;
command = self.commands.recv() => Woke::Command(command),
() = stopped(&mut self.stop) => Woke::Stop,
() = tokio::time::sleep_until(deadline) => Woke::Timer,
};
match woke {
Woke::Timer => return None,
Woke::Stop => {
return Some(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
});
}
Woke::Command(None) => {
return Some(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
});
}
Woke::Command(Some(command)) => match self.on_command(command).await {
Flow::Closed(closed) => return Some(closed),
Flow::Continue | Flow::Unbind(_) => {
if self.stopping {
return Some(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
});
}
}
},
Woke::Line(_) => {}
}
}
}
fn build_open(&mut self, intent: &OpenIntent) -> Result<StreamOpen, ProtocolError> {
let session = match (&self.session, intent) {
(Some(session), OpenIntent::Rebind | OpenIntent::Recover { .. }) => Some(session),
_ => None,
};
match (session, intent) {
(Some(session), OpenIntent::Recover { requested }) => {
Ok(StreamOpen::Bind(Box::new(BindSession {
session: Some(session.as_str().to_owned()),
recovery_from: Some(*requested),
content_length: self.options.content_length,
connection: self.options.connection,
reduce_head: None,
})))
}
(Some(session), OpenIntent::Rebind) => Ok(StreamOpen::Bind(Box::new(BindSession {
session: Some(session.as_str().to_owned()),
recovery_from: None,
content_length: self.options.content_length,
connection: self.options.connection,
reduce_head: None,
}))),
_ => {
let credentials = &self.options.credentials;
Ok(StreamOpen::Create(Box::new(CreateSession {
user: credentials.user.clone(),
password: credentials.password.clone(),
adapter_set: credentials.adapter_set.clone(),
requested_max_bandwidth: None,
content_length: self.options.content_length,
connection: self.options.connection,
reduce_head: None,
ttl: None,
})))
}
}
}
fn plan_after(&mut self, reason: &UnbindReason) -> Option<Reopen> {
if reason.is_clean_loop() {
self.backoff.reset();
self.refresh_streak = 0;
return Some(Reopen {
intent: OpenIntent::Rebind,
delay: reason.expected_delay().unwrap_or(Duration::ZERO),
});
}
if !matches!(reason, UnbindReason::ServerRefresh { .. }) {
self.refresh_streak = 0;
}
match reason {
UnbindReason::ConnectionFailed { .. } | UnbindReason::KeepaliveExpired { .. } => {
let delay = self.backoff.next_delay()?;
let intent = if self.session.is_some() {
OpenIntent::Recover {
requested: self.progressive,
}
} else {
OpenIntent::Create
};
Some(Reopen { intent, delay })
}
UnbindReason::Rejected { .. } => {
let delay = self.backoff.next_delay()?;
Some(Reopen {
intent: OpenIntent::Create,
delay,
})
}
UnbindReason::ServerRefresh { .. } => {
self.refresh_streak = self.refresh_streak.saturating_add(1);
if self.refresh_streak <= 1 {
self.backoff.reset();
return Some(Reopen {
intent: OpenIntent::Create,
delay: Duration::ZERO,
});
}
if let Some(max) = self.options.backoff.max_attempts
&& self.refresh_streak > max.get()
{
tracing::warn!(
streak = self.refresh_streak,
"the server kept asking for a fresh session; giving up"
);
return None;
}
let delay = self.backoff.next_delay()?;
tracing::warn!(
streak = self.refresh_streak,
?delay,
"repeated session refresh; delaying the next attempt"
);
Some(Reopen {
intent: OpenIntent::Create,
delay,
})
}
UnbindReason::ForcedByClient { .. }
| UnbindReason::ContentLengthReached { .. }
| UnbindReason::PollCycleExpired { .. }
| UnbindReason::Looped { .. } => Some(Reopen {
intent: OpenIntent::Rebind,
delay: Duration::ZERO,
}),
}
}
async fn pump(&mut self) -> Flow {
let mut consecutive_lines: u32 = 0;
loop {
let deadline = self.liveness.next_deadline();
let woke = if consecutive_lines >= MAX_LINES_PER_TURN {
consecutive_lines = 0;
tokio::select! {
biased;
command = self.commands.recv() => Woke::Command(command),
() = stopped(&mut self.stop) => Woke::Stop,
() = sleep_until_option(deadline) => Woke::Timer,
line = self.transport.next_line() => Woke::Line(line),
}
} else {
tokio::select! {
biased;
line = self.transport.next_line() => Woke::Line(line),
command = self.commands.recv() => Woke::Command(command),
() = stopped(&mut self.stop) => Woke::Stop,
() = sleep_until_option(deadline) => Woke::Timer,
}
};
consecutive_lines = if matches!(woke, Woke::Line(_)) {
consecutive_lines.saturating_add(1)
} else {
0
};
let flow = match woke {
Woke::Line(Some(Ok(line))) => self.on_line(line).await,
Woke::Line(Some(Err(error))) => Flow::Unbind(UnbindReason::ConnectionFailed {
detail: error.to_string(),
}),
Woke::Line(None) => self.on_stream_end(),
Woke::Command(Some(command)) => self.on_command(command).await,
Woke::Command(None) | Woke::Stop => Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
}),
Woke::Timer => self.on_timer().await,
};
if self.stopping {
return Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
});
}
match flow {
Flow::Continue => {}
other => return other,
}
}
}
fn on_stream_end(&mut self) -> Flow {
if self.destroy_requested {
return Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
});
}
if self.properties.is_polling {
return Flow::Unbind(UnbindReason::PollCycleExpired {
expected_delay: Duration::ZERO,
});
}
Flow::Unbind(UnbindReason::ConnectionFailed {
detail: "stream connection closed by the peer".to_owned(),
})
}
async fn on_timer(&mut self) -> Flow {
match self.liveness.due(Instant::now()) {
LivenessAction::Idle => Flow::Continue,
LivenessAction::SendHeartbeat => self.send_heartbeat().await,
LivenessAction::InboundStalled { budget } => match self.stream_state {
StreamState::Establishing => Flow::Unbind(UnbindReason::ConnectionFailed {
detail: format!("no response to the stream-opening request within {budget:?}"),
}),
StreamState::Bound => Flow::Unbind(UnbindReason::KeepaliveExpired { budget }),
},
}
}
async fn on_line(&mut self, line: String) -> Flow {
self.liveness.on_inbound(Instant::now());
match parse_line(&line) {
Ok(notification) => self.on_notification(notification).await,
Err(error) => {
tracing::warn!(%error, "unrecognized line preserved for the caller");
self.emit(SessionEvent::Unparsed { line, error }).await;
Flow::Continue
}
}
}
async fn on_notification(&mut self, notification: Notification) -> Flow {
match notification {
Notification::ConnectionOk {
session_id,
request_limit_bytes,
keep_alive_millis,
control_link,
} => {
self.on_conok(
SessionId::new(session_id),
request_limit_bytes,
Duration::from_millis(keep_alive_millis),
control_link,
)
.await;
Flow::Continue
}
Notification::ConnectionError { code, message } => {
let cause = ServerCause { code, message };
match classify(code) {
Recovery::Fatal => {
tracing::warn!(code, "session refused with a code that admits no retry");
Flow::Closed(SessionClosed::ByServer { cause })
}
Recovery::RecreateSession | Recovery::RetryLater => {
self.previous_session = self.session.take();
Flow::Unbind(UnbindReason::Rejected { cause })
}
}
}
Notification::End {
cause_code,
cause_message,
} => self.on_end(cause_code, cause_message),
Notification::Loop {
expected_delay_millis,
} => Flow::Unbind(self.loop_reason(Duration::from_millis(expected_delay_millis))),
Notification::Progressive { progressive } => {
self.on_prog(progressive).await;
Flow::Continue
}
Notification::RequestOk { request_id } => {
self.on_request_ok(request_id).await;
Flow::Continue
}
Notification::RequestError {
request_id,
code,
message,
} => {
self.on_request_error(request_id, ServerCause { code, message })
.await;
Flow::Continue
}
Notification::Error { code, message } => {
tracing::warn!(code, "server rejected a request it could not correlate");
self.emit(SessionEvent::ControlResponse {
request_id: None,
target: ControlTarget::Session,
outcome: ControlOutcome::Rejected {
cause: ServerCause { code, message },
},
})
.await;
Flow::Continue
}
Notification::Probe => {
tracing::trace!("probe");
Flow::Continue
}
Notification::NoOp { .. } => Flow::Continue,
Notification::WsOk => {
tracing::debug!("websocket establishment check acknowledged");
Flow::Continue
}
Notification::Sync { elapsed_seconds } => {
self.emit(SessionEvent::ServerInfo(ServerAnnouncement::Clock {
elapsed_seconds,
}))
.await;
Flow::Continue
}
Notification::ClientIp { address } => {
self.emit(SessionEvent::ServerInfo(ServerAnnouncement::ClientIp(
address,
)))
.await;
Flow::Continue
}
Notification::ServerName { name } => {
self.emit(SessionEvent::ServerInfo(ServerAnnouncement::Name(name)))
.await;
Flow::Continue
}
Notification::ConstraintsChanged { bandwidth } => {
self.emit(SessionEvent::ServerInfo(ServerAnnouncement::Bandwidth(
bandwidth,
)))
.await;
Flow::Continue
}
data => {
self.on_data(data).await;
Flow::Continue
}
}
}
async fn on_conok(
&mut self,
session_id: SessionId,
request_limit_bytes: u64,
keep_alive: Duration,
control_link: Option<String>,
) {
self.stream_state = StreamState::Bound;
self.liveness.on_bound(keep_alive, Instant::now());
self.negotiated_keep_alive = Some(keep_alive);
self.backoff.reset();
self.transport.set_control_link(control_link.as_deref());
let previous = self.session.take().or_else(|| self.previous_session.take());
let same_session = previous.as_ref() == Some(&session_id);
self.session = Some(session_id.clone());
self.previous_session = None;
let kind = match (&self.intent, same_session) {
(OpenIntent::Rebind, true) => BindKind::Rebound,
(OpenIntent::Recover { requested }, true) => BindKind::Recovering {
requested_progressive: *requested,
},
(OpenIntent::Rebind | OpenIntent::Recover { .. }, false) => {
tracing::warn!(
old = %previous.as_ref().map_or("<none>", SessionId::as_str),
new = %session_id,
"bind returned a different session id; treating it as a new session"
);
BindKind::Recreated { previous }
}
(OpenIntent::Create, _) => match previous {
Some(previous) => BindKind::Recreated {
previous: Some(previous),
},
None => BindKind::Created,
},
};
let fresh = matches!(kind, BindKind::Created | BindKind::Recreated { .. });
if fresh {
self.progressive = 0;
self.skip_remaining = 0;
self.registry.reset_for_new_session();
self.generation = self.generation.saturating_add(1);
self.retire_stale_pending().await;
}
tracing::info!(session = %session_id, ?kind, ?keep_alive, "session bound");
self.emit(SessionEvent::Bound(Box::new(BoundInfo {
session_id,
kind,
keep_alive,
request_limit_bytes,
control_link,
})))
.await;
self.issue_unbound_subscriptions().await;
}
async fn retire_stale_pending(&mut self) {
let generation = self.generation;
let stale: Vec<(String, PendingKind)> = self
.pending
.extract_if(|_, pending| pending.generation != generation)
.map(|(id, pending)| (id, pending.kind))
.collect();
if stale.is_empty() {
return;
}
tracing::debug!(
count = stale.len(),
"retiring control requests issued by a session that was replaced"
);
for (request_id, kind) in stale {
let target = kind.target();
if matches!(kind, PendingKind::ForceRebind) {
self.force_rebind_outstanding = false;
}
if let PendingKind::SendMessage { sequence, prog } = kind {
self.emit(SessionEvent::Message {
progressive: None,
sequence,
prog,
result: MessageResult::NotSent {
reason: "the session the message was sent on was replaced".to_owned(),
},
})
.await;
}
self.emit(SessionEvent::ControlResponse {
request_id: RequestId::try_new(request_id).ok(),
target,
outcome: ControlOutcome::NotSent {
reason: "the session the request was issued on was replaced".to_owned(),
},
})
.await;
}
}
fn on_end(&mut self, cause_code: i64, cause_message: String) -> Flow {
let cause = ServerCause {
code: cause_code,
message: cause_message,
};
if self.destroy_requested {
return Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: true,
cause: Some(cause),
});
}
match classify(cause_code) {
Recovery::RecreateSession => {
tracing::info!(code = cause_code, "server asked for a fresh session");
self.previous_session = self.session.take();
Flow::Unbind(UnbindReason::ServerRefresh { cause })
}
Recovery::RetryLater => {
self.previous_session = self.session.take();
Flow::Unbind(UnbindReason::Rejected { cause })
}
Recovery::Fatal => {
tracing::warn!(code = cause_code, "server ended the session");
Flow::Closed(SessionClosed::ByServer { cause })
}
}
}
fn loop_reason(&mut self, expected_delay: Duration) -> UnbindReason {
if self.force_rebind_outstanding {
self.force_rebind_outstanding = false;
return UnbindReason::ForcedByClient { expected_delay };
}
if self.properties.is_polling {
return UnbindReason::PollCycleExpired { expected_delay };
}
if self.properties.ends_on_content_length {
return UnbindReason::ContentLengthReached { expected_delay };
}
UnbindReason::Looped { expected_delay }
}
async fn on_prog(&mut self, resumed_at: u64) {
let requested = match self.intent {
OpenIntent::Recover { requested } if !self.prog_honoured => requested,
OpenIntent::Recover { .. } => {
tracing::warn!(
resumed_at,
"ignoring a second PROG on one stream connection"
);
return;
}
_ => {
tracing::warn!(
resumed_at,
"ignoring a PROG that answers no recovery request"
);
return;
}
};
self.prog_honoured = true;
let kind = match resumed_at.cmp(&requested) {
std::cmp::Ordering::Equal => RecoveryKind::Exact,
std::cmp::Ordering::Less => match requested.checked_sub(resumed_at) {
Some(count) => {
self.skip_remaining = count;
RecoveryKind::Duplicated { count }
}
None => RecoveryKind::Exact,
},
std::cmp::Ordering::Greater => match resumed_at.checked_sub(requested) {
Some(missing) => {
self.progressive = resumed_at;
tracing::warn!(
requested,
resumed_at,
missing,
"recovery resumed past the requested point"
);
RecoveryKind::Gap { missing }
}
None => RecoveryKind::Exact,
},
};
self.emit(SessionEvent::Recovered(RecoveryOutcome {
requested,
resumed_at,
kind,
}))
.await;
}
async fn on_data(&mut self, notification: Notification) {
self.refresh_streak = 0;
if self.skip_remaining > 0 {
if let Some(remaining) = self.skip_remaining.checked_sub(1) {
self.skip_remaining = remaining;
}
tracing::trace!(
remaining = self.skip_remaining,
"discarding a duplicate re-delivered by recovery"
);
return;
}
self.progressive = match self.progressive.checked_add(1) {
Some(next) => next,
None => {
tracing::error!(
"data-notification counter overflowed; recovery is no longer exact"
);
self.progressive
}
};
let progressive = self.progressive;
match notification {
Notification::MessageDone {
sequence,
prog,
response,
} => {
self.emit(SessionEvent::Message {
progressive: Some(progressive),
sequence,
prog: Some(prog),
result: MessageResult::Done { response },
})
.await;
}
Notification::MessageFailed {
sequence,
prog,
code,
message,
} => {
self.emit(SessionEvent::Message {
progressive: Some(progressive),
sequence,
prog: Some(prog),
result: MessageResult::Failed {
cause: ServerCause { code, message },
},
})
.await;
}
notification => {
if let Some(event) = self.interpret_for_subscription(progressive, ¬ification) {
self.emit(event).await;
}
}
}
}
fn interpret_for_subscription(
&mut self,
progressive: u64,
notification: &Notification,
) -> Option<SessionEvent> {
let Some((sub_id, key)) = self.resolve_subscription(notification) else {
tracing::debug!(?notification, "unattributed data notification");
return None;
};
let Some(entry) = self.registry.entries.get_mut(&key) else {
tracing::debug!(
key = key.get(),
"a notification named a retired subscription"
);
return None;
};
let outcome = entry.manager.handle(notification);
match notification {
Notification::SubscriptionOk { .. } | Notification::SubscriptionCommandOk { .. } => {
entry.state = SubscriptionState::Active;
entry.was_active = true;
}
Notification::Unsubscribed { .. } => {
self.registry.entries.remove(&key);
self.registry.by_sub_id.remove(&sub_id);
}
_ => {}
}
match outcome {
Ok(None) => None,
Ok(Some(event)) => Some(SessionEvent::Subscription {
progressive,
key,
outcome: Box::new(Ok(event)),
}),
Err(error) => Some(SessionEvent::Subscription {
progressive,
key,
outcome: Box::new(Err(error)),
}),
}
}
fn resolve_subscription(&self, notification: &Notification) -> Option<(u32, SubscriptionKey)> {
let sub_id = match notification {
Notification::Update {
subscription_id, ..
}
| Notification::SubscriptionOk {
subscription_id, ..
}
| Notification::SubscriptionCommandOk {
subscription_id, ..
}
| Notification::Unsubscribed { subscription_id }
| Notification::EndOfSnapshot {
subscription_id, ..
}
| Notification::ClearSnapshot {
subscription_id, ..
}
| Notification::Overflow {
subscription_id, ..
}
| Notification::SubscriptionReconfigured {
subscription_id, ..
} => *subscription_id,
_ => return None,
};
let key = self.registry.by_sub_id.get(&sub_id).copied()?;
Some((sub_id, key))
}
fn take_pending(&mut self, request_id: &str) -> Option<PendingKind> {
let pending = self.pending.remove(request_id)?;
if pending.generation == self.generation {
return Some(pending.kind);
}
tracing::debug!(
request_id,
issued_on = pending.generation,
current = self.generation,
"ignoring a late response from a session that was replaced"
);
None
}
async fn on_request_ok(&mut self, request_id: Option<String>) {
let Some(request_id) = request_id else {
tracing::trace!("heartbeat acknowledged");
return;
};
let target = match self.take_pending(&request_id) {
Some(kind) => {
tracing::debug!(request_id, ?kind, "control request accepted");
let target = kind.target();
if let PendingKind::Reconfigure { key, max_frequency } = kind
&& let Some(entry) = self.registry.entries.get_mut(&key)
{
entry.spec.requested_max_frequency =
Some(as_requested_frequency(max_frequency));
tracing::debug!(
key = key.get(),
"frequency reconfiguration accepted and committed to the desired state"
);
}
target
}
None => {
tracing::debug!(request_id, "response to an unknown or stale request");
ControlTarget::Session
}
};
let outcome = ControlOutcome::Accepted;
match RequestId::try_new(request_id.clone()) {
Ok(id) => {
self.emit(SessionEvent::ControlResponse {
request_id: Some(id),
target,
outcome,
})
.await;
}
Err(error) => {
tracing::warn!(%error, "server echoed a request id this client cannot represent");
self.emit(SessionEvent::ControlResponse {
request_id: None,
target,
outcome,
})
.await;
}
}
}
async fn on_request_error(&mut self, request_id: String, cause: ServerCause) {
let pending = self.take_pending(&request_id);
let target = pending
.as_ref()
.map_or(ControlTarget::Session, PendingKind::target);
if let Some(kind) = pending {
match kind {
PendingKind::Subscribe(key) => {
if let Some(entry) = self.registry.entries.remove(&key)
&& let Some(sub_id) = entry.wire
{
self.registry.by_sub_id.remove(&sub_id.get());
}
tracing::warn!(
key = key.get(),
code = cause.code,
"subscription refused by the server"
);
}
PendingKind::Unsubscribe(key) => {
if let Some(entry) = self.registry.entries.get_mut(&key) {
entry.state = if entry.was_active {
SubscriptionState::Active
} else {
SubscriptionState::Pending
};
}
}
PendingKind::Reconfigure { .. } => {}
PendingKind::SendMessage { sequence, prog } => {
self.emit(SessionEvent::Message {
progressive: None,
sequence,
prog,
result: MessageResult::Refused {
cause: cause.clone(),
},
})
.await;
}
PendingKind::ForceRebind => self.force_rebind_outstanding = false,
PendingKind::Destroy => self.destroy_requested = false,
}
}
let outcome = ControlOutcome::Rejected { cause };
let id = RequestId::try_new(request_id).ok();
self.emit(SessionEvent::ControlResponse {
request_id: id,
target,
outcome,
})
.await;
}
async fn on_command(&mut self, command: SessionCommand) -> Flow {
match command {
SessionCommand::Subscribe { key, spec } => {
let manager = spec.new_manager();
self.registry.entries.insert(
key,
RegistryEntry {
spec: *spec,
wire: None,
state: SubscriptionState::Pending,
was_active: false,
manager,
},
);
if self.stream_state == StreamState::Bound {
self.issue_subscription(key).await;
}
Flow::Continue
}
SessionCommand::Unsubscribe { key } => {
self.unsubscribe(key).await;
Flow::Continue
}
SessionCommand::Reconfigure { key, max_frequency } => {
self.reconfigure(key, max_frequency).await;
Flow::Continue
}
SessionCommand::SendMessage { message } => {
self.send_message(*message).await;
Flow::Continue
}
SessionCommand::ForceRebind { close_socket } => {
if self.stream_state != StreamState::Bound {
tracing::debug!("force_rebind ignored: no stream connection is bound");
return Flow::Continue;
}
let Some(session) = self.session.clone() else {
return Flow::Continue;
};
let Some(request_id) = self.next_request_id() else {
self.report_id_exhaustion(ControlTarget::Session).await;
return Flow::Continue;
};
self.force_rebind_outstanding = true;
let request = ForceRebind {
session: Some(session.as_str().to_owned()),
request_id: request_id.clone(),
polling_millis: None,
close_socket,
};
self.dispatch(&request, request_id, PendingKind::ForceRebind)
.await;
Flow::Continue
}
SessionCommand::Destroy { cause } => self.destroy(cause).await,
SessionCommand::Shutdown => Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
}),
}
}
async fn destroy(&mut self, cause: Option<DestroyCause>) -> Flow {
let Some(session) = self.session.clone() else {
return Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
});
};
let Some(request_id) = self.next_request_id() else {
self.report_id_exhaustion(ControlTarget::Session).await;
return Flow::Closed(SessionClosed::Internal {
reason: "the request-id space is exhausted".to_owned(),
});
};
self.destroy_requested = true;
let request = DestroySession {
session: Some(session.as_str().to_owned()),
request_id: request_id.clone(),
cause,
close_socket: None,
};
let sent = self
.dispatch(&request, request_id, PendingKind::Destroy)
.await;
if self.stream_state == StreamState::Bound && sent {
return Flow::Continue;
}
Flow::Closed(SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
})
}
async fn send_message(&mut self, message: OutgoingMessage) {
let (sequence, prog) = message.identity();
let session = match (self.stream_state, self.session.clone()) {
(StreamState::Bound, Some(session)) => session,
_ => {
tracing::debug!(?sequence, ?prog, "message not sent: no session is bound");
self.emit(SessionEvent::Message {
progressive: None,
sequence,
prog,
result: MessageResult::NotSent {
reason: "no session is bound to send the message on".to_owned(),
},
})
.await;
return;
}
};
let Some(request_id) = self.next_request_id() else {
self.report_id_exhaustion(ControlTarget::Session).await;
self.emit(SessionEvent::Message {
progressive: None,
sequence,
prog,
result: MessageResult::NotSent {
reason: "the request-id space is exhausted".to_owned(),
},
})
.await;
return;
};
let request = MessageSend {
session: Some(session.as_str().to_owned()),
request_id: request_id.clone(),
message: message.message,
sequence: message.sequence,
msg_prog: message.msg_prog,
max_wait_millis: message.max_wait_millis,
ack: None,
outcome: message.outcome,
};
let kind = PendingKind::SendMessage {
sequence: sequence.clone(),
prog,
};
if !self.dispatch(&request, request_id, kind).await {
self.emit(SessionEvent::Message {
progressive: None,
sequence,
prog,
result: MessageResult::NotSent {
reason: "the message request could not be sent".to_owned(),
},
})
.await;
}
}
async fn unsubscribe(&mut self, key: SubscriptionKey) {
let Some(entry) = self.registry.entries.get_mut(&key) else {
tracing::debug!(key = key.get(), "unsubscribe for an unknown subscription");
return;
};
let Some(sub_id) = entry.wire else {
self.registry.entries.remove(&key);
return;
};
entry.state = SubscriptionState::Removing;
let Some(session) = self.session.clone() else {
return;
};
let Some(request_id) = self.next_request_id() else {
self.report_id_exhaustion(ControlTarget::Subscription {
key,
operation: SubscriptionOperation::Unsubscribe,
})
.await;
return;
};
let request = Unsubscribe {
session: Some(session.as_str().to_owned()),
request_id: request_id.clone(),
subscription_id: sub_id,
ack: None,
};
self.dispatch(&request, request_id, PendingKind::Unsubscribe(key))
.await;
}
async fn reconfigure(&mut self, key: SubscriptionKey, max_frequency: MaxFrequencyLimit) {
let Some(entry) = self.registry.entries.get(&key) else {
return;
};
let Some(sub_id) = entry.wire else {
return;
};
let Some(session) = self.session.clone() else {
return;
};
let Some(request_id) = self.next_request_id() else {
self.report_id_exhaustion(ControlTarget::Subscription {
key,
operation: SubscriptionOperation::Reconfigure,
})
.await;
return;
};
let request = ReconfigureSubscription {
session: Some(session.as_str().to_owned()),
request_id: request_id.clone(),
subscription_id: sub_id,
requested_max_frequency: Some(max_frequency.clone()),
};
self.dispatch(
&request,
request_id,
PendingKind::Reconfigure { key, max_frequency },
)
.await;
}
async fn issue_unbound_subscriptions(&mut self) {
let keys: Vec<SubscriptionKey> = self
.registry
.entries
.iter()
.filter(|(_, entry)| entry.wire.is_none() && entry.state != SubscriptionState::Removing)
.map(|(key, _)| *key)
.collect();
if keys.is_empty() {
return;
}
let mut issued = Vec::with_capacity(keys.len());
for key in keys {
if let Some(entry) = self.issue_subscription(key).await {
issued.push(entry);
}
}
if !issued.is_empty() {
tracing::info!(count = issued.len(), "subscriptions re-established");
self.emit(SessionEvent::Resubscribed(issued)).await;
}
}
async fn issue_subscription(&mut self, key: SubscriptionKey) -> Option<ResubscribedEntry> {
let session = self.session.clone()?;
let sub_id = self.registry.allocate()?;
let entry = self.registry.entries.get_mut(&key)?;
entry.wire = Some(sub_id);
entry.state = SubscriptionState::Pending;
let spec = entry.spec.clone();
let previously_active = entry.was_active;
self.registry.by_sub_id.insert(sub_id.get(), key);
let Some(request_id) = self.next_request_id() else {
if let Some(entry) = self.registry.entries.get_mut(&key) {
entry.wire = None;
}
self.registry.by_sub_id.remove(&sub_id.get());
self.report_id_exhaustion(ControlTarget::Subscription {
key,
operation: SubscriptionOperation::Subscribe,
})
.await;
return None;
};
let request = Subscribe {
session: Some(session.as_str().to_owned()),
request_id: request_id.clone(),
subscription_id: sub_id,
group: spec.group.clone(),
schema: spec.schema.clone(),
mode: spec.mode,
data_adapter: spec.data_adapter.clone(),
selector: spec.selector.clone(),
requested_buffer_size: spec.requested_buffer_size,
requested_max_frequency: spec.requested_max_frequency.clone(),
snapshot: spec.snapshot,
ack: None,
};
if !self
.dispatch(&request, request_id, PendingKind::Subscribe(key))
.await
{
if let Some(entry) = self.registry.entries.get_mut(&key) {
entry.wire = None;
}
self.registry.by_sub_id.remove(&sub_id.get());
tracing::warn!(
key = key.get(),
"subscription could not be issued; it will be retried on the next bind"
);
return None;
}
Some(ResubscribedEntry {
key,
subscription_id: sub_id,
previously_active,
snapshot_requested: spec.wants_snapshot(),
})
}
fn next_request_id(&mut self) -> Option<RequestId> {
let next = self.next_request_id.checked_add(1)?;
let id = RequestId::from_progressive(self.next_request_id);
self.next_request_id = next;
Some(id)
}
async fn report_id_exhaustion(&mut self, target: ControlTarget) {
tracing::error!(
"the request-id space is exhausted; no further control request can be sent"
);
self.emit(SessionEvent::ControlResponse {
request_id: None,
target,
outcome: ControlOutcome::NotSent {
reason: "the request-id space is exhausted".to_owned(),
},
})
.await;
}
#[must_use]
const fn control_kind(&self) -> TransportKind {
if self.properties.control_shares_stream {
TransportKind::WebSocket
} else {
TransportKind::Http
}
}
async fn dispatch<R: TlcpRequest>(
&mut self,
request: &R,
request_id: RequestId,
kind: PendingKind,
) -> bool {
let target = kind.target();
let parameters = match request.encode_parameters(self.control_kind()) {
Ok(parameters) => parameters,
Err(error) => {
tracing::error!(%error, "cannot encode control request");
self.emit(SessionEvent::ControlResponse {
request_id: Some(request_id),
target,
outcome: ControlOutcome::NotSent {
reason: error.to_string(),
},
})
.await;
return false;
}
};
self.pending.insert(
request_id.as_str().to_owned(),
Pending {
kind,
generation: self.generation,
},
);
let encoded = EncodedRequest {
name: R::NAME,
path: R::PATH,
parameters,
};
match self.transport.send_control(encoded).await {
Ok(()) => {
self.liveness.on_outbound(Instant::now());
true
}
Err(error) => {
self.pending.remove(request_id.as_str());
tracing::warn!(%error, "control request could not be sent");
self.emit(SessionEvent::ControlResponse {
request_id: Some(request_id),
target,
outcome: ControlOutcome::NotSent {
reason: error.to_string(),
},
})
.await;
false
}
}
}
async fn send_heartbeat(&mut self) -> Flow {
let session = self.session.as_ref().map(|s| s.as_str().to_owned());
let request = Heartbeat { session };
let failure = match request.encode_parameters(self.control_kind()) {
Ok(parameters) => {
let encoded = EncodedRequest {
name: Heartbeat::NAME,
path: Heartbeat::PATH,
parameters,
};
match self.transport.send_control(encoded).await {
Ok(()) => None,
Err(error) => Some(error.to_string()),
}
}
Err(error) => Some(error.to_string()),
};
self.liveness.on_outbound(Instant::now());
let Some(detail) = failure else {
self.heartbeat_failures = 0;
return Flow::Continue;
};
self.heartbeat_failures = self.heartbeat_failures.saturating_add(1);
tracing::warn!(
detail,
attempts = self.heartbeat_failures,
"heartbeat could not be sent"
);
if self.heartbeat_failures >= MAX_HEARTBEAT_FAILURES {
self.heartbeat_failures = 0;
return Flow::Unbind(UnbindReason::ConnectionFailed {
detail: format!(
"the outbound heartbeat failed {MAX_HEARTBEAT_FAILURES} times in a row: {detail}"
),
});
}
Flow::Continue
}
async fn emit(&mut self, event: SessionEvent) {
let events = &self.events;
let stop = &mut self.stop;
let interrupted = tokio::select! {
biased;
permit = events.reserve() => match permit {
Ok(permit) => {
permit.send(event);
false
}
Err(_) => {
tracing::debug!("event receiver dropped; events are no longer delivered");
false
}
},
() = stopped(stop) => true,
};
if interrupted {
tracing::debug!(
"a stop was ordered while the event stream was full; the event was not delivered"
);
self.stopping = true;
}
}
}
async fn sleep_until_option(deadline: Option<Instant>) {
match deadline {
Some(instant) => tokio::time::sleep_until(instant).await,
None => std::future::pending().await,
}
}
#[cfg(test)]
mod tests {
#![allow(clippy::unwrap_used, clippy::expect_used)]
use std::collections::VecDeque;
use std::sync::Mutex;
use super::*;
use crate::protocol::request::ConnectionMode;
use crate::session::backoff::BackoffPolicy;
use crate::session::options::{Credentials, SessionOptions};
#[derive(Debug, Clone, PartialEq, Eq)]
enum Step {
Line(String),
Fail(&'static str),
End,
Silence,
AwaitControls(usize),
FloodUntilControls(&'static str, usize),
}
fn line(text: &str) -> Step {
Step::Line(text.to_owned())
}
fn transcript(text: &str) -> Vec<Step> {
text.lines()
.map(str::trim)
.filter(|l| !l.is_empty() && *l != "[…]")
.map(line)
.collect()
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum OpenRecord {
Create,
Bind {
session: Option<String>,
recovery_from: Option<u64>,
},
}
#[derive(Debug, Default)]
struct MockLog {
opens: Vec<OpenRecord>,
controls: Vec<EncodedRequest>,
control_links: Vec<Option<String>>,
closes: usize,
}
impl MockLog {
fn control_parameters(&self) -> Vec<String> {
self.controls
.iter()
.map(|request| request.parameters.clone())
.collect()
}
}
struct MockTransport {
properties: TransportProperties,
scripts: VecDeque<Vec<Step>>,
current: VecDeque<Step>,
control_failures: usize,
hangs_on_open: bool,
open_delay: Duration,
log: Arc<Mutex<MockLog>>,
}
impl MockTransport {
fn new(properties: TransportProperties, scripts: Vec<Vec<Step>>) -> Self {
Self {
properties,
scripts: scripts.into(),
current: VecDeque::new(),
control_failures: 0,
hangs_on_open: false,
open_delay: Duration::ZERO,
log: Arc::new(Mutex::new(MockLog::default())),
}
}
fn failing_controls(mut self, count: usize) -> Self {
self.control_failures = count;
self
}
const fn hanging_open(mut self) -> Self {
self.hangs_on_open = true;
self
}
const fn with_open_delay(mut self, delay: Duration) -> Self {
self.open_delay = delay;
self
}
}
impl Transport for MockTransport {
fn properties(&self) -> TransportProperties {
self.properties
}
fn set_control_link(&mut self, host: Option<&str>) {
self.log
.lock()
.expect("mock log poisoned")
.control_links
.push(host.map(str::to_owned));
}
async fn open_stream(&mut self, request: StreamOpen) -> Result<(), TransportError> {
let record = match request {
StreamOpen::Create(_) => OpenRecord::Create,
StreamOpen::Bind(bind) => OpenRecord::Bind {
session: bind.session.clone(),
recovery_from: bind.recovery_from,
},
};
self.log
.lock()
.expect("mock log poisoned")
.opens
.push(record);
if self.hangs_on_open {
return std::future::pending().await;
}
if !self.open_delay.is_zero() {
tokio::time::sleep(self.open_delay).await;
}
self.current = self.scripts.pop_front().unwrap_or_default().into();
Ok(())
}
async fn next_line(&mut self) -> Option<Result<String, TransportError>> {
loop {
let Some(step) = self.current.front().cloned() else {
return std::future::pending().await;
};
match step {
Step::Line(text) => {
self.current.pop_front();
return Some(Ok(text));
}
Step::Fail(reason) => {
self.current.pop_front();
return Some(Err(TransportError::ConnectionLost {
reason: reason.to_owned(),
}));
}
Step::End => {
self.current.pop_front();
return None;
}
Step::Silence => return std::future::pending().await,
Step::FloodUntilControls(text, wanted) => {
let sent = self.log.lock().expect("mock log poisoned").controls.len();
if sent >= wanted {
self.current.pop_front();
continue;
}
return Some(Ok(text.to_owned()));
}
Step::AwaitControls(wanted) => {
let sent = self.log.lock().expect("mock log poisoned").controls.len();
if sent >= wanted {
self.current.pop_front();
continue;
}
return std::future::pending().await;
}
}
}
}
async fn send_control(&mut self, request: EncodedRequest) -> Result<(), TransportError> {
let name = request.name;
self.log
.lock()
.expect("mock log poisoned")
.controls
.push(request);
if self.control_failures > 0
&& let Some(remaining) = self.control_failures.checked_sub(1)
{
self.control_failures = remaining;
return Err(TransportError::Send {
name,
reason: "control connection refused".to_owned(),
});
}
Ok(())
}
async fn close(&mut self) -> Result<(), TransportError> {
self.log.lock().expect("mock log poisoned").closes += 1;
Ok(())
}
}
const fn websocket() -> TransportProperties {
TransportProperties {
control_shares_stream: true,
ends_on_content_length: false,
is_polling: false,
}
}
const fn http_streaming() -> TransportProperties {
TransportProperties {
control_shares_stream: false,
ends_on_content_length: true,
is_polling: false,
}
}
const fn http_polling() -> TransportProperties {
TransportProperties {
control_shares_stream: false,
ends_on_content_length: false,
is_polling: true,
}
}
fn options() -> SessionOptions {
SessionOptions::default()
.with_credentials(Credentials {
user: None,
password: None,
adapter_set: Some("WELCOME".to_owned()),
})
.with_backoff(BackoffPolicy {
initial: Duration::from_millis(10),
max: Duration::from_millis(40),
max_attempts: NonZeroU32::new(4),
})
}
struct Outcome {
closed: SessionClosed,
events: Vec<SessionEvent>,
log: Arc<Mutex<MockLog>>,
}
impl Outcome {
fn opens(&self) -> Vec<OpenRecord> {
self.log.lock().expect("mock log poisoned").opens.clone()
}
fn controls(&self) -> Vec<String> {
self.log
.lock()
.expect("mock log poisoned")
.control_parameters()
}
fn closes(&self) -> usize {
self.log.lock().expect("mock log poisoned").closes
}
fn bound(&self) -> Vec<&BoundInfo> {
self.events
.iter()
.filter_map(|event| match event {
SessionEvent::Bound(info) => Some(info.as_ref()),
_ => None,
})
.collect()
}
fn unbinds(&self) -> Vec<&UnbindReason> {
self.events
.iter()
.filter_map(|event| match event {
SessionEvent::Unbound { reason, .. } => Some(reason),
_ => None,
})
.collect()
}
fn data(&self) -> Vec<(u64, SubscriptionKey, &SubscriptionOutcome)> {
self.events
.iter()
.filter_map(|event| match event {
SessionEvent::Subscription {
progressive,
key,
outcome,
} => Some((*progressive, *key, outcome.as_ref())),
_ => None,
})
.collect()
}
fn recoveries(&self) -> Vec<RecoveryOutcome> {
self.events
.iter()
.filter_map(|event| match event {
SessionEvent::Recovered(outcome) => Some(*outcome),
_ => None,
})
.collect()
}
fn resubscribes(&self) -> Vec<&Vec<ResubscribedEntry>> {
self.events
.iter()
.filter_map(|event| match event {
SessionEvent::Resubscribed(entries) => Some(entries),
_ => None,
})
.collect()
}
}
async fn drive(
properties: TransportProperties,
options: SessionOptions,
scripts: Vec<Vec<Step>>,
commands: Vec<SessionCommand>,
) -> Outcome {
drive_transport(MockTransport::new(properties, scripts), options, commands).await
}
async fn drive_transport(
transport: MockTransport,
options: SessionOptions,
commands: Vec<SessionCommand>,
) -> Outcome {
let log = Arc::clone(&transport.log);
let (driver, handle, mut events) =
connect(transport, options).expect("options must match the transport");
for command in commands {
handle
.commands
.send(command)
.await
.expect("the driver has not started yet");
}
let closed = driver.run().await;
drop(handle);
let mut collected = Vec::new();
while let Ok(event) = events.try_recv() {
collected.push(event);
}
Outcome {
closed,
events: collected,
log,
}
}
const SERVER_END: &str = "END,32,session closed on the Server side";
const F1_CONOK: &str = "CONOK,S1d7c802482843a26T5626355,50000,5000,*";
const F1_SESSION: &str = "S1d7c802482843a26T5626355";
const F6_CONOK: &str = "CONOK,Se939a67a9be2d336T3823582,50000,5000,*";
const F6_SESSION: &str = "Se939a67a9be2d336T3823582";
const F8_CONOK: &str = "CONOK,S22dee113e3f71b1fT4327493,50000,5000,*";
const F8_SESSION: &str = "S22dee113e3f71b1fT4327493";
fn subscribe() -> SessionCommand {
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
}
}
fn subscribed_transcript(text: &str) -> Vec<Step> {
let mut steps = transcript(text);
steps.insert(1, Step::AwaitControls(1));
steps
}
fn spec() -> SubscriptionSpec {
let mut spec = SubscriptionSpec::new(
"item2",
"stock_name time last_price",
SubscriptionMode::Merge,
);
spec.data_adapter = Some("STOCKS".to_owned());
spec
}
#[tokio::test(start_paused = true)]
async fn test_t1_create_session_binds_and_reports_a_new_session() {
let mut script = transcript(
"
CONOK,S1d7c802482843a26T5626355,50000,5000,*
SERVNAME,Lightstreamer HTTP Server
CLIENTIP,0:0:0:0:0:0:0:1
NOOP,sending placeholder data
CONS,unlimited
PROBE
PROBE
PROBE
",
);
script.push(line(SERVER_END));
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
let bound = outcome.bound();
assert_eq!(bound.len(), 1);
let info = bound.first().expect("one bind");
assert_eq!(info.session_id.as_str(), F1_SESSION);
assert_eq!(info.request_limit_bytes, 50_000);
assert_eq!(info.keep_alive, Duration::from_millis(5000));
assert_eq!(info.control_link, None);
assert_eq!(info.kind, BindKind::Created);
assert_eq!(outcome.opens(), vec![OpenRecord::Create]);
assert_eq!(outcome.closes(), 1);
}
#[tokio::test(start_paused = true)]
async fn test_t1_head_notifications_are_surfaced_but_not_counted() {
let mut script = transcript(
"
CONOK,S1d7c802482843a26T5626355,50000,5000,*
SERVNAME,Lightstreamer HTTP Server
CLIENTIP,0:0:0:0:0:0:0:1
NOOP,sending placeholder data
CONS,unlimited
PROBE
",
);
script.push(line(SERVER_END));
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
let infos = outcome
.events
.iter()
.filter(|event| matches!(event, SessionEvent::ServerInfo(_)))
.count();
assert_eq!(infos, 3, "SERVNAME, CLIENTIP and CONS reach the caller");
assert!(
outcome.data().is_empty(),
"none of them is a data notification"
);
}
#[tokio::test(start_paused = true)]
async fn test_t2_force_rebind_unbinds_and_rebinds_the_same_session() {
let first = vec![
line(F6_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("LOOP,0"),
Step::End,
];
let second = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![SessionCommand::ForceRebind {
close_socket: Some(true),
}],
)
.await;
assert_eq!(
outcome.unbinds(),
vec![&UnbindReason::ForcedByClient {
expected_delay: Duration::ZERO
}]
);
let controls = outcome.controls();
assert!(
controls
.iter()
.any(|parameters| parameters.contains("LS_op=force_rebind")
&& parameters.contains("LS_close_socket=true")),
"{controls:?}"
);
assert_eq!(
outcome.opens().get(1),
Some(&OpenRecord::Bind {
session: Some(F6_SESSION.to_owned()),
recovery_from: None,
})
);
}
#[tokio::test(start_paused = true)]
async fn test_t3_content_length_reached_rebinds_on_a_new_connection() {
let mut first = transcript(
"
CONOK,Se939a67a9be2d336T3823582,50000,5000,*
SUBOK,1,1,3
U,1,1,stock|15:57:51|16.27
U,1,1,||16.24
U,1,1,||16.11
U,1,1,|15:57:52|15.99
U,1,1,||15.93
LOOP,0
",
);
first.insert(1, Step::AwaitControls(1));
first.push(Step::End);
let mut second = transcript(
"
CONOK,Se939a67a9be2d336T3823582,50000,5000,*
NOOP,sending placeholder data
CONS,unlimited
U,1,1,|15:57:53|17.7
U,1,1,||17.81
",
);
second.push(line(SERVER_END));
let outcome = drive(
http_streaming(),
options(),
vec![first, second],
vec![subscribe()],
)
.await;
assert_eq!(
outcome.unbinds(),
vec![&UnbindReason::ContentLengthReached {
expected_delay: Duration::ZERO
}]
);
let bound = outcome.bound();
assert_eq!(bound.len(), 2);
assert_eq!(
bound.first().map(|info| info.session_id.clone()),
bound.get(1).map(|info| info.session_id.clone())
);
assert_eq!(
bound.get(1).map(|info| info.kind.clone()),
Some(BindKind::Rebound)
);
assert!(outcome.resubscribes().is_empty());
assert_eq!(
outcome
.data()
.iter()
.map(|(p, _, _)| *p)
.collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5, 6, 7, 8]
);
}
#[tokio::test(start_paused = true)]
async fn test_t4_poll_cycle_expired_rebinds_after_the_expected_delay() {
let first = vec![line(F1_CONOK), line("LOOP,2000"), Step::End];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let polling = options().with_connection(ConnectionMode::Polling {
polling_millis: 2000,
idle_millis: Some(10_000),
});
let outcome = drive(http_polling(), polling, vec![first, second], Vec::new()).await;
assert_eq!(
outcome.unbinds(),
vec![&UnbindReason::PollCycleExpired {
expected_delay: Duration::from_millis(2000)
}]
);
let retry_in = outcome.events.iter().find_map(|event| match event {
SessionEvent::Unbound { retry_in, .. } => *retry_in,
_ => None,
});
assert_eq!(retry_in, Some(Duration::from_millis(2000)));
}
#[tokio::test(start_paused = true)]
async fn test_t4_polling_stream_ending_without_a_loop_is_still_a_poll_cycle() {
let first = vec![line(F1_CONOK), Step::End];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let polling = options().with_connection(ConnectionMode::Polling {
polling_millis: 1000,
idle_millis: None,
});
let outcome = drive(http_polling(), polling, vec![first, second], Vec::new()).await;
assert_eq!(
outcome.unbinds(),
vec![&UnbindReason::PollCycleExpired {
expected_delay: Duration::ZERO
}]
);
assert_eq!(
outcome.opens().get(1),
Some(&OpenRecord::Bind {
session: Some(F1_SESSION.to_owned()),
recovery_from: None,
})
);
}
#[tokio::test(start_paused = true)]
async fn test_polling_establishment_budget_covers_a_withheld_response() {
let conok = "CONOK,S1,50000,15000,*";
let first = vec![line(conok), Step::End];
let second = vec![line(conok), line(SERVER_END)];
let polling = options().with_connection(ConnectionMode::Polling {
polling_millis: 5000,
idle_millis: Some(15_000),
});
assert_eq!(
polling.open_timeout,
Duration::from_secs(10),
"the fixture relies on the default 10s open timeout"
);
let transport = MockTransport::new(http_polling(), vec![first, second])
.with_open_delay(Duration::from_secs(13));
let outcome = drive_transport(transport, polling, Vec::new()).await;
assert_eq!(outcome.bound().len(), 2, "both poll cycles bound");
assert!(
!outcome
.unbinds()
.iter()
.any(|reason| matches!(reason, UnbindReason::ConnectionFailed { .. })),
"no cycle was aborted as a failed establishment: {:?}",
outcome.unbinds()
);
assert_eq!(
outcome.unbinds(),
vec![&UnbindReason::PollCycleExpired {
expected_delay: Duration::ZERO
}]
);
}
#[tokio::test(start_paused = true)]
async fn test_t3_and_t4_are_told_apart_by_declared_properties_only() {
let script = || {
vec![
vec![line(F1_CONOK), line("LOOP,0"), Step::End],
vec![line(F1_CONOK), line(SERVER_END)],
]
};
let a = drive(http_streaming(), options(), script(), Vec::new()).await;
assert!(matches!(
a.unbinds().first(),
Some(UnbindReason::ContentLengthReached { .. })
));
let b = drive(websocket(), options(), script(), Vec::new()).await;
assert!(matches!(
b.unbinds().first(),
Some(UnbindReason::Looped { .. })
));
}
#[tokio::test(start_paused = true)]
async fn test_t5_connection_failed_recovers_from_the_counted_progressive() {
let mut first = subscribed_transcript(F8_STREAM);
first.push(Step::Fail("connection reset"));
let second = vec![line(F8_CONOK), line("PROG,15"), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![subscribe()],
)
.await;
assert!(matches!(
outcome.unbinds().first(),
Some(UnbindReason::ConnectionFailed { .. })
));
assert_eq!(
outcome.opens().get(1),
Some(&OpenRecord::Bind {
session: Some(F8_SESSION.to_owned()),
recovery_from: Some(15),
})
);
assert_eq!(
outcome.recoveries(),
vec![RecoveryOutcome {
requested: 15,
resumed_at: 15,
kind: RecoveryKind::Exact,
}]
);
assert_eq!(
outcome.bound().get(1).map(|info| info.kind.clone()),
Some(BindKind::Recovering {
requested_progressive: 15
})
);
}
const F8_STREAM: &str = "
CONOK,S22dee113e3f71b1fT4327493,50000,5000,*
SERVNAME,Lightstreamer HTTP Server
CLIENTIP,0:0:0:0:0:0:0:1
NOOP,sending placeholder data
CONS,unlimited
PROBE
SUBOK,1,1,3
CONF,1,unlimited,filtered
PROBE
U,1,1,Ations Europe|15:55:08|14.81
U,1,1,||14.66
U,1,1,|15:55:09|14.62
U,1,1,||14.71
U,1,1,|15:55:10|14.63
U,1,1,||14.77
U,1,1,|15:55:11|
U,1,1,||14.61
U,1,1,|15:55:12|14.5
U,1,1,|15:55:13|14.64
U,1,1,|15:55:14|
SYNC,25
U,1,1,||14.74
U,1,1,||14.66
";
#[tokio::test(start_paused = true)]
async fn test_data_notification_count_matches_the_spec_worked_example() {
let mut script = subscribed_transcript(F8_STREAM);
script.push(line(SERVER_END));
let outcome = drive(websocket(), options(), vec![script], vec![subscribe()]).await;
let data = outcome.data();
assert_eq!(data.len(), 15, "1 SUBOK + 1 CONF + 13 U");
assert_eq!(data.last().map(|(p, _, _)| *p), Some(15));
assert!(
data.iter().all(|(_, _, outcome)| outcome.is_ok()),
"PROBE and SYNC are not data notifications, and every counted line decoded"
);
}
#[tokio::test(start_paused = true)]
async fn test_recovery_resuming_earlier_discards_the_duplicates() {
let mut first = subscribed_transcript(F8_STREAM);
first.push(Step::Fail("connection reset"));
let mut second = transcript(
"
CONOK,S22dee113e3f71b1fT4327493,50000,5000,*
NOOP,sending placeholder data
CONS,unlimited
PROG,11
U,1,1,|15:55:13|14.64
U,1,1,|15:55:14|
U,1,1,||14.74
U,1,1,||14.66
U,1,1,|15:55:15|17.7
U,1,1,||17.81
SYNC,5
",
);
second.push(line(SERVER_END));
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![subscribe()],
)
.await;
assert_eq!(
outcome.recoveries(),
vec![RecoveryOutcome {
requested: 15,
resumed_at: 11,
kind: RecoveryKind::Duplicated { count: 4 },
}]
);
let progressives: Vec<u64> = outcome.data().iter().map(|(p, _, _)| *p).collect();
assert_eq!(progressives.len(), 15 + 2);
assert_eq!(progressives.get(15), Some(&16));
assert_eq!(progressives.last(), Some(&17));
}
#[tokio::test(start_paused = true)]
async fn test_recovery_resuming_later_is_reported_as_a_gap() {
let first = vec![
line(F8_CONOK),
Step::AwaitControls(1),
line("SUBOK,1,1,3"),
line("U,1,1,a|b|c"),
Step::Fail("connection reset"),
];
let second = vec![
line(F8_CONOK),
line("PROG,9"),
line("U,1,1,d|e|f"),
line(SERVER_END),
];
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![subscribe()],
)
.await;
assert_eq!(
outcome.recoveries(),
vec![RecoveryOutcome {
requested: 2,
resumed_at: 9,
kind: RecoveryKind::Gap { missing: 7 },
}]
);
assert_eq!(outcome.data().last().map(|(p, _, _)| *p), Some(10));
}
#[tokio::test(start_paused = true)]
async fn test_t6_rebind_preserves_the_session_and_its_subscriptions() {
let first = vec![
line(F6_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
line("LOOP,0"),
Step::End,
];
let second = vec![line(F6_CONOK), line("U,1,1,x|y|z"), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
}],
)
.await;
assert_eq!(
outcome.bound().get(1).map(|info| info.kind.clone()),
Some(BindKind::Rebound)
);
let adds = outcome
.controls()
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.count();
assert_eq!(adds, 1, "the subscription must not be issued twice");
assert!(outcome.resubscribes().is_empty());
}
#[tokio::test(start_paused = true)]
async fn test_t7_destroy_while_bound_ends_with_the_servers_cause() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("END,31,Destroy invoked by client"),
Step::End,
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::Destroy { cause: None }],
)
.await;
assert_eq!(
outcome.closed,
SessionClosed::ByClient {
destroy_confirmed: true,
cause: Some(ServerCause {
code: 31,
message: "Destroy invoked by client".to_owned(),
}),
}
);
assert!(
outcome
.controls()
.iter()
.any(|parameters| parameters.contains("LS_op=destroy")),
"{:?}",
outcome.controls()
);
assert_eq!(outcome.closes(), 1);
}
#[tokio::test(start_paused = true)]
async fn test_t8_destroy_while_unbound_ends_without_waiting_for_an_end() {
let first = vec![line(F1_CONOK), line("LOOP,60000"), Step::End];
let outcome = drive(
websocket(),
options(),
vec![first],
vec![SessionCommand::Destroy { cause: None }],
)
.await;
assert_eq!(
outcome.closed,
SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
}
);
assert_eq!(outcome.opens(), vec![OpenRecord::Create]);
}
#[tokio::test(start_paused = true)]
async fn test_t9_timed_out_session_is_recreated_and_subscriptions_re_executed() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(2),
line("REQOK,1"),
line("REQOK,2"),
line("SUBOK,1,1,3"),
line("SUBOK,2,1,3"),
Step::Fail("connection reset"),
];
let second = vec![line("CONERR,20,Specified session not found"), Step::End];
let third = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second, third],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Subscribe {
key: SubscriptionKey(2),
spec: Box::new(SubscriptionSpec::new(
"item1",
"last_price",
SubscriptionMode::Merge,
)),
},
],
)
.await;
assert!(matches!(
outcome.unbinds().get(1),
Some(UnbindReason::Rejected {
cause: ServerCause { code: 20, .. }
})
));
assert_eq!(
outcome.bound().get(1).map(|info| info.kind.clone()),
Some(BindKind::Recreated {
previous: Some(SessionId::new(F1_SESSION)),
})
);
let resubscribed = outcome.resubscribes();
let entries = resubscribed.first().expect("a resubscribe batch");
assert_eq!(entries.len(), 2);
assert_eq!(
entries.iter().map(|e| e.key).collect::<Vec<_>>(),
vec![SubscriptionKey(1), SubscriptionKey(2)]
);
assert_eq!(
entries
.iter()
.map(|e| e.subscription_id.get())
.collect::<Vec<_>>(),
vec![1, 2]
);
assert!(entries.iter().all(|e| e.previously_active));
let adds = outcome
.controls()
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.count();
assert_eq!(adds, 4);
}
#[tokio::test(start_paused = true)]
async fn test_server_initiated_end_on_a_bound_session_is_definitive_loss() {
let script = vec![
line(F1_CONOK),
line("END,35,Another session was opened for this user"),
Step::End,
];
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
assert_eq!(
outcome.closed,
SessionClosed::ByServer {
cause: ServerCause {
code: 35,
message: "Another session was opened for this user".to_owned(),
},
}
);
assert_eq!(outcome.opens().len(), 1, "no retry after a fatal cause");
assert!(matches!(
outcome.events.last(),
Some(SessionEvent::Closed(SessionClosed::ByServer { .. }))
));
}
#[tokio::test(start_paused = true)]
async fn test_server_end_code_48_recreates_the_session_immediately() {
let first = vec![
line(F1_CONOK),
line("END,48,Maximum session duration reached"),
Step::End,
];
let second = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(websocket(), options(), vec![first, second], Vec::new()).await;
assert!(matches!(
outcome.unbinds().first(),
Some(UnbindReason::ServerRefresh {
cause: ServerCause { code: 48, .. }
})
));
let retry_in = outcome.events.iter().find_map(|event| match event {
SessionEvent::Unbound { retry_in, .. } => *retry_in,
_ => None,
});
assert_eq!(retry_in, Some(Duration::ZERO));
assert_eq!(outcome.opens().get(1), Some(&OpenRecord::Create));
}
#[tokio::test(start_paused = true)]
async fn test_conerr_authentication_failure_is_definitive_loss() {
let script = vec![line("CONERR,1,User/password check failed"), Step::End];
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
assert_eq!(
outcome.closed,
SessionClosed::ByServer {
cause: ServerCause {
code: 1,
message: "User/password check failed".to_owned(),
},
}
);
assert_eq!(outcome.opens().len(), 1);
}
#[tokio::test(start_paused = true)]
async fn test_conerr_retry_later_code_is_retried_and_succeeds() {
let first = vec![
line("CONERR,5,The Server is temporarily overloaded: retry later"),
Step::End,
];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let outcome = drive(websocket(), options(), vec![first, second], Vec::new()).await;
assert!(matches!(
outcome.unbinds().first(),
Some(UnbindReason::Rejected {
cause: ServerCause { code: 5, .. }
})
));
assert_eq!(outcome.bound().len(), 1);
assert_eq!(outcome.opens().len(), 2);
}
#[tokio::test(start_paused = true)]
async fn test_conerr_cluster_affinity_code_is_treated_as_permanent() {
let script = vec![
line("CONERR,21,Session ID not compatible with this Server instance"),
Step::End,
];
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
assert!(matches!(
outcome.closed,
SessionClosed::ByServer {
cause: ServerCause { code: 21, .. }
}
));
}
#[test]
fn test_classify_covers_every_code_the_spec_prescribes_an_action_for() {
assert_eq!(classify(4), Recovery::RecreateSession);
assert_eq!(classify(20), Recovery::RecreateSession);
assert_eq!(classify(48), Recovery::RecreateSession);
assert_eq!(classify(5), Recovery::RetryLater);
assert_eq!(classify(6), Recovery::RetryLater);
assert_eq!(classify(10), Recovery::RetryLater);
assert_eq!(classify(1), Recovery::Fatal);
assert_eq!(classify(21), Recovery::Fatal);
assert_eq!(classify(31), Recovery::Fatal);
assert_eq!(classify(9999), Recovery::Fatal);
assert_eq!(classify(0), Recovery::Fatal);
assert_eq!(classify(-7), Recovery::Fatal);
}
#[tokio::test(start_paused = true)]
async fn test_keepalive_expiry_forces_a_rebind_on_a_wedged_connection() {
let first = vec![line(F1_CONOK), Step::Silence];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let outcome = drive(websocket(), options(), vec![first, second], Vec::new()).await;
assert_eq!(
outcome.unbinds(),
vec![&UnbindReason::KeepaliveExpired {
budget: Duration::from_millis(8000)
}]
);
assert_eq!(
outcome.opens().get(1),
Some(&OpenRecord::Bind {
session: Some(F1_SESSION.to_owned()),
recovery_from: Some(0),
})
);
assert_eq!(outcome.closes(), 2);
}
#[tokio::test(start_paused = true)]
async fn test_open_timeout_abandons_an_unanswered_handshake() {
let first = vec![Step::Silence];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let outcome = drive(websocket(), options(), vec![first, second], Vec::new()).await;
assert!(matches!(
outcome.unbinds().first(),
Some(UnbindReason::ConnectionFailed { .. })
));
assert_eq!(outcome.bound().len(), 1);
}
#[tokio::test(start_paused = true)]
async fn test_reverse_heartbeat_is_sent_before_the_inactivity_commitment_expires() {
let first = vec![line(F1_CONOK), Step::Silence];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let committed = options().with_connection(ConnectionMode::Streaming {
inactivity_millis: Some(8000),
keepalive_millis: None,
send_sync: None,
});
let outcome = drive(websocket(), committed, vec![first, second], Vec::new()).await;
let controls = outcome.controls();
assert!(
controls
.iter()
.any(|parameters| parameters.contains("LS_session=")),
"the heartbeat declares the session it is listening to: {controls:?}"
);
assert!(!controls.is_empty(), "a heartbeat must have been sent");
assert!(matches!(
outcome.unbinds().first(),
Some(UnbindReason::KeepaliveExpired { .. })
));
}
#[tokio::test(start_paused = true)]
async fn test_a_quiet_but_probing_connection_is_never_declared_stalled() {
let mut script = vec![line(F1_CONOK)];
for _ in 0..20 {
script.push(line("PROBE"));
}
script.push(line(SERVER_END));
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
assert!(outcome.unbinds().is_empty());
assert!(matches!(outcome.closed, SessionClosed::ByServer { .. }));
}
#[tokio::test(start_paused = true)]
async fn test_control_responses_are_correlated_when_interleaved_with_data() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(2),
line("U,1,1,a|b|c"),
line("REQOK,2"),
line("U,1,1,d|e|f"),
line("REQOK,1"),
line("SUBOK,1,1,3"),
line(SERVER_END),
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Subscribe {
key: SubscriptionKey(2),
spec: Box::new(SubscriptionSpec::new(
"item1",
"last_price",
SubscriptionMode::Merge,
)),
},
],
)
.await;
let responses: Vec<String> = outcome
.events
.iter()
.filter_map(|event| match event {
SessionEvent::ControlResponse {
request_id: Some(id),
outcome: ControlOutcome::Accepted,
..
} => Some(id.as_str().to_owned()),
_ => None,
})
.collect();
assert_eq!(responses, vec!["2".to_owned(), "1".to_owned()]);
let controls = outcome.controls();
assert_eq!(
controls
.iter()
.filter(|parameters| parameters.contains("LS_reqId=1"))
.count(),
1,
"{controls:?}"
);
assert_eq!(
controls
.iter()
.filter(|parameters| parameters.contains("LS_reqId=2"))
.count(),
1,
"{controls:?}"
);
assert_eq!(outcome.data().len(), 3);
}
#[tokio::test(start_paused = true)]
async fn test_request_ids_are_never_reused_across_sessions() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::Fail("connection reset"),
];
let second = vec![line("CONERR,20,Specified session not found"), Step::End];
let third = vec![line(F6_CONOK), Step::AwaitControls(2), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second, third],
vec![SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
}],
)
.await;
let controls = outcome.controls();
assert_eq!(controls.len(), 2);
assert_eq!(
controls
.iter()
.filter(|parameters| parameters.contains("LS_reqId=1"))
.count(),
1,
"{controls:?}"
);
assert_eq!(
controls
.iter()
.filter(|parameters| parameters.contains("LS_reqId=2"))
.count(),
1,
"{controls:?}"
);
}
#[tokio::test(start_paused = true)]
async fn test_rejected_subscription_is_reported_and_not_re_issued() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQERR,1,23,Bad Field schema name"),
Step::Fail("connection reset"),
];
let second = vec![line("CONERR,20,Specified session not found"), Step::End];
let third = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second, third],
vec![SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
}],
)
.await;
let rejected = outcome.events.iter().any(|event| {
matches!(
event,
SessionEvent::ControlResponse {
outcome: ControlOutcome::Rejected {
cause: ServerCause { code: 23, .. }
},
..
}
)
});
assert!(rejected, "the refusal must reach the caller");
let adds = outcome
.controls()
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.count();
assert_eq!(adds, 1, "a refused subscription is not re-issued");
}
#[tokio::test(start_paused = true)]
async fn test_uncorrelatable_error_is_surfaced_without_a_request_id() {
let script = vec![
line(F1_CONOK),
line("ERROR,67,Malformed request"),
line(SERVER_END),
];
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
assert!(outcome.events.iter().any(|event| matches!(
event,
SessionEvent::ControlResponse {
request_id: None,
target: ControlTarget::Session,
outcome: ControlOutcome::Rejected {
cause: ServerCause { code: 67, .. }
},
}
)));
}
fn control_targets(outcome: &Outcome) -> Vec<(ControlTarget, ControlOutcome)> {
outcome
.events
.iter()
.filter_map(|event| match event {
SessionEvent::ControlResponse {
target, outcome, ..
} => Some((target.clone(), outcome.clone())),
_ => None,
})
.collect()
}
#[tokio::test(start_paused = true)]
async fn test_rejected_subscription_names_the_subscription_it_killed() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQERR,1,23,Bad Field schema name"),
line(SERVER_END),
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::Subscribe {
key: SubscriptionKey(7),
spec: Box::new(spec()),
}],
)
.await;
assert_eq!(
control_targets(&outcome),
vec![(
ControlTarget::Subscription {
key: SubscriptionKey(7),
operation: SubscriptionOperation::Subscribe,
},
ControlOutcome::Rejected {
cause: ServerCause {
code: 23,
message: "Bad Field schema name".to_owned(),
},
},
)]
);
}
#[tokio::test(start_paused = true)]
async fn test_rejected_unsubscribe_names_the_subscription_that_is_still_alive() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::AwaitControls(2),
line("REQERR,2,19,Specified subscription not found"),
line("U,1,1,still|coming|through"),
line(SERVER_END),
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Unsubscribe {
key: SubscriptionKey(1),
},
],
)
.await;
let rejection = control_targets(&outcome)
.into_iter()
.find(|(_, outcome)| matches!(outcome, ControlOutcome::Rejected { .. }));
assert_eq!(
rejection.map(|(target, _)| target),
Some(ControlTarget::Subscription {
key: SubscriptionKey(1),
operation: SubscriptionOperation::Unsubscribe,
})
);
let routed = outcome.data().iter().any(|(_, key, outcome)| {
*key == SubscriptionKey(1) && matches!(outcome, Ok(SubscriptionEvent::Update(_)))
});
assert!(
routed,
"the surviving subscription still routes its updates"
);
}
#[tokio::test(start_paused = true)]
async fn test_rejected_reconfigure_names_the_subscription_and_the_operation() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::AwaitControls(2),
line("REQERR,2,13,Reconfiguration not allowed"),
line(SERVER_END),
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Reconfigure {
key: SubscriptionKey(1),
max_frequency: MaxFrequencyLimit::Unlimited,
},
],
)
.await;
let rejection = control_targets(&outcome)
.into_iter()
.find(|(_, outcome)| matches!(outcome, ControlOutcome::Rejected { .. }));
assert_eq!(
rejection.map(|(target, _)| target),
Some(ControlTarget::Subscription {
key: SubscriptionKey(1),
operation: SubscriptionOperation::Reconfigure,
})
);
}
#[tokio::test(start_paused = true)]
async fn test_session_level_control_responses_target_the_session() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("END,31,Destroy invoked by client"),
Step::End,
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::Destroy { cause: None }],
)
.await;
assert_eq!(
control_targets(&outcome)
.into_iter()
.map(|(target, _)| target)
.collect::<Vec<_>>(),
vec![ControlTarget::Session]
);
}
#[tokio::test(start_paused = true)]
async fn test_subscription_whose_add_could_not_be_sent_is_re_issued_on_the_next_bind() {
let first = vec![line(F1_CONOK), Step::AwaitControls(1), Step::Fail("reset")];
let second = vec![line(F1_CONOK), Step::AwaitControls(2), line(SERVER_END)];
let transport = MockTransport::new(websocket(), vec![first, second]).failing_controls(1);
let outcome = drive_transport(
transport,
options(),
vec![SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
}],
)
.await;
assert_eq!(
control_targets(&outcome)
.into_iter()
.filter(|(_, outcome)| matches!(outcome, ControlOutcome::NotSent { .. }))
.map(|(target, _)| target)
.collect::<Vec<_>>(),
vec![ControlTarget::Subscription {
key: SubscriptionKey(1),
operation: SubscriptionOperation::Subscribe,
}]
);
let adds: Vec<String> = outcome
.controls()
.into_iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.collect();
assert_eq!(adds.len(), 2, "{adds:?}");
assert!(
adds.last().is_some_and(|last| last.contains("LS_subId=2")),
"{adds:?}"
);
}
fn messages(
outcome: &Outcome,
) -> Vec<(Option<u64>, MessageSequence, Option<u64>, MessageResult)> {
outcome
.events
.iter()
.filter_map(|event| match event {
SessionEvent::Message {
progressive,
sequence,
prog,
result,
} => Some((*progressive, sequence.clone(), *prog, result.clone())),
_ => None,
})
.collect()
}
#[tokio::test(start_paused = true)]
async fn test_message_send_is_dispatched_and_its_outcome_reported() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("MSGDONE,Orders_Sequence,3,Processed with ID 32652506"),
line(SERVER_END),
];
let sequence = SequenceName::try_new("Orders_Sequence").expect("a valid sequence name");
let message = OutgoingMessage::numbered("buy 100", NonZeroU32::new(3).expect("non-zero"))
.in_sequence(sequence);
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::SendMessage {
message: Box::new(message),
}],
)
.await;
let controls = outcome.controls();
assert_eq!(controls.len(), 1);
let sent = controls.first().expect("one control request");
assert!(sent.contains("LS_reqId=1"), "{sent}");
assert!(sent.contains("LS_sequence=Orders_Sequence"), "{sent}");
assert!(sent.contains("LS_msg_prog=3"), "{sent}");
assert!(sent.contains("LS_message=buy%20100"), "{sent}");
assert_eq!(
messages(&outcome),
vec![(
Some(1),
MessageSequence::Named("Orders_Sequence".to_owned()),
Some(3),
MessageResult::Done {
response: "Processed with ID 32652506".to_owned(),
},
)]
);
assert!(outcome.events.iter().any(|event| matches!(
event,
SessionEvent::ControlResponse {
outcome: ControlOutcome::Accepted,
..
}
)));
}
#[tokio::test(start_paused = true)]
async fn test_message_outcomes_are_counted_toward_the_recovery_progressive() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
line("MSGDONE,*,1,"),
line("U,1,1,a|b|c"),
Step::Fail("connection reset"),
];
let second = vec![line(F1_CONOK), line("PROG,3"), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![script, second],
vec![SessionCommand::SendMessage {
message: Box::new(OutgoingMessage::numbered(
"hello",
NonZeroU32::new(1).expect("non-zero"),
)),
}],
)
.await;
assert_eq!(
messages(&outcome)
.first()
.map(|(progressive, ..)| *progressive),
Some(Some(2)),
"the MSGDONE sits between the SUBOK and the U"
);
assert_eq!(
outcome.opens().get(1),
Some(&OpenRecord::Bind {
session: Some(F1_SESSION.to_owned()),
recovery_from: Some(3),
})
);
}
#[tokio::test(start_paused = true)]
async fn test_message_failure_is_reported_with_the_servers_cause() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line(
"MSGFAIL,Orders_Sequence,4,38,The specified progressive number has been skipped by timeout",
),
line(SERVER_END),
];
let sequence = SequenceName::try_new("Orders_Sequence").expect("a valid sequence name");
let message = OutgoingMessage::numbered("sell 50", NonZeroU32::new(4).expect("non-zero"))
.in_sequence(sequence);
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::SendMessage {
message: Box::new(message),
}],
)
.await;
assert_eq!(
messages(&outcome)
.first()
.map(|(_, _, _, result)| result.clone()),
Some(MessageResult::Failed {
cause: ServerCause {
code: 38,
message: "The specified progressive number has been skipped by timeout"
.to_owned(),
},
})
);
}
#[tokio::test(start_paused = true)]
async fn test_message_refused_at_submission_gets_exactly_one_outcome() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQERR,1,33,A message with this number has already been enqueued"),
line(SERVER_END),
];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::SendMessage {
message: Box::new(OutgoingMessage::numbered(
"duplicate",
NonZeroU32::new(2).expect("non-zero"),
)),
}],
)
.await;
let reported = messages(&outcome);
assert_eq!(reported.len(), 1, "exactly one outcome per message");
assert_eq!(
reported.first().map(|(_, _, _, result)| result.clone()),
Some(MessageResult::Refused {
cause: ServerCause {
code: 33,
message: "A message with this number has already been enqueued".to_owned(),
},
})
);
}
#[tokio::test(start_paused = true)]
async fn test_message_sent_while_unbound_fails_fast_and_is_never_buffered() {
let first = vec![line(F1_CONOK), line("LOOP,60000"), Step::End];
let outcome = drive(
websocket(),
options(),
vec![first],
vec![
SessionCommand::SendMessage {
message: Box::new(OutgoingMessage::numbered(
"too late",
NonZeroU32::new(1).expect("non-zero"),
)),
},
SessionCommand::Shutdown,
],
)
.await;
assert!(matches!(
messages(&outcome)
.first()
.map(|(_, _, _, result)| result.clone()),
Some(MessageResult::NotSent { .. })
));
assert!(outcome.controls().is_empty());
}
#[tokio::test(start_paused = true)]
async fn test_fire_and_forget_message_declines_its_outcome() {
let script = vec![line(F1_CONOK), Step::AwaitControls(1), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::SendMessage {
message: Box::new(OutgoingMessage::fire_and_forget("ping")),
}],
)
.await;
let controls = outcome.controls();
let sent = controls.first().expect("one control request");
assert!(sent.contains("LS_outcome=false"), "{sent}");
assert!(!sent.contains("LS_msg_prog"), "{sent}");
assert!(messages(&outcome).is_empty());
}
#[tokio::test(start_paused = true)]
async fn test_adr0005_recovery_and_reestablishment_are_distinguishable() {
let interrupted = || {
vec![
line(F8_CONOK),
line("SUBOK,1,1,3"),
Step::Fail("connection reset"),
]
};
let recovered = drive(
websocket(),
options(),
vec![
interrupted(),
vec![line(F8_CONOK), line("PROG,1"), line(SERVER_END)],
],
Vec::new(),
)
.await;
assert_eq!(
recovered.bound().get(1).map(|info| info.kind.clone()),
Some(BindKind::Recovering {
requested_progressive: 1
})
);
assert_eq!(
recovered.recoveries().first().map(|outcome| outcome.kind),
Some(RecoveryKind::Exact)
);
let replaced = drive(
websocket(),
options(),
vec![
interrupted(),
vec![line("CONERR,4,Recovery not possible"), Step::End],
vec![line(F1_CONOK), line(SERVER_END)],
],
Vec::new(),
)
.await;
assert_eq!(
replaced.bound().get(1).map(|info| info.kind.clone()),
Some(BindKind::Recreated {
previous: Some(SessionId::new(F8_SESSION)),
})
);
assert!(replaced.recoveries().is_empty());
}
#[tokio::test(start_paused = true)]
async fn test_adr0005_snapshot_re_delivery_is_visible_on_resubscribe() {
let mut with_snapshot = spec();
with_snapshot.snapshot = Some(Snapshot::On);
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::Fail("connection reset"),
];
let second = vec![line("CONERR,20,Specified session not found"), Step::End];
let third = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second, third],
vec![SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(with_snapshot),
}],
)
.await;
let resubscribed = outcome.resubscribes();
let entry = resubscribed
.first()
.and_then(|entries| entries.first())
.expect("one re-subscribed entry");
assert!(entry.snapshot_requested);
assert!(entry.previously_active);
}
#[tokio::test(start_paused = true)]
async fn test_retries_are_bounded_and_reported_as_definitive_loss() {
let failing = || vec![Step::Fail("connection refused")];
let options = options().with_backoff(BackoffPolicy {
initial: Duration::from_millis(10),
max: Duration::from_millis(20),
max_attempts: NonZeroU32::new(2),
});
let outcome = drive(
websocket(),
options,
vec![failing(), failing(), failing(), failing()],
Vec::new(),
)
.await;
assert!(matches!(
outcome.closed,
SessionClosed::RetriesExhausted { .. }
));
assert_eq!(outcome.opens().len(), 3);
assert_eq!(outcome.closes(), 4);
}
#[tokio::test(start_paused = true)]
async fn test_a_successful_bind_resets_the_retry_budget() {
let options = options().with_backoff(BackoffPolicy {
initial: Duration::from_millis(10),
max: Duration::from_millis(20),
max_attempts: NonZeroU32::new(1),
});
let outcome = drive(
websocket(),
options,
vec![
vec![line(F1_CONOK), Step::Fail("blip")],
vec![line(F1_CONOK), Step::Fail("blip")],
vec![line(F1_CONOK), line(SERVER_END)],
],
Vec::new(),
)
.await;
assert!(matches!(outcome.closed, SessionClosed::ByServer { .. }));
assert_eq!(outcome.bound().len(), 3);
}
#[tokio::test(start_paused = true)]
async fn test_shutdown_stops_the_driver_and_closes_the_transport() {
let script = vec![line(F1_CONOK), Step::Silence];
let outcome = drive(
websocket(),
options(),
vec![script],
vec![SessionCommand::Shutdown],
)
.await;
assert_eq!(
outcome.closed,
SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
}
);
assert_eq!(outcome.closes(), 1);
assert_eq!(outcome.opens().len(), 1);
}
#[tokio::test(start_paused = true)]
async fn test_dropping_every_handle_stops_the_driver() {
let transport = MockTransport::new(websocket(), vec![vec![line(F1_CONOK), Step::Silence]]);
let log = Arc::clone(&transport.log);
let (driver, handle, _events) =
connect(transport, options()).expect("options must match the transport");
drop(handle);
let closed = driver.run().await;
assert_eq!(
closed,
SessionClosed::ByClient {
destroy_confirmed: false,
cause: None,
}
);
assert_eq!(log.lock().expect("mock log poisoned").closes, 1);
}
#[tokio::test(start_paused = true)]
async fn test_connect_rejects_a_polling_mismatch_between_options_and_transport() {
let transport = MockTransport::new(http_polling(), Vec::new());
let result = connect(transport, options());
assert!(matches!(result, Err(SessionError::Configuration { .. })));
}
#[tokio::test(start_paused = true)]
async fn test_unknown_tag_is_surfaced_and_never_fatal() {
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("SUBOK,1,1,3"),
line("MPNREG,devid,adapter"),
line("U,1,1,a|b|c"),
line(SERVER_END),
];
let outcome = drive(websocket(), options(), vec![script], vec![subscribe()]).await;
assert!(
outcome
.events
.iter()
.any(|event| matches!(event, SessionEvent::Unparsed { .. }))
);
assert_eq!(
outcome
.data()
.iter()
.map(|(p, _, _)| *p)
.collect::<Vec<_>>(),
vec![1, 2]
);
assert!(matches!(outcome.closed, SessionClosed::ByServer { .. }));
}
#[tokio::test(start_paused = true)]
async fn test_control_link_is_handed_to_the_transport() {
let script = vec![
line("CONOK,S1,50000,5000,push2.example.com:8080"),
line(SERVER_END),
];
let outcome = drive(websocket(), options(), vec![script], Vec::new()).await;
assert_eq!(
outcome.log.lock().expect("mock log poisoned").control_links,
vec![Some("push2.example.com:8080".to_owned())]
);
assert_eq!(
outcome
.bound()
.first()
.and_then(|info| info.control_link.clone()),
Some("push2.example.com:8080".to_owned())
);
}
#[tokio::test(start_paused = true)]
async fn test_unsubscribe_removes_the_subscription_from_the_desired_set() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::AwaitControls(2),
line("REQOK,2"),
line("UNSUB,1"),
Step::Fail("connection reset"),
];
let second = vec![line("CONERR,20,Specified session not found"), Step::End];
let third = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second, third],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Unsubscribe {
key: SubscriptionKey(1),
},
],
)
.await;
let adds = outcome
.controls()
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.count();
assert_eq!(adds, 1, "an unsubscribed item is not re-established");
assert!(outcome.resubscribes().is_empty());
}
fn adds(outcome: &Outcome) -> usize {
outcome
.controls()
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.count()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn test_a_flooding_server_cannot_starve_a_destroy() {
let script = vec![
line(F1_CONOK),
Step::FloodUntilControls("PROBE", 1),
line("END,31,session destroyed"),
];
let transport = MockTransport::new(websocket(), vec![script]);
let log = Arc::clone(&transport.log);
let (driver, handle, events) = connect(transport, options()).expect("valid options");
handle
.destroy(None)
.await
.expect("the driver has not started yet");
let running = tokio::spawn(driver.run());
tokio::time::timeout(Duration::from_secs(10), running)
.await
.expect("an unbounded line stream must not postpone a command for ever")
.expect("the driver task did not panic");
let controls = log.lock().expect("mock log poisoned").control_parameters();
assert!(
controls
.iter()
.any(|parameters| parameters.contains("LS_op=destroy")),
"the destroy reached the wire: {controls:?}"
);
drop((handle, events));
}
#[tokio::test(start_paused = true)]
async fn test_a_stop_is_honoured_while_the_event_stream_is_full() {
let mut options = options();
options.event_capacity = std::num::NonZeroUsize::MIN;
let script = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("SUBOK,1,1,1"),
Step::FloodUntilControls("U,1,1,20.4", usize::MAX),
];
let transport = MockTransport::new(websocket(), vec![script]);
let log = Arc::clone(&transport.log);
let (driver, handle, events) = connect(transport, options).expect("valid options");
handle
.commands
.send(subscribe())
.await
.expect("the driver has not started yet");
let running = tokio::spawn(driver.run());
tokio::time::sleep(Duration::from_millis(50)).await;
handle.stop();
let closed = tokio::time::timeout(Duration::from_secs(5), running)
.await
.expect("a stop must not be subject to backpressure")
.expect("the driver task did not panic");
assert!(matches!(closed, SessionClosed::ByClient { .. }));
assert!(
log.lock().expect("mock log poisoned").closes >= 1,
"the transport is closed on every exit path"
);
drop(events);
}
#[tokio::test(start_paused = true)]
async fn test_an_establishment_that_hangs_is_abandoned_within_the_budget() {
let mut options = options();
options.open_timeout = Duration::from_secs(3);
let transport = MockTransport::new(websocket(), Vec::new()).hanging_open();
let outcome = tokio::time::timeout(
Duration::from_secs(120),
drive_transport(transport, options, Vec::new()),
)
.await
.expect("a hanging dial must not hang the session");
assert!(matches!(
outcome.closed,
SessionClosed::RetriesExhausted { .. }
));
assert!(
outcome
.unbinds()
.iter()
.all(|reason| matches!(reason, UnbindReason::ConnectionFailed { .. })),
"every attempt timed out: {:?}",
outcome.unbinds()
);
}
#[tokio::test(start_paused = true)]
async fn test_a_stop_is_honoured_while_a_dial_is_hanging() {
let mut options = options();
options.open_timeout = Duration::from_secs(3600);
let transport = MockTransport::new(websocket(), Vec::new()).hanging_open();
let (driver, handle, events) = connect(transport, options).expect("valid options");
let running = tokio::spawn(driver.run());
tokio::time::sleep(Duration::from_millis(50)).await;
handle.stop();
let closed = tokio::time::timeout(Duration::from_secs(5), running)
.await
.expect("a stop must interrupt an establishment")
.expect("the driver task did not panic");
assert!(matches!(closed, SessionClosed::ByClient { .. }));
drop(events);
}
#[tokio::test(start_paused = true)]
async fn test_repeated_heartbeat_failures_give_the_connection_up() {
let options = options().with_connection(ConnectionMode::Streaming {
inactivity_millis: Some(2000),
keepalive_millis: None,
send_sync: None,
});
let first = vec![line(F1_CONOK), Step::Silence];
let second = vec![line("CONERR,1,User/password check failed")];
let transport = MockTransport::new(websocket(), vec![first, second])
.failing_controls(usize::MAX);
let outcome = tokio::time::timeout(
Duration::from_secs(600),
drive_transport(transport, options, Vec::new()),
)
.await
.expect("a failing heartbeat must not spin for ever");
assert!(
outcome
.unbinds()
.iter()
.any(|reason| matches!(reason, UnbindReason::ConnectionFailed { .. })),
"{:?}",
outcome.unbinds()
);
let heartbeats = outcome
.controls()
.iter()
.filter(|parameters| parameters.contains("LS_session"))
.count();
assert!(
heartbeats <= usize::try_from(MAX_HEARTBEAT_FAILURES).unwrap_or(usize::MAX),
"at most {MAX_HEARTBEAT_FAILURES} attempts before giving up, got {heartbeats}"
);
}
#[tokio::test]
async fn test_identifier_spaces_are_exhausted_rather_than_reused() {
let (mut driver, handle, events) =
connect(MockTransport::new(websocket(), Vec::new()), options()).expect("valid options");
driver.next_request_id = u64::MAX;
assert!(
driver.next_request_id().is_none(),
"the last id is never handed out, because it cannot be advanced past"
);
handle.next_key.store(u64::MAX, Ordering::Relaxed);
assert!(matches!(
handle.allocate_key(),
Err(SessionError::Exhausted { .. })
));
drop((driver, events));
}
#[tokio::test(start_paused = true)]
async fn test_an_unsolicited_prog_cannot_move_the_recovery_baseline() {
let first = vec![
line(F1_CONOK),
line("U,1,1,20.4"),
line("PROG,1000"),
Step::Fail("connection reset"),
];
let second = vec![line(F1_CONOK), line(SERVER_END)];
let outcome = drive(websocket(), options(), vec![first, second], Vec::new()).await;
assert_eq!(
outcome.opens().get(1),
Some(&OpenRecord::Bind {
session: Some(F1_SESSION.to_owned()),
recovery_from: Some(1),
}),
"the recovery asks to resume from what was actually counted"
);
assert!(
outcome.recoveries().is_empty(),
"an unsolicited PROG reports no recovery outcome"
);
}
#[tokio::test(start_paused = true)]
async fn test_repeated_session_refreshes_are_delayed_and_eventually_given_up() {
let refresh = || vec![line(F1_CONOK), line("END,48,please reconnect")];
let scripts = (0..12).map(|_| refresh()).collect();
let outcome = tokio::time::timeout(
Duration::from_secs(600),
drive(websocket(), options(), scripts, Vec::new()),
)
.await
.expect("a refresh loop must terminate");
let delays: Vec<Option<Duration>> = outcome
.events
.iter()
.filter_map(|event| match event {
SessionEvent::Unbound { retry_in, .. } => Some(*retry_in),
_ => None,
})
.collect();
assert_eq!(
delays.first(),
Some(&Some(Duration::ZERO)),
"the first refresh is immediate, as the spec prescribes"
);
assert!(
delays
.iter()
.skip(1)
.any(|delay| matches!(delay, Some(delay) if !delay.is_zero())),
"later refreshes are delayed: {delays:?}"
);
assert!(
matches!(outcome.closed, SessionClosed::RetriesExhausted { .. }),
"a server that only ever refreshes is a definitive loss, got {:?}",
outcome.closed
);
}
#[tokio::test(start_paused = true)]
async fn test_a_productive_session_clears_the_refresh_streak() {
let productive = || {
vec![
line(F1_CONOK),
line("U,1,1,20.4"),
line("END,48,please reconnect"),
]
};
let scripts = vec![
productive(),
productive(),
productive(),
vec![line(F1_CONOK), line(SERVER_END)],
];
let outcome = drive(websocket(), options(), scripts, Vec::new()).await;
let delays: Vec<Option<Duration>> = outcome
.events
.iter()
.filter_map(|event| match event {
SessionEvent::Unbound { retry_in, .. } => Some(*retry_in),
_ => None,
})
.collect();
assert!(
delays.iter().all(|delay| *delay == Some(Duration::ZERO)),
"every refresh follows a productive session: {delays:?}"
);
assert!(matches!(outcome.closed, SessionClosed::ByServer { .. }));
}
#[tokio::test(start_paused = true)]
async fn test_a_late_rejection_from_a_replaced_session_changes_nothing() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("END,48,please reconnect"),
];
let second = vec![
line(F6_CONOK),
Step::AwaitControls(2),
line("REQERR,1,19,Specified subscription not found"),
line("U,1,1,20.4"),
line("END,48,please reconnect"),
];
let third = vec![line(F8_CONOK), Step::AwaitControls(3), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second, third],
vec![SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
}],
)
.await;
assert_eq!(
adds(&outcome),
3,
"the subscription survived a refusal aimed at a session that no longer exists"
);
assert!(
control_targets(&outcome).iter().any(|(target, outcome)| {
matches!(
target,
ControlTarget::Subscription {
operation: SubscriptionOperation::Subscribe,
..
}
) && matches!(outcome, ControlOutcome::NotSent { .. })
}),
"a retired request reports exactly one outcome: {:?}",
control_targets(&outcome)
);
}
#[tokio::test(start_paused = true)]
async fn test_a_message_pending_across_a_replacement_still_gets_one_outcome() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("END,48,please reconnect"),
];
let second = vec![line(F6_CONOK), line(SERVER_END)];
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![SessionCommand::SendMessage {
message: Box::new(OutgoingMessage::numbered("BUY 100", NonZeroU32::MIN)),
}],
)
.await;
let reported = messages(&outcome);
assert_eq!(reported.len(), 1, "exactly one report: {reported:?}");
assert!(matches!(
reported.first().map(|(_, _, _, result)| result),
Some(MessageResult::NotSent { .. })
));
}
#[tokio::test(start_paused = true)]
async fn test_an_accepted_reconfiguration_survives_a_recreated_session() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::AwaitControls(2),
line("REQOK,2"),
line("END,48,please reconnect"),
];
let second = vec![line(F6_CONOK), Step::AwaitControls(3), line(SERVER_END)];
let limit = MaxFrequencyLimit::Limited(
crate::protocol::request::DecimalNumber::try_new("2.5").expect("a decimal number"),
);
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Reconfigure {
key: SubscriptionKey(1),
max_frequency: limit,
},
],
)
.await;
let controls = outcome.controls();
let reissued = controls
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.nth(1)
.expect("the subscription was re-issued on the new session");
assert!(
reissued.contains("LS_requested_max_frequency=2.5"),
"the recreated subscription keeps the frequency the server granted: {reissued}"
);
}
#[tokio::test(start_paused = true)]
async fn test_a_refused_reconfiguration_leaves_the_desired_state_alone() {
let first = vec![
line(F1_CONOK),
Step::AwaitControls(1),
line("REQOK,1"),
line("SUBOK,1,1,3"),
Step::AwaitControls(2),
line("REQERR,2,25,Subscription is not unfiltered"),
line("END,48,please reconnect"),
];
let second = vec![line(F6_CONOK), Step::AwaitControls(3), line(SERVER_END)];
let limit = MaxFrequencyLimit::Limited(
crate::protocol::request::DecimalNumber::try_new("2.5").expect("a decimal number"),
);
let outcome = drive(
websocket(),
options(),
vec![first, second],
vec![
SessionCommand::Subscribe {
key: SubscriptionKey(1),
spec: Box::new(spec()),
},
SessionCommand::Reconfigure {
key: SubscriptionKey(1),
max_frequency: limit,
},
],
)
.await;
let controls = outcome.controls();
let reissued = controls
.iter()
.filter(|parameters| parameters.contains("LS_op=add"))
.nth(1)
.expect("the subscription was re-issued on the new session");
assert!(
!reissued.contains("LS_requested_max_frequency"),
"a refused change is not desired state: {reissued}"
);
}
}