use crate::kv::KVTransaction;
use crate::types::{Key, Value};
use crate::Result;
use serde::{Deserialize, Serialize};
const JOURNAL_PREFIX: &[u8] = b"\x00alopex/range-change/";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum RangeChangePayload {
UpsertRow {
row_key: Vec<u8>,
encoded_row: Vec<u8>,
},
DeleteRow {
row_key: Vec<u8>,
tombstone: Vec<u8>,
},
UpsertIndex {
index_id: u32,
index_key: Vec<u8>,
row_key: Vec<u8>,
},
DeleteIndex {
index_id: u32,
index_key: Vec<u8>,
row_key: Vec<u8>,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RangeChangeRecord {
pub range_id: String,
pub generation: u64,
pub epoch: u64,
pub predecessor_epoch: Option<u64>,
pub replay_id: String,
pub payload: Vec<RangeChangePayload>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RangeChangeJournalCapability {
Supported,
Unavailable,
}
pub fn stage_range_change<'a, T: KVTransaction<'a>>(
transaction: &mut T,
record: &RangeChangeRecord,
) -> Result<Key> {
let key = journal_key(record);
let value = bincode::serialize(record).expect("RangeChangeRecord serialization is infallible");
transaction.put(key.clone(), value)?;
Ok(key)
}
pub fn decode_range_change(value: &Value) -> Result<RangeChangeRecord> {
bincode::deserialize(value).map_err(|error| {
crate::Error::InvalidFormat(format!("invalid range-change journal: {error}"))
})
}
pub fn journal_key(record: &RangeChangeRecord) -> Key {
let mut key = JOURNAL_PREFIX.to_vec();
key.extend_from_slice(record.range_id.as_bytes());
key.push(0);
key.extend_from_slice(&record.generation.to_be_bytes());
key.extend_from_slice(&record.epoch.to_be_bytes());
key.push(0);
key.extend_from_slice(record.replay_id.as_bytes());
key
}
#[cfg(test)]
mod tests {
use super::*;
use crate::kv::memory::MemoryKV;
use crate::kv::{KVStore, KVTransaction};
use crate::types::TxnMode;
fn record() -> RangeChangeRecord {
RangeChangeRecord {
range_id: "range-a".into(),
generation: 1,
epoch: 2,
predecessor_epoch: Some(1),
replay_id: "replay-1".into(),
payload: vec![RangeChangePayload::UpsertRow {
row_key: b"row-a".to_vec(),
encoded_row: b"value-a".to_vec(),
}],
}
}
#[test]
fn journal_and_data_share_one_local_commit_boundary() {
let store = MemoryKV::new();
let record = record();
let key = journal_key(&record);
let mut transaction = store.begin(TxnMode::ReadWrite).unwrap();
transaction
.put(b"row-a".to_vec(), b"value-a".to_vec())
.unwrap();
stage_range_change(&mut transaction, &record).unwrap();
transaction.commit_self().unwrap();
let mut reader = store.begin(TxnMode::ReadOnly).unwrap();
assert_eq!(
reader.get(&b"row-a".to_vec()).unwrap(),
Some(b"value-a".to_vec())
);
assert_eq!(
decode_range_change(&reader.get(&key).unwrap().unwrap()).unwrap(),
record
);
}
#[test]
fn rollback_exposes_neither_data_nor_journal_entry() {
let store = MemoryKV::new();
let record = record();
let key = journal_key(&record);
let mut transaction = store.begin(TxnMode::ReadWrite).unwrap();
transaction
.put(b"row-a".to_vec(), b"value-a".to_vec())
.unwrap();
stage_range_change(&mut transaction, &record).unwrap();
transaction.rollback_self().unwrap();
let mut reader = store.begin(TxnMode::ReadOnly).unwrap();
assert!(reader.get(&b"row-a".to_vec()).unwrap().is_none());
assert!(reader.get(&key).unwrap().is_none());
}
#[test]
fn failed_commit_exposes_neither_data_nor_journal_entry() {
let store = MemoryKV::new_with_limit(Some(0));
let record = record();
let key = journal_key(&record);
let mut transaction = store.begin(TxnMode::ReadWrite).unwrap();
transaction
.put(b"row-a".to_vec(), b"value-a".to_vec())
.unwrap();
stage_range_change(&mut transaction, &record).unwrap();
assert!(transaction.commit_self().is_err());
let mut reader = store.begin(TxnMode::ReadOnly).unwrap();
assert!(reader.get(&b"row-a".to_vec()).unwrap().is_none());
assert!(reader.get(&key).unwrap().is_none());
}
}