use std::sync::{
Arc,
atomic::{AtomicI64, Ordering},
};
use parking_lot::RwLock;
use super::virtual_sublog_replay_state::ReadSessionWaiter;
pub struct ReplicaReadSessionContext {
session_version: i64,
maximum_session_sequence_number: i64,
last_hash: i64,
last_virtual_sublog_idx: i32,
cached_sublog_max: Arc<RwLock<Vec<AtomicI64>>>,
waiter: Arc<ReadSessionWaiter>,
}
impl Default for ReplicaReadSessionContext {
fn default() -> Self {
Self {
session_version: -1,
maximum_session_sequence_number: 0,
last_hash: 0,
last_virtual_sublog_idx: -1,
cached_sublog_max: Arc::new(RwLock::new(Vec::new())),
waiter: Arc::new(ReadSessionWaiter::new()),
}
}
}
impl ReplicaReadSessionContext {
#[inline]
pub fn session_version(&self) -> i64 {
self.session_version
}
#[inline]
pub fn set_session_version(&mut self, version: i64) {
self.session_version = version;
}
#[inline]
pub fn maximum_session_sequence_number(&self) -> i64 {
self.maximum_session_sequence_number
}
#[inline]
pub fn advance_maximum_session_sequence_number(&mut self, value: i64) {
self.maximum_session_sequence_number = self.maximum_session_sequence_number.max(value);
}
#[inline]
pub fn set_maximum_session_sequence_number(&mut self, value: i64) {
self.maximum_session_sequence_number = value;
}
#[inline]
pub fn last_hash(&self) -> i64 {
self.last_hash
}
#[inline]
pub fn set_last_hash(&mut self, hash: i64) {
self.last_hash = hash;
}
#[inline]
pub fn last_virtual_sublog_idx(&self) -> i32 {
self.last_virtual_sublog_idx
}
#[inline]
pub fn set_last_virtual_sublog_idx(&mut self, idx: i32) {
self.last_virtual_sublog_idx = idx;
}
pub fn waiter(&self) -> Arc<ReadSessionWaiter> {
Arc::clone(&self.waiter)
}
pub fn get_power_of_two_size(value: usize) -> usize {
value.max(1).next_power_of_two()
}
pub fn expand_key_hash_cache(&self, key_count: usize) {
let new_size = Self::get_power_of_two_size(key_count);
let mut cache = self.cached_sublog_max.write();
if cache.len() >= new_size {
return;
}
*cache = (0..new_size).map(|_| AtomicI64::new(0)).collect();
}
pub fn shrink_key_hash_cache(&self, key_count: usize) {
let mut cache = self.cached_sublog_max.write();
let new_size = if key_count == 0 {
0
} else {
Self::get_power_of_two_size(key_count) / 2
};
cache.truncate(new_size);
}
pub fn cached_sublog_max(&self, idx: usize) -> i64 {
self
.cached_sublog_max
.read()
.get(idx)
.map_or(0, |slot| slot.load(Ordering::Acquire))
}
pub fn set_cached_sublog_max(&self, idx: usize, value: i64) {
let cache = self.cached_sublog_max.read();
if let Some(slot) = cache.get(idx) {
slot.store(value, Ordering::Release);
}
}
pub fn reset_cached_sublog_max(&self) {
for slot in self.cached_sublog_max.read().iter() {
slot.store(0, Ordering::Release);
}
}
pub fn cached_len(&self) -> usize {
self.cached_sublog_max.read().len()
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn power_of_two_sizes() {
assert_eq!(ReplicaReadSessionContext::get_power_of_two_size(0), 1);
assert_eq!(ReplicaReadSessionContext::get_power_of_two_size(1), 1);
assert_eq!(ReplicaReadSessionContext::get_power_of_two_size(5), 8);
assert_eq!(ReplicaReadSessionContext::get_power_of_two_size(16), 16);
}
#[test]
fn expand_shrink_and_cache_roundtrip() {
let ctx = ReplicaReadSessionContext::default();
assert_eq!(ctx.session_version(), -1);
ctx.expand_key_hash_cache(5);
assert_eq!(ctx.cached_len(), 8);
ctx.set_cached_sublog_max(3, 99);
assert_eq!(ctx.cached_sublog_max(3), 99);
assert_eq!(ctx.cached_sublog_max(100), 0);
ctx.shrink_key_hash_cache(5);
assert_eq!(ctx.cached_len(), 4);
assert_eq!(ctx.cached_sublog_max(3), 99);
ctx.shrink_key_hash_cache(3);
assert_eq!(ctx.cached_len(), 2);
assert_eq!(ctx.cached_sublog_max(3), 0);
ctx.reset_cached_sublog_max();
ctx.set_cached_sublog_max(1, 7);
ctx.reset_cached_sublog_max();
assert_eq!(ctx.cached_sublog_max(1), 0);
}
#[test]
fn session_state_defaults_and_setters() {
let mut ctx = ReplicaReadSessionContext::default();
assert_eq!(ctx.last_virtual_sublog_idx(), -1);
ctx.set_session_version(3);
ctx.set_last_hash(0xdead);
ctx.set_last_virtual_sublog_idx(2);
ctx.set_maximum_session_sequence_number(10);
ctx.advance_maximum_session_sequence_number(7);
assert_eq!(ctx.maximum_session_sequence_number(), 10, "单调推进");
ctx.advance_maximum_session_sequence_number(15);
assert_eq!(ctx.maximum_session_sequence_number(), 15);
assert_eq!(ctx.last_hash(), 0xdead);
assert_eq!(ctx.last_virtual_sublog_idx(), 2);
let _ = ctx.waiter();
}
}