#![allow(dead_code)]
use std::collections::{HashMap, VecDeque};
#[cfg(feature = "sync-sender-qwp-ws")]
use std::io::Write;
use std::sync::Arc;
use std::sync::atomic::{AtomicU8, AtomicUsize, Ordering};
use std::time::{Duration, Instant};
#[cfg(feature = "sync-sender-qwp-ws")]
use rand::Rng;
use crate::error;
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::conf::{QwpWsConfig, QwpWsEndpoint};
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ingress::tls::TlsSettings;
use crate::{Error, ErrorCode};
#[cfg(feature = "sync-sender-qwp-ws")]
use super::qwp_ws::{
QwpWsConnectKind, QwpWsConnectRoundSuccess, QwpWsHostHealthTracker, TrafficGate, WsFrameRead,
WsFrameReader, WsStream, connect_qwp_ws_endpoint_round, qwp_ws_configured_endpoints,
write_binary_frame, write_ping_frame,
};
use super::qwp_ws_codec::{self as codec, PipelinedResponse};
use super::qwp_ws_ownership::{QwpWsErrorCategory, QwpWsErrorPolicy, QwpWsSenderError};
use super::qwp_ws_queue::{
OutboundFrame, OutboundFrameView, QueueError, QwpReceipt, QwpReceiptStatus, SentFrame,
};
use super::qwp_ws_sfa_catchup::{
CatchUpEntryTooLarge, CatchUpFrameBuildError, CatchUpStreamError, SentDictMirror,
frame_delta_start,
};
#[cfg(test)]
use super::qwp_ws_sfa_queue::SfaMemoryQueueOptions;
use super::qwp_ws_sfa_queue::{
SfaCleanupFailure, SfaFrameQueue, SfaProducer, SfaProgressView, SfaSendCursor,
SfaStorageFinish, SfaStorageResult, SfaStorageStep,
};
#[cfg(feature = "sync-sender-qwp-ws")]
use crate::ws::frame::OPCODE_BINARY;
pub(crate) const DEFAULT_EVENT_CAPACITY: usize = 1024;
pub(crate) const DEFAULT_MAX_FRAME_REJECTIONS: usize = 4;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct ReconnectPolicy {
max_duration: Duration,
initial_backoff: Duration,
max_backoff: Duration,
}
impl ReconnectPolicy {
pub(crate) fn bounded(
max_duration: Duration,
initial_backoff: Duration,
max_backoff: Duration,
) -> Self {
Self {
max_duration,
initial_backoff,
max_backoff,
}
}
fn no_backoff(max_duration: Duration) -> Self {
Self {
max_duration,
initial_backoff: Duration::ZERO,
max_backoff: Duration::ZERO,
}
}
pub(crate) fn max_duration(&self) -> Duration {
self.max_duration
}
pub(crate) fn initial_backoff(&self) -> Duration {
self.initial_backoff
}
pub(crate) fn max_backoff(&self) -> Duration {
self.max_backoff
}
}
#[derive(Debug, Default, Clone, Copy)]
struct PoisonFrameTracker {
fsn: Option<u64>,
completed_fsn: Option<u64>,
rejection_count: usize,
first_strike_at: Option<Instant>,
}
impl PoisonFrameTracker {
fn clear(&mut self) {
*self = Self::default();
}
fn record_failure(
&mut self,
fsn: u64,
completed_fsn: Option<u64>,
limit: usize,
min_escalation_window: Duration,
now: Instant,
) -> bool {
if self.fsn == Some(fsn) && self.completed_fsn == completed_fsn {
self.rejection_count = self.rejection_count.saturating_add(1);
} else {
self.fsn = Some(fsn);
self.completed_fsn = completed_fsn;
self.rejection_count = 1;
self.first_strike_at = Some(now);
}
if self.rejection_count < limit {
return false;
}
min_escalation_window.is_zero()
|| self.first_strike_at.is_some_and(|first_strike_at| {
now.saturating_duration_since(first_strike_at) >= min_escalation_window
})
}
fn strikes(&self) -> usize {
self.rejection_count
}
}
#[cfg(test)]
#[derive(Debug)]
pub(crate) struct QwpWsCoreTestHarness<Q, T> {
store: QwpWsPublicationStore<Q>,
send_core: QwpWsSendCore<T>,
}
#[derive(Debug)]
pub(crate) struct QwpWsSendCore<T> {
transport: T,
send_cursor: SendCursor,
dict_mirror: SentDictMirror,
catch_up_pending: bool,
catch_up_retry_strikes: usize,
durable_ack: Option<DurableAckTracker>,
reconnect_policy: ReconnectPolicy,
pending_reconnect: Option<QwpWsReconnectState>,
poison_tracker: PoisonFrameTracker,
max_frame_rejections: usize,
poison_min_escalation_window: Duration,
zero_progress_role_recycles: usize,
completed_at_last_role_recycle: Option<u64>,
sends_on_connection: u64,
}
#[derive(Debug)]
pub(crate) enum QwpWsSendProgress {
Outcome(DriveOutcome),
TransportFailure(TransportFailure),
}
#[derive(Debug)]
pub(crate) enum QwpWsHotSendProgress {
NoResponse {
frame: SentFrame,
replayed: bool,
},
Response {
frame: SentFrame,
replayed: bool,
response: TransportResponse,
},
TransportFailure {
frame: SentFrame,
replayed: bool,
failure: TransportFailure,
},
}
#[derive(Debug)]
pub(crate) struct QwpWsHotResponseProgress {
pub(crate) outcome: DriveOutcome,
pub(crate) events: Vec<DriverEvent>,
pub(crate) ok_fsn: Option<u64>,
}
impl QwpWsHotResponseProgress {
fn idle() -> Self {
Self {
outcome: DriveOutcome::Idle,
events: Vec::new(),
ok_fsn: None,
}
}
fn from_optional_event(outcome: DriveOutcome, event: Option<DriverEvent>) -> Self {
let events = event.into_iter().collect();
Self {
outcome,
events,
ok_fsn: None,
}
}
fn with_ok_fsn(mut self, fsn: u64) -> Self {
self.ok_fsn = Some(fsn);
self
}
}
#[derive(Debug)]
pub(crate) enum QwpWsTransportFailureAction {
Reconnect {
reason: ReconnectReason,
initial_error: Error,
pace: Duration,
},
Terminal(Error),
}
#[derive(Debug)]
pub(crate) enum QwpWsReconnectStep {
Reconnected { reason: ReconnectReason },
RetryAfter { sleep_for: Duration },
Terminal(Error),
}
#[derive(Debug)]
pub(crate) struct QwpWsReconnectState {
policy: ReconnectPolicy,
context: &'static str,
reason: ReconnectReason,
started: Instant,
deadline: Option<Instant>,
backoff: Duration,
pace_first_attempt: Option<Duration>,
last_error: Error,
attempts: usize,
}
impl QwpWsReconnectState {
fn new(
policy: ReconnectPolicy,
context: &'static str,
reason: ReconnectReason,
initial_error: Error,
) -> Self {
let started = Instant::now();
Self {
policy,
context,
reason,
started,
deadline: started.checked_add(policy.max_duration),
backoff: policy.initial_backoff,
pace_first_attempt: None,
last_error: initial_error,
attempts: 0,
}
}
pub(crate) fn with_pace(mut self, pace: Duration) -> Self {
if !pace.is_zero() {
self.deadline = self
.deadline
.and_then(|deadline| deadline.checked_add(pace));
self.pace_first_attempt = Some(pace);
}
self
}
pub(crate) fn deadline(&self) -> Option<Instant> {
self.deadline
}
pub(crate) fn initial_backoff(&self) -> Duration {
self.policy.initial_backoff
}
fn deadline_expired(&self) -> bool {
reconnect_deadline_expired(self.deadline)
}
pub(crate) fn retry_budget_exhausted_error(&self) -> Error {
retry_budget_exhausted_error(
self.context,
self.attempts,
self.started,
Some(self.last_error.clone()),
)
}
pub(crate) fn next_after_retryable_terminal(&self, err: Error) -> Self {
Self::new(self.policy, self.context, self.reason, err)
}
pub(crate) fn take_first_attempt_pace(&mut self) -> Option<Duration> {
self.pace_first_attempt.take()
}
fn record_retryable_error(&mut self, err: Error) -> Duration {
let role_reject = is_qwp_ws_role_reject_error(&err);
self.last_error = err;
let sleep_for =
reconnect_sleep_duration(role_reject, self.policy.initial_backoff, self.backoff);
self.backoff = if role_reject {
self.policy.initial_backoff
} else {
double_duration(self.backoff).min(self.policy.max_backoff)
};
sleep_for
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum PublicationState {
Open,
Closing,
Terminal,
}
#[derive(Debug, Clone)]
pub(crate) struct PublicationLifecycle {
state: Arc<AtomicU8>,
}
const PUBLICATION_OPEN: u8 = 0;
const PUBLICATION_CLOSING: u8 = 1;
const PUBLICATION_TERMINAL: u8 = 2;
impl PublicationState {
fn from_raw(state: u8) -> Self {
match state {
PUBLICATION_OPEN => Self::Open,
PUBLICATION_CLOSING => Self::Closing,
PUBLICATION_TERMINAL => Self::Terminal,
_ => Self::Terminal,
}
}
}
impl PublicationLifecycle {
fn new() -> Self {
Self {
state: Arc::new(AtomicU8::new(PUBLICATION_OPEN)),
}
}
pub(crate) fn load(&self) -> PublicationState {
PublicationState::from_raw(self.state.load(Ordering::Acquire))
}
pub(crate) fn begin_close(&self) {
let _ = self.state.compare_exchange(
PUBLICATION_OPEN,
PUBLICATION_CLOSING,
Ordering::AcqRel,
Ordering::Acquire,
);
}
pub(crate) fn terminalize(&self) -> PublicationState {
PublicationState::from_raw(self.state.swap(PUBLICATION_TERMINAL, Ordering::AcqRel))
}
pub(crate) fn is_terminal(&self) -> bool {
self.load() == PublicationState::Terminal
}
}
#[derive(Debug, Default, Clone, Copy)]
pub(crate) struct QwpWsCounters {
pub total_frames_sent: u64,
pub total_frames_replayed: u64,
pub total_acks: u64,
pub total_reconnect_attempts: u64,
pub total_reconnects_succeeded: u64,
pub total_server_errors: u64,
}
impl From<QwpWsCounters> for super::qwp_ws_ownership::QwpWsTotals {
fn from(counters: QwpWsCounters) -> Self {
Self {
frames_sent: counters.total_frames_sent,
frames_replayed: counters.total_frames_replayed,
acks: counters.total_acks,
reconnect_attempts: counters.total_reconnect_attempts,
reconnects_succeeded: counters.total_reconnects_succeeded,
server_errors: counters.total_server_errors,
}
}
}
#[derive(Debug)]
pub(crate) struct QwpWsPublicationStore<Q = SfaFrameQueue> {
queue: Q,
events: DriverEventRing,
lifecycle: PublicationLifecycle,
terminal_error: Option<Error>,
terminal_sender_error: Option<QwpWsSenderError>,
last_server_error: Option<QwpServerError>,
rejected_frames: VecDeque<QwpRejectedFrame>,
sender_errors: SenderErrorLog,
rejection_sink: Option<Arc<crate::ingress::rejection_events::RejectionEventSource>>,
counters: QwpWsCounters,
}
impl<Q: PublicationLog> QwpWsPublicationStore<Q> {
pub(crate) fn new(queue: Q, event_capacity: usize) -> Self {
Self {
queue,
events: DriverEventRing::new(event_capacity),
lifecycle: PublicationLifecycle::new(),
terminal_error: None,
terminal_sender_error: None,
last_server_error: None,
rejected_frames: VecDeque::new(),
sender_errors: SenderErrorLog::new(event_capacity),
rejection_sink: None,
counters: QwpWsCounters::default(),
}
}
pub(crate) fn set_rejection_sink(
&mut self,
sink: Option<Arc<crate::ingress::rejection_events::RejectionEventSource>>,
) {
self.rejection_sink = sink;
}
pub(crate) fn counters(&self) -> QwpWsCounters {
self.counters
}
pub(crate) fn record_reconnect_attempt(&mut self) {
self.counters.total_reconnect_attempts += 1;
}
pub(crate) fn lifecycle(&self) -> PublicationLifecycle {
self.lifecycle.clone()
}
pub(crate) fn check_durability(&self) -> Result<(), DriverError> {
self.queue.check_durability()
}
pub(crate) fn storage_maintenance_in_flight(&self) -> Result<bool, DriverError> {
self.queue.storage_maintenance_in_flight()
}
pub(crate) fn try_submit(&mut self, payload: &[u8]) -> Result<QwpReceipt, DriverError> {
match self.lifecycle.load() {
PublicationState::Open => {}
PublicationState::Closing => return Err(DriverError::Closing),
PublicationState::Terminal => return Err(DriverError::Terminal),
}
let receipt = self.queue.try_publish(payload)?;
self.push_event(DriverEvent::Published { fsn: receipt.fsn });
Ok(receipt)
}
pub(crate) fn take_producer(&mut self) -> Option<SfaProducer> {
self.queue.take_producer()
}
pub(crate) fn progress_view(&self) -> SfaProgressView {
self.queue.progress_view()
}
pub(crate) fn record_sent_frame(
&mut self,
send_cursor: &mut SendCursor,
frame: SentFrame,
) -> Result<(), DriverError> {
let replayed = send_cursor.commit_sent(frame)?;
self.record_sent_event(frame, replayed);
Ok(())
}
pub(crate) fn record_sent_event(&mut self, frame: SentFrame, replayed: bool) {
self.counters.total_frames_sent += 1;
if replayed {
self.counters.total_frames_replayed += 1;
}
self.push_event(DriverEvent::Sent {
fsn: frame.fsn,
wire_seq: frame.wire_seq,
});
}
pub(crate) fn record_completed_through_event(&mut self, fsn: u64, wire_seq: u64) {
self.queue.persist_completed_fsn(fsn);
self.push_event(DriverEvent::CompletedThrough { fsn, wire_seq });
}
pub(crate) fn record_driver_event(&mut self, event: DriverEvent) {
self.push_event(event);
}
pub(crate) fn mark_terminal(&mut self, error: Option<Error>) {
self.try_mark_terminal(error);
}
fn try_mark_terminal(&mut self, error: Option<Error>) -> bool {
if self.lifecycle.load() == PublicationState::Terminal {
return false;
}
self.terminal_error = error;
let previous = self.lifecycle.terminalize();
debug_assert_ne!(previous, PublicationState::Terminal);
self.push_event(DriverEvent::Terminal);
true
}
pub(crate) fn is_terminal(&self) -> bool {
self.lifecycle.is_terminal()
}
fn receipt_status(
&self,
send_cursor: &SendCursor,
durable_ack: Option<&DurableAckTracker>,
receipt: QwpReceipt,
) -> QwpReceiptStatus {
let status = self.queue.receipt_status(receipt);
if self.is_terminal() && status.is_pending() {
return QwpReceiptStatus::Terminal { fsn: receipt.fsn };
}
if matches!(status, QwpReceiptStatus::Published { .. })
&& let Some(wire_seq) = send_cursor.wire_seq_for_fsn(receipt.fsn)
{
return QwpReceiptStatus::Sent {
fsn: receipt.fsn,
wire_seq,
};
}
if matches!(status, QwpReceiptStatus::Published { .. })
&& let Some(wire_seq) =
durable_ack.and_then(|tracker| tracker.pending_wire_seq_for_fsn(receipt.fsn))
{
return QwpReceiptStatus::Sent {
fsn: receipt.fsn,
wire_seq,
};
}
status
}
pub(crate) fn take_storage_maintenance_step(
&mut self,
) -> Result<Option<SfaStorageStep>, DriverError> {
self.queue
.take_storage_maintenance_step(self.lifecycle.load() == PublicationState::Open)
}
pub(crate) fn finish_storage_maintenance(
&mut self,
result: SfaStorageResult,
) -> Result<SfaStorageFinish, DriverError> {
self.queue
.finish_storage_maintenance(result, self.lifecycle.load() == PublicationState::Open)
}
pub(crate) fn complete_storage_maintenance(&mut self) -> Result<(), DriverError> {
self.queue.complete_storage_maintenance()
}
pub(crate) fn record_storage_cleanup_failure(
&mut self,
failure: SfaCleanupFailure,
) -> Result<(), DriverError> {
self.queue.record_storage_cleanup_failure(failure)
}
fn clear_unresolved_rejected_frames(&mut self) {
let completed_fsn = self.queue.completed_fsn();
self.rejected_frames.retain(|rejected| {
completed_fsn.is_some_and(|completed_fsn| rejected.fsn <= completed_fsn)
});
}
fn record_rejected_frame(
&mut self,
fsn: u64,
wire_seq: u64,
error: QwpServerError,
policy: QwpWsErrorPolicy,
) -> QwpWsSenderError {
self.last_server_error = Some(error.clone());
let sender_error = sender_error_for_qwp_error(&error, wire_seq, fsn, policy);
if self.sender_errors.capacity() != 0 {
if self.rejected_frames.len() == self.sender_errors.capacity() {
self.rejected_frames.pop_front();
}
self.rejected_frames.push_back(QwpRejectedFrame {
fsn,
wire_seq,
error,
});
}
self.push_sender_error(sender_error.clone());
sender_error
}
fn record_reject_error(
&mut self,
fsn: u64,
wire_seq: u64,
error: QwpServerError,
policy: QwpWsErrorPolicy,
) -> QwpWsSenderError {
let sender_error = sender_error_for_qwp_error(&error, wire_seq, fsn, policy);
self.last_server_error = Some(error);
self.push_sender_error(sender_error.clone());
sender_error
}
fn record_terminal_sender_error(
&mut self,
sender_error: QwpWsSenderError,
terminal_error: Error,
last_server_error: Option<QwpServerError>,
) -> Error {
if self.lifecycle.load() == PublicationState::Terminal {
return self.terminal_error.clone().unwrap_or(terminal_error);
}
if let Some(last_server_error) = last_server_error {
self.last_server_error = Some(last_server_error);
}
self.terminal_sender_error = Some(sender_error.clone());
self.sender_errors.push(sender_error.clone());
let committed = self.try_mark_terminal(Some(terminal_error.clone()));
debug_assert!(committed);
if let Some(sink) = &self.rejection_sink {
sink.publish(sender_error);
}
terminal_error
}
fn delivery_status(
&self,
send_cursor: &SendCursor,
durable_ack: Option<&DurableAckTracker>,
receipt: QwpReceipt,
) -> Result<Option<DeliveryOutcome>, DriverError> {
match self.receipt_status(send_cursor, durable_ack, receipt) {
QwpReceiptStatus::Completed { .. } => Ok(Some(DeliveryOutcome::Completed)),
QwpReceiptStatus::Terminal { .. } => Ok(Some(DeliveryOutcome::Terminal)),
QwpReceiptStatus::Published { .. } | QwpReceiptStatus::Sent { .. } => Ok(None),
QwpReceiptStatus::Unknown { fsn } => Err(DriverError::UnknownReceipt { fsn }),
}
}
pub(crate) fn set_closing(&mut self) {
self.lifecycle.begin_close();
}
pub(crate) fn all_published_receipts_resolved(&self) -> bool {
match self.queue.published_fsn() {
None => true,
Some(published_fsn) => self
.queue
.completed_fsn()
.is_some_and(|completed_fsn| completed_fsn >= published_fsn),
}
}
pub(crate) fn published_fsn(&self) -> Option<u64> {
self.queue.published_fsn()
}
pub(crate) fn completed_fsn(&self) -> Option<u64> {
self.queue.completed_fsn()
}
pub(crate) fn close_queue(&mut self) -> Result<(), DriverError> {
self.queue.close()
}
pub(crate) fn record_protocol_violation(
&mut self,
close_code: Option<u16>,
reason: String,
) -> Error {
let from_fsn = self
.queue
.completed_fsn()
.map_or(0, |fsn| fsn.saturating_add(1));
let to_fsn = self.queue.published_fsn().unwrap_or(from_fsn).max(from_fsn);
let message = match (close_code, reason.is_empty()) {
(Some(close_code), true) => format!("ws-close[{close_code}]"),
(Some(close_code), false) => format!("ws-close[{close_code}]: {reason}"),
(None, true) => "WebSocket protocol violation".to_string(),
(None, false) => reason,
};
let sender_error = QwpWsSenderError {
category: QwpWsErrorCategory::ProtocolViolation,
applied_policy: QwpWsErrorPolicy::Terminal,
status: None,
message: Some(message.clone()),
message_sequence: None,
from_fsn,
to_fsn,
};
let terminal_error = error::fmt!(
ServerRejection,
"QWP/WebSocket protocol violation: {message}"
)
.with_qwp_ws_rejection(sender_error.clone());
self.record_terminal_sender_error(sender_error, terminal_error, None)
}
pub(crate) fn poll_sender_error(&mut self) -> Option<QwpWsSenderError> {
self.sender_errors.poll()
}
pub(crate) fn poll_sender_error_notification(&mut self) -> Option<QwpWsSenderError> {
self.sender_errors.poll_notification()
}
pub(crate) fn sender_errors_dropped_total(&self) -> u64 {
self.sender_errors.dropped_total()
}
pub(crate) fn poll_event(&mut self) -> Option<DriverEvent> {
self.events.pop()
}
pub(crate) fn events_dropped_total(&self) -> u64 {
self.events.dropped_total()
}
pub(crate) fn terminal_error(&self) -> Option<&Error> {
self.terminal_error.as_ref()
}
pub(crate) fn terminal_sender_error(&self) -> Option<&QwpWsSenderError> {
self.terminal_sender_error.as_ref()
}
pub(crate) fn last_server_error(&self) -> Option<&QwpServerError> {
self.last_server_error.as_ref()
}
pub(crate) fn rejected_frame(&self, receipt: QwpReceipt) -> Option<&QwpRejectedFrame> {
self.rejected_frames
.iter()
.find(|rejected| rejected.fsn == receipt.fsn)
}
fn push_event(&mut self, event: DriverEvent) {
self.events.push(event);
}
fn push_sender_error(&mut self, error: QwpWsSenderError) {
self.sender_errors.push(error.clone());
if let Some(sink) = &self.rejection_sink {
sink.publish(error);
}
}
}
impl<T: QwpWsCoreTransport> QwpWsSendCore<T> {
pub(crate) fn new(transport: T, reconnect_policy: ReconnectPolicy) -> Self {
Self::new_with_durable_ack(transport, reconnect_policy, false)
}
pub(crate) fn new_with_durable_ack(
transport: T,
reconnect_policy: ReconnectPolicy,
durable_ack: bool,
) -> Self {
Self::new_with_durable_ack_and_rejection_limit(
transport,
reconnect_policy,
durable_ack,
DEFAULT_MAX_FRAME_REJECTIONS,
Duration::ZERO,
)
}
pub(crate) fn new_with_durable_ack_and_rejection_limit(
transport: T,
reconnect_policy: ReconnectPolicy,
durable_ack: bool,
max_frame_rejections: usize,
poison_min_escalation_window: Duration,
) -> Self {
Self {
transport,
send_cursor: SendCursor::new(),
dict_mirror: SentDictMirror::new(false),
catch_up_pending: false,
catch_up_retry_strikes: 0,
durable_ack: durable_ack.then(DurableAckTracker::new),
reconnect_policy,
pending_reconnect: None,
poison_tracker: PoisonFrameTracker::default(),
max_frame_rejections,
poison_min_escalation_window,
zero_progress_role_recycles: 0,
completed_at_last_role_recycle: None,
sends_on_connection: 0,
}
}
fn apply_response<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
response: TransportResponse,
defer_reconnect: bool,
) -> Result<DriveOutcome, DriverError> {
match response {
TransportResponse::Ack { wire_seq } => {
store.counters.total_acks += 1;
self.complete_ack_through(store, wire_seq)
}
TransportResponse::DurableOk {
wire_seq,
table_seq_txns,
} => {
store.counters.total_acks += 1;
if self.durable_ack.is_some() {
self.apply_durable_ok(store, wire_seq, table_seq_txns)
} else {
self.complete_ack_through(store, wire_seq)
}
}
TransportResponse::DurableAck { table_seq_txns } => {
store.counters.total_acks += 1;
let Some(tracker) = self.durable_ack.as_mut() else {
return Ok(DriveOutcome::Idle);
};
tracker.apply_ack(table_seq_txns);
self.complete_ready_durable(store)
}
TransportResponse::Reject { wire_seq, error } => {
store.counters.total_server_errors += 1;
let policy = server_error_policy(error.status);
let Some((fsn, _effect_wire_seq)) =
self.send_cursor.reject_fsn_for_wire_seq(wire_seq)?
else {
return self.record_presend_reject(
store,
wire_seq,
error,
policy,
defer_reconnect,
);
};
if self.reject_target_already_accounted(store, fsn) {
store.record_reject_error(fsn, wire_seq, error, policy);
return Ok(DriveOutcome::Idle);
}
if !self.reject_target_can_complete(store, fsn) {
store.last_server_error = Some(error.clone());
self.reject_gap_protocol_error(store, fsn);
return Ok(DriveOutcome::Terminal);
}
if policy == QwpWsErrorPolicy::Terminal {
let sender_error = sender_error_for_qwp_error(&error, wire_seq, fsn, policy);
let terminal_error =
server_rejection_error(error.error.clone(), sender_error.clone());
store.record_terminal_sender_error(sender_error, terminal_error, Some(error));
return Ok(DriveOutcome::Terminal);
}
if policy != QwpWsErrorPolicy::RetriableOther
&& self.rejected_head_is_poison(store, fsn)
{
let strikes = self.poison_tracker.strikes();
let reason = if error.message.is_empty() {
format!(
"QWP/WebSocket frame fsn {fsn} was rejected {strikes} times without ACK progress"
)
} else {
format!(
"QWP/WebSocket frame fsn {fsn} was rejected {} times without ACK \
progress; last server error: {}",
strikes, error.message
)
};
store.last_server_error = Some(error);
store.record_protocol_violation(None, reason);
return Ok(DriveOutcome::Terminal);
}
let error_for_reconnect = error.error.clone();
let reconnect_reason = reconnect_reason_for_policy(policy);
let pace = self.reconnect_pace_for_reject_policy(store, policy);
let sender_error = store.record_rejected_frame(fsn, wire_seq, error, policy);
store.push_event(DriverEvent::Rejected { fsn, wire_seq });
let initial_error = server_rejection_error(error_for_reconnect, sender_error);
self.pending_reconnect = Some(
self.begin_reconnect(
"QWP/WebSocket reconnect after server rejection",
reconnect_reason,
initial_error,
)
.with_pace(pace),
);
self.continue_or_defer_reconnect(store, defer_reconnect)
}
}
}
fn reject_target_already_accounted<Q: PublicationLog>(
&self,
store: &QwpWsPublicationStore<Q>,
fsn: u64,
) -> bool {
self.durable_ack
.as_ref()
.is_some_and(|tracker| tracker.pending_wire_seq_for_fsn(fsn).is_some())
|| store
.queue
.completed_fsn()
.is_some_and(|completed_fsn| fsn <= completed_fsn)
}
fn reject_target_can_complete<Q: PublicationLog>(
&self,
store: &QwpWsPublicationStore<Q>,
fsn: u64,
) -> bool {
let Some(oldest) = store.queue.oldest_unresolved_fsn() else {
return false;
};
if fsn == oldest {
return true;
}
self.durable_ack
.as_ref()
.is_some_and(|tracker| tracker.pending_prefix_covers(oldest, fsn))
}
fn reject_gap_protocol_error<Q: PublicationLog>(
&self,
store: &mut QwpWsPublicationStore<Q>,
fsn: u64,
) -> Error {
let oldest = store.queue.oldest_unresolved_fsn();
store.record_protocol_violation(
None,
match oldest {
Some(oldest) => format!(
"QWP/WebSocket reject response for fsn {fsn} skipped unresolved fsn {oldest}"
),
None => {
format!("QWP/WebSocket reject response for fsn {fsn} has no unresolved frame")
}
},
)
}
fn rejected_head_is_poison<Q: PublicationLog>(
&mut self,
store: &QwpWsPublicationStore<Q>,
fsn: u64,
) -> bool {
if store.queue.oldest_unresolved_fsn() != Some(fsn) {
return false;
}
let now = Instant::now();
self.poison_tracker.record_failure(
fsn,
store.queue.completed_fsn(),
self.max_frame_rejections,
self.poison_min_escalation_window,
now,
)
}
fn record_presend_reject<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
wire_seq: u64,
error: QwpServerError,
policy: QwpWsErrorPolicy,
defer_reconnect: bool,
) -> Result<DriveOutcome, DriverError> {
let from_fsn = store
.queue
.completed_fsn()
.map_or(0, |fsn| fsn.saturating_add(1));
let to_fsn = store
.queue
.published_fsn()
.unwrap_or(from_fsn)
.max(from_fsn);
let sender_error =
sender_error_for_qwp_error_span(&error, wire_seq, from_fsn, to_fsn, policy);
if policy == QwpWsErrorPolicy::Terminal {
let terminal_error = server_rejection_error(error.error.clone(), sender_error.clone());
store.record_terminal_sender_error(sender_error, terminal_error, Some(error));
Ok(DriveOutcome::Terminal)
} else {
if policy != QwpWsErrorPolicy::RetriableOther
&& let Some(oldest) = store.queue.oldest_unresolved_fsn()
&& self.rejected_head_is_poison(store, oldest)
{
let strikes = self.poison_tracker.strikes();
let reason = if error.message.is_empty() {
format!(
"QWP/WebSocket symbol-dictionary catch-up before fsn {oldest} was \
rejected {strikes} times without progress"
)
} else {
format!(
"QWP/WebSocket symbol-dictionary catch-up before fsn {oldest} was \
rejected {} times without progress; last server error: {}",
strikes, error.message
)
};
store.last_server_error = Some(error);
store.push_sender_error(sender_error);
store.record_protocol_violation(None, reason);
return Ok(DriveOutcome::Terminal);
}
let initial_error = server_rejection_error(error.error.clone(), sender_error.clone());
store.last_server_error = Some(error);
store.push_sender_error(sender_error);
let pace = self.reconnect_pace_for_reject_policy(store, policy);
self.pending_reconnect = Some(
self.begin_reconnect(
"QWP/WebSocket reconnect after server rejection",
reconnect_reason_for_policy(policy),
initial_error,
)
.with_pace(pace),
);
self.continue_or_defer_reconnect(store, defer_reconnect)
}
}
fn reconnect_pace_for_reject_policy<Q: PublicationLog>(
&mut self,
store: &QwpWsPublicationStore<Q>,
policy: QwpWsErrorPolicy,
) -> Duration {
match policy {
QwpWsErrorPolicy::Retriable => {
self.reconnect_pace_for_strikes(self.poison_tracker.strikes().max(1))
}
QwpWsErrorPolicy::RetriableOther => self.role_recycle_pace(store),
QwpWsErrorPolicy::Terminal => Duration::ZERO,
}
}
fn role_recycle_pace<Q: PublicationLog>(
&mut self,
store: &QwpWsPublicationStore<Q>,
) -> Duration {
let completed = store.queue.completed_fsn();
if completed != self.completed_at_last_role_recycle {
self.zero_progress_role_recycles = 0;
self.completed_at_last_role_recycle = completed;
}
let level = self.zero_progress_role_recycles;
self.zero_progress_role_recycles = level.saturating_add(1);
if level == 0 {
Duration::ZERO
} else {
self.reconnect_pace_for_strikes(level)
}
}
fn reconnect_pace_for_strikes(&self, strikes: usize) -> Duration {
let dose = pace_dose(
self.reconnect_policy.initial_backoff,
self.reconnect_policy.max_backoff,
strikes,
);
pace_jitter_duration(dose)
}
fn complete_ack_through<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
wire_seq: u64,
) -> Result<DriveOutcome, DriverError> {
let Some((fsn, ack_wire_seq)) = self.send_cursor.ack_fsn_for_wire_seq(wire_seq)? else {
return Ok(DriveOutcome::Idle);
};
self.complete_through(store, fsn, ack_wire_seq)
}
fn apply_durable_ok<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
wire_seq: u64,
table_seq_txns: Vec<TableSeqTxn>,
) -> Result<DriveOutcome, DriverError> {
let Some((fsn, ack_wire_seq)) = self.send_cursor.ack_fsn_for_wire_seq(wire_seq)? else {
return Ok(DriveOutcome::Idle);
};
if self
.durable_ack
.as_ref()
.is_some_and(|tracker| tracker.pending_wire_seq_for_fsn(fsn).is_some())
{
self.send_cursor.ack_through(fsn);
return Ok(DriveOutcome::Idle);
}
if store
.queue
.completed_fsn()
.is_some_and(|completed_fsn| fsn <= completed_fsn)
{
self.send_cursor.ack_through(fsn);
return Ok(DriveOutcome::Acked {
wire_seq: ack_wire_seq,
});
}
self.send_cursor.ack_through(fsn);
let tracker = self.durable_ack.as_mut().expect("durable ACK mode");
tracker.enqueue_ok(ack_wire_seq, fsn, table_seq_txns);
self.complete_ready_durable(store)
}
fn complete_ready_durable<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
let mut last_resolved = None;
while let Some(resolved) = self
.durable_ack
.as_mut()
.and_then(DurableAckTracker::pop_ready)
{
last_resolved = Some(resolved);
}
match last_resolved {
Some(resolved) => self.complete_through(store, resolved.fsn, resolved.wire_seq),
None => Ok(DriveOutcome::Idle),
}
}
fn complete_through<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
fsn: u64,
wire_seq: u64,
) -> Result<DriveOutcome, DriverError> {
let progress = store.progress_view();
if progress.completion_reaches_published(fsn) {
self.send_cursor.release_sfa_cursor();
}
let advanced = progress
.complete_through_fsn(fsn)
.map_err(DriverError::from)?;
self.send_cursor.ack_through(fsn);
if advanced {
self.poison_tracker.clear();
store.record_completed_through_event(fsn, wire_seq);
}
Ok(DriveOutcome::Acked { wire_seq })
}
fn close_publication_queue<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<(), DriverError> {
self.send_cursor.release_sfa_cursor();
store.close_queue()
}
pub(crate) fn next_outbound_sfa_frame(
&mut self,
progress: &SfaProgressView,
) -> Result<Option<OutboundFrame>, DriverError> {
progress.next_outbound_frame(&mut self.send_cursor)
}
pub(crate) fn send_frame(
&mut self,
outbound: OutboundFrame,
) -> (SentFrame, Result<TransportSendResult, TransportFailure>) {
let frame = outbound.sent_frame();
self.sends_on_connection = self.sends_on_connection.saturating_add(1);
let transport = &mut self.transport;
let dict_mirror = &mut self.dict_mirror;
let result = outbound.with_view(|view| {
let payload = view.payload;
let sent = transport.send_frame(view);
if sent.is_ok() {
let _mirrored = dict_mirror.accumulate(payload);
}
sent
});
(frame, result)
}
pub(crate) fn guard_dict_not_torn(&self, payload: &[u8]) -> Result<(), Error> {
let Some(delta_start) = frame_delta_start(payload) else {
return Ok(()); };
if !self.dict_mirror.is_enabled() {
if delta_start > 0 {
return Err(error::fmt!(
StoreResendRequired,
"QWP/WebSocket store-and-forward: a stored frame bases at symbol \
id {} but delta symbol-dictionary mode is disabled on this \
connection, so the dictionary it depends on cannot be \
re-registered; the affected data must be resent",
delta_start
));
}
return Ok(());
}
let registered = u64::from(self.dict_mirror.count());
if delta_start > registered {
return Err(error::fmt!(
StoreResendRequired,
"QWP/WebSocket store-and-forward symbol dictionary is torn: a stored \
frame references symbol id {} but only {} dictionary entries were \
recovered (a host crash lost entries the frame depends on); the \
affected data must be resent",
delta_start,
registered
));
}
if self.dict_mirror.conflicts_with(payload) {
return Err(error::fmt!(
StoreResendRequired,
"QWP/WebSocket store-and-forward symbol dictionary is torn: a stored \
frame redefines an already-registered symbol id to a different \
symbol (the recovered dictionary disagrees with the queued frames, \
most often a host crash that tore the side-file); the affected data \
must be resent"
));
}
Ok(())
}
pub(crate) fn enable_delta_dict(&mut self, seed_entries: &[u8], seed_count: u32) {
self.dict_mirror = SentDictMirror::new(true);
let _seeded = self.dict_mirror.seed(seed_entries, seed_count);
self.catch_up_pending = !self.dict_mirror.is_empty();
}
pub(crate) fn enable_delta_dict_owned(&mut self, seed_entries: Vec<u8>, seed_count: u32) {
self.dict_mirror = SentDictMirror::new(true);
self.dict_mirror.seed_owned(seed_entries, seed_count);
self.catch_up_pending = !self.dict_mirror.is_empty();
}
fn emit_dict_catch_up(&mut self) -> Result<u64, DictCatchUpError> {
let cap = self.transport.server_max_batch_size();
let version = self.transport.negotiated_qwp_version();
let transport = &mut self.transport;
let mut wire_seq = 0u64;
self.dict_mirror
.for_each_catch_up_frame(cap, version, |frame_bytes| {
let view = OutboundFrameView {
fsn: 0,
wire_seq,
payload: frame_bytes,
};
wire_seq += 1;
transport.send_frame(view).map(|_| ())
})
.map_err(|e| match e {
CatchUpStreamError::EntryTooLarge(e) => DictCatchUpError::EntryTooLarge(e),
CatchUpStreamError::FrameBuild(e) => DictCatchUpError::FrameBuild(e),
CatchUpStreamError::Emit(failure) => DictCatchUpError::Transport(failure),
})
}
pub(crate) fn drive_catch_up(&mut self) -> Result<(), CatchUpDriveError> {
if !self.catch_up_pending {
return Ok(());
}
self.catch_up_pending = false;
match self.emit_dict_catch_up() {
Ok(catch_up_frames) => {
self.send_cursor.begin_catch_up(catch_up_frames);
self.catch_up_retry_strikes = 0;
Ok(())
}
Err(DictCatchUpError::Transport(failure)) => Err(CatchUpDriveError::Transport(failure)),
Err(DictCatchUpError::EntryTooLarge(e)) => {
let pace = self.next_catch_up_retry_pace();
Err(CatchUpDriveError::RetryConnection {
error: error::fmt!(
BatchTooLarge,
"QWP/WebSocket symbol dictionary entry ({} bytes) exceeds the server \
batch cap ({} bytes) during reconnect catch-up; queued data is \
preserved while another connection is tried",
e.entry_bytes,
e.budget
),
pace,
})
}
Err(DictCatchUpError::FrameBuild(CatchUpFrameBuildError::AllocationFailed)) => {
let pace = self.next_catch_up_retry_pace();
Err(CatchUpDriveError::RetryConnection {
error: error::fmt!(
SocketError,
"QWP/WebSocket reconnect catch-up could not allocate a \
symbol-dictionary frame; queued data is preserved while a \
fresh connection is tried"
),
pace,
})
}
Err(DictCatchUpError::FrameBuild(CatchUpFrameBuildError::PayloadTooLarge)) => {
Err(CatchUpDriveError::Terminal(error::fmt!(
BatchTooLarge,
"QWP/WebSocket reconnect catch-up could not build a symbol-dictionary \
frame because its payload exceeds the protocol limit; queued data \
is preserved but requires resend"
)))
}
}
}
fn next_catch_up_retry_pace(&mut self) -> Duration {
self.catch_up_retry_strikes = self.catch_up_retry_strikes.saturating_add(1);
self.reconnect_pace_for_strikes(self.catch_up_retry_strikes)
}
pub(crate) fn finish_send_result<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
frame: SentFrame,
send_result: TransportSendResult,
) -> Result<QwpWsSendProgress, DriverError> {
match self.finish_send_result_hot(frame, send_result)? {
QwpWsHotSendProgress::NoResponse { frame, replayed } => {
store.record_sent_event(frame, replayed);
Ok(QwpWsSendProgress::Outcome(DriveOutcome::Sent(frame)))
}
QwpWsHotSendProgress::Response {
frame,
replayed,
response,
} => {
store.record_sent_event(frame, replayed);
self.apply_response(store, response, false)
.map(QwpWsSendProgress::Outcome)
}
QwpWsHotSendProgress::TransportFailure {
frame,
replayed,
failure,
} => {
store.record_sent_event(frame, replayed);
Ok(QwpWsSendProgress::TransportFailure(failure))
}
}
}
pub(crate) fn finish_send_result_hot(
&mut self,
frame: SentFrame,
send_result: TransportSendResult,
) -> Result<QwpWsHotSendProgress, DriverError> {
let replayed = self.send_cursor.commit_sent(frame)?;
match send_result {
TransportSendResult::NoResponse => {
Ok(QwpWsHotSendProgress::NoResponse { frame, replayed })
}
TransportSendResult::Response(response) => Ok(QwpWsHotSendProgress::Response {
frame,
replayed,
response,
}),
TransportSendResult::Failure(failure) => Ok(QwpWsHotSendProgress::TransportFailure {
frame,
replayed,
failure,
}),
}
}
pub(crate) fn try_poll_response(&mut self) -> Result<TransportPoll, TransportFailure> {
self.transport.try_poll_response()
}
pub(crate) fn send_durable_ack_keepalive_if_due(
&mut self,
durable_ack_pending: bool,
) -> Result<bool, TransportFailure> {
self.transport
.send_durable_ack_keepalive_if_due(durable_ack_pending)
}
pub(crate) fn finish_response<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
response: TransportResponse,
) -> Result<DriveOutcome, DriverError> {
self.apply_response(store, response, false)
}
pub(crate) fn finish_response_defer_reconnect<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
response: TransportResponse,
) -> Result<DriveOutcome, DriverError> {
self.apply_response(store, response, true)
}
pub(crate) fn finish_ack_response_sfa(
&mut self,
progress: &SfaProgressView,
wire_seq: u64,
) -> Result<QwpWsHotResponseProgress, DriverError> {
let Some((fsn, ack_wire_seq)) = self.send_cursor.ack_fsn_for_wire_seq(wire_seq)? else {
return Ok(QwpWsHotResponseProgress::idle());
};
if progress.completion_reaches_published(fsn) {
self.send_cursor.release_sfa_cursor();
}
let advanced = progress
.complete_through_fsn(fsn)
.map_err(DriverError::from)?;
self.send_cursor.ack_through(fsn);
if advanced {
self.poison_tracker.clear();
}
let event = advanced.then_some(DriverEvent::CompletedThrough {
fsn,
wire_seq: ack_wire_seq,
});
Ok(QwpWsHotResponseProgress::from_optional_event(
DriveOutcome::Acked {
wire_seq: ack_wire_seq,
},
event,
)
.with_ok_fsn(fsn))
}
pub(crate) fn finish_durable_ok_response_sfa(
&mut self,
progress: &SfaProgressView,
wire_seq: u64,
table_seq_txns: Vec<TableSeqTxn>,
) -> Result<QwpWsHotResponseProgress, DriverError> {
if self.durable_ack.is_none() {
return self.finish_ack_response_sfa(progress, wire_seq);
}
let Some((fsn, ack_wire_seq)) = self.send_cursor.ack_fsn_for_wire_seq(wire_seq)? else {
return Ok(QwpWsHotResponseProgress::idle());
};
if self
.durable_ack
.as_ref()
.is_some_and(|tracker| tracker.pending_wire_seq_for_fsn(fsn).is_some())
{
self.send_cursor.ack_through(fsn);
return Ok(QwpWsHotResponseProgress::idle().with_ok_fsn(fsn));
}
if progress
.completed_fsn()
.is_some_and(|completed_fsn| fsn <= completed_fsn)
{
self.send_cursor.ack_through(fsn);
return Ok(QwpWsHotResponseProgress {
outcome: DriveOutcome::Acked {
wire_seq: ack_wire_seq,
},
events: Vec::new(),
ok_fsn: Some(fsn),
});
}
self.send_cursor.ack_through(fsn);
let tracker = self.durable_ack.as_mut().expect("durable ACK mode");
tracker.enqueue_ok(ack_wire_seq, fsn, table_seq_txns);
Ok(self.complete_ready_durable_sfa(progress)?.with_ok_fsn(fsn))
}
pub(crate) fn finish_durable_ack_response_sfa(
&mut self,
progress: &SfaProgressView,
table_seq_txns: Vec<TableSeqTxn>,
) -> Result<QwpWsHotResponseProgress, DriverError> {
let Some(tracker) = self.durable_ack.as_mut() else {
return Ok(QwpWsHotResponseProgress::idle());
};
tracker.apply_ack(table_seq_txns);
self.complete_ready_durable_sfa(progress)
}
fn complete_ready_durable_sfa(
&mut self,
progress: &SfaProgressView,
) -> Result<QwpWsHotResponseProgress, DriverError> {
let mut last_resolved = None;
while let Some(resolved) = self
.durable_ack
.as_mut()
.and_then(DurableAckTracker::pop_ready)
{
last_resolved = Some(resolved);
}
let Some(resolved) = last_resolved else {
return Ok(QwpWsHotResponseProgress::idle());
};
if progress.completion_reaches_published(resolved.fsn) {
self.send_cursor.release_sfa_cursor();
}
let advanced = progress
.complete_through_fsn(resolved.fsn)
.map_err(DriverError::from)?;
self.send_cursor.ack_through(resolved.fsn);
if advanced {
self.poison_tracker.clear();
}
let event = advanced.then_some(DriverEvent::CompletedThrough {
fsn: resolved.fsn,
wire_seq: resolved.wire_seq,
});
Ok(QwpWsHotResponseProgress::from_optional_event(
DriveOutcome::Acked {
wire_seq: resolved.wire_seq,
},
event,
))
}
pub(crate) fn receipt_status<Q: PublicationLog>(
&self,
store: &QwpWsPublicationStore<Q>,
receipt: QwpReceipt,
) -> QwpReceiptStatus {
store.receipt_status(&self.send_cursor, self.durable_ack.as_ref(), receipt)
}
pub(crate) fn delivery_status<Q: PublicationLog>(
&self,
store: &QwpWsPublicationStore<Q>,
receipt: QwpReceipt,
) -> Result<Option<DeliveryOutcome>, DriverError> {
store.delivery_status(&self.send_cursor, self.durable_ack.as_ref(), receipt)
}
pub(crate) fn has_pending_durable_ack(&self) -> bool {
self.durable_ack
.as_ref()
.is_some_and(DurableAckTracker::has_pending)
}
pub(crate) fn transport_failure_action<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
failure: TransportFailure,
) -> QwpWsTransportFailureAction {
match failure {
TransportFailure::Disconnect(initial_error) => QwpWsTransportFailureAction::Reconnect {
reason: ReconnectReason::Disconnect,
initial_error,
pace: Duration::ZERO,
},
TransportFailure::ServerClose(initial_error) => {
if self.sends_on_connection == 0 {
return QwpWsTransportFailureAction::Reconnect {
reason: ReconnectReason::Disconnect,
initial_error,
pace: Duration::ZERO,
};
}
if let Some(error) = self.server_close_poison_error(store) {
QwpWsTransportFailureAction::Terminal(error)
} else {
let pace = self.reconnect_pace_for_strikes(self.poison_tracker.strikes());
QwpWsTransportFailureAction::Reconnect {
reason: ReconnectReason::Disconnect,
initial_error,
pace,
}
}
}
TransportFailure::Retryable(initial_error) => QwpWsTransportFailureAction::Reconnect {
reason: ReconnectReason::RetryableFailure,
initial_error,
pace: Duration::ZERO,
},
TransportFailure::Terminal(error) => {
store.mark_terminal(Some(error.clone()));
QwpWsTransportFailureAction::Terminal(error)
}
TransportFailure::ProtocolViolation { close_code, reason } => {
let error = store.record_protocol_violation(close_code, reason);
QwpWsTransportFailureAction::Terminal(error)
}
}
}
fn server_close_poison_error<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Option<Error> {
let fsn = store.queue.oldest_unresolved_fsn()?;
let now = Instant::now();
if self.poison_tracker.record_failure(
fsn,
store.queue.completed_fsn(),
self.max_frame_rejections,
self.poison_min_escalation_window,
now,
) {
let strikes = self.poison_tracker.strikes();
return Some(store.record_protocol_violation(
None,
format!(
"QWP/WebSocket frame fsn {fsn} was closed {strikes} times without ACK progress"
),
));
}
None
}
pub(crate) fn restart_connection(
&mut self,
reason: ReconnectReason,
) -> Result<(), DriverError> {
self.transport.restart_connection(reason)
}
pub(crate) fn begin_reconnect(
&self,
context: &'static str,
reason: ReconnectReason,
initial_error: Error,
) -> QwpWsReconnectState {
QwpWsReconnectState::new(self.reconnect_policy, context, reason, initial_error)
}
pub(crate) fn has_pending_reconnect(&self) -> bool {
self.pending_reconnect.is_some()
}
pub(crate) fn take_pending_reconnect(&mut self) -> Option<QwpWsReconnectState> {
self.pending_reconnect.take()
}
pub(crate) fn reconnect_once(
&mut self,
reconnect: &mut QwpWsReconnectState,
) -> Result<QwpWsReconnectStep, DriverError> {
if reconnect.deadline_expired() {
return Ok(QwpWsReconnectStep::Terminal(
reconnect.retry_budget_exhausted_error(),
));
}
reconnect.attempts += 1;
match self.restart_connection(reconnect.reason) {
Ok(()) => Ok(QwpWsReconnectStep::Reconnected {
reason: reconnect.reason,
}),
Err(err) => match reconnect_attempt_error(err) {
Ok(err) if reconnect_error_is_terminal(&err) => {
Ok(QwpWsReconnectStep::Terminal(err))
}
Ok(err) => {
let sleep_for = reconnect.record_retryable_error(err);
Ok(QwpWsReconnectStep::RetryAfter { sleep_for })
}
Err(err) => Err(err),
},
}
}
pub(crate) fn finish_reconnect_success<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
reason: ReconnectReason,
) -> DriveOutcome {
self.pending_reconnect = None;
self.sends_on_connection = 0;
self.send_cursor.restart(&store.queue);
if self.durable_ack.is_some() {
store.clear_unresolved_rejected_frames();
}
if let Some(tracker) = self.durable_ack.as_mut() {
tracker.reset();
}
self.catch_up_pending = self.dict_mirror.is_enabled() && !self.dict_mirror.is_empty();
store.counters.total_reconnects_succeeded += 1;
store.push_event(DriverEvent::Reconnected { reason });
DriveOutcome::Reconnected { reason }
}
fn continue_or_defer_reconnect<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
defer_reconnect: bool,
) -> Result<DriveOutcome, DriverError> {
if defer_reconnect {
let deadline = self
.pending_reconnect
.as_ref()
.and_then(QwpWsReconnectState::deadline);
Ok(DriveOutcome::ReconnectDelay {
sleep_for: Duration::ZERO,
deadline,
})
} else {
self.continue_reconnect(store)
}
}
fn continue_reconnect<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
let Some(mut reconnect) = self.pending_reconnect.take() else {
return Ok(DriveOutcome::Idle);
};
if let Some(sleep_for) = reconnect.take_first_attempt_pace() {
let deadline = reconnect.deadline();
self.pending_reconnect = Some(reconnect);
return Ok(DriveOutcome::ReconnectDelay {
sleep_for,
deadline,
});
}
store.record_reconnect_attempt();
match self.reconnect_once(&mut reconnect)? {
QwpWsReconnectStep::Reconnected { reason } => {
Ok(self.finish_reconnect_success(store, reason))
}
QwpWsReconnectStep::RetryAfter { sleep_for } => {
let deadline = reconnect.deadline();
self.pending_reconnect = Some(reconnect);
Ok(DriveOutcome::ReconnectDelay {
sleep_for,
deadline,
})
}
QwpWsReconnectStep::Terminal(error) => {
if reconnect_error_is_terminal(&error) {
store.mark_terminal(Some(error));
Ok(DriveOutcome::Terminal)
} else {
let sleep_for = reconnect.policy.initial_backoff();
let next = QwpWsReconnectState::new(
reconnect.policy,
reconnect.context,
reconnect.reason,
error,
);
let deadline = next.deadline();
self.pending_reconnect = Some(next);
Ok(DriveOutcome::ReconnectDelay {
sleep_for,
deadline,
})
}
}
}
}
pub(crate) fn drive_once<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
if store.is_terminal() {
return Ok(DriveOutcome::Terminal);
}
if self.pending_reconnect.is_some() {
return self.continue_reconnect(store);
}
let mut outcome = DriveOutcome::Idle;
if let Some(send_outcome) = self.drive_send_available(store)? {
if drive_outcome_stops_tick(send_outcome) {
return Ok(send_outcome);
}
if send_outcome != DriveOutcome::Idle {
outcome = send_outcome;
}
}
let receive = self.drive_receive_ready_until_idle(store)?;
if drive_outcome_stops_tick(receive) {
return Ok(receive);
}
if receive != DriveOutcome::Idle
&& (outcome == DriveOutcome::Idle || receive != DriveOutcome::Progress)
{
outcome = receive;
}
if self.drive_storage_once(store)? && outcome == DriveOutcome::Idle {
outcome = DriveOutcome::Progress;
}
if outcome == DriveOutcome::Idle {
outcome = self.drive_durable_ack_keepalive_once(store)?;
}
Ok(outcome)
}
pub(crate) fn close_drain_ready_once<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<CloseOutcome, DriverError> {
store.set_closing();
if store.is_terminal() {
return Ok(CloseOutcome::Terminal);
}
if store.all_published_receipts_resolved() {
self.close_publication_queue(store)?;
return Ok(CloseOutcome::Drained);
}
match self.drive_once(store)? {
DriveOutcome::Terminal => return Ok(CloseOutcome::Terminal),
DriveOutcome::ReconnectDelay {
sleep_for,
deadline,
} => {
return Ok(CloseOutcome::Waiting {
sleep_for,
deadline,
});
}
_ => {}
}
if store.is_terminal() {
Ok(CloseOutcome::Terminal)
} else if store.all_published_receipts_resolved() {
self.close_publication_queue(store)?;
Ok(CloseOutcome::Drained)
} else {
Ok(CloseOutcome::Timeout)
}
}
pub(crate) fn close_drain_ready_step<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<CloseStepOutcome, DriverError> {
store.set_closing();
if store.is_terminal() {
return Ok(CloseStepOutcome::Terminal);
}
if store.all_published_receipts_resolved() {
self.close_publication_queue(store)?;
return Ok(CloseStepOutcome::Drained);
}
let outcome = self.drive_once(store)?;
match outcome {
DriveOutcome::Terminal => return Ok(CloseStepOutcome::Terminal),
DriveOutcome::ReconnectDelay {
sleep_for,
deadline,
} => {
return Ok(CloseStepOutcome::Waiting {
sleep_for,
deadline,
});
}
_ => {}
}
if store.is_terminal() {
return Ok(CloseStepOutcome::Terminal);
}
if store.all_published_receipts_resolved() {
self.close_publication_queue(store)?;
return Ok(CloseStepOutcome::Drained);
}
if outcome == DriveOutcome::Idle {
Ok(CloseStepOutcome::Idle)
} else {
Ok(CloseStepOutcome::Progress)
}
}
pub(crate) fn drive_storage_once<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<bool, DriverError> {
let Some(step) = store.take_storage_maintenance_step()? else {
return Ok(false);
};
let changed_before_io = step.changes_queue_before_io();
let result = match step.perform() {
Ok(result) => result,
Err(err) => {
store.complete_storage_maintenance()?;
return Err(err.into());
}
};
let finish = match store.finish_storage_maintenance(result) {
Ok(finish) => finish,
Err(err) => {
store.complete_storage_maintenance()?;
return Err(err);
}
};
let changed = changed_before_io || finish.did_change();
if let Some(cleanup) = finish.into_cleanup()
&& let Some(failure) = cleanup.perform()
&& let Err(err) = store.record_storage_cleanup_failure(failure)
{
store.complete_storage_maintenance()?;
return Err(err);
}
store.complete_storage_maintenance()?;
Ok(changed)
}
pub(crate) fn drive_send_once<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
if store.is_terminal() {
return Ok(DriveOutcome::Terminal);
}
if self.pending_reconnect.is_some() {
return self.continue_reconnect(store);
}
Ok(self
.drive_send_available(store)?
.unwrap_or(DriveOutcome::Idle))
}
pub(crate) fn drive_receive_once<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
if store.is_terminal() {
return Ok(DriveOutcome::Terminal);
}
if self.pending_reconnect.is_some() {
return self.continue_reconnect(store);
}
let response = self.try_poll_response();
match response {
Ok(TransportPoll::Response(response)) => self.finish_polled_response(store, response),
Ok(TransportPoll::Progress) => Ok(DriveOutcome::Progress),
Ok(TransportPoll::Idle) => Ok(DriveOutcome::Idle),
Err(failure) => self.apply_transport_failure(store, failure),
}
}
fn drive_receive_ready_until_idle<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
if store.is_terminal() {
return Ok(DriveOutcome::Terminal);
}
if self.pending_reconnect.is_some() {
return self.continue_reconnect(store);
}
let mut outcome = DriveOutcome::Idle;
loop {
match self.try_poll_response() {
Ok(TransportPoll::Response(response)) => {
let response_outcome = self.finish_polled_response(store, response)?;
if drive_outcome_stops_tick(response_outcome) {
return Ok(response_outcome);
}
if response_outcome != DriveOutcome::Idle {
outcome = response_outcome;
}
}
Ok(TransportPoll::Progress) => {
if outcome == DriveOutcome::Idle {
outcome = DriveOutcome::Progress;
}
}
Ok(TransportPoll::Idle) => return Ok(outcome),
Err(failure) => return self.apply_transport_failure(store, failure),
}
}
}
fn finish_polled_response<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
response: TransportResponse,
) -> Result<DriveOutcome, DriverError> {
let outcome = self.finish_response(store, response)?;
Ok(if outcome == DriveOutcome::Idle {
DriveOutcome::Progress
} else {
outcome
})
}
fn drive_durable_ack_keepalive_once<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<DriveOutcome, DriverError> {
let durable_ack_pending = self.has_pending_durable_ack();
match self.send_durable_ack_keepalive_if_due(durable_ack_pending) {
Ok(true) | Ok(false) => Ok(DriveOutcome::Idle),
Err(failure) => self.apply_transport_failure(store, failure),
}
}
fn drive_send_available<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
) -> Result<Option<DriveOutcome>, DriverError> {
match self.drive_catch_up() {
Ok(()) => {}
Err(CatchUpDriveError::Transport(failure)) => {
return Ok(Some(self.apply_transport_failure(store, failure)?));
}
Err(CatchUpDriveError::Terminal(err)) => {
store.mark_terminal(Some(err));
return Ok(Some(DriveOutcome::Terminal));
}
Err(CatchUpDriveError::RetryConnection { error, pace }) => {
self.pending_reconnect = Some(
self.begin_reconnect(
"QWP/WebSocket reconnect after catch-up build failure",
ReconnectReason::RetryableFailure,
error,
)
.with_pace(pace),
);
return Ok(Some(self.continue_reconnect(store)?));
}
}
let progress = store.progress_view();
let Some(outbound) = self.next_outbound_sfa_frame(&progress)? else {
return Ok(None);
};
if let Err(err) = outbound.with_view(|view| self.guard_dict_not_torn(view.payload)) {
store.mark_terminal(Some(err));
return Ok(Some(DriveOutcome::Terminal));
}
let (frame, send_result) = self.send_frame(outbound);
let send_result = match send_result {
Ok(result) => result,
Err(failure) => {
return Ok(Some(self.apply_transport_failure(store, failure)?));
}
};
match self.finish_send_result(store, frame, send_result)? {
QwpWsSendProgress::Outcome(outcome) => Ok(Some(outcome)),
QwpWsSendProgress::TransportFailure(failure) => {
Ok(Some(self.apply_transport_failure(store, failure)?))
}
}
}
fn apply_transport_failure<Q: PublicationLog>(
&mut self,
store: &mut QwpWsPublicationStore<Q>,
failure: TransportFailure,
) -> Result<DriveOutcome, DriverError> {
match self.transport_failure_action(store, failure) {
QwpWsTransportFailureAction::Reconnect {
reason,
initial_error,
pace,
} => {
self.pending_reconnect = Some(
self.begin_reconnect("QWP/WebSocket reconnect", reason, initial_error)
.with_pace(pace),
);
self.continue_reconnect(store)
}
QwpWsTransportFailureAction::Terminal(error) => {
self.pending_reconnect = None;
store.mark_terminal(Some(error));
Ok(DriveOutcome::Terminal)
}
}
}
}
#[cfg(test)]
impl QwpWsCoreTestHarness<SfaFrameQueue, FakeOrderedServer> {
pub(crate) fn new(
options: SfaMemoryQueueOptions,
server: FakeOrderedServer,
) -> Result<Self, DriverError> {
let queue = SfaFrameQueue::open_memory(options)?;
Ok(Self {
store: QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY),
send_core: QwpWsSendCore::new(server, ReconnectPolicy::no_backoff(Duration::MAX)),
})
}
}
#[cfg(test)]
impl<Q: PublicationLog, T: QwpWsCoreTransport> QwpWsCoreTestHarness<Q, T> {
pub(crate) fn from_queue(queue: Q, transport: T) -> Self {
Self {
store: QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY),
send_core: QwpWsSendCore::new(transport, ReconnectPolicy::no_backoff(Duration::MAX)),
}
}
pub(crate) fn from_queue_with_reconnect_policy(
queue: Q,
transport: T,
reconnect_policy: ReconnectPolicy,
durable_ack: bool,
) -> Self {
Self {
store: QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY),
send_core: QwpWsSendCore::new_with_durable_ack(
transport,
reconnect_policy,
durable_ack,
),
}
}
pub(crate) fn from_queue_with_rejection_limit_and_window(
queue: Q,
transport: T,
max_frame_rejections: usize,
poison_min_escalation_window: Duration,
) -> Self {
Self {
store: QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY),
send_core: QwpWsSendCore::new_with_durable_ack_and_rejection_limit(
transport,
ReconnectPolicy::no_backoff(Duration::MAX),
false,
max_frame_rejections,
poison_min_escalation_window,
),
}
}
pub(crate) fn from_queue_with_event_capacity(
queue: Q,
transport: T,
event_capacity: usize,
) -> Self {
Self {
store: QwpWsPublicationStore::new(queue, event_capacity),
send_core: QwpWsSendCore::new(transport, ReconnectPolicy::no_backoff(Duration::MAX)),
}
}
fn from_queue_with_durable_ack(queue: Q, transport: T) -> Self {
Self {
store: QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY),
send_core: QwpWsSendCore::new_with_durable_ack(
transport,
ReconnectPolicy::no_backoff(Duration::MAX),
true,
),
}
}
pub(crate) fn try_submit(&mut self, payload: &[u8]) -> Result<QwpReceipt, DriverError> {
self.store.try_submit(payload)
}
pub(crate) fn set_closing(&mut self) {
self.store.set_closing();
}
pub(crate) fn submit_with_drive_limit(
&mut self,
payload: &[u8],
max_drive_steps: usize,
) -> Result<QwpReceipt, DriverError> {
let mut drive_steps = 0;
loop {
if self.store.is_terminal() {
return Err(DriverError::Terminal);
}
match self.store.try_submit(payload) {
Ok(receipt) => return Ok(receipt),
Err(DriverError::Queue(
QueueError::FrameCapacityFull { .. }
| QueueError::ByteCapacityFull { .. }
| QueueError::StorageSpareNotReady { .. }
| QueueError::StorageSegmentCapFull { .. },
)) if drive_steps < max_drive_steps => {
self.drive_once()?;
drive_steps += 1;
}
Err(DriverError::Queue(
err @ (QueueError::FrameCapacityFull { .. }
| QueueError::ByteCapacityFull { .. }
| QueueError::StorageSpareNotReady { .. }
| QueueError::StorageSegmentCapFull { .. }),
)) => {
return Err(DriverError::SubmitTimedOut {
backpressure: Some(err),
});
}
Err(
err @ (DriverError::Queue(
QueueError::InvalidCapacity
| QueueError::EmptyPayload
| QueueError::PayloadExceedsByteCapacity { .. }
| QueueError::NoUnsentFrame
| QueueError::ProtocolAckWithoutConnection
| QueueError::ProtocolAckBeyondSent { .. }
| QueueError::ProtocolAckedUnsentFrame { .. }
| QueueError::ProtocolRejectWithoutConnection
| QueueError::ProtocolRejectBeyondSent { .. }
| QueueError::ProtocolRejectedUnsentFrame { .. }
| QueueError::OutboundFrameUnavailable { .. }
| QueueError::SequenceOverflow,
)
| DriverError::Transport(_)
| DriverError::Storage(_)
| DriverError::SubmitTimedOut { .. }
| DriverError::Terminal
| DriverError::Closing
| DriverError::UnknownReceipt { .. }),
) => return Err(err),
}
}
}
pub(crate) fn submit_with_drive_deadline(
&mut self,
payload: &[u8],
append_deadline: Duration,
) -> Result<QwpReceipt, DriverError> {
let deadline = Instant::now().checked_add(append_deadline);
loop {
if self.store.is_terminal() {
return Err(DriverError::Terminal);
}
match self.store.try_submit(payload) {
Ok(receipt) => return Ok(receipt),
Err(err) => {
let Some(backpressure) = driver_backpressure_queue(&err) else {
return Err(err);
};
if drive_deadline_expired(deadline) {
return Err(DriverError::SubmitTimedOut {
backpressure: Some(backpressure),
});
}
if self.drive_once()? == DriveOutcome::Idle {
sleep_until_drive_deadline(deadline);
}
}
}
}
}
pub(crate) fn drive_once(&mut self) -> Result<DriveOutcome, DriverError> {
self.send_core.drive_once(&mut self.store)
}
pub(crate) fn drive_send_once(&mut self) -> Result<DriveOutcome, DriverError> {
self.send_core.drive_send_once(&mut self.store)
}
pub(crate) fn drive_receive_once(&mut self) -> Result<DriveOutcome, DriverError> {
self.send_core.drive_receive_once(&mut self.store)
}
pub(crate) fn delivery_status(
&self,
receipt: QwpReceipt,
) -> Result<Option<DeliveryOutcome>, DriverError> {
self.send_core.delivery_status(&self.store, receipt)
}
pub(crate) fn close_drain_steps(
&mut self,
max_drive_steps: usize,
) -> Result<CloseOutcome, DriverError> {
self.store.set_closing();
for _ in 0..max_drive_steps {
if self.store.is_terminal() {
return Ok(CloseOutcome::Terminal);
}
if self.store.all_published_receipts_resolved() {
self.send_core.close_publication_queue(&mut self.store)?;
return Ok(CloseOutcome::Drained);
}
if self.drive_once()? == DriveOutcome::Terminal {
return Ok(CloseOutcome::Terminal);
}
}
if self.store.is_terminal() {
Ok(CloseOutcome::Terminal)
} else if self.store.all_published_receipts_resolved() {
self.send_core.close_publication_queue(&mut self.store)?;
Ok(CloseOutcome::Drained)
} else {
Ok(CloseOutcome::Timeout)
}
}
pub(crate) fn close_drain_ready_once(&mut self) -> Result<CloseOutcome, DriverError> {
self.send_core.close_drain_ready_once(&mut self.store)
}
pub(crate) fn close_drain_ready_step(&mut self) -> Result<CloseStepOutcome, DriverError> {
self.send_core.close_drain_ready_step(&mut self.store)
}
pub(crate) fn receipt_status(&self, receipt: QwpReceipt) -> QwpReceiptStatus {
self.send_core.receipt_status(&self.store, receipt)
}
pub(crate) fn poll_event(&mut self) -> Option<DriverEvent> {
self.store.poll_event()
}
pub(crate) fn events_dropped_total(&self) -> u64 {
self.store.events_dropped_total()
}
pub(crate) fn terminal_error(&self) -> Option<&Error> {
self.store.terminal_error()
}
pub(crate) fn terminal_sender_error(&self) -> Option<&QwpWsSenderError> {
self.store.terminal_sender_error()
}
pub(crate) fn published_fsn(&self) -> Option<u64> {
self.store.queue.published_fsn()
}
pub(crate) fn acked_fsn(&self) -> Option<u64> {
self.store.queue.completed_fsn()
}
pub(crate) fn poll_sender_error(&mut self) -> Option<QwpWsSenderError> {
self.store.poll_sender_error()
}
pub(crate) fn poll_sender_error_notification(&mut self) -> Option<QwpWsSenderError> {
self.store.poll_sender_error_notification()
}
pub(crate) fn sender_errors_dropped_total(&self) -> u64 {
self.store.sender_errors_dropped_total()
}
pub(crate) fn counters(&self) -> QwpWsCounters {
self.store.counters()
}
pub(crate) fn last_server_error(&self) -> Option<&QwpServerError> {
self.store.last_server_error()
}
pub(crate) fn rejected_frame(&self, receipt: QwpReceipt) -> Option<&QwpRejectedFrame> {
self.store.rejected_frame(receipt)
}
pub(crate) fn is_terminal(&self) -> bool {
self.store.is_terminal()
}
pub(crate) fn into_parts(self) -> (QwpWsPublicationStore<Q>, QwpWsSendCore<T>) {
(self.store, self.send_core)
}
}
pub(super) fn reconnect_attempt_error(err: DriverError) -> Result<Error, DriverError> {
match err {
DriverError::Transport(err) | DriverError::Storage(err) => Ok(err),
err => Err(err),
}
}
pub(crate) fn reconnect_error_is_terminal(err: &Error) -> bool {
if is_qwp_ws_role_reject_error(err) {
return false;
}
matches!(
err.code(),
ErrorCode::AuthError
| ErrorCode::ConfigError
| ErrorCode::ProtocolVersionError
| ErrorCode::StoreResendRequired
| ErrorCode::SymbolDictFull
)
}
pub(super) fn is_qwp_ws_role_reject_error(err: &Error) -> bool {
err.qwp_ws_role_reject().is_some()
}
fn reconnect_deadline_expired(deadline: Option<Instant>) -> bool {
deadline.is_some_and(|deadline| Instant::now() >= deadline)
}
fn double_duration(duration: Duration) -> Duration {
duration.checked_mul(2).unwrap_or(Duration::MAX)
}
fn pace_dose(initial_backoff: Duration, max_backoff: Duration, strikes: usize) -> Duration {
if initial_backoff.is_zero() {
return Duration::ZERO;
}
let shift = strikes.max(1).saturating_sub(1).min(6) as u32;
initial_backoff
.checked_mul(1u32 << shift)
.unwrap_or(Duration::MAX)
.min(max_backoff)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn pace_jitter_duration(dose: Duration) -> Duration {
let dose_nanos = dose.as_nanos().min(u128::from(u64::MAX)) as u64;
if dose_nanos == 0 {
return dose;
}
let extra = rand::rng().random_range(0..dose_nanos);
Duration::from_nanos(dose_nanos.saturating_add(extra))
}
#[cfg(not(feature = "sync-sender-qwp-ws"))]
fn pace_jitter_duration(dose: Duration) -> Duration {
dose
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn centered_jitter_duration(base: Duration) -> Duration {
let base_nanos = base.as_nanos().min(u128::from(u64::MAX)) as u64;
if base_nanos == 0 {
return base;
}
let extra = rand::rng().random_range(0..base_nanos);
Duration::from_nanos((base_nanos / 2).saturating_add(extra))
}
#[cfg(not(feature = "sync-sender-qwp-ws"))]
fn centered_jitter_duration(base: Duration) -> Duration {
base
}
pub(super) fn reconnect_sleep_duration(
role_reject: bool,
initial_backoff: Duration,
backoff: Duration,
) -> Duration {
if role_reject {
initial_backoff
} else {
centered_jitter_duration(backoff)
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
pub(crate) fn reconnect_backoff_step(
err: &Error,
initial_backoff: Duration,
max_backoff: Duration,
backoff: Duration,
) -> (Duration, Duration) {
let role_reject = is_qwp_ws_role_reject_error(err);
let sleep_for = reconnect_sleep_duration(role_reject, initial_backoff, backoff);
let next_backoff = if role_reject {
initial_backoff
} else {
double_duration(backoff).min(max_backoff)
};
(sleep_for, next_backoff)
}
pub(super) fn retry_budget_exhausted_error(
context: &str,
attempts: usize,
started: Instant,
last_error: Option<Error>,
) -> Error {
let elapsed_ms = started.elapsed().as_millis();
let code = last_error
.as_ref()
.map_or(ErrorCode::SocketError, |err| err.code());
let last_error_msg = last_error
.as_ref()
.map_or_else(|| "none".to_string(), |err| err.msg().to_string());
let qwp_ws_rejection = last_error
.as_ref()
.and_then(|err| err.qwp_ws_rejection().cloned());
let qwp_ws_role_reject = last_error
.as_ref()
.and_then(|err| err.qwp_ws_role_reject().cloned());
let mut err = Error::new(
code,
format!(
"{context} retry budget exhausted [attempts={attempts}, elapsed_ms={elapsed_ms}, last_error={last_error_msg}]"
),
);
if let Some(rejection) = qwp_ws_rejection {
err = err.with_qwp_ws_rejection(rejection);
}
if let Some(role_reject) = qwp_ws_role_reject {
err = err.with_qwp_ws_role_reject(role_reject);
}
err
}
pub(crate) trait PublicationLog {
fn try_publish(&mut self, payload: &[u8]) -> Result<QwpReceipt, DriverError>;
fn take_producer(&mut self) -> Option<SfaProducer> {
None
}
fn progress_view(&self) -> SfaProgressView;
fn check_durability(&self) -> Result<(), DriverError> {
Ok(())
}
fn storage_maintenance_in_flight(&self) -> Result<bool, DriverError> {
Ok(false)
}
fn take_storage_maintenance_step(
&mut self,
_allow_create: bool,
) -> Result<Option<SfaStorageStep>, DriverError> {
Ok(None)
}
fn finish_storage_maintenance(
&mut self,
_result: SfaStorageResult,
_allow_install: bool,
) -> Result<SfaStorageFinish, DriverError> {
Ok(SfaStorageFinish::unchanged())
}
fn complete_storage_maintenance(&mut self) -> Result<(), DriverError> {
Ok(())
}
fn record_storage_cleanup_failure(
&mut self,
_failure: SfaCleanupFailure,
) -> Result<(), DriverError> {
Ok(())
}
fn oldest_unresolved_fsn(&self) -> Option<u64>;
fn persist_completed_fsn(&mut self, _fsn: u64) {}
fn close(&mut self) -> Result<(), DriverError> {
Ok(())
}
fn receipt_status(&self, receipt: QwpReceipt) -> QwpReceiptStatus;
fn published_fsn(&self) -> Option<u64>;
fn completed_fsn(&self) -> Option<u64>;
}
#[derive(Debug)]
pub(crate) struct SendCursor {
fsn_at_zero: Option<u64>,
next_fsn: Option<u64>,
replay_target_fsn: Option<u64>,
next_wire_seq: u64,
catch_up_offset: u64,
last_sent_wire_seq: Option<u64>,
in_flight: InFlightRun,
sfa_cursor: Option<SfaSendCursor>,
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
struct InFlightRun {
front_fsn: u64,
front_wire_seq: u64,
len: usize,
}
impl InFlightRun {
fn len(&self) -> usize {
self.len
}
fn clear(&mut self) {
self.len = 0;
}
fn push(&mut self, frame: &SentFrame) {
if self.len == 0 {
self.front_fsn = frame.fsn;
self.front_wire_seq = frame.wire_seq;
}
self.len += 1;
}
fn ack_through(&mut self, acked_fsn: u64) {
if self.len == 0 || acked_fsn < self.front_fsn {
return;
}
let dropped = (acked_fsn - self.front_fsn + 1).min(self.len as u64);
self.front_fsn += dropped;
self.front_wire_seq += dropped;
self.len -= dropped as usize;
}
fn wire_seq_for_fsn(&self, fsn: u64) -> Option<u64> {
if self.len == 0 {
return None;
}
let offset = fsn.checked_sub(self.front_fsn)?;
if offset >= self.len as u64 {
return None;
}
Some(self.front_wire_seq + offset)
}
}
impl SendCursor {
fn new() -> Self {
Self {
fsn_at_zero: None,
next_fsn: None,
replay_target_fsn: None,
next_wire_seq: 0,
catch_up_offset: 0,
last_sent_wire_seq: None,
in_flight: InFlightRun::default(),
sfa_cursor: None,
}
}
pub(crate) fn sfa_cursor_mut(&mut self) -> &mut Option<SfaSendCursor> {
&mut self.sfa_cursor
}
fn release_sfa_cursor(&mut self) {
self.sfa_cursor = None;
}
pub(crate) fn peek_next_frame_from_oldest(
&mut self,
oldest_unresolved_fsn: Option<u64>,
) -> Result<Option<(u64, u64)>, DriverError> {
let fsn = match self.next_fsn {
Some(fsn) => fsn,
None => {
let Some(fsn) = oldest_unresolved_fsn else {
return Ok(None);
};
self.fsn_at_zero = Some(fsn);
self.next_fsn = Some(fsn);
fsn
}
};
Ok(Some((fsn, self.next_wire_seq)))
}
fn commit_sent(&mut self, frame: SentFrame) -> Result<bool, DriverError> {
if self.next_fsn != Some(frame.fsn) || self.next_wire_seq != frame.wire_seq {
return Err(DriverError::Queue(QueueError::OutboundFrameUnavailable {
fsn: frame.fsn,
wire_seq: frame.wire_seq,
}));
}
let replayed = self
.replay_target_fsn
.is_some_and(|target_fsn| frame.fsn <= target_fsn);
self.next_fsn = Some(
frame
.fsn
.checked_add(1)
.ok_or(DriverError::Queue(QueueError::SequenceOverflow))?,
);
self.next_wire_seq = self
.next_wire_seq
.checked_add(1)
.ok_or(DriverError::Queue(QueueError::SequenceOverflow))?;
self.last_sent_wire_seq = Some(frame.wire_seq);
self.in_flight.push(&frame);
if self
.replay_target_fsn
.is_some_and(|target_fsn| frame.fsn >= target_fsn)
{
self.replay_target_fsn = None;
}
Ok(replayed)
}
fn reject_fsn_for_wire_seq(&self, wire_seq: u64) -> Result<Option<(u64, u64)>, DriverError> {
let Some(fsn_at_zero) = self.fsn_at_zero else {
return Ok(None);
};
let Some(last_sent_wire_seq) = self.last_sent_wire_seq else {
return Ok(None);
};
let effect_wire_seq = wire_seq.min(last_sent_wire_seq);
let Some(real_delta) = effect_wire_seq.checked_sub(self.catch_up_offset) else {
return Ok(None); };
let fsn = fsn_at_zero
.checked_add(real_delta)
.ok_or(DriverError::Queue(QueueError::SequenceOverflow))?;
Ok(Some((fsn, effect_wire_seq)))
}
fn ack_fsn_for_wire_seq(&self, wire_seq: u64) -> Result<Option<(u64, u64)>, DriverError> {
let Some(fsn_at_zero) = self.fsn_at_zero else {
return Ok(None);
};
let Some(last_sent_wire_seq) = self.last_sent_wire_seq else {
return Ok(None);
};
let ack_wire_seq = wire_seq.min(last_sent_wire_seq);
let Some(real_delta) = ack_wire_seq.checked_sub(self.catch_up_offset) else {
return Ok(None);
};
let fsn = fsn_at_zero
.checked_add(real_delta)
.ok_or(DriverError::Queue(QueueError::SequenceOverflow))?;
Ok(Some((fsn, ack_wire_seq)))
}
fn ack_through(&mut self, acked_fsn: u64) {
self.in_flight.ack_through(acked_fsn);
}
fn restart<Q: PublicationLog>(&mut self, log: &Q) {
self.in_flight.clear();
self.fsn_at_zero = log.oldest_unresolved_fsn();
self.next_fsn = self.fsn_at_zero;
self.replay_target_fsn = match (self.fsn_at_zero, log.published_fsn()) {
(Some(oldest_unresolved), Some(published)) if oldest_unresolved <= published => {
Some(published)
}
_ => None,
};
self.next_wire_seq = 0;
self.catch_up_offset = 0;
self.last_sent_wire_seq = None;
self.sfa_cursor = None;
}
fn begin_catch_up(&mut self, catch_up_frames: u64) {
self.catch_up_offset = catch_up_frames;
self.next_wire_seq = catch_up_frames;
}
fn wire_seq_for_fsn(&self, fsn: u64) -> Option<u64> {
self.in_flight.wire_seq_for_fsn(fsn)
}
}
pub(crate) trait QwpWsCoreTransport {
fn try_poll_response(&mut self) -> Result<TransportPoll, TransportFailure>;
fn send_durable_ack_keepalive_if_due(
&mut self,
_durable_ack_pending: bool,
) -> Result<bool, TransportFailure> {
Ok(false)
}
fn send_frame(
&mut self,
frame: OutboundFrameView<'_>,
) -> Result<TransportSendResult, TransportFailure>;
fn restart_connection(&mut self, _reason: ReconnectReason) -> Result<(), DriverError> {
Ok(())
}
fn server_max_batch_size(&self) -> usize {
0
}
fn negotiated_qwp_version(&self) -> u8 {
1
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[derive(Debug, Default)]
struct PendingWireSequenceRun {
front: u64,
len: u64,
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl PendingWireSequenceRun {
fn is_empty(&self) -> bool {
self.len == 0
}
fn clear(&mut self) {
self.len = 0;
}
fn push(&mut self, wire_seq: u64) {
if self.len == 0 {
self.front = wire_seq;
} else {
debug_assert_eq!(self.front.checked_add(self.len), Some(wire_seq));
}
self.len = self
.len
.checked_add(1)
.expect("pending wire-sequence count overflow");
}
fn complete_through(&mut self, wire_seq: u64) {
if self.len == 0 || wire_seq < self.front {
return;
}
let dropped = wire_seq
.saturating_sub(self.front)
.saturating_add(1)
.min(self.len);
self.len -= dropped;
if self.len == 0 {
self.front = 0;
} else {
self.front += dropped;
}
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
pub(crate) struct BlockingQwpWsTransport {
endpoints: Arc<[QwpWsEndpoint]>,
previous_idx: Option<usize>,
tracker: QwpWsHostHealthTracker,
use_tls: bool,
tls_settings: Option<TlsSettings>,
connect_kind: QwpWsConnectKind,
qwp_ws: QwpWsConfig,
auth_header: Option<String>,
negotiated_version: u8,
server_max_batch_size: Arc<AtomicUsize>,
traffic_gate: Option<Arc<TrafficGate>>,
stream: WsStream,
reader: WsFrameReader,
send_buf: Vec<u8>,
pending_wire_sequences: PendingWireSequenceRun,
last_durable_keepalive_ping: Option<Instant>,
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl BlockingQwpWsTransport {
#[allow(clippy::too_many_arguments)]
pub(crate) fn connect(
host: impl Into<String>,
port: impl Into<String>,
use_tls: bool,
tls_settings: Option<TlsSettings>,
connect_kind: QwpWsConnectKind,
qwp_ws: QwpWsConfig,
auth_header: Option<String>,
server_max_batch_size: Arc<AtomicUsize>,
traffic_gate: Option<Arc<TrafficGate>>,
) -> crate::Result<Self> {
let host = host.into();
let port = port.into();
let endpoints = qwp_ws_configured_endpoints(&host, &port, &qwp_ws);
let mut tracker = QwpWsHostHealthTracker::new(endpoints.len());
let mut previous_idx = None;
let connected = connect_qwp_ws_endpoint_round(
&endpoints,
&mut tracker,
&mut previous_idx,
None,
use_tls,
tls_settings.clone(),
connect_kind,
&qwp_ws,
auth_header.as_deref(),
qwp_ws.conn_events.as_deref(),
traffic_gate.as_deref(),
)?;
Ok(Self::from_connected(
endpoints,
tracker,
use_tls,
tls_settings,
connect_kind,
qwp_ws,
auth_header,
server_max_batch_size,
traffic_gate,
connected,
))
}
#[allow(clippy::too_many_arguments)]
pub(super) fn from_connected(
endpoints: Arc<[QwpWsEndpoint]>,
tracker: QwpWsHostHealthTracker,
use_tls: bool,
tls_settings: Option<TlsSettings>,
connect_kind: QwpWsConnectKind,
qwp_ws: QwpWsConfig,
auth_header: Option<String>,
server_max_batch_size: Arc<AtomicUsize>,
traffic_gate: Option<Arc<TrafficGate>>,
connected: QwpWsConnectRoundSuccess,
) -> Self {
server_max_batch_size.store(connected.server_max_batch_size, Ordering::Release);
let transport = Self {
endpoints,
previous_idx: Some(connected.endpoint_idx),
tracker,
use_tls,
tls_settings,
connect_kind,
qwp_ws,
auth_header,
negotiated_version: connected.negotiated_version,
server_max_batch_size,
traffic_gate,
stream: connected.stream,
reader: WsFrameReader::with_initial_input(connected.leftover),
send_buf: Vec::with_capacity(16 * 1024),
pending_wire_sequences: PendingWireSequenceRun::default(),
last_durable_keepalive_ping: None,
};
transport.emit_connect_succeeded();
transport
}
pub(crate) fn negotiated_version(&self) -> u8 {
self.negotiated_version
}
fn reconnect(&mut self, reason: ReconnectReason) -> Result<(), DriverError> {
if let Some(gate) = self.traffic_gate.as_deref() {
gate.clear();
}
if matches!(self.connect_kind, QwpWsConnectKind::Foreground)
&& let Some(events) = self.qwp_ws.conn_events.as_deref()
&& let Some(idx) = self.previous_idx
&& let Some(endpoint) = self.endpoints.get(idx)
{
events.disconnected(&endpoint.host, &endpoint.port);
}
let connected = connect_qwp_ws_endpoint_round(
&self.endpoints,
&mut self.tracker,
&mut self.previous_idx,
Some(reason),
self.use_tls,
self.tls_settings.clone(),
self.connect_kind.for_reconnect(),
&self.qwp_ws,
self.auth_header.as_deref(),
self.qwp_ws.conn_events.as_deref(),
self.traffic_gate.as_deref(),
)
.map_err(DriverError::Transport)?;
self.previous_idx = Some(connected.endpoint_idx);
self.stream = connected.stream;
self.negotiated_version = connected.negotiated_version;
self.server_max_batch_size
.store(connected.server_max_batch_size, Ordering::Release);
self.reader = WsFrameReader::with_initial_input(connected.leftover);
self.send_buf.clear();
self.pending_wire_sequences.clear();
self.last_durable_keepalive_ping = None;
self.emit_connect_succeeded();
Ok(())
}
fn emit_connect_succeeded(&self) {
if matches!(self.connect_kind, QwpWsConnectKind::Foreground)
&& let Some(events) = self.qwp_ws.conn_events.as_deref()
&& let Some(idx) = self.previous_idx
&& let Some(endpoint) = self.endpoints.get(idx)
{
events.connect_succeeded(&endpoint.host, &endpoint.port);
}
}
fn complete_pending_through(&mut self, sequence: u64) {
self.pending_wire_sequences.complete_through(sequence);
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl Drop for BlockingQwpWsTransport {
fn drop(&mut self) {
if let Some(gate) = self.traffic_gate.as_deref() {
gate.clear();
}
}
}
fn decode_transport_response(
payload: &[u8],
) -> Result<Option<TransportResponse>, TransportFailure> {
match codec::parse_pipelined_response(payload) {
Ok(PipelinedResponse::Ok { sequence }) => {
Ok(Some(TransportResponse::Ack { wire_seq: sequence }))
}
Ok(PipelinedResponse::DurableAck) => Ok(None),
Ok(PipelinedResponse::Error(error)) => {
let wire_seq = error.sequence;
let server_error = QwpServerError::from(error);
Ok(Some(TransportResponse::Reject {
wire_seq,
error: server_error,
}))
}
Err(err) => Err(TransportFailure::Retryable(err)),
}
}
fn decode_durable_transport_response(
payload: &[u8],
) -> Result<Option<TransportResponse>, TransportFailure> {
let mut table_seq_txns = Vec::new();
let mut handler = |_, table: &str, seq_txn| {
table_seq_txns.push(TableSeqTxn {
table: table.to_string(),
seq_txn,
});
Ok(())
};
match codec::parse_pipelined_response_with_table_handler(payload, Some(&mut handler)) {
Ok(PipelinedResponse::Ok { sequence }) => Ok(Some(TransportResponse::DurableOk {
wire_seq: sequence,
table_seq_txns,
})),
Ok(PipelinedResponse::DurableAck) => {
Ok(Some(TransportResponse::DurableAck { table_seq_txns }))
}
Ok(PipelinedResponse::Error(error)) => {
let wire_seq = error.sequence;
let server_error = QwpServerError::from(error);
Ok(Some(TransportResponse::Reject {
wire_seq,
error: server_error,
}))
}
Err(err) => Err(TransportFailure::Retryable(err)),
}
}
#[cfg(feature = "sync-sender-qwp-ws")]
impl QwpWsCoreTransport for BlockingQwpWsTransport {
fn try_poll_response(&mut self) -> Result<TransportPoll, TransportFailure> {
if self.pending_wire_sequences.is_empty() && !*self.qwp_ws.request_durable_ack {
return Ok(TransportPoll::Idle);
}
let read = self
.reader
.try_read_one(&mut self.stream, &mut self.send_buf)
.map_err(|err| match err {
super::qwp_ws::WsMessageError::Close(close) if close.is_orderly() => {
TransportFailure::Disconnect(close.into_error())
}
super::qwp_ws::WsMessageError::Close(close) => {
TransportFailure::ServerClose(close.into_error())
}
super::qwp_ws::WsMessageError::ProtocolViolation(reason) => {
TransportFailure::ProtocolViolation {
close_code: None,
reason,
}
}
super::qwp_ws::WsMessageError::Error(err) => match err.code() {
ErrorCode::SocketError => TransportFailure::Disconnect(err),
ErrorCode::AuthError | ErrorCode::ProtocolVersionError => {
TransportFailure::Terminal(err)
}
_ => TransportFailure::Terminal(err),
},
})?;
let WsFrameRead::Message { opcode } = read else {
return Ok(match read {
WsFrameRead::Progress => TransportPoll::Progress,
WsFrameRead::Idle => TransportPoll::Idle,
WsFrameRead::Message { .. } => unreachable!(),
});
};
if opcode != OPCODE_BINARY {
self.reader.clear_message();
return Err(TransportFailure::ProtocolViolation {
close_code: None,
reason: "QWP/WebSocket server response was not a binary frame".to_string(),
});
}
let response = if *self.qwp_ws.request_durable_ack {
decode_durable_transport_response(self.reader.message())
} else {
decode_transport_response(self.reader.message())
};
self.reader.clear_message();
let response = response?;
if let Some(
TransportResponse::Ack { wire_seq }
| TransportResponse::DurableOk { wire_seq, .. }
| TransportResponse::Reject { wire_seq, .. },
) = &response
{
self.complete_pending_through(*wire_seq);
}
Ok(match response {
Some(response) => TransportPoll::Response(response),
None => TransportPoll::Progress,
})
}
fn send_durable_ack_keepalive_if_due(
&mut self,
durable_ack_pending: bool,
) -> Result<bool, TransportFailure> {
if !durable_ack_pending || !*self.qwp_ws.request_durable_ack {
return Ok(false);
}
let interval = *self.qwp_ws.durable_ack_keepalive_interval;
if interval.is_zero() {
return Ok(false);
}
if self
.last_durable_keepalive_ping
.is_some_and(|sent_at| sent_at.elapsed() < interval)
{
return Ok(false);
}
write_ping_frame(&mut self.stream, &mut self.send_buf, b"").map_err(|io| {
TransportFailure::Disconnect(error::fmt!(
SocketError,
"Could not send WebSocket durable ACK keepalive PING: {}",
io
))
})?;
self.stream.flush().map_err(|io| {
TransportFailure::Disconnect(error::fmt!(
SocketError,
"Could not flush WebSocket durable ACK keepalive PING: {}",
io
))
})?;
self.last_durable_keepalive_ping = Some(Instant::now());
Ok(true)
}
fn send_frame(
&mut self,
frame: OutboundFrameView<'_>,
) -> Result<TransportSendResult, TransportFailure> {
write_binary_frame(&mut self.stream, &mut self.send_buf, frame.payload).map_err(|io| {
TransportFailure::Disconnect(error::fmt!(
SocketError,
"Could not send WebSocket frame: {}",
io
))
})?;
self.stream.flush().map_err(|io| {
TransportFailure::Disconnect(error::fmt!(
SocketError,
"Could not flush WebSocket frame: {}",
io
))
})?;
self.pending_wire_sequences.push(frame.wire_seq);
Ok(TransportSendResult::NoResponse)
}
fn restart_connection(&mut self, reason: ReconnectReason) -> Result<(), DriverError> {
self.reconnect(reason)
}
fn server_max_batch_size(&self) -> usize {
self.server_max_batch_size.load(Ordering::Acquire)
}
fn negotiated_qwp_version(&self) -> u8 {
self.negotiated_version
}
}
enum DictCatchUpError {
Transport(TransportFailure),
EntryTooLarge(CatchUpEntryTooLarge),
FrameBuild(CatchUpFrameBuildError),
}
pub(crate) enum CatchUpDriveError {
Transport(TransportFailure),
RetryConnection { error: Error, pace: Duration },
Terminal(Error),
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum DriverError {
Queue(QueueError),
Transport(Error),
Storage(Error),
SubmitTimedOut { backpressure: Option<QueueError> },
Terminal,
Closing,
UnknownReceipt { fsn: u64 },
}
impl From<QueueError> for DriverError {
fn from(value: QueueError) -> Self {
Self::Queue(value)
}
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct QwpServerError {
pub(crate) status: u8,
pub(crate) message: String,
pub(crate) error: Error,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) struct QwpRejectedFrame {
pub(crate) fsn: u64,
pub(crate) wire_seq: u64,
pub(crate) error: QwpServerError,
}
impl From<codec::PipelinedError> for QwpServerError {
fn from(value: codec::PipelinedError) -> Self {
Self {
status: value.status,
message: value.message,
error: value.err,
}
}
}
fn server_error_category(status: u8) -> QwpWsErrorCategory {
match status {
codec::WS_STATUS_SCHEMA_MISMATCH => QwpWsErrorCategory::SchemaMismatch,
codec::WS_STATUS_PARSE_ERROR => QwpWsErrorCategory::ParseError,
codec::WS_STATUS_INTERNAL_ERROR => QwpWsErrorCategory::InternalError,
codec::WS_STATUS_SECURITY_ERROR => QwpWsErrorCategory::SecurityError,
codec::WS_STATUS_WRITE_ERROR => QwpWsErrorCategory::WriteError,
codec::WS_STATUS_NOT_WRITABLE => QwpWsErrorCategory::NotWritable,
_ => QwpWsErrorCategory::Unknown,
}
}
fn server_error_policy(status: u8) -> QwpWsErrorPolicy {
match status {
codec::WS_STATUS_WRITE_ERROR | codec::WS_STATUS_INTERNAL_ERROR => {
QwpWsErrorPolicy::Retriable
}
codec::WS_STATUS_NOT_WRITABLE => QwpWsErrorPolicy::RetriableOther,
codec::WS_STATUS_SCHEMA_MISMATCH
| codec::WS_STATUS_PARSE_ERROR
| codec::WS_STATUS_SECURITY_ERROR => QwpWsErrorPolicy::Terminal,
_ => QwpWsErrorPolicy::Retriable,
}
}
fn reconnect_reason_for_policy(policy: QwpWsErrorPolicy) -> ReconnectReason {
match policy {
QwpWsErrorPolicy::Retriable => ReconnectReason::RetryableFailure,
QwpWsErrorPolicy::RetriableOther => ReconnectReason::NotWritable,
QwpWsErrorPolicy::Terminal => ReconnectReason::RetryableFailure,
}
}
fn sender_error_for_qwp_error(
error: &QwpServerError,
wire_seq: u64,
fsn: u64,
applied_policy: QwpWsErrorPolicy,
) -> QwpWsSenderError {
sender_error_for_qwp_error_span(error, wire_seq, fsn, fsn, applied_policy)
}
fn sender_error_for_qwp_error_span(
error: &QwpServerError,
wire_seq: u64,
from_fsn: u64,
to_fsn: u64,
applied_policy: QwpWsErrorPolicy,
) -> QwpWsSenderError {
QwpWsSenderError {
category: server_error_category(error.status),
applied_policy,
status: Some(error.status),
message: (!error.message.is_empty()).then(|| error.message.clone()),
message_sequence: Some(wire_seq),
from_fsn,
to_fsn,
}
}
fn server_rejection_error(error: Error, sender_error: QwpWsSenderError) -> Error {
Error::new(ErrorCode::ServerRejection, error.msg().to_string())
.with_qwp_ws_rejection(sender_error)
}
#[cfg(test)]
fn fake_transport_error(message: &'static str) -> Error {
Error::new(ErrorCode::SocketError, message)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum DriveOutcome {
Idle,
Sent(SentFrame),
Acked {
wire_seq: u64,
},
Rejected {
fsn: u64,
wire_seq: u64,
},
Reconnected {
reason: ReconnectReason,
},
ReconnectDelay {
sleep_for: Duration,
deadline: Option<Instant>,
},
Progress,
Terminal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum DriverEvent {
Published { fsn: u64 },
Sent { fsn: u64, wire_seq: u64 },
CompletedThrough { fsn: u64, wire_seq: u64 },
Rejected { fsn: u64, wire_seq: u64 },
Reconnected { reason: ReconnectReason },
Terminal,
}
fn drive_outcome_stops_tick(outcome: DriveOutcome) -> bool {
matches!(
outcome,
DriveOutcome::Terminal | DriveOutcome::ReconnectDelay { .. }
)
}
fn driver_backpressure_queue(err: &DriverError) -> Option<QueueError> {
match err {
DriverError::Queue(
err @ (QueueError::FrameCapacityFull { .. }
| QueueError::ByteCapacityFull { .. }
| QueueError::StorageSpareNotReady { .. }
| QueueError::StorageSegmentCapFull { .. }),
) => Some(*err),
_ => None,
}
}
fn drive_deadline_expired(deadline: Option<Instant>) -> bool {
deadline.is_some_and(|deadline| deadline <= Instant::now())
}
fn sleep_until_drive_deadline(deadline: Option<Instant>) {
let sleep_for = match deadline {
Some(deadline) => {
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return;
}
remaining.min(Duration::from_millis(10))
}
None => Duration::from_millis(10),
};
std::thread::sleep(sleep_for);
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ReconnectReason {
Disconnect,
RetryableFailure,
NotWritable,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum DeliveryOutcome {
Completed,
Terminal,
Timeout,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CloseOutcome {
Drained,
Waiting {
sleep_for: Duration,
deadline: Option<Instant>,
},
Timeout,
Terminal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum CloseStepOutcome {
Drained,
Terminal,
Waiting {
sleep_for: Duration,
deadline: Option<Instant>,
},
Progress,
Idle,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum TransportSendResult {
NoResponse,
Response(TransportResponse),
Failure(TransportFailure),
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum TransportFailure {
Disconnect(Error),
ServerClose(Error),
Retryable(Error),
Terminal(Error),
ProtocolViolation {
close_code: Option<u16>,
reason: String,
},
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum TransportPoll {
Response(TransportResponse),
Progress,
Idle,
}
#[derive(Debug, Clone, PartialEq)]
pub(crate) enum TransportResponse {
Ack {
wire_seq: u64,
},
DurableOk {
wire_seq: u64,
table_seq_txns: Vec<TableSeqTxn>,
},
DurableAck {
table_seq_txns: Vec<TableSeqTxn>,
},
Reject {
wire_seq: u64,
error: QwpServerError,
},
}
#[cfg(test)]
pub(crate) type FakeServerResponse = TransportResponse;
#[derive(Debug)]
struct DriverEventRing {
events: VecDeque<DriverEvent>,
capacity: usize,
dropped_total: u64,
}
#[derive(Debug)]
struct SenderErrorLog {
errors: VecDeque<QwpWsSenderError>,
capacity: usize,
dropped_total: u64,
first_seq: u64,
next_seq: u64,
poll_next_seq: u64,
notification_next_seq: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct TableSeqTxn {
pub(crate) table: String,
pub(crate) seq_txn: i64,
}
#[derive(Debug)]
struct DurableAckTracker {
table_watermarks: HashMap<String, i64>,
pending: VecDeque<PendingDurableFrame>,
}
#[derive(Debug)]
enum PendingDurableFrame {
Ok {
wire_seq: u64,
fsn: u64,
table_seq_txns: Vec<TableSeqTxn>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct DurableCompletion {
wire_seq: u64,
fsn: u64,
}
impl DurableAckTracker {
fn new() -> Self {
Self {
table_watermarks: HashMap::new(),
pending: VecDeque::new(),
}
}
fn reset(&mut self) {
self.table_watermarks.clear();
self.pending.clear();
}
fn has_pending(&self) -> bool {
!self.pending.is_empty()
}
fn enqueue_ok(&mut self, wire_seq: u64, fsn: u64, table_seq_txns: Vec<TableSeqTxn>) {
self.pending.push_back(PendingDurableFrame::Ok {
wire_seq,
fsn,
table_seq_txns,
});
}
fn apply_ack(&mut self, table_seq_txns: Vec<TableSeqTxn>) {
for entry in table_seq_txns {
match self.table_watermarks.get_mut(&entry.table) {
Some(current) if entry.seq_txn > *current => {
*current = entry.seq_txn;
}
Some(_) => {}
None => {
self.table_watermarks.insert(entry.table, entry.seq_txn);
}
}
}
}
fn pending_wire_seq_for_fsn(&self, fsn: u64) -> Option<u64> {
let front = self.pending.front()?;
let back = self.pending.back()?;
if fsn < front.fsn() || fsn > back.fsn() {
return None;
}
let offset = usize::try_from(fsn - front.fsn()).ok()?;
if let Some(entry) = self.pending.get(offset)
&& entry.fsn() == fsn
{
return Some(entry.wire_seq());
}
self.pending
.binary_search_by_key(&fsn, PendingDurableFrame::fsn)
.ok()
.and_then(|index| self.pending.get(index))
.map(PendingDurableFrame::wire_seq)
}
fn pending_prefix_covers(&self, start_fsn: u64, end_before_fsn: u64) -> bool {
let mut next_fsn = start_fsn;
for entry in &self.pending {
if next_fsn >= end_before_fsn {
return true;
}
let PendingDurableFrame::Ok { fsn, .. } = entry;
if *fsn < next_fsn {
continue;
}
next_fsn = match fsn.checked_add(1) {
Some(next_fsn) => next_fsn,
None => return false,
}
}
next_fsn >= end_before_fsn
}
fn pop_ready(&mut self) -> Option<DurableCompletion> {
if !self
.pending
.front()
.is_some_and(|entry| entry.is_covered_by(&self.table_watermarks))
{
return None;
}
match self.pending.pop_front().unwrap() {
PendingDurableFrame::Ok { wire_seq, fsn, .. } => {
Some(DurableCompletion { wire_seq, fsn })
}
}
}
}
impl PendingDurableFrame {
fn is_covered_by(&self, table_watermarks: &HashMap<String, i64>) -> bool {
match self {
PendingDurableFrame::Ok { table_seq_txns, .. } => table_seq_txns.iter().all(|entry| {
table_watermarks
.get(&entry.table)
.is_some_and(|watermark| *watermark >= entry.seq_txn)
}),
}
}
fn fsn(&self) -> u64 {
match self {
PendingDurableFrame::Ok { fsn, .. } => *fsn,
}
}
fn wire_seq(&self) -> u64 {
match self {
PendingDurableFrame::Ok { wire_seq, .. } => *wire_seq,
}
}
}
impl SenderErrorLog {
fn new(capacity: usize) -> Self {
Self {
errors: VecDeque::with_capacity(capacity),
capacity,
dropped_total: 0,
first_seq: 0,
next_seq: 0,
poll_next_seq: 0,
notification_next_seq: 0,
}
}
fn push(&mut self, error: QwpWsSenderError) {
let seq = self.next_seq;
self.next_seq = self.next_seq.saturating_add(1);
if self.capacity == 0 {
self.dropped_total += 1;
self.first_seq = self.next_seq;
self.poll_next_seq = self.poll_next_seq.max(self.first_seq);
self.notification_next_seq = self.notification_next_seq.max(self.first_seq);
return;
}
self.discard_consumed_prefix();
if self.errors.len() == self.capacity {
self.errors.pop_front();
self.first_seq = self.first_seq.saturating_add(1);
self.dropped_total += 1;
self.poll_next_seq = self.poll_next_seq.max(self.first_seq);
self.notification_next_seq = self.notification_next_seq.max(self.first_seq);
}
debug_assert_eq!(self.first_seq + self.errors.len() as u64, seq);
self.errors.push_back(error);
}
fn poll(&mut self) -> Option<QwpWsSenderError> {
let (next_seq, error) = self.poll_at(self.poll_next_seq);
self.poll_next_seq = next_seq;
self.discard_consumed_prefix();
error
}
fn poll_notification(&mut self) -> Option<QwpWsSenderError> {
let (next_seq, error) = self.poll_at(self.notification_next_seq);
self.notification_next_seq = next_seq;
self.discard_consumed_prefix();
error
}
fn poll_at(&self, mut next_seq: u64) -> (u64, Option<QwpWsSenderError>) {
next_seq = next_seq.max(self.first_seq);
if next_seq >= self.next_seq {
return (next_seq, None);
}
let index = (next_seq - self.first_seq) as usize;
match self.errors.get(index) {
Some(error) => (next_seq.saturating_add(1), Some(error.clone())),
None => (self.next_seq, None),
}
}
fn discard_consumed_prefix(&mut self) {
let keep_from = self.poll_next_seq.min(self.notification_next_seq);
while self.first_seq < keep_from && !self.errors.is_empty() {
self.errors.pop_front();
self.first_seq = self.first_seq.saturating_add(1);
}
}
fn capacity(&self) -> usize {
self.capacity
}
fn dropped_total(&self) -> u64 {
self.dropped_total
}
}
fn sender_error_overlaps(error: &QwpWsSenderError, from_fsn: u64, to_fsn: u64) -> bool {
error.from_fsn <= to_fsn && error.to_fsn >= from_fsn
}
impl DriverEventRing {
fn new(capacity: usize) -> Self {
Self {
events: VecDeque::with_capacity(capacity),
capacity,
dropped_total: 0,
}
}
fn push(&mut self, event: DriverEvent) {
if self.capacity == 0 {
self.dropped_total += 1;
return;
}
if self.events.len() == self.capacity {
self.events.pop_front();
self.dropped_total += 1;
}
self.events.push_back(event);
}
fn pop(&mut self) -> Option<DriverEvent> {
self.events.pop_front()
}
fn dropped_total(&self) -> u64 {
self.dropped_total
}
}
#[cfg(test)]
#[derive(Debug)]
pub(crate) struct FakeOrderedServer {
send_results: VecDeque<FakeSendResult>,
poll_responses: VecDeque<TransportResponse>,
default_send_result: FakeSendResult,
sent_frames: Vec<SentFrame>,
sent_payloads: Vec<Vec<u8>>,
server_max_batch_size: usize,
}
#[cfg(test)]
impl FakeOrderedServer {
pub(crate) fn no_response() -> Self {
Self::with_default(FakeSendResult::NoResponse)
}
pub(crate) fn ack_each_send() -> Self {
Self::with_default(FakeSendResult::AckSent)
}
pub(crate) fn scripted(send_results: impl IntoIterator<Item = FakeSendResult>) -> Self {
let mut server = Self::no_response();
server.send_results.extend(send_results);
server
}
pub(crate) fn push_response(&mut self, response: TransportResponse) {
self.poll_responses.push_back(response);
}
pub(crate) fn sent_frames(&self) -> &[SentFrame] {
&self.sent_frames
}
pub(crate) fn sent_payloads(&self) -> &[Vec<u8>] {
&self.sent_payloads
}
pub(crate) fn with_batch_cap(mut self, cap: usize) -> Self {
self.server_max_batch_size = cap;
self
}
fn with_default(default_send_result: FakeSendResult) -> Self {
Self {
send_results: VecDeque::new(),
poll_responses: VecDeque::new(),
default_send_result,
sent_frames: Vec::new(),
sent_payloads: Vec::new(),
server_max_batch_size: 0,
}
}
}
#[cfg(test)]
impl QwpWsCoreTransport for FakeOrderedServer {
fn try_poll_response(&mut self) -> Result<TransportPoll, TransportFailure> {
Ok(self
.poll_responses
.pop_front()
.map_or(TransportPoll::Idle, TransportPoll::Response))
}
fn send_frame(
&mut self,
frame: OutboundFrameView<'_>,
) -> Result<TransportSendResult, TransportFailure> {
self.sent_payloads.push(frame.payload.to_vec());
let frame = frame.sent_frame();
let wire_seq = frame.wire_seq;
self.sent_frames.push(frame);
let result = self
.send_results
.pop_front()
.unwrap_or(self.default_send_result);
Ok(match result {
FakeSendResult::NoResponse => TransportSendResult::NoResponse,
FakeSendResult::AckSent => {
TransportSendResult::Response(TransportResponse::Ack { wire_seq })
}
FakeSendResult::AckWire { wire_seq } => {
TransportSendResult::Response(TransportResponse::Ack { wire_seq })
}
FakeSendResult::RejectWire { wire_seq } => {
TransportSendResult::Response(TransportResponse::Reject {
wire_seq,
error: QwpServerError {
status: codec::WS_STATUS_WRITE_ERROR,
message: "fake write error".to_string(),
error: Error::new(ErrorCode::ServerFlushError, "fake write error"),
},
})
}
FakeSendResult::RejectWireNotWritable { wire_seq } => {
TransportSendResult::Response(TransportResponse::Reject {
wire_seq,
error: QwpServerError {
status: codec::WS_STATUS_NOT_WRITABLE,
message: "replica is read-only".to_string(),
error: Error::new(ErrorCode::ServerFlushError, "replica is read-only"),
},
})
}
FakeSendResult::Disconnect => TransportSendResult::Failure(
TransportFailure::Disconnect(fake_transport_error("fake disconnect")),
),
FakeSendResult::RetryableFailure => TransportSendResult::Failure(
TransportFailure::Retryable(fake_transport_error("fake retryable failure")),
),
FakeSendResult::TerminalFailure => TransportSendResult::Failure(
TransportFailure::Terminal(fake_transport_error("fake terminal failure")),
),
})
}
fn server_max_batch_size(&self) -> usize {
self.server_max_batch_size
}
}
#[cfg(test)]
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum FakeSendResult {
NoResponse,
AckSent,
AckWire {
wire_seq: u64,
},
RejectWire {
wire_seq: u64,
},
RejectWireNotWritable {
wire_seq: u64,
},
Disconnect,
RetryableFailure,
TerminalFailure,
}
#[cfg(test)]
mod tests {
use super::super::qwp_ws_publisher::QwpWsReplayEncoder;
use super::*;
use crate::ingress::buffer::{QwpWsColumnarBuffer, QwpWsEncodeScratch, SymbolGlobalDict};
use crate::ingress::rejection_events::RejectionEventSource;
use crate::ingress::{Buffer, QwpWsErrorCategory, QwpWsErrorPolicy, TimestampNanos};
use std::sync::atomic::AtomicBool;
fn terminal_latch_asserting_sink(
lifecycle: PublicationLifecycle,
) -> (Arc<RejectionEventSource>, Arc<AtomicBool>) {
let callback_ran = Arc::new(AtomicBool::new(false));
let callback_ran_in_handler = Arc::clone(&callback_ran);
let sink = RejectionEventSource::inline_for_test(crate::ingress::QwpWsErrorHandler::new(
move |error| {
if error.applied_policy == QwpWsErrorPolicy::Terminal {
assert_eq!(
lifecycle.load(),
PublicationState::Terminal,
"terminal rejection callback ran before the terminal latch"
);
callback_ran_in_handler.store(true, Ordering::Release);
}
},
));
(Arc::new(sink), callback_ran)
}
#[test]
fn catch_up_offset_maps_replay_acks_and_ignores_catch_up_acks() {
let mut cursor = SendCursor::new();
cursor.fsn_at_zero = Some(5);
cursor.next_fsn = Some(5);
cursor.begin_catch_up(3);
assert_eq!(
cursor.next_wire_seq, 3,
"replay resumes past the catch-up frames"
);
cursor
.commit_sent(SentFrame {
fsn: 5,
wire_seq: 3,
payload_len: 0,
})
.unwrap();
cursor
.commit_sent(SentFrame {
fsn: 6,
wire_seq: 4,
payload_len: 0,
})
.unwrap();
assert_eq!(cursor.ack_fsn_for_wire_seq(0).unwrap(), None);
assert_eq!(cursor.ack_fsn_for_wire_seq(2).unwrap(), None);
assert_eq!(cursor.ack_fsn_for_wire_seq(3).unwrap(), Some((5, 3)));
assert_eq!(cursor.ack_fsn_for_wire_seq(4).unwrap(), Some((6, 4)));
assert_eq!(cursor.reject_fsn_for_wire_seq(1).unwrap(), None);
assert_eq!(cursor.reject_fsn_for_wire_seq(4).unwrap(), Some((6, 4)));
}
#[test]
fn no_catch_up_leaves_wire_mapping_unchanged() {
let mut cursor = SendCursor::new();
cursor.fsn_at_zero = Some(10);
cursor.next_fsn = Some(10);
cursor
.commit_sent(SentFrame {
fsn: 10,
wire_seq: 0,
payload_len: 0,
})
.unwrap();
cursor
.commit_sent(SentFrame {
fsn: 11,
wire_seq: 1,
payload_len: 0,
})
.unwrap();
assert_eq!(cursor.ack_fsn_for_wire_seq(0).unwrap(), Some((10, 0)));
assert_eq!(cursor.ack_fsn_for_wire_seq(1).unwrap(), Some((11, 1)));
}
#[test]
fn hot_final_ack_releases_the_sfa_cursor() {
let mut queue = SfaFrameQueue::open_memory(SfaMemoryQueueOptions {
segment_size_bytes: 128,
max_bytes: 256,
})
.unwrap();
queue.try_submit(b"frame").unwrap();
let progress = queue.progress_view();
let mut core = QwpWsSendCore::new(
FakeOrderedServer::no_response(),
ReconnectPolicy::no_backoff(Duration::MAX),
);
let outbound = core.next_outbound_sfa_frame(&progress).unwrap().unwrap();
core.send_cursor.commit_sent(outbound.sent_frame()).unwrap();
assert!(core.send_cursor.sfa_cursor.is_some());
let completion = core.finish_ack_response_sfa(&progress, 0).unwrap();
assert_eq!(completion.outcome, DriveOutcome::Acked { wire_seq: 0 });
assert!(core.send_cursor.sfa_cursor.is_none());
assert_eq!(progress.completed_fsn(), Some(0));
}
fn write_frame_varint(out: &mut Vec<u8>, mut value: u64) {
while value > 0x7F {
out.push(((value & 0x7F) as u8) | 0x80);
value >>= 7;
}
out.push(value as u8);
}
fn read_frame_varint(buf: &[u8], mut pos: usize) -> (u64, usize) {
let mut value = 0u64;
let mut shift = 0;
loop {
let b = buf[pos];
pos += 1;
value |= u64::from(b & 0x7F) << shift;
if b & 0x80 == 0 {
return (value, pos);
}
shift += 7;
}
}
fn make_delta_frame(delta_start: u64, symbols: &[&[u8]]) -> Vec<u8> {
let mut payload = Vec::new();
write_frame_varint(&mut payload, delta_start);
write_frame_varint(&mut payload, symbols.len() as u64);
for s in symbols {
write_frame_varint(&mut payload, s.len() as u64);
payload.extend_from_slice(s);
}
payload.push(0xAB); let mut frame = Vec::new();
frame.extend_from_slice(b"QWP1");
frame.push(1); frame.push(0x08); frame.extend_from_slice(&1u16.to_le_bytes()); frame.extend_from_slice(&(payload.len() as u32).to_le_bytes());
frame.extend_from_slice(&payload);
frame
}
fn assert_catch_up_frame(frame: &[u8], expected: &[&[u8]]) {
assert_eq!(&frame[0..4], b"QWP1", "catch-up frame magic");
assert_eq!(
frame[5] & 0x08,
0x08,
"catch-up frame carries the delta flag"
);
assert_eq!(
u16::from_le_bytes([frame[6], frame[7]]),
0,
"catch-up frame is table-less"
);
let (delta_start, p) = read_frame_varint(frame, 12);
assert_eq!(delta_start, 0, "catch-up re-registers from id 0");
let (count, mut p) = read_frame_varint(frame, p);
assert_eq!(count as usize, expected.len(), "catch-up entry count");
for exp in expected {
let (len, after) = read_frame_varint(frame, p);
let end = after + len as usize;
assert_eq!(&frame[after..end], *exp, "catch-up entry bytes");
p = end;
}
assert_eq!(p, frame.len(), "no trailing bytes after the dictionary");
}
#[test]
fn reconnect_emits_full_dict_catch_up_before_replay() {
let frame = make_delta_frame(0, &[b"AAPL", b"GOOG"]);
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[], 0);
driver.try_submit(&frame).unwrap();
driver.drive_send_once().unwrap();
assert_eq!(driver.send_core.dict_mirror.count(), 2);
assert_eq!(driver.send_core.transport.sent_payloads().len(), 1);
let outcome = driver
.send_core
.finish_reconnect_success(&mut driver.store, ReconnectReason::Disconnect);
assert!(matches!(outcome, DriveOutcome::Reconnected { .. }));
assert!(driver.send_core.catch_up_pending);
driver.drive_send_once().unwrap();
assert!(!driver.send_core.catch_up_pending, "catch-up consumed");
let payloads = driver.send_core.transport.sent_payloads();
assert_eq!(payloads.len(), 3, "catch-up frame + replay emitted");
assert_catch_up_frame(&payloads[1], &[b"AAPL", b"GOOG"]);
assert_eq!(payloads[2], frame, "the original frame replays verbatim");
}
#[test]
fn recovery_seed_emits_catch_up_on_first_connect_before_replay() {
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[1, b'a', 1, b'b'], 2);
assert!(
driver.send_core.catch_up_pending,
"a recovered seed must arm the first-connect catch-up"
);
let frame = make_delta_frame(2, &[b"c"]);
driver.try_submit(&frame).unwrap();
driver.drive_send_once().unwrap();
assert!(
!driver.send_core.catch_up_pending,
"catch-up consumed on the first send"
);
let payloads = driver.send_core.transport.sent_payloads();
assert_eq!(payloads.len(), 2, "catch-up frame precedes the replay");
assert_catch_up_frame(&payloads[0], &[b"a", b"b"]);
assert_eq!(payloads[1], frame, "the stored frame replays verbatim");
let counters = driver.counters();
assert_eq!(counters.total_frames_sent, 1);
assert_eq!(
counters.total_frames_replayed, 0,
"initial recovery sends are not post-reconnect replays"
);
}
#[test]
fn replay_totals_use_the_reconnect_time_publication_boundary() {
let mut driver = driver(FakeOrderedServer::no_response());
driver.try_submit(b"sent-before-reconnect").unwrap();
driver.drive_send_once().unwrap();
driver.try_submit(b"published-before-reconnect").unwrap();
let counters = driver.counters();
assert_eq!(counters.total_frames_sent, 1);
assert_eq!(counters.total_frames_replayed, 0);
assert_eq!(
driver
.send_core
.finish_reconnect_success(&mut driver.store, ReconnectReason::Disconnect),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver.try_submit(b"published-after-reconnect").unwrap();
driver.drive_send_once().unwrap();
let counters = driver.counters();
assert_eq!(counters.total_frames_sent, 4);
assert_eq!(counters.total_frames_replayed, 2);
assert_eq!(
driver
.send_core
.finish_reconnect_success(&mut driver.store, ReconnectReason::Disconnect),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
let counters = driver.counters();
assert_eq!(counters.total_frames_sent, 7);
assert_eq!(counters.total_frames_replayed, 5);
}
#[test]
fn torn_dict_guard_rejects_a_conflicting_symbol_redefinition() {
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[1, b'a'], 1);
assert!(
driver
.send_core
.dict_mirror
.accumulate(&make_delta_frame(1, &[b"b"])),
"folding a small frame cannot fail"
);
assert_eq!(driver.send_core.dict_mirror.count(), 2);
let same = make_delta_frame(1, &[b"b"]);
assert!(driver.send_core.guard_dict_not_torn(&same).is_ok());
let conflicting = make_delta_frame(1, &[b"c"]);
let err = driver
.send_core
.guard_dict_not_torn(&conflicting)
.unwrap_err();
assert_eq!(err.code(), ErrorCode::StoreResendRequired);
}
#[test]
fn torn_dict_guard_rejects_a_frame_referencing_lost_symbols() {
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[1, b'a', 1, b'b'], 2);
let ok = make_delta_frame(2, &[b"c"]);
assert!(driver.send_core.guard_dict_not_torn(&ok).is_ok());
let torn = make_delta_frame(5, &[b"z"]);
let err = driver.send_core.guard_dict_not_torn(&torn).unwrap_err();
assert_eq!(err.code(), ErrorCode::StoreResendRequired);
assert!(
err.msg().contains("dictionary is torn"),
"msg: {}",
err.msg()
);
}
#[test]
fn torn_dict_guard_in_full_dict_mode_passes_dense_frames_but_rejects_stray_deltas() {
let driver = driver(FakeOrderedServer::no_response());
let dense = make_delta_frame(0, &[b"z"]);
assert!(driver.send_core.guard_dict_not_torn(&dense).is_ok());
let stray = make_delta_frame(9, &[b"z"]);
let err = driver.send_core.guard_dict_not_torn(&stray).unwrap_err();
assert_eq!(err.code(), ErrorCode::StoreResendRequired);
assert!(
err.msg()
.contains("delta symbol-dictionary mode is disabled"),
"msg: {}",
err.msg()
);
}
#[test]
fn torn_dict_guard_accepts_a_dense_frame_while_the_mirror_is_enabled() {
{
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[1, b'a', 1, b'b'], 2);
assert_eq!(driver.send_core.dict_mirror.count(), 2);
let dense = make_delta_frame(0, &[b"a", b"b", b"c"]);
assert!(
driver.send_core.guard_dict_not_torn(&dense).is_ok(),
"a self-sufficient dense frame re-shipping the mirrored prefix is safe"
);
assert!(
driver.send_core.dict_mirror.accumulate(&dense),
"folding a small frame cannot fail"
);
assert_eq!(
driver.send_core.dict_mirror.count(),
3,
"accumulate folds only the [M, K) suffix, keeping the mirror in lockstep"
);
}
{
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[1, b'a', 1, b'b'], 2);
let conflicting = make_delta_frame(0, &[b"a", b"X", b"c"]);
let err = driver
.send_core
.guard_dict_not_torn(&conflicting)
.unwrap_err();
assert_eq!(err.code(), ErrorCode::StoreResendRequired);
}
}
#[test]
fn reconnect_error_is_terminal_treats_store_resend_required_as_terminal() {
assert!(
reconnect_error_is_terminal(&Error::new(
ErrorCode::StoreResendRequired,
"the affected data must be resent"
)),
"StoreResendRequired must be terminal so the reconnect loop stops"
);
assert!(
!reconnect_error_is_terminal(&Error::new(ErrorCode::SocketError, "connection reset")),
"a transient SocketError must stay retryable"
);
}
#[test]
fn reconnect_catch_up_splits_across_frames_when_server_caps_batch() {
let syms: Vec<Vec<u8>> = (0..5).map(|i| format!("sy{i:02}").into_bytes()).collect();
let mut seed = Vec::new();
for s in &syms {
write_frame_varint(&mut seed, s.len() as u64);
seed.extend_from_slice(s);
}
let cap = 12 + 16 + 12;
let mut driver = driver(FakeOrderedServer::no_response().with_batch_cap(cap));
driver.send_core.enable_delta_dict(&seed, syms.len() as u32);
assert!(driver.send_core.catch_up_pending);
let replay = make_delta_frame(5, &[b"new"]);
driver.try_submit(&replay).unwrap();
driver.drive_send_once().unwrap();
assert!(!driver.send_core.catch_up_pending, "catch-up consumed");
let payloads = driver.send_core.transport.sent_payloads();
assert!(
payloads.len() > 2,
"expected a multi-frame catch-up split + replay, got {}",
payloads.len()
);
let (catch_up, replay_sent) = payloads.split_at(payloads.len() - 1);
assert!(catch_up.len() >= 2, "catch-up must span multiple frames");
let mut expected_start = 0u64;
let mut reassembled: Vec<Vec<u8>> = Vec::new();
for f in catch_up {
assert!(
f.len() <= cap,
"catch-up frame ({}) exceeds cap {}",
f.len(),
cap
);
assert_eq!(&f[0..4], b"QWP1");
assert_eq!(f[5] & 0x08, 0x08, "catch-up carries the delta flag");
assert_eq!(
u16::from_le_bytes([f[6], f[7]]),
0,
"catch-up is table-less"
);
let (start, p) = read_frame_varint(f, 12);
assert_eq!(start, expected_start, "catch-up ranges are contiguous");
let (count, mut p) = read_frame_varint(f, p);
for _ in 0..count {
let (len, after) = read_frame_varint(f, p);
reassembled.push(f[after..after + len as usize].to_vec());
p = after + len as usize;
}
expected_start += count;
}
assert_eq!(
reassembled, syms,
"catch-up re-registers every symbol gap-free"
);
assert_eq!(
replay_sent[0], replay,
"the queued frame replays verbatim after the catch-up"
);
}
#[test]
fn catch_up_entry_exceeding_batch_cap_reconnects_without_terminalizing() {
let big = vec![b'x'; 64];
let mut seed = Vec::new();
write_frame_varint(&mut seed, big.len() as u64);
seed.extend_from_slice(&big);
let mut driver = driver(FakeOrderedServer::no_response().with_batch_cap(20));
driver.send_core.enable_delta_dict(&seed, 1);
assert!(driver.send_core.catch_up_pending);
let outcome = driver.drive_send_once().unwrap();
assert!(matches!(
outcome,
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
));
assert!(!driver.store.is_terminal());
assert!(
driver.send_core.catch_up_pending,
"successful reconnect must re-arm catch-up"
);
assert!(
driver.send_core.transport.sent_payloads().is_empty(),
"oversized catch-up entry must not be sent"
);
assert_eq!(driver.send_core.catch_up_retry_strikes, 1);
driver.send_core.transport.server_max_batch_size = 0;
driver.drive_send_once().unwrap();
assert_eq!(driver.send_core.transport.sent_payloads().len(), 1);
assert_eq!(driver.send_core.catch_up_retry_strikes, 0);
}
#[test]
fn drive_marks_terminal_when_a_stored_frame_outruns_the_recovered_dictionary() {
let mut driver = driver(FakeOrderedServer::no_response());
driver.send_core.enable_delta_dict(&[1, b'a', 1, b'b'], 2);
let torn = make_delta_frame(5, &[b"z"]);
driver.try_submit(&torn).unwrap();
let outcome = driver.drive_send_once().unwrap();
assert!(matches!(outcome, DriveOutcome::Terminal));
let err = driver
.store
.terminal_error()
.expect("terminal error recorded");
assert!(
err.msg().contains("dictionary is torn"),
"msg: {}",
err.msg()
);
}
#[test]
fn presend_reject_of_the_catch_up_escalates_to_terminal_not_an_infinite_loop() {
let mut driver = driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"queued-frame").unwrap();
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 },
"the queued frame the catch-up exists to unblock",
);
for _ in 0..DEFAULT_MAX_FRAME_REJECTIONS - 1 {
let outcome = driver
.send_core
.record_presend_reject(
&mut driver.store,
0,
write_error("replica rejects writes"),
QwpWsErrorPolicy::Retriable,
true,
)
.unwrap();
assert_ne!(outcome, DriveOutcome::Terminal, "must not terminate early");
}
let outcome = driver
.send_core
.record_presend_reject(
&mut driver.store,
0,
write_error("replica rejects writes"),
QwpWsErrorPolicy::Retriable,
true,
)
.unwrap();
assert_eq!(
outcome,
DriveOutcome::Terminal,
"a persistently-rejected catch-up must escalate, not loop forever",
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
let terminal = driver
.store
.terminal_error()
.expect("terminal error recorded");
assert!(
terminal.msg().contains("catch-up"),
"terminal error names the catch-up cause: {}",
terminal.msg(),
);
}
use std::collections::VecDeque;
#[cfg(feature = "sync-sender-qwp-ws")]
use std::fs::File;
#[cfg(feature = "sync-sender-qwp-ws")]
use std::io::{Read, Write};
#[cfg(feature = "sync-sender-qwp-ws")]
use std::net::TcpListener;
#[cfg(feature = "sync-sender-qwp-ws")]
use std::sync::mpsc;
#[cfg(feature = "sync-sender-qwp-ws")]
use std::thread;
#[cfg(feature = "sync-sender-qwp-ws")]
use std::time::Duration;
#[cfg(feature = "sync-sender-qwp-ws")]
use rustls::{ServerConfig, StreamOwned, server::ServerConnection};
#[cfg(feature = "sync-sender-qwp-ws")]
use rustls_pki_types::pem::PemObject;
#[cfg(feature = "sync-sender-qwp-ws")]
use rustls_pki_types::{CertificateDer, PrivateKeyDer};
#[derive(Debug)]
struct DelayedPollAckServer {
polls_before_ack: usize,
ack_sent: bool,
sent_frames: Vec<SentFrame>,
}
impl DelayedPollAckServer {
fn new(polls_before_ack: usize) -> Self {
Self {
polls_before_ack,
ack_sent: false,
sent_frames: Vec::new(),
}
}
}
impl QwpWsCoreTransport for DelayedPollAckServer {
fn try_poll_response(&mut self) -> Result<TransportPoll, TransportFailure> {
if self.sent_frames.is_empty() || self.ack_sent {
return Ok(TransportPoll::Idle);
}
if self.polls_before_ack > 0 {
self.polls_before_ack -= 1;
return Ok(TransportPoll::Idle);
}
self.ack_sent = true;
Ok(TransportPoll::Response(TransportResponse::Ack {
wire_seq: self.sent_frames[0].wire_seq,
}))
}
fn send_frame(
&mut self,
frame: OutboundFrameView<'_>,
) -> Result<TransportSendResult, TransportFailure> {
self.sent_frames.push(frame.sent_frame());
Ok(TransportSendResult::NoResponse)
}
}
fn options(_max_frames: usize, max_bytes: usize) -> SfaMemoryQueueOptions {
SfaMemoryQueueOptions {
segment_size_bytes: 256,
max_bytes,
}
}
fn memory_queue(options: SfaMemoryQueueOptions) -> SfaFrameQueue {
SfaFrameQueue::open_memory(options).unwrap()
}
type FakeDriver = QwpWsCoreTestHarness<SfaFrameQueue, FakeOrderedServer>;
fn driver(server: FakeOrderedServer) -> FakeDriver {
QwpWsCoreTestHarness::new(options(8, 1024), server).unwrap()
}
#[derive(Debug)]
enum PublishTestError {
Encode(crate::Error),
Driver(DriverError),
}
impl From<DriverError> for PublishTestError {
fn from(value: DriverError) -> Self {
Self::Driver(value)
}
}
fn publish_qwp<Q, T>(
driver: &mut QwpWsCoreTestHarness<Q, T>,
encoder: &mut QwpWsReplayEncoder,
buffer: &QwpWsColumnarBuffer,
) -> Result<QwpReceipt, PublishTestError>
where
Q: PublicationLog,
T: QwpWsCoreTransport,
{
let payload = encoder.encode(buffer).map_err(PublishTestError::Encode)?;
Ok(driver.try_submit(payload)?)
}
fn wait_for_delivery<Q, T>(
driver: &mut QwpWsCoreTestHarness<Q, T>,
receipt: QwpReceipt,
timeout: Duration,
) -> Result<DeliveryOutcome, DriverError>
where
Q: PublicationLog,
T: QwpWsCoreTransport,
{
let deadline = Instant::now() + timeout;
loop {
if let Some(outcome) = driver.delivery_status(receipt)? {
return Ok(outcome);
}
let remaining = deadline.saturating_duration_since(Instant::now());
if remaining.is_zero() {
return Ok(DeliveryOutcome::Timeout);
}
if driver.drive_once()? == DriveOutcome::Idle {
std::thread::sleep(remaining.min(Duration::from_micros(100)));
}
}
}
fn durable_driver(server: FakeOrderedServer) -> FakeDriver {
QwpWsCoreTestHarness::from_queue_with_durable_ack(memory_queue(options(8, 1024)), server)
}
fn durable_driver_with_options(
options: SfaMemoryQueueOptions,
server: FakeOrderedServer,
) -> FakeDriver {
QwpWsCoreTestHarness::from_queue_with_durable_ack(memory_queue(options), server)
}
fn driver_with_event_capacity(server: FakeOrderedServer, event_capacity: usize) -> FakeDriver {
QwpWsCoreTestHarness::from_queue_with_event_capacity(
memory_queue(options(8, 1024)),
server,
event_capacity,
)
}
fn sender_error(fsn: u64) -> QwpWsSenderError {
QwpWsSenderError {
category: QwpWsErrorCategory::WriteError,
applied_policy: QwpWsErrorPolicy::Retriable,
status: Some(codec::WS_STATUS_WRITE_ERROR),
message: Some(format!("error {fsn}")),
message_sequence: Some(fsn),
from_fsn: fsn,
to_fsn: fsn,
}
}
const QWP_WS_COLUMNAR_BENCH_BATCH_SIZE: usize = 1000;
fn qwp_ws_columnar_bench_rows() -> usize {
std::env::var("QWP_WS_COLUMNAR_BENCH_ROWS")
.ok()
.and_then(|value| value.parse::<usize>().ok())
.filter(|rows| *rows > 0)
.unwrap_or(20_000_000)
}
fn fill_qwp_ws_columnar_benchmark_batch(buffer: &mut Buffer, batch_idx: usize, rows: usize) {
let symbols = [
"SYM000", "SYM001", "SYM002", "SYM003", "SYM004", "SYM005", "SYM006", "SYM007",
];
let venues = ["ldn", "nyc", "ams", "fra", "sin", "hkg", "tyo", "sfo"];
for row_idx in 0..rows {
let seq = (batch_idx * QWP_WS_COLUMNAR_BENCH_BATCH_SIZE + row_idx) as i64;
buffer
.table("trades")
.unwrap()
.symbol("sym", symbols[row_idx & 7])
.unwrap()
.column_i64("qty", seq)
.unwrap()
.column_f64("px", 100.0 + (seq & 1023) as f64)
.unwrap()
.column_str("venue", venues[row_idx & 7])
.unwrap()
.column_ts("event_ts", TimestampNanos::new(seq))
.unwrap()
.at(TimestampNanos::new(seq))
.unwrap();
}
}
fn qwp_buffer(sym: &str, qty: i64, ts: i64) -> Buffer {
let mut buffer = Buffer::qwp_ws_with_max_name_len(127);
buffer
.table("trades")
.unwrap()
.symbol("sym", sym)
.unwrap()
.column_i64("qty", qty)
.unwrap()
.column_f64("px", 100.0 + qty as f64)
.unwrap()
.at(TimestampNanos::new(ts))
.unwrap();
buffer
}
fn replay_payload(buffer: &Buffer) -> Vec<u8> {
let mut scratch = QwpWsEncodeScratch::new();
let mut global_dict = SymbolGlobalDict::new();
buffer
.as_qwp_ws()
.unwrap()
.encode_ws_replay_message(&mut scratch, &mut global_dict, 1)
.unwrap();
scratch.message
}
fn qwp_error_payload(status: u8, sequence: u64, message: &str) -> Vec<u8> {
let msg = message.as_bytes();
let mut payload = Vec::with_capacity(11 + msg.len());
payload.push(status);
payload.extend_from_slice(&sequence.to_le_bytes());
payload.extend_from_slice(&(msg.len() as u16).to_le_bytes());
payload.extend_from_slice(msg);
payload
}
fn qwp_durable_ack_payload(entries: &[(&str, i64)]) -> Vec<u8> {
let mut payload = Vec::new();
payload.push(codec::WS_STATUS_DURABLE_ACK);
append_table_seq_txns(&mut payload, entries);
payload
}
fn qwp_ok_payload_with_table_entries(sequence: u64, entries: &[(&str, i64)]) -> Vec<u8> {
let mut payload = Vec::new();
payload.push(codec::WS_STATUS_OK);
payload.extend_from_slice(&sequence.to_le_bytes());
append_table_seq_txns(&mut payload, entries);
payload
}
fn append_table_seq_txns(payload: &mut Vec<u8>, entries: &[(&str, i64)]) {
payload.extend_from_slice(&(entries.len() as u16).to_le_bytes());
for (table, seq_txn) in entries {
payload.extend_from_slice(&(table.len() as u16).to_le_bytes());
payload.extend_from_slice(table.as_bytes());
payload.extend_from_slice(&seq_txn.to_le_bytes());
}
}
fn table_seq_txns(entries: &[(&str, i64)]) -> Vec<TableSeqTxn> {
entries
.iter()
.map(|(table, seq_txn)| TableSeqTxn {
table: (*table).to_string(),
seq_txn: *seq_txn,
})
.collect()
}
fn schema_mismatch_error(message: &str) -> QwpServerError {
QwpServerError {
status: codec::WS_STATUS_SCHEMA_MISMATCH,
message: message.to_string(),
error: Error::new(ErrorCode::InvalidApiCall, message),
}
}
fn write_error(message: &str) -> QwpServerError {
QwpServerError {
status: codec::WS_STATUS_WRITE_ERROR,
message: message.to_string(),
error: Error::new(ErrorCode::ServerFlushError, message),
}
}
fn not_writable_error(message: &str) -> QwpServerError {
QwpServerError {
status: codec::WS_STATUS_NOT_WRITABLE,
message: message.to_string(),
error: Error::new(ErrorCode::ServerFlushError, message),
}
}
fn drain_events<Q: PublicationLog, T: QwpWsCoreTransport>(
driver: &mut QwpWsCoreTestHarness<Q, T>,
) -> Vec<DriverEvent> {
let mut events = Vec::new();
while let Some(event) = driver.poll_event() {
events.push(event);
}
events
}
#[derive(Debug)]
struct TestTransport {
send_results: VecDeque<Result<TransportSendResult, TransportFailure>>,
poll_results: VecDeque<Result<TransportPoll, TransportFailure>>,
keepalive_results: VecDeque<Result<bool, TransportFailure>>,
restart_results: VecDeque<Result<(), DriverError>>,
keepalive_attempts: usize,
keepalive_pending_args: Vec<bool>,
restart_attempts: usize,
restart_reasons: Vec<ReconnectReason>,
sent_frames: Vec<SentFrame>,
sent_payloads: Vec<Vec<u8>>,
}
impl TestTransport {
fn scripted(
send_results: impl IntoIterator<Item = Result<TransportSendResult, TransportFailure>>,
) -> Self {
Self {
send_results: send_results.into_iter().collect(),
poll_results: VecDeque::new(),
keepalive_results: VecDeque::new(),
restart_results: VecDeque::new(),
keepalive_attempts: 0,
keepalive_pending_args: Vec::new(),
restart_attempts: 0,
restart_reasons: Vec::new(),
sent_frames: Vec::new(),
sent_payloads: Vec::new(),
}
}
fn with_restart_results(
mut self,
restart_results: impl IntoIterator<Item = Result<(), DriverError>>,
) -> Self {
self.restart_results = restart_results.into_iter().collect();
self
}
fn with_poll_results(
mut self,
poll_results: impl IntoIterator<Item = Result<Option<TransportResponse>, TransportFailure>>,
) -> Self {
self.poll_results = poll_results
.into_iter()
.map(|result| {
result.map(|response| {
response.map_or(TransportPoll::Idle, TransportPoll::Response)
})
})
.collect();
self
}
fn with_poll_events(
mut self,
poll_results: impl IntoIterator<Item = Result<TransportPoll, TransportFailure>>,
) -> Self {
self.poll_results = poll_results.into_iter().collect();
self
}
fn with_keepalive_results(
mut self,
keepalive_results: impl IntoIterator<Item = Result<bool, TransportFailure>>,
) -> Self {
self.keepalive_results = keepalive_results.into_iter().collect();
self
}
}
impl QwpWsCoreTransport for TestTransport {
fn try_poll_response(&mut self) -> Result<TransportPoll, TransportFailure> {
self.poll_results
.pop_front()
.unwrap_or(Ok(TransportPoll::Idle))
}
fn send_durable_ack_keepalive_if_due(
&mut self,
durable_ack_pending: bool,
) -> Result<bool, TransportFailure> {
self.keepalive_attempts += 1;
self.keepalive_pending_args.push(durable_ack_pending);
if !durable_ack_pending {
return Ok(false);
}
self.keepalive_results.pop_front().unwrap_or(Ok(false))
}
fn send_frame(
&mut self,
frame: OutboundFrameView<'_>,
) -> Result<TransportSendResult, TransportFailure> {
self.sent_frames.push(frame.sent_frame());
self.sent_payloads.push(frame.payload.to_vec());
self.send_results
.pop_front()
.unwrap_or(Ok(TransportSendResult::NoResponse))
}
fn restart_connection(&mut self, reason: ReconnectReason) -> Result<(), DriverError> {
self.restart_attempts += 1;
self.restart_reasons.push(reason);
self.restart_results.pop_front().unwrap_or(Ok(()))
}
}
#[test]
fn publisher_sends_replay_payload_to_driver_transport() {
let buffer = qwp_buffer("SYM_001", 7, 1_000);
let expected = replay_payload(&buffer);
let driver = QwpWsCoreTestHarness::from_queue(
memory_queue(options(8, 4096)),
TestTransport::scripted([Ok(TransportSendResult::NoResponse)]),
);
let mut driver = driver;
let mut encoder = QwpWsReplayEncoder::new(1);
let receipt = publish_qwp(&mut driver, &mut encoder, buffer.as_qwp_ws().unwrap()).unwrap();
let outcome = driver.drive_once().unwrap();
assert_eq!(receipt, QwpReceipt { fsn: 0 });
assert_eq!(
outcome,
DriveOutcome::Sent(SentFrame {
fsn: 0,
wire_seq: 0,
payload_len: expected.len(),
})
);
assert_eq!(driver.send_core.transport.sent_payloads, vec![expected]);
}
#[test]
fn publisher_rejects_empty_buffer_without_publication() {
let buffer = Buffer::qwp_ws_with_max_name_len(127);
let driver = QwpWsCoreTestHarness::from_queue(
memory_queue(options(8, 4096)),
TestTransport::scripted([]),
);
let mut driver = driver;
let mut encoder = QwpWsReplayEncoder::new(1);
let err = publish_qwp(&mut driver, &mut encoder, buffer.as_qwp_ws().unwrap()).unwrap_err();
match err {
PublishTestError::Encode(err) => {
assert_eq!(err.code(), crate::ErrorCode::InvalidApiCall);
assert_eq!(err.msg(), "Cannot submit an empty QWP/WebSocket buffer.");
}
PublishTestError::Driver(err) => panic!("unexpected driver error: {err:?}"),
}
assert!(driver.send_core.transport.sent_payloads.is_empty());
}
#[test]
#[ignore]
fn publisher_memory_sfa_zero_alloc_after_warmup() {
use crate::alloc_counter;
let buffer = qwp_buffer("SYM_001", 7, 1_000);
let queue = memory_queue(options(8, 4096));
let driver = QwpWsCoreTestHarness::from_queue(queue, FakeOrderedServer::ack_each_send());
let mut driver = driver;
let mut encoder = QwpWsReplayEncoder::new(1);
for _ in 0..4 {
let receipt =
publish_qwp(&mut driver, &mut encoder, buffer.as_qwp_ws().unwrap()).unwrap();
assert_eq!(
wait_for_delivery(&mut driver, receipt, Duration::from_secs(5)).unwrap(),
DeliveryOutcome::Completed
);
}
alloc_counter::start_counting();
let receipt = publish_qwp(&mut driver, &mut encoder, buffer.as_qwp_ws().unwrap()).unwrap();
let alloc_count = alloc_counter::stop_counting();
assert_eq!(receipt, QwpReceipt { fsn: 4 });
assert_eq!(
alloc_count, 0,
"Expected zero allocations for warmed QWP/WebSocket memory publication, got {alloc_count}"
);
}
#[test]
#[ignore = "performance benchmark"]
fn qwp_ws_columnar_memory_publication_benchmark() {
let rows = qwp_ws_columnar_bench_rows();
let batches = rows.div_ceil(QWP_WS_COLUMNAR_BENCH_BATCH_SIZE);
let mut buffer = Buffer::qwp_ws_with_max_name_len(127);
let queue = memory_queue(options(8, 1 << 20));
let driver = QwpWsCoreTestHarness::from_queue(queue, FakeOrderedServer::ack_each_send());
let mut driver = driver;
let mut encoder = QwpWsReplayEncoder::new(1);
fill_qwp_ws_columnar_benchmark_batch(&mut buffer, 0, QWP_WS_COLUMNAR_BENCH_BATCH_SIZE);
let receipt = publish_qwp(&mut driver, &mut encoder, buffer.as_qwp_ws().unwrap()).unwrap();
assert_eq!(
wait_for_delivery(&mut driver, receipt, Duration::from_secs(5)).unwrap(),
DeliveryOutcome::Completed
);
buffer.clear();
let started = std::time::Instant::now();
let mut published_rows = 0usize;
for batch_idx in 0..batches {
let rows_in_batch = (rows - published_rows).min(QWP_WS_COLUMNAR_BENCH_BATCH_SIZE);
fill_qwp_ws_columnar_benchmark_batch(&mut buffer, batch_idx, rows_in_batch);
let receipt =
publish_qwp(&mut driver, &mut encoder, buffer.as_qwp_ws().unwrap()).unwrap();
assert_eq!(
wait_for_delivery(&mut driver, receipt, Duration::from_secs(5)).unwrap(),
DeliveryOutcome::Completed
);
buffer.clear();
published_rows += rows_in_batch;
}
let elapsed = started.elapsed();
eprintln!(
"qwp_ws_columnar_memory_publication_benchmark rows={} batch_size={} end_to_end_ms={} rows_per_sec={:.2}",
rows,
QWP_WS_COLUMNAR_BENCH_BATCH_SIZE,
elapsed.as_millis(),
rows as f64 / elapsed.as_secs_f64()
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn tls_certs_dir() -> std::path::PathBuf {
let mut certs_dir = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"));
certs_dir.pop();
certs_dir.push("tls_certs");
certs_dir
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn tls_server_config() -> std::sync::Arc<ServerConfig> {
let certs_dir = tls_certs_dir();
let cert_file = File::open(certs_dir.join("server.crt")).unwrap();
let private_key_file = File::open(certs_dir.join("server.key")).unwrap();
let certs = CertificateDer::pem_reader_iter(cert_file)
.collect::<Result<Vec<_>, _>>()
.unwrap();
let private_key = PrivateKeyDer::from_pem_reader(private_key_file).unwrap();
std::sync::Arc::new(
ServerConfig::builder()
.with_no_client_auth()
.with_single_cert(certs, private_key)
.unwrap(),
)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn tls_client_settings() -> TlsSettings {
let cert_file = File::open(tls_certs_dir().join("server_rootCA.pem")).unwrap();
let certs = CertificateDer::pem_reader_iter(cert_file)
.collect::<Result<Vec<_>, _>>()
.unwrap();
TlsSettings::PemFile(certs)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn read_request_until_blank<S: Read>(stream: &mut S) -> std::io::Result<String> {
let mut bytes = Vec::new();
let mut tmp = [0u8; 256];
loop {
let n = stream.read(&mut tmp)?;
if n == 0 {
break;
}
bytes.extend_from_slice(&tmp[..n]);
if bytes.windows(4).any(|window| window == b"\r\n\r\n") {
break;
}
}
Ok(String::from_utf8_lossy(&bytes).to_string())
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn header_value(request: &str, name: &str) -> String {
request
.split("\r\n")
.find_map(|line| {
let (key, value) = line.split_once(':')?;
key.trim()
.eq_ignore_ascii_case(name)
.then(|| value.trim().to_string())
})
.unwrap()
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn read_client_frame<S: Read>(stream: &mut S) -> std::io::Result<(u8, Vec<u8>)> {
let mut hdr = [0u8; 2];
stream.read_exact(&mut hdr)?;
let opcode = hdr[0] & 0x0f;
assert_ne!(hdr[1] & 0x80, 0, "client WebSocket frame must be masked");
let len_short = hdr[1] & 0x7f;
let payload_len = match len_short {
126 => {
let mut bytes = [0u8; 2];
stream.read_exact(&mut bytes)?;
u16::from_be_bytes(bytes) as usize
}
127 => {
let mut bytes = [0u8; 8];
stream.read_exact(&mut bytes)?;
u64::from_be_bytes(bytes) as usize
}
len => len as usize,
};
let mut mask = [0u8; 4];
stream.read_exact(&mut mask)?;
let mut payload = vec![0u8; payload_len];
stream.read_exact(&mut payload)?;
for (index, byte) in payload.iter_mut().enumerate() {
*byte ^= mask[index & 3];
}
Ok((opcode, payload))
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn write_server_binary_frame<S: Write>(stream: &mut S, payload: &[u8]) -> std::io::Result<()> {
let mut frame = vec![0x82];
match payload.len() {
len @ 0..=125 => frame.push(len as u8),
len @ 126..=0xffff => {
frame.push(126);
frame.extend_from_slice(&(len as u16).to_be_bytes());
}
len => {
frame.push(127);
frame.extend_from_slice(&(len as u64).to_be_bytes());
}
}
frame.extend_from_slice(payload);
stream.write_all(&frame)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn write_ok_response<S: Write>(stream: &mut S, wire_seq: u64) -> std::io::Result<()> {
let mut ok = vec![0u8];
ok.extend_from_slice(&wire_seq.to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(stream, &ok)
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn upgrade_qwp_ws_test_connection<S: Read + Write>(stream: &mut S, durable_ack: bool) {
let request = read_request_until_blank(stream).unwrap();
let accept =
crate::ws::crypto::compute_accept(&header_value(&request, "Sec-WebSocket-Key"));
if durable_ack {
assert_eq!(header_value(&request, "X-QWP-Request-Durable-Ack"), "true");
}
let durable_ack_header = if durable_ack {
"X-QWP-Durable-Ack: enabled\r\n"
} else {
""
};
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
{durable_ack_header}\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn connect_blocking_test_transport(
port: u16,
durable_ack: bool,
traffic_gate: Option<Arc<TrafficGate>>,
) -> BlockingQwpWsTransport {
let durable_ack_conf = if durable_ack {
"request_durable_ack=on;durable_ack_keepalive_interval_millis=60000;"
} else {
""
};
let builder = crate::ingress::SenderBuilder::from_conf(format!(
"ws::addr=127.0.0.1:{port};{durable_ack_conf}"
))
.unwrap();
let qwp_ws = builder.qwp_ws.unwrap();
BlockingQwpWsTransport::connect(
"127.0.0.1",
port.to_string(),
false,
None,
QwpWsConnectKind::Foreground,
qwp_ws,
None,
Arc::new(AtomicUsize::new(0)),
traffic_gate,
)
.unwrap()
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn serve_qwp_ws_connection<S: Read + Write>(
stream: &mut S,
frames: usize,
payload_tx: mpsc::Sender<Vec<u8>>,
) {
upgrade_qwp_ws_test_connection(stream, false);
for wire_seq in 0..frames {
let (opcode, payload) = read_client_frame(stream).unwrap();
assert_eq!(opcode, OPCODE_BINARY);
payload_tx.send(payload).unwrap();
write_ok_response(stream, wire_seq as u64).unwrap();
}
let mut sink = [0u8; 1024];
while matches!(stream.read(&mut sink), Ok(n) if n > 0) {}
}
#[cfg(feature = "sync-sender-qwp-ws")]
fn spawn_real_qwp_ws_server(
use_tls: bool,
frames: usize,
) -> (String, u16, mpsc::Receiver<Vec<u8>>) {
let host = if use_tls { "localhost" } else { "127.0.0.1" };
let listener = TcpListener::bind((host, 0)).unwrap();
let port = listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
if use_tls {
let server = ServerConnection::new(tls_server_config()).unwrap();
let mut stream = StreamOwned::new(server, stream);
serve_qwp_ws_connection(&mut stream, frames, payload_tx);
} else {
let mut stream = stream;
serve_qwp_ws_connection(&mut stream, frames, payload_tx);
}
});
(host.to_string(), port, payload_rx)
}
#[test]
fn submit_returns_before_ack() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
drain_events(&mut driver),
vec![DriverEvent::Published { fsn: 0 }]
);
}
#[test]
fn drive_send_once_sends_without_polling_for_ack() {
let transport = TestTransport::scripted([
Ok(TransportSendResult::NoResponse),
Ok(TransportSendResult::NoResponse),
])
.with_poll_results([Err(TransportFailure::Terminal(fake_transport_error(
"should not poll",
)))]);
let mut driver =
QwpWsCoreTestHarness::from_queue(memory_queue(options(8, 1024)), transport);
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert_eq!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame {
fsn: first.fsn,
wire_seq: 0,
payload_len: b"first".len(),
})
);
assert_eq!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame {
fsn: second.fsn,
wire_seq: 1,
payload_len: b"second".len(),
})
);
assert_eq!(driver.drive_send_once().unwrap(), DriveOutcome::Idle);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Sent {
fsn: first.fsn,
wire_seq: 0,
}
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Sent {
fsn: second.fsn,
wire_seq: 1,
}
);
assert_eq!(
driver.send_core.transport.sent_frames.as_slice(),
&[
SentFrame {
fsn: first.fsn,
wire_seq: 0,
payload_len: b"first".len(),
},
SentFrame {
fsn: second.fsn,
wire_seq: 1,
payload_len: b"second".len(),
},
]
);
}
#[test]
fn drive_send_once_applies_immediate_transport_response() {
let transport =
TestTransport::scripted([Ok(TransportSendResult::Response(TransportResponse::Ack {
wire_seq: 0,
}))]);
let mut driver =
QwpWsCoreTestHarness::from_queue(memory_queue(options(8, 1024)), transport);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_send_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: receipt.fsn }
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn sender_qwp_ws_round_trip_delivers_replay_payload() {
let (host, port, payload_rx) = spawn_real_qwp_ws_server(false, 1);
let mut buffer = qwp_buffer("SYM_REAL", 42, 42_000);
let expected = replay_payload(&buffer);
let conf = format!("ws::addr={host}:{port};");
let mut sender = crate::ingress::Sender::from_conf(&conf).unwrap();
sender.flush(&mut buffer).unwrap();
sender.close_drain().unwrap();
assert_eq!(
payload_rx.recv_timeout(Duration::from_secs(5)).unwrap(),
expected
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn sender_qwp_wss_round_trip_delivers_replay_payload() {
let (host, port, payload_rx) = spawn_real_qwp_ws_server(true, 1);
let mut buffer = qwp_buffer("SYM_SECURE", 7, 7_000);
let expected = replay_payload(&buffer);
let ca_path = tls_certs_dir().join("server_rootCA.pem");
let conf = format!("wss::addr={host}:{port};tls_roots={};", ca_path.display());
let mut sender = crate::ingress::Sender::from_conf(&conf).unwrap();
sender.flush(&mut buffer).unwrap();
sender.close_drain().unwrap();
assert_eq!(
payload_rx.recv_timeout(Duration::from_secs(5)).unwrap(),
expected
);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn blocking_transport_emits_durable_keepalive_ping() {
let listener = TcpListener::bind(("127.0.0.1", 0)).unwrap();
let port = listener.local_addr().unwrap().port();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_qwp_ws_test_connection(&mut stream, true);
read_client_frame(&mut stream).unwrap()
});
let mut transport = connect_blocking_test_transport(port, true, None);
assert!(transport.send_durable_ack_keepalive_if_due(true).unwrap());
assert!(!transport.send_durable_ack_keepalive_if_due(true).unwrap());
let (opcode, payload) = server.join().unwrap();
assert_eq!(opcode, crate::ws::frame::OPCODE_PING);
assert!(payload.is_empty());
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn blocking_transport_reconnect_registers_replacement_with_traffic_gate() {
let listener = TcpListener::bind(("127.0.0.1", 0)).unwrap();
let port = listener.local_addr().unwrap().port();
let server = thread::spawn(move || {
let mut streams = Vec::new();
for _ in 0..2 {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_qwp_ws_test_connection(&mut stream, false);
streams.push(stream);
}
streams
});
let traffic_gate = Arc::new(TrafficGate::default());
let mut transport =
connect_blocking_test_transport(port, false, Some(Arc::clone(&traffic_gate)));
transport
.restart_connection(ReconnectReason::Disconnect)
.unwrap();
let _server_streams = server.join().unwrap();
transport
.stream
.set_timeouts(Some(Duration::from_millis(100)), None)
.unwrap();
let reader = thread::spawn(move || {
let deadline = std::time::Instant::now() + Duration::from_secs(10);
let mut byte = [0u8; 1];
loop {
match transport.stream.read(&mut byte) {
Err(err)
if matches!(
err.kind(),
std::io::ErrorKind::TimedOut | std::io::ErrorKind::WouldBlock
) && std::time::Instant::now() < deadline => {}
result => return result,
}
}
});
traffic_gate.shutdown().unwrap();
match reader.join().unwrap() {
Ok(0) => {}
Err(err)
if matches!(
err.kind(),
std::io::ErrorKind::Interrupted
| std::io::ErrorKind::ConnectionAborted
| std::io::ErrorKind::ConnectionReset
| std::io::ErrorKind::NotConnected
| std::io::ErrorKind::BrokenPipe
) => {}
Err(err) if matches!(err.raw_os_error(), Some(10004) | Some(10058)) => {}
result => panic!("replacement socket read was not interrupted: {result:?}"),
}
}
#[test]
fn submit_success_queues_published_event_but_failed_submit_queues_nothing() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
assert_eq!(
driver.try_submit(b""),
Err(DriverError::Queue(QueueError::EmptyPayload))
);
assert_eq!(driver.poll_event(), None);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.poll_event(),
Some(DriverEvent::Published { fsn: receipt.fsn })
);
assert_eq!(driver.poll_event(), None);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
}
#[test]
fn empty_submit_returns_api_error_without_receipt() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
assert_eq!(
driver.try_submit(b""),
Err(DriverError::Queue(QueueError::EmptyPayload))
);
assert_eq!(driver.poll_event(), None);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(receipt, QwpReceipt { fsn: 0 });
}
#[test]
fn blocking_submit_drives_until_local_capacity_frees() {
let mut driver = QwpWsCoreTestHarness::new(
SfaMemoryQueueOptions {
segment_size_bytes: 48,
max_bytes: 96,
},
FakeOrderedServer::ack_each_send(),
)
.unwrap();
let first = driver.try_submit(b"aaaaaaaaaa").unwrap();
driver.try_submit(b"bbbbbbbbbb").unwrap();
let third = driver.submit_with_drive_limit(b"cccccccccc", 16).unwrap();
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(third, QwpReceipt { fsn: 2 });
}
#[test]
fn blocking_submit_times_out_when_capacity_does_not_free() {
let mut driver = QwpWsCoreTestHarness::new(
SfaMemoryQueueOptions {
segment_size_bytes: 48,
max_bytes: 96,
},
FakeOrderedServer::no_response(),
)
.unwrap();
driver.try_submit(b"aaaaaaaaaa").unwrap();
driver.try_submit(b"bbbbbbbbbb").unwrap();
assert!(matches!(
driver.submit_with_drive_limit(b"cccccccccc", 1),
Err(DriverError::SubmitTimedOut {
backpressure: Some(QueueError::StorageSegmentCapFull { .. })
})
));
}
#[test]
fn blocking_submit_deadline_continues_past_fixed_step_budget() {
let queue = memory_queue(SfaMemoryQueueOptions {
segment_size_bytes: 48,
max_bytes: 96,
});
let mut driver = QwpWsCoreTestHarness::from_queue(queue, DelayedPollAckServer::new(20));
let first = driver.try_submit(b"aaaaaaaaaa").unwrap();
driver.try_submit(b"bbbbbbbbbb").unwrap();
let third = driver
.submit_with_drive_deadline(b"cccccccccc", Duration::from_secs(2))
.unwrap();
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(third, QwpReceipt { fsn: 2 });
}
#[test]
fn blocking_submit_deadline_can_expire_before_driving() {
let mut driver = QwpWsCoreTestHarness::new(
SfaMemoryQueueOptions {
segment_size_bytes: 48,
max_bytes: 96,
},
FakeOrderedServer::no_response(),
)
.unwrap();
driver.try_submit(b"aaaaaaaaaa").unwrap();
driver.try_submit(b"bbbbbbbbbb").unwrap();
assert!(matches!(
driver.submit_with_drive_deadline(b"cccccccccc", Duration::ZERO),
Err(DriverError::SubmitTimedOut {
backpressure: Some(QueueError::StorageSegmentCapFull { .. })
})
));
assert!(driver.send_core.transport.sent_frames().is_empty());
}
#[test]
fn blocking_submit_propagates_non_backpressure_errors() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
assert_eq!(
driver.submit_with_drive_limit(b"", 1),
Err(DriverError::Queue(QueueError::EmptyPayload))
);
}
fn sent(fsn: u64, wire_seq: u64) -> SentFrame {
SentFrame {
fsn,
wire_seq,
payload_len: 1,
}
}
#[test]
fn in_flight_run_tracks_a_contiguous_window_without_per_frame_state() {
let mut run = InFlightRun::default();
assert_eq!(run.len(), 0);
assert_eq!(run.wire_seq_for_fsn(0), None);
for (fsn, wire_seq) in [(10, 4), (11, 5), (12, 6)] {
run.push(&sent(fsn, wire_seq));
}
assert_eq!(run.len(), 3);
assert_eq!(run.wire_seq_for_fsn(10), Some(4));
assert_eq!(run.wire_seq_for_fsn(12), Some(6));
assert_eq!(run.wire_seq_for_fsn(9), None, "below the run");
assert_eq!(run.wire_seq_for_fsn(13), None, "past the run");
run.ack_through(9);
assert_eq!(run.len(), 3);
assert_eq!(run.wire_seq_for_fsn(10), Some(4));
run.ack_through(11);
assert_eq!(run.len(), 1);
assert_eq!(run.wire_seq_for_fsn(11), None, "acked frames leave the run");
assert_eq!(run.wire_seq_for_fsn(12), Some(6));
run.ack_through(99);
assert_eq!(run.len(), 0);
assert_eq!(run.wire_seq_for_fsn(12), None);
run.push(&sent(40, 7));
assert_eq!(run.len(), 1);
assert_eq!(run.wire_seq_for_fsn(40), Some(7));
assert_eq!(run.wire_seq_for_fsn(13), None);
run.clear();
assert_eq!(run.len(), 0);
assert_eq!(run.wire_seq_for_fsn(40), None);
}
#[cfg(feature = "sync-sender-qwp-ws")]
#[test]
fn pending_wire_sequence_run_tracks_cumulative_acks() {
let mut run = PendingWireSequenceRun::default();
assert!(run.is_empty());
for wire_seq in 4..=6 {
run.push(wire_seq);
}
assert_eq!((run.front, run.len), (4, 3));
run.complete_through(3);
assert_eq!((run.front, run.len), (4, 3));
run.complete_through(5);
assert_eq!((run.front, run.len), (6, 1));
run.complete_through(99);
assert!(run.is_empty());
run.push(u64::MAX);
run.complete_through(u64::MAX);
assert!(run.is_empty());
run.push(0);
run.clear();
assert!(run.is_empty());
}
#[test]
fn durable_pending_lookup_handles_contiguous_and_cumulative_gaps() {
let mut tracker = DurableAckTracker::new();
tracker.enqueue_ok(40, 100, Vec::new());
tracker.enqueue_ok(41, 101, Vec::new());
tracker.enqueue_ok(43, 103, Vec::new());
assert_eq!(tracker.pending_wire_seq_for_fsn(100), Some(40));
assert_eq!(tracker.pending_wire_seq_for_fsn(101), Some(41));
assert_eq!(tracker.pending_wire_seq_for_fsn(102), None);
assert_eq!(tracker.pending_wire_seq_for_fsn(103), Some(43));
assert_eq!(tracker.pending_wire_seq_for_fsn(99), None);
assert_eq!(tracker.pending_wire_seq_for_fsn(104), None);
}
#[test]
fn drive_once_streams_without_in_flight_cap() {
let mut driver =
QwpWsCoreTestHarness::new(options(8, 4096), FakeOrderedServer::no_response()).unwrap();
driver.try_submit(b"a").unwrap();
driver.try_submit(b"b").unwrap();
driver.try_submit(b"c").unwrap();
for _ in 0..3 {
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
}
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Idle);
assert_eq!(driver.send_core.transport.sent_frames().len(), 3);
}
#[test]
fn send_event_agrees_with_sent_receipt_status() {
let mut driver = driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
]
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0,
}
);
}
#[test]
fn reconnect_policy_retries_failed_reconnect_until_success() {
let transport = TestTransport::scripted([Ok(TransportSendResult::Failure(
TransportFailure::Disconnect(fake_transport_error("disconnect before reconnect")),
))])
.with_restart_results([
Err(DriverError::Transport(fake_transport_error(
"reconnect failed once",
))),
Ok(()),
]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
ReconnectPolicy::no_backoff(Duration::from_secs(1)),
false,
);
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::ReconnectDelay {
sleep_for: Duration::ZERO,
..
}
));
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect,
}
);
assert_eq!(driver.send_core.transport.restart_attempts, 2);
assert_eq!(
driver.send_core.transport.sent_payloads,
vec![b"payload".to_vec()]
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Reconnected {
reason: ReconnectReason::Disconnect,
},
]
);
}
fn paced_policy() -> ReconnectPolicy {
ReconnectPolicy::bounded(
Duration::from_secs(10),
Duration::from_millis(100),
Duration::from_secs(1),
)
}
fn assert_pace_range(actual: Duration, min: Duration, max: Duration) {
assert!(
actual >= min && actual < max,
"pace {actual:?} outside [{min:?}, {max:?})"
);
}
#[test]
fn retriable_nack_below_threshold_paces_before_first_reconnect_attempt() {
let transport = TestTransport::scripted([Ok(TransportSendResult::Response(
TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed"),
},
))])
.with_restart_results([Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
paced_policy(),
false,
);
driver.try_submit(b"payload").unwrap();
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => panic!("expected paced reconnect delay, got {other:?}"),
}
assert_eq!(driver.send_core.transport.restart_attempts, 0);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(driver.send_core.transport.restart_attempts, 1);
}
#[test]
fn second_same_frame_strike_doubles_pace_dose() {
let transport = TestTransport::scripted([
Ok(TransportSendResult::Response(TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed once"),
})),
Ok(TransportSendResult::Response(TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed twice"),
})),
])
.with_restart_results([Ok(()), Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
paced_policy(),
false,
);
driver.try_submit(b"payload").unwrap();
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => panic!("expected first paced reconnect delay, got {other:?}"),
}
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(200),
Duration::from_millis(400),
);
}
other => panic!("expected second paced reconnect delay, got {other:?}"),
}
assert_eq!(driver.send_core.transport.restart_attempts, 1);
}
#[test]
fn paced_reject_stops_receive_drain_before_buffered_ack() {
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
FakeOrderedServer::no_response(),
paced_policy(),
false,
);
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed"),
});
driver
.send_core
.transport
.push_response(TransportResponse::Ack { wire_seq: 1 });
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => panic!("expected paced reconnect delay, got {other:?}"),
}
assert_eq!(driver.send_core.transport.poll_responses.len(), 1);
assert_eq!(driver.store.completed_fsn(), None);
driver.send_core.transport.poll_responses.clear();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
for expected_fsn in [0, 1] {
match driver.drive_once().unwrap() {
DriveOutcome::Sent(frame) => assert_eq!(frame.fsn, expected_fsn),
other => panic!("expected replay of fsn {expected_fsn}, got {other:?}"),
}
}
driver
.send_core
.transport
.push_response(TransportResponse::Ack { wire_seq: 1 });
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn not_writable_first_recycle_immediate_then_zero_progress_recycles_escalate() {
let not_writable_reject = || {
Ok(TransportSendResult::Response(TransportResponse::Reject {
wire_seq: 0,
error: not_writable_error("replica access is read-only"),
}))
};
let transport = TestTransport::scripted([
not_writable_reject(),
not_writable_reject(),
not_writable_reject(),
])
.with_restart_results([Ok(()), Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
paced_policy(),
false,
);
driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::NotWritable
}
);
assert_eq!(driver.send_core.transport.restart_attempts, 1);
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => panic!("expected first zero-progress pace, got {other:?}"),
}
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::NotWritable
}
);
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(200),
Duration::from_millis(400),
);
}
other => panic!("expected doubled zero-progress pace, got {other:?}"),
}
assert_eq!(driver.send_core.transport.restart_attempts, 2);
}
#[test]
fn ack_progress_resets_not_writable_recycle_pace() {
let transport = TestTransport::scripted([
Ok(TransportSendResult::Response(TransportResponse::Reject {
wire_seq: 0,
error: not_writable_error("replica access is read-only"),
})),
Ok(TransportSendResult::Response(TransportResponse::Ack {
wire_seq: 0,
})),
Ok(TransportSendResult::Response(TransportResponse::Reject {
wire_seq: 1,
error: not_writable_error("replica access is read-only"),
})),
Ok(TransportSendResult::Response(TransportResponse::Reject {
wire_seq: 0,
error: not_writable_error("replica access is read-only"),
})),
])
.with_restart_results([Ok(()), Ok(()), Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
paced_policy(),
false,
);
driver.try_submit(b"first").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::NotWritable
}
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(driver.acked_fsn(), Some(0));
driver.try_submit(b"second").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::NotWritable
},
"the first recycle after ACK progress must be immediate again"
);
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => {
panic!("expected zero-progress pace to restart at the initial dose, got {other:?}")
}
}
}
#[test]
fn presend_reject_paces_with_min_one_strike() {
let transport = TestTransport::scripted([]).with_restart_results([Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
paced_policy(),
false,
);
driver.try_submit(b"payload").unwrap();
match driver
.send_core
.finish_response(
&mut driver.store,
TransportResponse::Reject {
wire_seq: 0,
error: write_error("pre-send reject"),
},
)
.unwrap()
{
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => panic!("expected paced reconnect delay, got {other:?}"),
}
assert_eq!(driver.send_core.transport.restart_attempts, 0);
}
#[test]
fn zero_initial_backoff_keeps_legacy_immediate_reconnect() {
let transport = TestTransport::scripted([Ok(TransportSendResult::Response(
TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed"),
},
))])
.with_restart_results([Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
ReconnectPolicy::no_backoff(Duration::MAX),
false,
);
driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(driver.send_core.transport.restart_attempts, 1);
}
#[test]
fn server_close_below_threshold_paces_reconnect() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_events([Err(TransportFailure::ServerClose(fake_transport_error(
"server close",
)))])
.with_restart_results([Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
paced_policy(),
false,
);
driver.try_submit(b"payload").unwrap();
match driver.drive_once().unwrap() {
DriveOutcome::ReconnectDelay { sleep_for, .. } => {
assert_pace_range(
sleep_for,
Duration::from_millis(100),
Duration::from_millis(200),
);
}
other => panic!("expected paced reconnect delay, got {other:?}"),
}
assert_eq!(driver.send_core.transport.restart_attempts, 0);
}
#[test]
fn role_reject_reconnect_sleep_uses_fixed_initial_backoff() {
let initial = Duration::from_millis(100);
let current = Duration::from_millis(400);
assert_eq!(reconnect_sleep_duration(true, initial, current), initial);
assert_eq!(
reconnect_sleep_duration(false, initial, Duration::ZERO),
Duration::ZERO
);
}
#[test]
fn centered_jitter_reconnect_sleep_scatters_around_backoff() {
let unused_initial = Duration::from_millis(100);
for base_ms in [2u64, 80, 100, 1_000, 5_000] {
let base = Duration::from_millis(base_ms);
for _ in 0..10_000 {
let d = reconnect_sleep_duration(false, unused_initial, base);
assert!(
d >= base / 2 && d < base + base / 2,
"centered-jitter sleep {d:?} outside [base/2, 3*base/2) for base={base:?}"
);
}
}
}
#[test]
fn pace_dose_doubles_per_strike_and_caps_at_max_backoff() {
let initial = Duration::from_millis(100);
let max = Duration::from_millis(1000);
assert_eq!(pace_dose(initial, max, 0), Duration::from_millis(100));
assert_eq!(pace_dose(initial, max, 1), Duration::from_millis(100));
assert_eq!(pace_dose(initial, max, 2), Duration::from_millis(200));
assert_eq!(pace_dose(initial, max, 3), Duration::from_millis(400));
assert_eq!(pace_dose(initial, max, 4), Duration::from_millis(800));
assert_eq!(pace_dose(initial, max, 5), Duration::from_millis(1000));
assert_eq!(pace_dose(initial, max, 99), Duration::from_millis(1000));
assert_eq!(
pace_dose(Duration::ZERO, max, 4),
Duration::ZERO,
"zero initial backoff preserves legacy immediate reconnects"
);
}
#[test]
fn pace_jitter_stays_within_one_to_two_dose() {
for dose_ms in [1u64, 80, 100, 1_000] {
let dose = Duration::from_millis(dose_ms);
for _ in 0..10_000 {
let d = pace_jitter_duration(dose);
assert!(
d >= dose && d < dose + dose,
"paced jitter {d:?} outside [dose, 2*dose) for dose={dose:?}"
);
}
}
}
#[test]
fn poison_tracker_dwell_holds_escalation_until_window_elapses() {
let mut tracker = PoisonFrameTracker::default();
let started = Instant::now();
let window = Duration::from_secs(5);
assert!(!tracker.record_failure(7, None, 2, window, started));
assert!(!tracker.record_failure(7, None, 2, window, started + Duration::from_secs(1)));
assert_eq!(tracker.strikes(), 2);
assert!(tracker.record_failure(7, None, 2, window, started + window));
assert_eq!(tracker.strikes(), 3);
}
#[test]
fn poison_tracker_dwell_restamps_when_suspect_key_changes() {
let mut tracker = PoisonFrameTracker::default();
let started = Instant::now();
let window = Duration::from_secs(5);
assert!(!tracker.record_failure(7, None, 2, window, started));
assert!(!tracker.record_failure(8, None, 2, window, started + Duration::from_secs(10)));
assert_eq!(tracker.strikes(), 1);
assert!(!tracker.record_failure(8, None, 2, window, started + Duration::from_secs(11)));
assert!(tracker.record_failure(8, None, 2, window, started + Duration::from_secs(15)));
}
#[test]
fn poison_tracker_window_zero_escalates_exactly_at_threshold() {
let mut tracker = PoisonFrameTracker::default();
let started = Instant::now();
assert!(!tracker.record_failure(7, None, 2, Duration::ZERO, started));
assert!(tracker.record_failure(7, None, 2, Duration::ZERO, started));
}
#[test]
fn poison_tracker_clear_resets_first_strike_timestamp() {
let mut tracker = PoisonFrameTracker::default();
let started = Instant::now();
let window = Duration::from_secs(5);
assert!(!tracker.record_failure(7, None, 2, window, started));
tracker.clear();
assert_eq!(tracker.strikes(), 0);
assert!(!tracker.record_failure(7, None, 1, window, started + Duration::from_secs(10)));
}
#[test]
fn reconnect_terminal_classification_keeps_protocol_version_errors_terminal() {
let durable_ack_mismatch = Error::new(
ErrorCode::ProtocolVersionError,
"server did not enable durable ACK",
);
assert!(reconnect_error_is_terminal(&durable_ack_mismatch));
let retryable_upgrade_version_error =
Error::new(ErrorCode::SocketError, "unsupported X-QWP-Version");
assert!(!reconnect_error_is_terminal(
&retryable_upgrade_version_error
));
}
#[test]
fn reconnect_terminal_classification_stops_on_a_full_symbol_dictionary() {
let dict_full = Error::new(
ErrorCode::SymbolDictFull,
"QWP/WS connection-scoped symbol dictionary reached its 2000000-entry cap",
);
assert!(reconnect_error_is_terminal(&dict_full));
let torn_dict = Error::new(
ErrorCode::StoreResendRequired,
"corrupt persisted symbol dictionary: duplicate entry at index 1",
);
assert!(reconnect_error_is_terminal(&torn_dict));
}
#[test]
fn reconnect_terminal_classification_retries_structured_role_rejects() {
let all_role_rejected = Error::new(
ErrorCode::ProtocolVersionError,
"server did not enable durable ACK; all endpoints rejected by role",
)
.with_qwp_ws_role_reject(crate::ingress::QwpWsRoleReject::new("REPLICA", None));
assert!(!reconnect_error_is_terminal(&all_role_rejected));
}
#[test]
fn transport_write_failure_does_not_commit_sent_receipt() {
let transport = TestTransport::scripted([Err(TransportFailure::Disconnect(
fake_transport_error("write failed"),
))]);
let mut driver =
QwpWsCoreTestHarness::from_queue(memory_queue(options(8, 1024)), transport);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect,
}
);
assert_eq!(
driver.send_core.transport.sent_payloads,
vec![b"payload".to_vec()]
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Reconnected {
reason: ReconnectReason::Disconnect,
},
]
);
}
#[test]
fn transport_poll_failure_enters_reconnect_policy() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_results([Err(TransportFailure::Disconnect(fake_transport_error(
"poll failed",
)))]);
let mut driver =
QwpWsCoreTestHarness::from_queue(memory_queue(options(8, 1024)), transport);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect,
}
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Reconnected {
reason: ReconnectReason::Disconnect,
},
]
);
}
#[test]
fn raw_qwp_schema_and_write_errors_decode_as_rejections() {
for (status, expected_code) in [
(codec::WS_STATUS_SCHEMA_MISMATCH, ErrorCode::InvalidApiCall),
(codec::WS_STATUS_WRITE_ERROR, ErrorCode::ServerFlushError),
] {
let payload = qwp_error_payload(status, 42, "server says no");
let response = decode_transport_response(&payload).unwrap().unwrap();
match response {
TransportResponse::Reject { wire_seq, error } => {
assert_eq!(wire_seq, 42);
assert_eq!(error.status, status);
assert_eq!(error.message, "server says no");
assert_eq!(error.error.code(), expected_code);
assert!(error.error.msg().contains("server says no"));
}
other => panic!("unexpected response: {other:?}"),
}
}
}
#[test]
#[ignore = "uses the process-global allocation counter"]
fn non_durable_ok_decode_with_table_entries_zero_alloc_after_warmup() {
use crate::alloc_counter;
let payload = qwp_ok_payload_with_table_entries(7, &[("table_a", 42), ("table_b", -7)]);
assert!(matches!(
decode_transport_response(&payload).unwrap(),
Some(TransportResponse::Ack { wire_seq: 7 })
));
alloc_counter::start_counting();
let response = decode_transport_response(&payload).unwrap();
let alloc_count = alloc_counter::stop_counting();
assert!(matches!(
response,
Some(TransportResponse::Ack { wire_seq: 7 })
));
assert_eq!(
alloc_count, 0,
"Expected zero allocations for non-durable OK decode, got {alloc_count}"
);
}
#[test]
fn durable_decode_preserves_ok_and_durable_ack_table_entries() {
let ok_payload = qwp_ok_payload_with_table_entries(7, &[("table_a", 42)]);
assert_eq!(
decode_durable_transport_response(&ok_payload).unwrap(),
Some(TransportResponse::DurableOk {
wire_seq: 7,
table_seq_txns: table_seq_txns(&[("table_a", 42)])
})
);
let durable_ack_payload = qwp_durable_ack_payload(&[("table_a", 42), ("table_b", 99)]);
assert_eq!(
decode_durable_transport_response(&durable_ack_payload).unwrap(),
Some(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("table_a", 42), ("table_b", 99)])
})
);
}
#[test]
fn malformed_qwp_response_frames_are_retryable_failures() {
let failure = decode_transport_response(&[codec::WS_STATUS_OK]).unwrap_err();
assert_retryable_transport_failure(failure, "QWP OK response truncated");
let failure =
decode_durable_transport_response(&[codec::WS_STATUS_DURABLE_ACK]).unwrap_err();
assert_retryable_transport_failure(failure, "QWP durable ACK response truncated");
}
fn assert_retryable_transport_failure(failure: TransportFailure, expected_message: &str) {
match failure {
TransportFailure::Retryable(error) => assert!(
error.msg().contains(expected_message),
"expected error message containing {expected_message:?}, got {:?}",
error.msg()
),
other => panic!("expected retryable transport failure, got {other:?}"),
}
}
#[test]
fn durable_pending_ok_triggers_ready_keepalive() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_results([
Ok(Some(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
})),
Ok(None),
])
.with_keepalive_results([Ok(false), Ok(true)]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_durable_ack(
memory_queue(options(8, 1024)),
transport,
);
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Idle);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0,
}
);
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Idle);
assert_eq!(driver.send_core.transport.keepalive_attempts, 2);
assert_eq!(
driver.send_core.transport.keepalive_pending_args,
vec![true, true]
);
}
#[test]
fn durable_keepalive_disconnect_reports_reconnect() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_results([
Ok(Some(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
})),
Ok(None),
])
.with_keepalive_results([Err(TransportFailure::Disconnect(fake_transport_error(
"keepalive disconnect",
)))])
.with_restart_results([Ok(())]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_durable_ack(
memory_queue(options(8, 1024)),
transport,
);
driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
assert_eq!(driver.send_core.transport.keepalive_attempts, 1);
assert_eq!(driver.send_core.transport.restart_attempts, 1);
}
#[test]
fn durable_keepalive_terminal_failure_reports_terminal() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_results([
Ok(Some(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
})),
Ok(None),
])
.with_keepalive_results([Err(TransportFailure::Terminal(fake_transport_error(
"terminal keepalive failure",
)))]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_durable_ack(
memory_queue(options(8, 1024)),
transport,
);
driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert!(driver.is_terminal());
assert_eq!(driver.send_core.transport.keepalive_attempts, 1);
}
#[test]
fn consumed_durable_ok_is_progress_and_does_not_send_keepalive() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_results([
Ok(Some(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
})),
Ok(None),
])
.with_keepalive_results([Ok(true)]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_durable_ack(
memory_queue(options(8, 1024)),
transport,
);
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0,
}
);
assert_eq!(driver.send_core.transport.keepalive_attempts, 0);
}
#[test]
fn drive_ready_drains_control_progress_and_durable_ack_until_idle() {
let transport = TestTransport::scripted([Ok(TransportSendResult::NoResponse)])
.with_poll_events([
Ok(TransportPoll::Response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
})),
Ok(TransportPoll::Progress),
Ok(TransportPoll::Response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
})),
Ok(TransportPoll::Idle),
])
.with_keepalive_results([Ok(true)]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_durable_ack(
memory_queue(options(8, 1024)),
transport,
);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert!(matches!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
));
assert_eq!(driver.send_core.transport.keepalive_attempts, 0);
}
#[test]
fn durable_ok_releases_send_window_without_completing_receipt() {
let mut server = FakeOrderedServer::no_response();
server.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
let mut driver = durable_driver_with_options(options(4, 1024), server);
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(driver.acked_fsn(), None);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0
}
);
assert_eq!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame {
fsn: 1,
wire_seq: 1,
payload_len: 6,
})
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Sent {
fsn: 1,
wire_seq: 1
}
);
}
#[test]
fn durable_ack_covering_pending_ok_advances_acked_fsn() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(_)
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(driver.acked_fsn(), Some(0));
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
}
#[test]
fn future_durable_ok_wire_sequence_clamps_to_highest_sent_like_java() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { wire_seq: 0, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 99,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0
}
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
}
#[test]
fn future_durable_reject_wire_sequence_reconnects_from_highest_sent_like_java_v2() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 99,
error: write_error("write failed"),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.message_sequence, Some(99));
assert_eq!(error.from_fsn, 1);
assert_eq!(error.to_fsn, 1);
}
#[test]
fn late_future_durable_reject_for_acked_frame_reports_error_without_reblocking_tracker() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(driver.poll_sender_error(), None);
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 99,
error: write_error("late write failure"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Retriable);
assert_eq!(error.message_sequence, Some(99));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
}
#[test]
fn future_durable_reject_for_pending_ok_reports_error_without_duplicate_tracker_entry() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 99,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0
}
);
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 99,
error: write_error("late write failure"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0
}
);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Retriable);
assert_eq!(error.message_sequence, Some(99));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(driver.poll_sender_error(), None);
}
#[test]
fn durable_ack_does_not_skip_earlier_pending_ok_gap() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { wire_seq: 0, .. })
));
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { wire_seq: 1, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("a", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("b", 20)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("b", 20)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(driver.acked_fsn(), None);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("a", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn durable_empty_ok_waits_behind_prior_non_empty_ok() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: Vec::new(),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(driver.acked_fsn(), None);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn stale_durable_ack_watermark_does_not_move_tracker_backwards() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("trades", 12)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 9)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Sent {
fsn: 1,
wire_seq: 1
}
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 12)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn durable_ack_before_ok_drains_when_ok_later_arrives() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
}
#[test]
fn stale_durable_ok_after_completion_does_not_reblock_tracker() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 99)]),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("trades", 12)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 12)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn durable_reconnect_clears_pending_ok_tracking_for_replay() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame {
fsn: 0,
wire_seq: 0,
payload_len: 7,
})
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0
}
);
assert_eq!(
driver
.send_core
.finish_reconnect_success(&mut driver.store, ReconnectReason::Disconnect),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame {
fsn: 0,
wire_seq: 0,
payload_len: 7,
})
);
}
#[test]
fn durable_reconnect_clears_unresolved_reject_so_replayed_ok_can_complete() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 1,
error: write_error("write failed"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
assert_eq!(
driver
.send_core
.finish_reconnect_success(&mut driver.store, ReconnectReason::Disconnect),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 1, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("quotes", 20)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("quotes", 20)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
assert_eq!(driver.acked_fsn(), Some(1));
}
#[test]
fn durable_reconnect_clears_unresolved_reject_so_replayed_reject_retries_again() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 1,
error: write_error("write failed"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver
.send_core
.finish_reconnect_success(&mut driver.store, ReconnectReason::Disconnect),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 1, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 1,
error: write_error("write failed again"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
assert_eq!(driver.acked_fsn(), Some(0));
}
#[test]
fn durable_retryable_reject_preserves_prior_pending_ok_for_replay() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
let third = driver.try_submit(b"third").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 1,
error: write_error("write failed"),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
assert!(driver.receipt_status(third).is_pending());
assert_eq!(driver.acked_fsn(), None);
assert_eq!(
driver
.last_server_error()
.map(|err| (err.status, err.message.as_str())),
Some((codec::WS_STATUS_WRITE_ERROR, "write failed"))
);
assert_eq!(
driver.poll_sender_error().map(|err| err.applied_policy),
Some(QwpWsErrorPolicy::Retriable)
);
}
#[test]
fn durable_reject_after_cumulative_pending_ok_reconnects_without_gap_violation() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
let third = driver.try_submit(b"third").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 2,
error: write_error("write failed"),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert!(!driver.is_terminal());
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
assert!(driver.receipt_status(third).is_pending());
assert_eq!(
driver.poll_sender_error().map(|err| err.applied_policy),
Some(QwpWsErrorPolicy::Retriable)
);
}
#[test]
fn durable_consecutive_retryable_rejects_preserve_replay_tail() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
let third = driver.try_submit(b"third").unwrap();
let fourth = driver.try_submit(b"fourth").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 1,
error: write_error("write failed"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert!(matches!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
));
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
assert!(driver.receipt_status(third).is_pending());
assert!(driver.receipt_status(fourth).is_pending());
assert_eq!(driver.acked_fsn(), None);
}
#[test]
fn stale_durable_ok_after_retryable_reject_does_not_complete_before_replay() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed"),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
let first_error = driver.poll_sender_error().unwrap();
assert_eq!(first_error.message_sequence, Some(0));
assert_eq!(first_error.message.as_deref(), Some("write failed"));
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Published { fsn: 1 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Sent {
fsn: 1,
wire_seq: 1,
},
DriverEvent::Rejected {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Reconnected {
reason: ReconnectReason::RetryableFailure,
},
]
);
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 99)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 1, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("quotes", 20)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("quotes", 20)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn stale_durable_reject_before_replay_reconnects_without_completion() {
let mut driver = durable_driver(FakeOrderedServer::no_response());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_send_once().unwrap();
driver.drive_send_once().unwrap();
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed"),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
let first_error = driver.poll_sender_error().unwrap();
assert_eq!(first_error.message_sequence, Some(0));
assert_eq!(first_error.message.as_deref(), Some("write failed"));
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 0,
error: write_error("write failed again"),
});
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
let stale_error = driver.poll_sender_error().unwrap();
assert_eq!(stale_error.message_sequence, Some(0));
assert_eq!(stale_error.message.as_deref(), Some("write failed again"));
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 1, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 0,
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("trades", 10)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableOk {
wire_seq: 1,
table_seq_txns: table_seq_txns(&[("quotes", 20)]),
});
driver
.send_core
.transport
.push_response(TransportResponse::DurableAck {
table_seq_txns: table_seq_txns(&[("quotes", 20)]),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn server_error_policy_matches_java_v2_defaults() {
for (status, category, policy) in [
(
codec::WS_STATUS_SCHEMA_MISMATCH,
QwpWsErrorCategory::SchemaMismatch,
QwpWsErrorPolicy::Terminal,
),
(
codec::WS_STATUS_PARSE_ERROR,
QwpWsErrorCategory::ParseError,
QwpWsErrorPolicy::Terminal,
),
(
codec::WS_STATUS_SECURITY_ERROR,
QwpWsErrorCategory::SecurityError,
QwpWsErrorPolicy::Terminal,
),
(
codec::WS_STATUS_WRITE_ERROR,
QwpWsErrorCategory::WriteError,
QwpWsErrorPolicy::Retriable,
),
(
codec::WS_STATUS_INTERNAL_ERROR,
QwpWsErrorCategory::InternalError,
QwpWsErrorPolicy::Retriable,
),
(
codec::WS_STATUS_NOT_WRITABLE,
QwpWsErrorCategory::NotWritable,
QwpWsErrorPolicy::RetriableOther,
),
(
0x7f,
QwpWsErrorCategory::Unknown,
QwpWsErrorPolicy::Retriable,
),
] {
assert_eq!(server_error_category(status), category);
assert_eq!(server_error_policy(status), policy);
}
}
#[test]
fn not_writable_reject_uses_role_reconnect_reason_and_preserves_receipt() {
let transport = TestTransport::scripted([Ok(TransportSendResult::Response(
TransportResponse::Reject {
wire_seq: 0,
error: not_writable_error("replica access is read-only"),
},
))]);
let mut driver =
QwpWsCoreTestHarness::from_queue(memory_queue(options(8, 1024)), transport);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once(),
Ok(DriveOutcome::Reconnected {
reason: ReconnectReason::NotWritable
})
);
assert_eq!(driver.send_core.transport.restart_attempts, 1);
assert_eq!(
driver.send_core.transport.restart_reasons,
[ReconnectReason::NotWritable]
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), None);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::NotWritable);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::RetriableOther);
assert_eq!(error.status, Some(codec::WS_STATUS_NOT_WRITABLE));
assert_eq!(
error.message.as_deref(),
Some("replica access is read-only")
);
}
#[test]
fn raw_qwp_error_statuses_decode_with_expected_error_codes() {
for (status, expected_code) in [
(codec::WS_STATUS_PARSE_ERROR, ErrorCode::InvalidApiCall),
(codec::WS_STATUS_INTERNAL_ERROR, ErrorCode::ServerFlushError),
(codec::WS_STATUS_SECURITY_ERROR, ErrorCode::AuthError),
(codec::WS_STATUS_NOT_WRITABLE, ErrorCode::ServerFlushError),
(0x7f, ErrorCode::ServerFlushError),
] {
let payload = qwp_error_payload(status, 7, "fatal server error");
let response = decode_transport_response(&payload).unwrap().unwrap();
match response {
TransportResponse::Reject { wire_seq, error } => {
assert_eq!(wire_seq, 7);
assert_eq!(error.status, status);
assert_eq!(error.error.code(), expected_code);
assert!(error.error.msg().contains("fatal server error"));
}
other => panic!("unexpected response: {other:?}"),
}
}
}
#[test]
fn future_ack_wire_sequence_clamps_to_highest_sent_like_java() {
let mut driver = driver(FakeOrderedServer::scripted([FakeSendResult::AckWire {
wire_seq: 99,
}]));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.drive_once(), Ok(DriveOutcome::Acked { wire_seq: 0 }));
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
}
#[test]
fn future_reject_wire_sequence_reconnects_from_highest_sent_like_java_v2() {
let mut driver = driver(FakeOrderedServer::scripted([FakeSendResult::RejectWire {
wire_seq: 99,
}]));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once(),
Ok(DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
})
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), None);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Retriable);
assert_eq!(error.message_sequence, Some(99));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
}
#[test]
fn late_future_reject_for_acked_frame_reports_error_without_changing_receipt() {
let mut driver = driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(SentFrame { fsn: 0, .. })
));
driver
.send_core
.transport
.push_response(TransportResponse::Ack { wire_seq: 0 });
assert_eq!(
driver.drive_receive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(driver.poll_sender_error(), None);
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 99,
error: write_error("late write failure"),
});
assert_eq!(driver.drive_receive_once().unwrap(), DriveOutcome::Progress);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), Some(0));
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Retriable);
assert_eq!(error.message_sequence, Some(99));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
}
#[test]
fn wait_drives_until_receipt_acked() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_once().unwrap();
assert_eq!(
driver.delivery_status(receipt).unwrap(),
Some(DeliveryOutcome::Completed)
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
}
#[test]
fn coalesced_ack_completes_multiple_receipts() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::NoResponse,
FakeSendResult::AckWire { wire_seq: 1 },
]));
let first = driver.try_submit(b"a").unwrap();
let second = driver.try_submit(b"b").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn cumulative_ack_event_agrees_with_wait_and_receipt_status() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::NoResponse,
FakeSendResult::AckWire { wire_seq: 1 },
]));
let first = driver.try_submit(b"a").unwrap();
let second = driver.try_submit(b"b").unwrap();
driver.drive_once().unwrap();
driver.drive_once().unwrap();
assert_eq!(
driver.delivery_status(second).unwrap(),
Some(DeliveryOutcome::Completed)
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Published { fsn: 1 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Sent {
fsn: 1,
wire_seq: 1,
},
DriverEvent::CompletedThrough {
fsn: 1,
wire_seq: 1,
},
]
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn stale_completion_response_does_not_emit_duplicate_progress_event() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::NoResponse,
FakeSendResult::AckWire { wire_seq: 1 },
]));
let first = driver.try_submit(b"a").unwrap();
let second = driver.try_submit(b"b").unwrap();
driver.drive_once().unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 1 }
);
driver
.send_core
.transport
.push_response(FakeServerResponse::Ack { wire_seq: 0 });
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Published { fsn: 1 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Sent {
fsn: 1,
wire_seq: 1,
},
DriverEvent::CompletedThrough {
fsn: 1,
wire_seq: 1,
},
]
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
#[test]
fn disconnect_replays_from_oldest_unresolved_with_zero_based_wire_sequence() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::NoResponse,
FakeSendResult::Disconnect,
FakeSendResult::AckSent,
]));
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(SentFrame {
fsn: 0,
wire_seq: 0,
payload_len: 5
})
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::Disconnect
}
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.send_core.transport.sent_frames(),
&[
SentFrame {
fsn: 0,
wire_seq: 0,
payload_len: 5,
},
SentFrame {
fsn: 1,
wire_seq: 1,
payload_len: 6,
},
SentFrame {
fsn: 0,
wire_seq: 0,
payload_len: 5,
},
]
);
}
#[test]
fn retryable_failure_does_not_complete_receipt() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::RetryableFailure,
FakeSendResult::AckSent,
]));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
}
#[test]
fn reconnect_event_agrees_with_republished_receipt_status() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::RetryableFailure,
]));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure,
}
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Reconnected {
reason: ReconnectReason::RetryableFailure,
},
]
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
}
#[test]
fn reconnect_policy_exhaustion_keeps_sf_receipts_replayable() {
let transport = TestTransport::scripted([Ok(TransportSendResult::Failure(
TransportFailure::Retryable(fake_transport_error("retryable outage")),
))])
.with_restart_results([Err(DriverError::Transport(fake_transport_error(
"reconnect failed once",
)))]);
let mut driver = QwpWsCoreTestHarness::from_queue_with_reconnect_policy(
memory_queue(options(8, 1024)),
transport,
ReconnectPolicy::bounded(
Duration::from_millis(1),
Duration::from_millis(10),
Duration::from_millis(10),
),
false,
);
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::ReconnectDelay { .. }
));
std::thread::sleep(Duration::from_millis(20));
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::ReconnectDelay { .. }
));
assert_eq!(driver.send_core.transport.restart_attempts, 1);
assert_eq!(driver.terminal_error(), None);
assert!(driver.receipt_status(first).is_pending());
assert!(driver.receipt_status(second).is_pending());
assert!(driver.try_submit(b"third").is_ok());
}
#[test]
fn terminal_event_is_emitted_once_on_transition() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::TerminalFailure,
]));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Terminal,
]
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
assert_eq!(
driver.delivery_status(receipt).unwrap(),
Some(DeliveryOutcome::Terminal)
);
}
#[test]
fn lifecycle_terminalizes_when_store_terminalizes() {
let queue = memory_queue(options(4, 1024));
let mut store = QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY);
let lifecycle = store.lifecycle();
store.mark_terminal(Some(Error::new(ErrorCode::SocketError, "terminal")));
assert_eq!(lifecycle.load(), PublicationState::Terminal);
assert!(store.is_terminal());
}
#[test]
fn first_terminal_error_wins_over_late_structured_diagnostic() {
let queue = memory_queue(options(4, 1024));
let mut store = QwpWsPublicationStore::new(queue, DEFAULT_EVENT_CAPACITY);
store.mark_terminal(Some(Error::new(ErrorCode::SocketError, "first terminal")));
let returned = store.record_protocol_violation(None, "late violation".to_string());
assert_eq!(returned.msg(), "first terminal");
assert_eq!(store.terminal_error().unwrap().msg(), "first terminal");
assert_eq!(store.terminal_sender_error(), None);
assert_eq!(store.poll_sender_error(), None);
}
#[test]
fn close_drain_drives_until_all_published_receipts_resolve() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert_eq!(driver.close_drain_steps(2).unwrap(), CloseOutcome::Drained);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
assert_eq!(driver.try_submit(b"third"), Err(DriverError::Closing));
}
#[test]
fn close_drain_timeout_keeps_existing_receipt_observable() {
let mut driver = driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.close_drain_steps(1).unwrap(), CloseOutcome::Timeout);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0,
}
);
assert_eq!(driver.delivery_status(receipt).unwrap(), None);
assert_eq!(driver.try_submit(b"next"), Err(DriverError::Closing));
}
#[test]
fn close_drain_clamps_future_ack_and_drains() {
let mut server = FakeOrderedServer::no_response();
server.push_response(FakeServerResponse::Ack { wire_seq: 1 });
let mut driver = driver(server);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.close_drain_steps(2).unwrap(), CloseOutcome::Drained);
assert_eq!(driver.try_submit(b"next"), Err(DriverError::Closing));
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::CompletedThrough {
fsn: 0,
wire_seq: 0,
}
]
);
}
#[test]
fn close_drain_reports_terminal_failure() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::TerminalFailure,
]));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.close_drain_steps(1).unwrap(), CloseOutcome::Terminal);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(driver.try_submit(b"next"), Err(DriverError::Terminal));
}
#[test]
fn terminal_failure_marks_unresolved_receipts_and_rejects_future_submit() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::NoResponse,
FakeSendResult::TerminalFailure,
]));
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(
driver.terminal_error().map(Error::msg),
Some("fake terminal failure")
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Terminal { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Terminal { fsn: 1 }
);
assert_eq!(
driver.delivery_status(first).unwrap(),
Some(DeliveryOutcome::Terminal)
);
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(driver.try_submit(b"third"), Err(DriverError::Terminal));
}
#[test]
fn retriable_reject_after_prior_ack_reconnects_and_preserves_replay_tail() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::AckWire { wire_seq: 0 },
FakeSendResult::RejectWire { wire_seq: 1 },
]));
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
let third = driver.try_submit(b"third").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
assert_eq!(
driver
.last_server_error()
.map(|err| (err.status, err.message.as_str())),
Some((codec::WS_STATUS_WRITE_ERROR, "fake write error"))
);
assert_eq!(
driver.receipt_status(third),
QwpReceiptStatus::Published { fsn: 2 }
);
}
#[test]
fn reject_gap_terminalizes_without_completing_unresolved_lower_frame() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::NoResponse,
FakeSendResult::RejectWire { wire_seq: 1 },
]));
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Terminal { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Terminal { fsn: 1 }
);
let terminal_error = driver.terminal_sender_error().unwrap();
assert_eq!(
terminal_error.category,
QwpWsErrorCategory::ProtocolViolation
);
assert!(terminal_error.message.as_deref().is_some_and(|message| {
message.contains("reject response for fsn 1 skipped unresolved fsn 0")
}));
}
#[test]
fn retriable_rejection_event_agrees_with_pending_replay_status() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::AckWire { wire_seq: 0 },
FakeSendResult::RejectWire { wire_seq: 1 },
]));
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
let third = driver.try_submit(b"third").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
driver.drive_once().unwrap();
assert_eq!(driver.delivery_status(second).unwrap(), None);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Published { fsn: 0 },
DriverEvent::Published { fsn: 1 },
DriverEvent::Published { fsn: 2 },
DriverEvent::Sent {
fsn: 0,
wire_seq: 0,
},
DriverEvent::CompletedThrough {
fsn: 0,
wire_seq: 0,
},
DriverEvent::Sent {
fsn: 1,
wire_seq: 1,
},
DriverEvent::Rejected {
fsn: 1,
wire_seq: 1,
},
DriverEvent::Reconnected {
reason: ReconnectReason::RetryableFailure,
},
]
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Published { fsn: 1 }
);
assert_eq!(
driver.receipt_status(third),
QwpReceiptStatus::Published { fsn: 2 }
);
}
#[test]
fn retriable_reject_replays_from_rejected_frame_before_later_receipt() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::AckWire { wire_seq: 0 },
FakeSendResult::RejectWire { wire_seq: 1 },
FakeSendResult::AckWire { wire_seq: 2 },
]));
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
let third = driver.try_submit(b"third").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Acked { wire_seq: 0 }
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
assert_eq!(
driver.receipt_status(third),
QwpReceiptStatus::Published { fsn: 2 }
);
}
#[test]
fn wait_keeps_receipt_pending_after_retriable_reject_until_replay_ack() {
let mut driver = driver(FakeOrderedServer::scripted([FakeSendResult::RejectWire {
wire_seq: 0,
}]));
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_once().unwrap();
assert_eq!(driver.delivery_status(receipt).unwrap(), None);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
}
#[test]
fn retriable_server_error_is_pollable() {
let mut driver = driver(FakeOrderedServer::scripted([FakeSendResult::RejectWire {
wire_seq: 0,
}]));
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_once().unwrap();
assert_eq!(driver.delivery_status(receipt).unwrap(), None);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::WriteError);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Retriable);
assert_eq!(error.status, Some(codec::WS_STATUS_WRITE_ERROR));
assert_eq!(error.message.as_deref(), Some("fake write error"));
assert_eq!(error.message_sequence, Some(0));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
assert_eq!(driver.terminal_sender_error(), None);
assert_eq!(driver.poll_sender_error(), None);
}
#[test]
fn retriable_reject_poison_terminal_carries_last_server_error() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::RejectWire { wire_seq: 0 },
FakeSendResult::RejectWire { wire_seq: 1 },
FakeSendResult::RejectWire { wire_seq: 2 },
FakeSendResult::RejectWire { wire_seq: 3 },
]));
let (sink, callback_ran) = terminal_latch_asserting_sink(driver.store.lifecycle());
driver.store.set_rejection_sink(Some(sink));
let receipt = driver.try_submit(b"payload").unwrap();
let mut outcome = driver.drive_once().unwrap();
for _ in 0..8 {
if outcome == DriveOutcome::Terminal {
break;
}
outcome = driver.drive_once().unwrap();
}
assert_eq!(outcome, DriveOutcome::Terminal);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
let terminal_error = driver.terminal_sender_error().unwrap();
assert_eq!(
terminal_error.category,
QwpWsErrorCategory::ProtocolViolation
);
assert!(
terminal_error.message.as_deref().is_some_and(|message| {
message.contains("was rejected 4 times without ACK progress")
&& message.contains("fake write error")
}),
"poison terminal must carry the last server error, got: {:?}",
terminal_error.message
);
assert!(callback_ran.load(Ordering::Acquire));
}
#[test]
fn role_reject_never_terminalizes_and_never_strikes() {
let mut driver = driver(FakeOrderedServer::scripted([
FakeSendResult::RejectWireNotWritable { wire_seq: 0 },
FakeSendResult::RejectWireNotWritable { wire_seq: 1 },
FakeSendResult::RejectWireNotWritable { wire_seq: 2 },
FakeSendResult::RejectWireNotWritable { wire_seq: 3 },
FakeSendResult::RejectWireNotWritable { wire_seq: 4 },
FakeSendResult::RejectWireNotWritable { wire_seq: 5 },
]));
let receipt = driver.try_submit(b"payload").unwrap();
for _ in 0..12 {
let outcome = driver.drive_once().unwrap();
assert_ne!(
outcome,
DriveOutcome::Terminal,
"a role reject must never terminalize"
);
}
assert!(!driver.is_terminal());
assert_eq!(
driver.send_core.poison_tracker.strikes(),
0,
"role rejects are strike-exempt"
);
let status = driver.receipt_status(receipt);
assert!(
matches!(
status,
QwpReceiptStatus::Published { fsn: 0 } | QwpReceiptStatus::Sent { fsn: 0, .. }
),
"the queued frame must stay replayable, got {status:?}"
);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::NotWritable);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::RetriableOther);
}
#[test]
fn poison_dwell_driver_holds_terminal_until_window_elapses() {
let mut driver = QwpWsCoreTestHarness::from_queue_with_rejection_limit_and_window(
memory_queue(options(8, 1024)),
FakeOrderedServer::scripted([
FakeSendResult::RejectWire { wire_seq: 0 },
FakeSendResult::RejectWire { wire_seq: 1 },
FakeSendResult::RejectWire { wire_seq: 2 },
]),
2,
Duration::from_secs(2),
);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
std::thread::sleep(Duration::from_millis(2300));
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
let terminal_error = driver.terminal_sender_error().unwrap();
assert!(
terminal_error.message.as_deref().is_some_and(|message| {
message.contains("was rejected 3 times without ACK progress")
&& message.contains("fake write error")
}),
"poison terminal must report actual strike count, got: {:?}",
terminal_error.message
);
}
#[test]
fn server_close_before_any_send_does_not_strike() {
let mut driver = QwpWsCoreTestHarness::from_queue_with_rejection_limit_and_window(
memory_queue(options(8, 1024)),
FakeOrderedServer::no_response(),
1,
Duration::ZERO,
);
driver.try_submit(b"payload").unwrap();
match driver.send_core.transport_failure_action(
&mut driver.store,
TransportFailure::ServerClose(fake_transport_error("server close")),
) {
QwpWsTransportFailureAction::Reconnect { reason, .. } => {
assert_eq!(reason, ReconnectReason::Disconnect);
}
other => panic!("expected reconnect, got {other:?}"),
}
assert_eq!(driver.send_core.poison_tracker.strikes(), 0);
assert!(!driver.is_terminal());
}
#[test]
fn server_close_after_send_still_strikes() {
let mut driver = QwpWsCoreTestHarness::from_queue_with_rejection_limit_and_window(
memory_queue(options(8, 1024)),
FakeOrderedServer::no_response(),
1,
Duration::ZERO,
);
driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_send_once().unwrap(),
DriveOutcome::Sent(_)
));
match driver.send_core.transport_failure_action(
&mut driver.store,
TransportFailure::ServerClose(fake_transport_error("server close")),
) {
QwpWsTransportFailureAction::Terminal(err) => {
assert!(
err.msg()
.contains("was closed 1 times without ACK progress"),
"unexpected terminal error: {}",
err.msg()
);
}
other => panic!("expected terminal, got {other:?}"),
}
assert!(driver.is_terminal());
}
#[test]
fn sender_error_log_cursors_are_independent() {
let mut log = SenderErrorLog::new(2);
let error = sender_error(0);
log.push(error.clone());
assert_eq!(log.poll_notification(), Some(error.clone()));
assert_eq!(log.poll(), Some(error));
assert_eq!(log.poll_notification(), None);
assert_eq!(log.poll(), None);
assert_eq!(log.dropped_total(), 0);
}
#[test]
fn sender_error_log_drop_count_is_unified_across_cursors() {
let mut log = SenderErrorLog::new(1);
let first = sender_error(0);
let second = sender_error(1);
log.push(first.clone());
assert_eq!(log.poll(), Some(first));
log.push(second.clone());
assert_eq!(log.dropped_total(), 1);
assert_eq!(log.poll_notification(), Some(second.clone()));
assert_eq!(log.poll(), Some(second));
}
#[test]
fn terminal_qwp_server_error_is_pollable_after_terminalization() {
let mut driver = driver(FakeOrderedServer::no_response());
let (sink, callback_ran) = terminal_latch_asserting_sink(driver.store.lifecycle());
driver.store.set_rejection_sink(Some(sink));
let receipt = driver.try_submit(b"payload").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Sent(_)
));
driver
.send_core
.transport
.push_response(TransportResponse::Reject {
wire_seq: 0,
error: QwpServerError {
status: codec::WS_STATUS_PARSE_ERROR,
message: "bad payload".to_string(),
error: error::fmt!(InvalidApiCall, "QWP parse error: bad payload"),
},
});
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Terminal);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
let terminal_error = driver.terminal_sender_error().unwrap();
assert_eq!(terminal_error.category, QwpWsErrorCategory::ParseError);
assert_eq!(terminal_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(terminal_error.status, Some(codec::WS_STATUS_PARSE_ERROR));
assert_eq!(terminal_error.message.as_deref(), Some("bad payload"));
assert_eq!(terminal_error.message_sequence, Some(0));
assert_eq!(terminal_error.from_fsn, 0);
assert_eq!(terminal_error.to_fsn, 0);
let latched_error = driver.terminal_error().unwrap();
assert_eq!(latched_error.code(), ErrorCode::ServerRejection);
assert_eq!(latched_error.qwp_ws_rejection(), Some(terminal_error));
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::ParseError);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(error.status, Some(codec::WS_STATUS_PARSE_ERROR));
assert_eq!(error.message.as_deref(), Some("bad payload"));
assert_eq!(error.message_sequence, Some(0));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
assert_eq!(driver.terminal_sender_error(), Some(&error));
assert!(callback_ran.load(Ordering::Acquire));
}
#[test]
fn protocol_violation_error_records_unresolved_fsn_span() {
let mut driver = driver(FakeOrderedServer::no_response());
let (sink, callback_ran) = terminal_latch_asserting_sink(driver.store.lifecycle());
driver.store.set_rejection_sink(Some(sink));
driver.try_submit(b"first").unwrap();
driver.try_submit(b"second").unwrap();
assert_eq!(
driver
.send_core
.apply_transport_failure(
&mut driver.store,
TransportFailure::ProtocolViolation {
close_code: Some(1002),
reason: "bad frame".to_string(),
},
)
.unwrap(),
DriveOutcome::Terminal
);
let terminal_error = driver.terminal_sender_error().unwrap();
assert_eq!(
terminal_error.category,
QwpWsErrorCategory::ProtocolViolation
);
assert_eq!(terminal_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(terminal_error.status, None);
assert_eq!(terminal_error.message_sequence, None);
assert_eq!(terminal_error.from_fsn, 0);
assert_eq!(terminal_error.to_fsn, 1);
assert!(
terminal_error
.message
.as_deref()
.unwrap()
.contains("bad frame")
);
assert!(callback_ran.load(Ordering::Acquire));
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::ProtocolViolation);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(error.status, None);
assert_eq!(error.message_sequence, None);
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 1);
assert!(error.message.unwrap().contains("bad frame"));
}
#[test]
fn sender_error_overflow_drops_oldest_error() {
let mut driver = driver_with_event_capacity(
FakeOrderedServer::scripted([
FakeSendResult::RejectWire { wire_seq: 0 },
FakeSendResult::RejectWire { wire_seq: 1 },
]),
1,
);
driver.try_submit(b"first").unwrap();
driver.try_submit(b"second").unwrap();
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
));
assert!(matches!(
driver.drive_once().unwrap(),
DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
}
));
assert_eq!(driver.sender_errors_dropped_total(), 1);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.from_fsn, 0);
assert_eq!(driver.poll_sender_error(), None);
}
#[test]
fn delivery_status_returns_completed_for_completed_receipts() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_once().unwrap();
assert_eq!(
driver.delivery_status(receipt).unwrap(),
Some(DeliveryOutcome::Completed)
);
}
#[test]
fn completed_receipt_is_observable_without_redriving() {
let mut driver = driver(FakeOrderedServer::ack_each_send());
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_once().unwrap();
assert_eq!(
driver.delivery_status(receipt).unwrap(),
Some(DeliveryOutcome::Completed)
);
assert_eq!(driver.drive_once().unwrap(), DriveOutcome::Idle);
}
#[test]
fn delivery_status_unknown_receipt_is_api_error() {
let driver = driver(FakeOrderedServer::ack_each_send());
assert_eq!(
driver.delivery_status(QwpReceipt { fsn: 99 }),
Err(DriverError::UnknownReceipt { fsn: 99 })
);
}
#[test]
fn pending_receipt_status_after_drive_without_ack() {
let mut driver = driver(FakeOrderedServer::no_response());
let receipt = driver.try_submit(b"payload").unwrap();
driver.drive_once().unwrap();
assert_eq!(driver.delivery_status(receipt).unwrap(), None);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Sent {
fsn: 0,
wire_seq: 0
}
);
}
#[test]
fn drive_receive_once_ignores_ack_before_send_like_java() {
let mut server = FakeOrderedServer::no_response();
server.push_response(FakeServerResponse::Ack { wire_seq: 0 });
let mut driver = driver(server);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.drive_receive_once(), Ok(DriveOutcome::Progress));
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), None);
assert_eq!(driver.poll_sender_error(), None);
}
#[test]
fn drive_receive_once_terminalizes_presend_schema_reject_without_ack_advance() {
let mut server = FakeOrderedServer::no_response();
server.push_response(FakeServerResponse::Reject {
wire_seq: 42,
error: schema_mismatch_error("pre-send schema mismatch"),
});
let mut driver = driver(server);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.drive_receive_once(), Ok(DriveOutcome::Terminal));
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), None);
assert!(driver.terminal_sender_error().is_some());
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::SchemaMismatch);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(error.status, Some(codec::WS_STATUS_SCHEMA_MISMATCH));
assert_eq!(error.message.as_deref(), Some("pre-send schema mismatch"));
assert_eq!(error.message_sequence, Some(42));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
}
#[test]
fn drive_receive_once_reconnects_presend_retryable_reject_without_ack_advance() {
let mut server = FakeOrderedServer::no_response();
server.push_response(FakeServerResponse::Reject {
wire_seq: 42,
error: write_error("pre-send write failure"),
});
let mut driver = driver(server);
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(
driver.drive_receive_once(),
Ok(DriveOutcome::Reconnected {
reason: ReconnectReason::RetryableFailure
})
);
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Published { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), None);
assert_eq!(driver.terminal_sender_error(), None);
let error = driver.poll_sender_error().unwrap();
assert_eq!(error.category, QwpWsErrorCategory::WriteError);
assert_eq!(error.applied_policy, QwpWsErrorPolicy::Retriable);
assert_eq!(error.status, Some(codec::WS_STATUS_WRITE_ERROR));
assert_eq!(error.message.as_deref(), Some("pre-send write failure"));
assert_eq!(error.message_sequence, Some(42));
assert_eq!(error.from_fsn, 0);
assert_eq!(error.to_fsn, 0);
}
#[test]
fn drive_receive_once_terminalizes_presend_terminal_reject_without_ack_advance_like_java() {
let mut server = FakeOrderedServer::no_response();
server.push_response(FakeServerResponse::Reject {
wire_seq: 7,
error: QwpServerError {
status: codec::WS_STATUS_PARSE_ERROR,
message: "bad pre-send payload".to_string(),
error: error::fmt!(InvalidApiCall, "QWP parse error: bad pre-send payload"),
},
});
let mut driver = driver(server);
let (sink, callback_ran) = terminal_latch_asserting_sink(driver.store.lifecycle());
driver.store.set_rejection_sink(Some(sink));
let receipt = driver.try_submit(b"payload").unwrap();
assert_eq!(driver.drive_receive_once(), Ok(DriveOutcome::Terminal));
assert_eq!(
driver.receipt_status(receipt),
QwpReceiptStatus::Terminal { fsn: 0 }
);
assert_eq!(driver.acked_fsn(), None);
let terminal_error = driver.terminal_sender_error().unwrap();
assert_eq!(terminal_error.category, QwpWsErrorCategory::ParseError);
assert_eq!(terminal_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(terminal_error.message_sequence, Some(7));
assert_eq!(terminal_error.from_fsn, 0);
assert_eq!(terminal_error.to_fsn, 0);
assert!(callback_ran.load(Ordering::Acquire));
}
#[test]
fn event_ring_overflow_drops_old_events_without_corrupting_receipt_status() {
let mut driver = driver_with_event_capacity(FakeOrderedServer::ack_each_send(), 2);
let first = driver.try_submit(b"first").unwrap();
let second = driver.try_submit(b"second").unwrap();
driver.drive_once().unwrap();
driver.drive_once().unwrap();
assert_eq!(
driver.delivery_status(first).unwrap(),
Some(DeliveryOutcome::Completed)
);
assert_eq!(
driver.delivery_status(second).unwrap(),
Some(DeliveryOutcome::Completed)
);
assert_eq!(driver.events_dropped_total(), 4);
assert_eq!(
drain_events(&mut driver),
vec![
DriverEvent::Sent {
fsn: 1,
wire_seq: 1,
},
DriverEvent::CompletedThrough {
fsn: 1,
wire_seq: 1,
},
]
);
assert_eq!(
driver.receipt_status(first),
QwpReceiptStatus::Completed { fsn: 0 }
);
assert_eq!(
driver.receipt_status(second),
QwpReceiptStatus::Completed { fsn: 1 }
);
}
}