use std::sync::Arc;
use wdev::Device;
use crate::aof::{
aof_processor::{AofProcessor, AofReplayError, ReplayTarget},
garnet_append_only_file::GarnetAppendOnlyFile,
recover::recover_log_driver::RecoverLogDriver,
};
impl RecoverLogDriver {
pub async fn replay_page<D: Device>(
&self,
processor: &AofProcessor,
aof: &Arc<GarnetAppendOnlyFile>,
page_entries: &[&[u8]],
page_start_address: i64,
replay_task_idx: usize,
target: &ReplayTarget<'_, '_, D>,
) -> Result<u64, AofReplayError> {
let virtual_per_sublog = aof.virtual_sublog_count() / aof.log().size().max(1);
let mut applied = 0u64;
for (entry_idx, entry) in page_entries.iter().enumerate() {
let entry_address = page_start_address + entry_idx as i64;
let virtual_sublog_idx = self.physical_sublog_idx() * virtual_per_sublog + replay_task_idx;
if processor
.can_replay(entry, replay_task_idx, entry_address)
.is_some_and(|(owned, _)| owned)
{
processor
.process_aof_record_internal(virtual_sublog_idx, entry, true, entry_address, target)
.await?;
applied += 1;
}
}
Ok(applied)
}
pub async fn recover_replay_task_async<D: Device>(
&self,
processor: &AofProcessor,
aof: &Arc<GarnetAppendOnlyFile>,
page_entries: &[&[u8]],
page_start_address: i64,
replay_task_idx: usize,
target: &ReplayTarget<'_, '_, D>,
) -> Result<u64, AofReplayError> {
self
.replay_page(
processor,
aof,
page_entries,
page_start_address,
replay_task_idx,
target,
)
.await
}
}