use wdev::Device;
use super::{
aof_chunked_record_reader::ChunkedAccumulator,
aof_processor::{AofProcessor, AofReplayError, PreparedParameters, ReplayTarget},
};
use crate::aof::aof_entry_type::AofEntryType;
pub async fn process_chunked_record<D: Device>(
processor: &AofProcessor,
virtual_sublog_idx: usize,
acc: ChunkedAccumulator,
as_replica: bool,
log_address_sequence_number: i64,
target: &ReplayTarget<'_, '_, D>,
) -> Result<bool, AofReplayError> {
let buffered = processor.coordinator().buffer_chunk_operation(
virtual_sublog_idx,
acc.session_id,
super::replaycoordinator::aof_replay_context::ReplayOperation::Chunk(Box::new(acc.clone())),
);
if buffered {
return Ok(false);
}
replay_op_dispatch_chunk(
processor,
virtual_sublog_idx,
&acc,
as_replica,
log_address_sequence_number,
target,
)
.await?;
Ok(false)
}
pub async fn replay_op_dispatch_chunk<D: Device>(
processor: &AofProcessor,
virtual_sublog_idx: usize,
acc: &ChunkedAccumulator,
_as_replica: bool,
log_address_sequence_number: i64,
target: &ReplayTarget<'_, '_, D>,
) -> Result<(), AofReplayError> {
if (processor.append_only_file().log().size() > 1
|| processor.append_only_file().virtual_sublog_count() > 1)
&& let Some(manager) = processor.read_consistency_manager()
{
let sequence_number = if acc.sequence_number != 0 {
acc.sequence_number
} else {
log_address_sequence_number
};
manager.update_virtual_sublog_key_sequence_number(
virtual_sublog_idx,
acc.key_hash,
sequence_number,
);
}
if !processor.begin_replay_op(processor.should_skip_record_chunk(
virtual_sublog_idx,
acc,
target.store_version,
)) {
return Ok(());
}
replay_chunk(processor, virtual_sublog_idx, acc.clone(), target).await
}
pub async fn replay_chunk<D: Device>(
processor: &AofProcessor,
_virtual_sublog_idx: usize,
acc: ChunkedAccumulator,
target: &ReplayTarget<'_, '_, D>,
) -> Result<(), AofReplayError> {
let prepared = PreparedParameters {
key: acc.key_span().to_vec(),
key_hash: acc.key_hash,
payload: build_payload(&acc),
};
processor
.replay_op(acc.op_type, prepared, false, target)
.await
}
fn build_payload(acc: &ChunkedAccumulator) -> Vec<u8> {
let mut payload = Vec::new();
match acc.op_type {
AofEntryType::ObjectStoreUpsert | AofEntryType::UnifiedStoreObjectUpsert => {
let value: Vec<u8> = acc.get_value_sequence().concat();
payload.extend_from_slice(&(value.len() as u32).to_le_bytes());
payload.extend_from_slice(&value);
payload.extend_from_slice(acc.input_span());
}
AofEntryType::StoreUpsert | AofEntryType::UnifiedStoreStringUpsert => {
payload.extend_from_slice(&(acc.value_span().len() as u32).to_le_bytes());
payload.extend_from_slice(acc.value_span());
payload.extend_from_slice(acc.input_span());
}
_ => {
payload.extend_from_slice(acc.input_span());
}
}
payload
}
#[cfg(test)]
mod tests {
use super::*;
use crate::aof::aof_header::AofChunkHeader;
#[test]
fn payload_shape_per_op_type() {
let chunk_header = AofChunkHeader {
overflow_key_length: 1,
overflow_value_length: 2,
input_length: 0,
object_id: 1,
key_hash: 1,
};
let mut acc = ChunkedAccumulator::new(AofEntryType::StoreUpsert, &chunk_header);
assert!(acc.feed(b"k"));
assert!(acc.feed(b"vv"));
let payload = build_payload(&acc);
assert_eq!(payload, vec![2, 0, 0, 0, b'v', b'v']);
let del_header = AofChunkHeader {
overflow_key_length: 1,
overflow_value_length: 0,
input_length: 0,
object_id: 2,
key_hash: 1,
};
let mut del = ChunkedAccumulator::new(AofEntryType::StoreDelete, &del_header);
assert!(del.feed(b"k"));
assert!(build_payload(&del).is_empty());
}
}