use crate::RetryReason;
use std::collections::VecDeque;
use std::sync::Arc;
#[derive(Clone, Default)]
pub(crate) struct SendBatchTestHook {
inject_first_message_reasons: Arc<std::sync::Mutex<VecDeque<Option<Vec<RetryReason>>>>>,
fed_retry_reasons: Arc<std::sync::Mutex<Vec<(String, usize)>>>,
kill_on_read_by_name: Arc<std::sync::Mutex<Option<(String, usize)>>>,
}
#[allow(
clippy::expect_used,
reason = "test-support code: a panic is how a test reports failure"
)]
#[cfg_attr(
not(feature = "server-tests"),
expect(
dead_code,
reason = "the arming and reading half of the hook only has callers among the tests that need a live Redis"
)
)]
impl SendBatchTestHook {
pub(crate) fn new() -> Self {
Self::default()
}
pub(crate) fn push_injection(&self, reasons: Option<Vec<RetryReason>>) {
self.inject_first_message_reasons
.lock()
.expect("send batch test hook mutex poisoned")
.push_back(reasons);
}
pub(crate) fn fed_retry_reasons(&self) -> Vec<(String, usize)> {
self.fed_retry_reasons
.lock()
.expect("send batch test hook mutex poisoned")
.clone()
}
pub(super) fn take_injection(&self) -> Option<Vec<RetryReason>> {
self.inject_first_message_reasons
.lock()
.expect("send batch test hook mutex poisoned")
.pop_front()
.flatten()
}
pub(super) fn record_fed(&self, command_name: String, num_reasons: usize) {
self.fed_retry_reasons
.lock()
.expect("send batch test hook mutex poisoned")
.push((command_name, num_reasons));
}
pub(crate) fn arm_kill_on_read_for(&self, command_name: &str, num_reads: usize) {
*self
.kill_on_read_by_name
.lock()
.expect("send batch test hook mutex poisoned") =
Some((command_name.to_owned(), num_reads));
}
pub(super) fn take_kill_on_read_for(&self, command_name: &str) -> Option<usize> {
let mut guard = self
.kill_on_read_by_name
.lock()
.expect("send batch test hook mutex poisoned");
if guard.as_ref().is_some_and(|(name, _)| name == command_name) {
return guard.take().map(|(_, num_reads)| num_reads);
}
None
}
}
#[derive(Clone, Default)]
pub(crate) struct QueueMetricsTestHook {
messages_to_send_high_water: Arc<std::sync::atomic::AtomicUsize>,
messages_to_receive_high_water: Arc<std::sync::atomic::AtomicUsize>,
queued_commands: Arc<std::sync::atomic::AtomicUsize>,
queued_commands_high_water: Arc<std::sync::atomic::AtomicUsize>,
pub_sub_delivered: Arc<std::sync::atomic::AtomicUsize>,
pub_sub_delivery_failed: Arc<std::sync::atomic::AtomicUsize>,
pub_sub_delivered_bytes: Arc<std::sync::atomic::AtomicUsize>,
push_delivered: Arc<std::sync::atomic::AtomicUsize>,
push_delivery_failed: Arc<std::sync::atomic::AtomicUsize>,
push_delivered_bytes: Arc<std::sync::atomic::AtomicUsize>,
read_wave_high_water: Arc<std::sync::atomic::AtomicUsize>,
write_wave_high_water: Arc<std::sync::atomic::AtomicUsize>,
}
#[cfg_attr(
not(feature = "server-tests"),
expect(
dead_code,
reason = "the arming and reading half of the hook only has callers among the tests that need a live Redis"
)
)]
impl QueueMetricsTestHook {
pub(crate) fn new() -> Self {
Self::default()
}
pub(crate) fn messages_to_send_high_water(&self) -> usize {
self.messages_to_send_high_water
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn messages_to_receive_high_water(&self) -> usize {
self.messages_to_receive_high_water
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn queued_commands(&self) -> usize {
self.queued_commands
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn queued_commands_high_water(&self) -> usize {
self.queued_commands_high_water
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn pub_sub_delivered(&self) -> usize {
self.pub_sub_delivered
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn pub_sub_delivery_failed(&self) -> usize {
self.pub_sub_delivery_failed
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn pub_sub_delivered_bytes(&self) -> usize {
self.pub_sub_delivered_bytes
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn push_delivered(&self) -> usize {
self.push_delivered
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn push_delivery_failed(&self) -> usize {
self.push_delivery_failed
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn push_delivered_bytes(&self) -> usize {
self.push_delivered_bytes
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn read_wave_high_water(&self) -> usize {
self.read_wave_high_water
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(crate) fn write_wave_high_water(&self) -> usize {
self.write_wave_high_water
.load(std::sync::atomic::Ordering::Relaxed)
}
pub(super) fn record_read_wave(&self, handled: usize) {
self.read_wave_high_water
.fetch_max(handled, std::sync::atomic::Ordering::Relaxed);
}
pub(super) fn record_write_wave(&self, handled: usize) {
self.write_wave_high_water
.fetch_max(handled, std::sync::atomic::Ordering::Relaxed);
}
pub(super) fn record_queue_depths(
&self,
to_send: usize,
to_receive: usize,
queued_commands: usize,
) {
self.messages_to_send_high_water
.fetch_max(to_send, std::sync::atomic::Ordering::Relaxed);
self.messages_to_receive_high_water
.fetch_max(to_receive, std::sync::atomic::Ordering::Relaxed);
self.queued_commands
.store(queued_commands, std::sync::atomic::Ordering::Relaxed);
self.queued_commands_high_water
.fetch_max(queued_commands, std::sync::atomic::Ordering::Relaxed);
}
pub(super) fn record_pub_sub_delivered(&self, bytes: usize) {
self.pub_sub_delivered
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.pub_sub_delivered_bytes
.fetch_add(bytes, std::sync::atomic::Ordering::Relaxed);
}
pub(super) fn record_pub_sub_delivery_failed(&self) {
self.pub_sub_delivery_failed
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
pub(super) fn record_push_delivered(&self, bytes: usize) {
self.push_delivered
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
self.push_delivered_bytes
.fetch_add(bytes, std::sync::atomic::Ordering::Relaxed);
}
pub(super) fn record_push_delivery_failed(&self) {
self.push_delivery_failed
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}