use std::{
sync::atomic::{AtomicI64, Ordering},
time::Duration,
};
use super::{
replay_align_barrier::ReplayAlignBarrier,
replica_read_session_context::ReplicaReadSessionContext,
virtual_sublog_replay_state::VirtualSublogReplayState,
};
use crate::aof::garnet_log::GarnetLog;
pub struct ReadConsistencyManager {
current_version: AtomicI64,
physical_sublog_count: usize,
replay_task_count: usize,
vsrs: Vec<VirtualSublogReplayState>,
proactive_replay_drift_check_enabled: bool,
replay_drift_threshold: i64,
reactive_replay_drift_check_enabled: bool,
replay_drift_interval: i64,
pub replay_barrier: ReplayAlignBarrier,
}
impl ReadConsistencyManager {
pub fn new(
current_version: i64,
physical_sublog_count: usize,
replay_task_count: usize,
replay_drift_threshold: i64,
replay_drift_check_freq: i64,
) -> Self {
let virtual_sublog_count = physical_sublog_count.max(1) * replay_task_count.max(1);
let window_length = (replay_drift_check_freq.max(0) * replay_drift_threshold.max(0)).max(1);
let proactive =
replay_drift_check_freq > 0 && replay_drift_threshold >= 0 && virtual_sublog_count > 1;
let reactive = replay_drift_threshold >= 0 && virtual_sublog_count > 1;
let vsrs = (0..virtual_sublog_count)
.map(|idx| {
VirtualSublogReplayState::new(if proactive {
(idx as i64) * window_length
} else {
i64::MAX
})
})
.collect();
let replay_drift_interval = window_length * virtual_sublog_count as i64;
Self {
current_version: AtomicI64::new(current_version),
physical_sublog_count: physical_sublog_count.max(1),
replay_task_count: replay_task_count.max(1),
vsrs,
proactive_replay_drift_check_enabled: proactive,
replay_drift_threshold,
reactive_replay_drift_check_enabled: reactive,
replay_drift_interval,
replay_barrier: ReplayAlignBarrier::new(virtual_sublog_count, None),
}
}
#[inline]
pub fn current_version(&self) -> i64 {
self.current_version.load(Ordering::Acquire)
}
pub fn virtual_sublog_count(&self) -> usize {
self.vsrs.len()
}
#[inline]
pub fn key_hash(&self, key: &[u8]) -> i64 {
GarnetLog::hash(key)
}
#[inline]
pub fn virtual_sublog_idx_of_hash(&self, hash: i64) -> usize {
((hash as u64) % (self.physical_sublog_count as u64)) as usize * self.replay_task_count
+ ((hash as u64) / (self.physical_sublog_count as u64) % (self.replay_task_count as u64))
as usize
}
#[inline]
pub fn get_virtual_sublog_idx(&self, sublog_idx: usize, replay_idx: usize) -> usize {
sublog_idx * self.replay_task_count + replay_idx
}
#[inline]
fn vsr(&self, virtual_sublog_idx: usize) -> &VirtualSublogReplayState {
&self.vsrs[virtual_sublog_idx.min(self.vsrs.len() - 1)]
}
pub fn get_key_sequence_number(&self, key: &[u8], frontier: bool) -> i64 {
let hash = GarnetLog::hash(key);
if frontier {
self.get_sublog_frontier_sequence_number(hash)
} else {
self.get_key_sequence_number_by_hash(hash)
}
}
pub fn get_key_sequence_number_by_hash(&self, hash: i64) -> i64 {
self
.vsr(self.virtual_sublog_idx_of_hash(hash))
.get_key_sequence_number(hash)
}
pub fn get_sublog_frontier_sequence_number(&self, hash: i64) -> i64 {
self
.vsr(self.virtual_sublog_idx_of_hash(hash))
.get_frontier_sequence_number(hash)
}
pub fn get_physical_sublog_max_replayed_sequence_number(&self) -> Vec<i64> {
(0..self.physical_sublog_count)
.map(|physical_sublog_idx| self.get_physical_sublog_max(physical_sublog_idx))
.collect()
}
pub fn get_physical_sublog_max(&self, physical_sublog_idx: usize) -> i64 {
let start_idx = self.get_virtual_sublog_idx(physical_sublog_idx, 0);
(0..self.replay_task_count)
.map(|rt| self.vsr(start_idx + rt).max())
.max()
.unwrap_or(0)
}
pub fn get_physical_sublog_max_sequence_vector(&self) -> String {
(0..self.physical_sublog_count)
.map(|idx| self.get_physical_sublog_max(idx).to_string())
.collect::<Vec<_>>()
.join(",")
}
pub fn get_physical_sublog_max_drift_sequence_vector(&self) -> String {
let maxes: Vec<i64> = (0..self.physical_sublog_count)
.map(|idx| self.get_physical_sublog_max(idx))
.collect();
let max_sequence_number = maxes.iter().copied().max().unwrap_or(0);
maxes
.iter()
.map(|m| (max_sequence_number - m).to_string())
.collect::<Vec<_>>()
.join(",")
}
pub fn update_physical_sublog_max_sequence_number(
&self,
physical_sublog_idx: usize,
sequence_number: i64,
) {
let start_idx = self.get_virtual_sublog_idx(physical_sublog_idx, 0);
for rt in 0..self.replay_task_count {
self
.vsr(start_idx + rt)
.update_max_sequence_number(sequence_number);
}
}
pub fn advance_virtual_sublog_time(&self, virtual_sublog_idx: usize, sequence_number: i64) {
let vsr = self.vsr(virtual_sublog_idx);
vsr.update_max_sequence_number(sequence_number);
self
.replay_barrier
.signal_arrival(virtual_sublog_idx, vsr.max());
}
pub fn update_virtual_sublog_max_sequence_number(
&self,
virtual_sublog_idx: usize,
sequence_number: i64,
) {
self
.vsr(virtual_sublog_idx)
.update_max_sequence_number(sequence_number);
}
pub fn update_virtual_sublog_key_sequence_number(
&self,
virtual_sublog_idx: usize,
key_hash: i64,
sequence_number: i64,
) {
let vsr = self.vsr(virtual_sublog_idx);
vsr.update_max_sequence_number(sequence_number);
if self.proactive_replay_drift_check_enabled
&& sequence_number >= vsr.next_drift_check_window_lower_bound()
{
let interval = self.replay_drift_interval;
let mut next = vsr.next_drift_check_window_lower_bound() + interval;
if next <= sequence_number {
next += ((sequence_number - next) / interval + 1) * interval;
}
vsr.set_next_drift_check_window_lower_bound(next);
self.bound_replay_drift();
}
self
.replay_barrier
.signal_arrival_and_wait(virtual_sublog_idx, vsr.max());
vsr.update_key_sequence_number(key_hash, sequence_number);
}
pub fn update_key_sequence_number_by_hash(&self, key_hash: i64, sequence_number: i64) {
let idx = self.virtual_sublog_idx_of_hash(key_hash);
self
.vsr(idx)
.update_key_sequence_number(key_hash, sequence_number);
}
pub fn check_consistency_manager_version(
&self,
replica_read_session_context: &mut ReplicaReadSessionContext,
) {
if replica_read_session_context.session_version() != self.current_version() {
replica_read_session_context.set_session_version(self.current_version());
replica_read_session_context.set_last_virtual_sublog_idx(-1);
replica_read_session_context.set_maximum_session_sequence_number(0);
replica_read_session_context.reset_cached_sublog_max();
}
}
pub fn verify_key_freshness(
&self,
key_hash: i64,
replica_read_session_context: &mut ReplicaReadSessionContext,
timeout: Duration,
) {
let virtual_sublog_idx = self.virtual_sublog_idx_of_hash(key_hash);
let last_idx = replica_read_session_context.last_virtual_sublog_idx();
let init_or_same_sublog = last_idx == -1 || last_idx as usize == virtual_sublog_idx;
let mssn = replica_read_session_context.maximum_session_sequence_number();
let vsr = self.vsr(virtual_sublog_idx);
vsr.prefetch_key_sequence_number(key_hash);
if !init_or_same_sublog
&& mssn >= replica_read_session_context.cached_sublog_max(virtual_sublog_idx)
{
let sketch_max_value = vsr.max();
replica_read_session_context.set_cached_sublog_max(virtual_sublog_idx, sketch_max_value);
if mssn >= sketch_max_value {
self.bound_replay_drift();
if !vsr.wait_for_sequence_number(mssn, &replica_read_session_context.waiter(), timeout) {
}
replica_read_session_context.set_cached_sublog_max(virtual_sublog_idx, vsr.max());
}
}
replica_read_session_context.set_last_virtual_sublog_idx(virtual_sublog_idx as i32);
replica_read_session_context.set_last_hash(key_hash);
}
pub fn bound_replay_drift(&self) {
if !self.reactive_replay_drift_check_enabled {
return;
}
if self.replay_barrier.in_progress() {
return;
}
let mut min_frontier = i64::MAX;
let mut max_frontier = i64::MIN;
for vsr in &self.vsrs {
let frontier = vsr.max();
min_frontier = min_frontier.min(frontier);
max_frontier = max_frontier.max(frontier);
}
if max_frontier - min_frontier <= self.replay_drift_threshold {
return;
}
self.replay_barrier.try_open_round(max_frontier);
}
pub fn pre_single_key_consistent_read(
&self,
hash: i64,
replica_read_session_context: &mut ReplicaReadSessionContext,
timeout: Duration,
) {
self.check_consistency_manager_version(replica_read_session_context);
self.verify_key_freshness(hash, replica_read_session_context, timeout);
}
pub fn post_single_key_consistent_read(
&self,
replica_read_session_context: &mut ReplicaReadSessionContext,
) {
let key_sequence_number =
self.get_key_sequence_number_by_hash(replica_read_session_context.last_hash());
replica_read_session_context.advance_maximum_session_sequence_number(key_sequence_number);
}
pub fn pre_batch_key_consistent_read(
&self,
key: &[u8],
batch_read_context: &mut ReplicaReadSessionContext,
timeout: Duration,
) -> i64 {
let hash = GarnetLog::hash(key);
self.verify_key_freshness(hash, batch_read_context, timeout);
let key_sequence_number = self.get_key_sequence_number_by_hash(batch_read_context.last_hash());
batch_read_context.advance_maximum_session_sequence_number(key_sequence_number);
hash
}
pub fn post_batch_key_consistent_read_validate(
&self,
hash: i64,
batch_read_context: &ReplicaReadSessionContext,
) -> bool {
let key_sequence_number = self.get_key_sequence_number_by_hash(hash);
key_sequence_number <= batch_read_context.maximum_session_sequence_number()
}
}
#[cfg(test)]
mod tests {
use super::*;
fn manager() -> ReadConsistencyManager {
ReadConsistencyManager::new(1, 2, 2, -1, 0)
}
#[test]
fn version_check_resets_context_once() {
let m = manager();
let mut ctx = ReplicaReadSessionContext::default();
m.check_consistency_manager_version(&mut ctx);
assert_eq!(ctx.session_version(), 1);
ctx.set_maximum_session_sequence_number(50);
m.check_consistency_manager_version(&mut ctx);
assert_eq!(ctx.maximum_session_sequence_number(), 50, "同版本不重置");
}
#[test]
fn update_and_read_key_sequence_number() {
let m = manager();
let key = b"rk";
assert_eq!(m.get_key_sequence_number(key, false), 0);
m.update_key_sequence_number_by_hash(GarnetLog::hash(key), 11);
assert_eq!(m.get_key_sequence_number(key, false), 11);
assert!(m.get_key_sequence_number(key, true) >= 11);
}
#[test]
fn physical_sublog_max_vector_and_drift() {
let m = manager();
m.update_physical_sublog_max_sequence_number(0, 30);
m.update_physical_sublog_max_sequence_number(1, 10);
assert_eq!(m.get_physical_sublog_max(0), 30);
assert_eq!(m.get_physical_sublog_max(1), 10);
assert_eq!(m.get_physical_sublog_max_sequence_vector(), "30,10");
assert_eq!(m.get_physical_sublog_max_drift_sequence_vector(), "0,20");
assert_eq!(
m.get_physical_sublog_max_replayed_sequence_number(),
vec![30, 10]
);
}
#[test]
fn virtual_sublog_routing_matches_formula() {
let m = manager();
for hash in [0i64, 1, 12345, i64::MAX, i64::MIN] {
let expected = ((hash as u64) % 2) as usize * 2 + ((hash as u64) / 2 % 2) as usize;
assert_eq!(m.virtual_sublog_idx_of_hash(hash), expected);
assert!(m.virtual_sublog_idx_of_hash(hash) < 4);
}
assert_eq!(m.get_virtual_sublog_idx(1, 2), 4);
}
#[test]
fn consistent_read_protocol_single_key() {
let m = manager();
let key = b"proto";
let hash = GarnetLog::hash(key);
m.update_key_sequence_number_by_hash(hash, 5);
let mut ctx = ReplicaReadSessionContext::default();
m.pre_single_key_consistent_read(hash, &mut ctx, Duration::from_millis(10));
m.post_single_key_consistent_read(&mut ctx);
assert!(ctx.maximum_session_sequence_number() >= 5);
}
#[test]
fn batch_protocol_cross_sublog_wait_and_validate() {
let m = manager();
let k1 = b"b1";
let k2 = b"b2";
let (h1, h2) = (GarnetLog::hash(k1), GarnetLog::hash(k2));
if m.virtual_sublog_idx_of_hash(h1) == m.virtual_sublog_idx_of_hash(h2) {
return;
}
m.update_key_sequence_number_by_hash(h1, 3);
let mut ctx = ReplicaReadSessionContext::default();
let got1 = m.pre_batch_key_consistent_read(k1, &mut ctx, Duration::from_millis(10));
m.update_virtual_sublog_key_sequence_number(m.virtual_sublog_idx_of_hash(h2), h2, 8);
let got2 = m.pre_batch_key_consistent_read(k2, &mut ctx, Duration::from_millis(500));
assert!(m.post_batch_key_consistent_read_validate(got1, &ctx));
assert!(m.post_batch_key_consistent_read_validate(got2, &ctx));
assert!(ctx.maximum_session_sequence_number() >= 3);
}
#[test]
fn drift_bounding_opens_barrier_round() {
let m = ReadConsistencyManager::new(1, 2, 1, 5, 0);
assert!(!m.replay_barrier.in_progress());
m.update_virtual_sublog_max_sequence_number(0, 100);
m.bound_replay_drift();
assert!(m.replay_barrier.in_progress(), "漂移 100 应开轮");
m.bound_replay_drift();
assert!(m.replay_barrier.in_progress());
}
#[test]
fn advance_virtual_sublog_time_signals_barrier() {
let m = ReadConsistencyManager::new(1, 1, 2, 5, 0);
m.replay_barrier.try_open_round(10);
m.advance_virtual_sublog_time(0, 12);
m.advance_virtual_sublog_time(1, 12);
assert!(!m.replay_barrier.in_progress());
}
}