use std::ops::Range;
use anyhow::Result;
use velo::EventManager;
use crate::BlockId;
use kvbm_common::LogicalLayoutHandle;
use kvbm_physical::transfer::TransferCompleteNotification;
use super::CollectiveOps;
pub struct StubCollectiveOps {
events: EventManager,
rank: usize,
world_size: usize,
}
impl StubCollectiveOps {
pub fn new(events: EventManager, rank: usize, world_size: usize) -> Self {
Self {
events,
rank,
world_size,
}
}
pub fn single_worker(events: EventManager) -> Self {
Self::new(events, 0, 1)
}
}
impl CollectiveOps for StubCollectiveOps {
fn broadcast(
&self,
src: LogicalLayoutHandle,
dst: LogicalLayoutHandle,
src_block_ids: &[BlockId],
dst_block_ids: &[BlockId],
layer_range: Option<Range<usize>>,
) -> Result<TransferCompleteNotification> {
tracing::warn!(
rank = self.rank,
world_size = self.world_size,
?src,
?dst,
num_src_blocks = src_block_ids.len(),
num_dst_blocks = dst_block_ids.len(),
?layer_range,
"StubCollectiveOps::broadcast called - completing immediately without actual transfer"
);
let event = self.events.new_event()?;
let handle = event.handle();
event.trigger()?;
let awaiter = self.events.awaiter(handle)?;
Ok(TransferCompleteNotification::from_awaiter(awaiter))
}
fn rank(&self) -> usize {
self.rank
}
fn world_size(&self) -> usize {
self.world_size
}
}