use std::sync::{
Arc,
atomic::{AtomicU64, Ordering},
};
use wdev::Device;
use crate::aof::{
aof_processor::{AofProcessor, AofReplayError, ReplayTarget},
garnet_append_only_file::GarnetAppendOnlyFile,
};
pub struct RecoverLogDriver {
physical_sublog_idx: usize,
start_address: i64,
until_address: i64,
until_sequence_number: i64,
replayed_record_count: AtomicU64,
}
impl RecoverLogDriver {
pub fn physical_sublog_idx(&self) -> usize {
self.physical_sublog_idx
}
pub fn new(
physical_sublog_idx: usize,
start_address: i64,
until_address: i64,
until_sequence_number: i64,
) -> Self {
Self {
physical_sublog_idx,
start_address,
until_address,
until_sequence_number,
replayed_record_count: AtomicU64::new(0),
}
}
pub fn replayed_record_count(&self) -> u64 {
self.replayed_record_count.load(Ordering::Acquire)
}
pub fn throttle(&self) {}
pub async fn run<D: Device>(
&self,
processor: &AofProcessor,
aof: &Arc<GarnetAppendOnlyFile>,
target: &ReplayTarget<'_, '_, D>,
) -> Result<u64, AofReplayError> {
if self.start_address == self.until_address {
return Ok(0);
}
let records = aof.log().scan_single(
self.physical_sublog_idx,
self.start_address,
self.until_address,
);
for record in &records {
let entry = record.payload.as_slice();
if let Some((true, _sequence_number)) =
processor.skip_replay(entry, self.until_sequence_number, record.address)
{
break;
}
let virtual_sublog_idx = self.physical_sublog_idx * self.virtual_sublog_per_sublog(aof);
processor
.process_aof_record_internal(virtual_sublog_idx, entry, true, record.address, target)
.await?;
self.replayed_record_count.fetch_add(1, Ordering::AcqRel);
}
Ok(self.replayed_record_count.load(Ordering::Acquire))
}
fn virtual_sublog_per_sublog(&self, aof: &Arc<GarnetAppendOnlyFile>) -> usize {
aof.virtual_sublog_count() / aof.log().size().max(1)
}
pub async fn consume<D: Device>(
&self,
processor: &AofProcessor,
aof: &Arc<GarnetAppendOnlyFile>,
entry: &[u8],
current_address: i64,
replay_task_idx: usize,
target: &ReplayTarget<'_, '_, D>,
) -> Result<(), AofReplayError> {
let virtual_sublog_idx =
self.physical_sublog_idx * self.virtual_sublog_per_sublog(aof) + replay_task_idx;
processor
.process_aof_record_internal(virtual_sublog_idx, entry, true, current_address, target)
.await?;
self.replayed_record_count.fetch_add(1, Ordering::AcqRel);
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub async fn create_and_run_intra_page_parallel_replay_tasks<D: Device>(
&self,
processor: &AofProcessor,
aof: &Arc<GarnetAppendOnlyFile>,
page_entries: &[&[u8]],
page_start_address: i64,
entry_stride: usize,
target: &ReplayTarget<'_, '_, D>,
) -> Result<(), AofReplayError> {
let replay_task_count = aof.virtual_sublog_count() / aof.log().size().max(1);
for (entry_idx, entry) in page_entries.iter().enumerate() {
let entry_address = page_start_address + (entry_idx * entry_stride) as i64;
let virtual_sublog_idx = self.physical_sublog_idx * self.virtual_sublog_per_sublog(aof)
+ entry_idx % replay_task_count.max(1);
if processor
.can_replay(entry, virtual_sublog_idx, entry_address)
.is_some_and(|(owned, _)| owned)
{
processor
.process_aof_record_internal(virtual_sublog_idx, entry, true, entry_address, target)
.await?;
self.replayed_record_count.fetch_add(1, Ordering::AcqRel);
}
}
Ok(())
}
}