use std::sync::Arc;
use anyhow::Result;
use crate::{BlockId, G2, G3, worker::group::ParallelWorkers};
use kvbm_common::LogicalLayoutHandle;
use kvbm_logical::{blocks::ImmutableBlock, manager::BlockManager};
use kvbm_physical::transfer::TransferOptions;
use super::blocks::BlockHolder;
pub struct StagingResult {
pub new_g2_blocks: Vec<ImmutableBlock<G2>>,
}
pub async fn stage_g3_to_g2(
g3_blocks: &BlockHolder<G3>,
g2_manager: &BlockManager<G2>,
parallel_worker: &dyn ParallelWorkers,
) -> Result<StagingResult> {
if g3_blocks.is_empty() {
return Ok(StagingResult {
new_g2_blocks: Vec::new(),
});
}
let src_ids: Vec<BlockId> = g3_blocks.blocks().iter().map(|b| b.block_id()).collect();
let dst_blocks = g2_manager
.allocate_blocks(src_ids.len())
.ok_or_else(|| anyhow::anyhow!("Failed to allocate G2 blocks"))?;
let dst_ids: Vec<BlockId> = dst_blocks.iter().map(|b| b.block_id()).collect();
let notification = parallel_worker.execute_local_transfer(
LogicalLayoutHandle::G3,
LogicalLayoutHandle::G2,
Arc::from(src_ids),
Arc::from(dst_ids),
TransferOptions::default(),
)?;
notification.await?;
let new_g2_blocks: Vec<ImmutableBlock<G2>> = dst_blocks
.into_iter()
.zip(g3_blocks.blocks().iter())
.map(|(dst, src)| {
let complete = dst
.stage(src.sequence_hash(), g2_manager.block_size())
.expect("block size mismatch");
g2_manager.register_block(complete)
})
.collect();
Ok(StagingResult { new_g2_blocks })
}