use super::{
transaction_manager::{TransactionManager, TxnState},
txn_key_entry::{LockType, TxnKeyEntries},
txn_key_entry_comparison::TxnKeyEntryComparison,
};
use crate::{
resp::resp_server_session::RespServerSession, storage::session::storage_session::StoreType,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct TxnKeySpec {
pub first_idx: usize,
pub last_idx: i64,
pub step: usize,
pub read_only: bool,
}
impl TxnKeySpec {
pub fn new(first_idx: usize, last_idx: i64, step: usize, read_only: bool) -> Self {
Self {
first_idx,
last_idx,
step,
read_only,
}
}
}
#[derive(Debug, Clone)]
pub struct TxnCommandKeys {
pub store_type: StoreType,
pub key_specs: Vec<TxnKeySpec>,
}
impl TransactionManager {
pub(crate) fn register_key_lock(
key_entries: &mut TxnKeyEntries,
perform_writes: &mut bool,
key: &[u8],
lock_type: LockType,
) {
*perform_writes |= lock_type == LockType::Exclusive;
key_entries.add_key(TxnKeyEntryComparison::key_hash(key), lock_type);
}
pub fn save_key_entry_to_lock(&mut self, key: &[u8], lock_type: LockType) {
Self::register_key_lock(
&mut self.key_entries,
&mut self.perform_writes,
key,
lock_type,
);
}
pub fn reset_cache_slot_verification_result(&mut self) {
if self.cluster_enabled {
}
}
pub fn write_cached_slot_verification_message(&self, _output: &mut Vec<u8>) {
if self.cluster_enabled {
}
}
pub fn verify_key_ownership(
&mut self,
_session: &RespServerSession,
_key: &[u8],
lock_type: LockType,
) {
if !self.cluster_enabled || self.is_replaying {
return;
}
let _ = lock_type == LockType::Shared;
}
pub fn abort_key_ownership(&mut self) {
self.state = TxnState::Aborted;
}
pub fn lock_keys(&mut self, session: &RespServerSession, command_keys: &TxnCommandKeys) {
if command_keys.key_specs.is_empty() {
return;
}
self.add_transaction_store_type(command_keys.store_type);
for key_spec in &command_keys.key_specs {
let last_idx = if key_spec.last_idx < 0 {
(session.parse_state.count as i64 + key_spec.last_idx).max(0)
} else {
key_spec.last_idx.min(session.parse_state.count as i64)
};
let mut curr_idx = key_spec.first_idx;
while curr_idx <= last_idx as usize && curr_idx < session.parse_state.count {
let key = session.parse_state.get_arg_slice_by_ref(curr_idx);
let key_bytes = key.as_slice();
let lock_type = if key_spec.read_only {
LockType::Shared
} else {
LockType::Exclusive
};
self.save_key_entry_to_lock(key_bytes, lock_type);
self.save_key_arg_slice(key_bytes);
curr_idx += key_spec.step.max(1);
}
}
}
}
#[cfg(test)]
mod tests {
use std::sync::Arc;
use super::{
super::{
transaction_manager::{TransactionStoreTypes, TxnState},
watch_version_map::WatchVersionMap,
},
*,
};
use crate::{arg_slice::ArgSlice, resp::resp_server_session::RespServerSession};
fn manager() -> TransactionManager {
TransactionManager::new(Arc::new(WatchVersionMap::new(64)), None, false)
}
fn session_with_args(args: &[&[u8]]) -> (RespServerSession, Vec<u8>) {
let mut buffer: Vec<u8> = Vec::new();
for arg in args {
buffer.extend_from_slice(arg);
}
let mut slices = Vec::with_capacity(args.len());
let mut offset = 0usize;
for arg in args {
slices.push(ArgSlice::new(
unsafe { buffer.as_ptr().add(offset) },
arg.len(),
));
offset += arg.len();
}
let mut session = RespServerSession::default();
session.parse_state.initialize_with_args(&slices);
(session, buffer)
}
#[test]
fn save_key_entry_marks_perform_writes_on_exclusive() {
let mut txn = manager();
txn.save_key_entry_to_lock(b"a", LockType::Shared);
assert!(!txn.perform_writes);
txn.save_key_entry_to_lock(b"b", LockType::Exclusive);
assert!(txn.perform_writes);
assert_eq!(txn.key_entries.count(), 2);
}
#[test]
fn lock_keys_expands_window_and_registers_keys() {
let mut txn = manager();
let (session, _buffer) = session_with_args(&[b"k1", b"k2", b"k3"]);
let keys = TxnCommandKeys {
store_type: StoreType::Main,
key_specs: vec![TxnKeySpec::new(0, 2, 1, false)],
};
txn.lock_keys(&session, &keys);
assert_eq!(txn.key_entries.count(), 3);
assert!(txn.perform_writes);
assert!(txn.txn_keys.is_empty());
assert!(txn.store_types.contains(TransactionStoreTypes::Main));
}
#[test]
fn lock_keys_negative_index_pins_to_last_arg() {
let mut txn = manager();
let (session, _buffer) = session_with_args(&[b"k1", b"k2", b"k3"]);
let keys = TxnCommandKeys {
store_type: StoreType::Main,
key_specs: vec![TxnKeySpec::new(2, -1, 1, true)],
};
txn.lock_keys(&session, &keys);
assert_eq!(txn.key_entries.count(), 1);
assert!(!txn.perform_writes);
}
#[test]
fn verify_key_ownership_noop_without_cluster() {
let mut txn = manager();
let session = RespServerSession::default();
txn.verify_key_ownership(&session, b"k", LockType::Shared);
assert_eq!(txn.state, TxnState::None);
}
}