use super::*;
use crate::KvbmRuntime;
use crate::collectives::CollectiveOps;
use anyhow::Result;
use std::sync::Arc;
#[allow(dead_code)]
pub struct ReplicatedDataWorker {
inner: Arc<PhysicalWorker>,
runtime: Arc<KvbmRuntime>,
collective: Arc<dyn CollectiveOps>,
}
#[allow(dead_code)]
impl ReplicatedDataWorker {
pub fn new(
worker: Arc<PhysicalWorker>, runtime: Arc<KvbmRuntime>,
collective: Arc<dyn CollectiveOps>,
) -> Self {
Self {
inner: worker,
runtime,
collective,
}
}
pub fn inner(&self) -> &PhysicalWorker {
&self.inner
}
pub fn rank(&self) -> usize {
self.inner.rank().expect("Worker must have a rank")
}
#[expect(unused_variables)]
fn broadcast(
&self,
xfer_completion: TransferCompleteNotification,
dst: LogicalLayoutHandle,
dst_block_ids: Arc<[BlockId]>,
options: kvbm_physical::transfer::TransferOptions,
) -> Result<TransferCompleteNotification> {
unimplemented!()
}
}
impl WorkerTransfers for ReplicatedDataWorker {
fn execute_local_transfer(
&self,
src: LogicalLayoutHandle,
dst: LogicalLayoutHandle,
src_block_ids: Arc<[BlockId]>,
dst_block_ids: Arc<[BlockId]>,
options: kvbm_physical::transfer::TransferOptions,
) -> Result<TransferCompleteNotification> {
let is_rank0 = self.rank() == 0;
let use_bcast = dst == LogicalLayoutHandle::G1;
if src == LogicalLayoutHandle::G1 && dst == LogicalLayoutHandle::G1 {
return self.inner.execute_local_transfer(
src,
dst,
src_block_ids,
dst_block_ids.clone(),
options,
);
}
if !is_rank0 && !use_bcast {
return Ok(TransferCompleteNotification::completed());
} else if is_rank0 {
let xfer_completion = self.inner.execute_local_transfer(
src,
dst,
src_block_ids,
dst_block_ids.clone(),
options.clone(),
)?;
if use_bcast {
self.broadcast(xfer_completion, dst, dst_block_ids, options)
} else {
Ok(xfer_completion)
}
} else {
let xfer_completion = TransferCompleteNotification::completed();
self.broadcast(xfer_completion, dst, dst_block_ids, options)
}
}
#[expect(unused_variables)]
fn execute_remote_onboard(
&self,
src: RemoteDescriptor,
dst: LogicalLayoutHandle,
dst_block_ids: Arc<[BlockId]>,
options: kvbm_physical::transfer::TransferOptions,
) -> Result<TransferCompleteNotification> {
unimplemented!()
}
#[expect(unused_variables)]
fn execute_remote_offload(
&self,
src: LogicalLayoutHandle,
src_block_ids: Arc<[BlockId]>,
dst: RemoteDescriptor,
options: kvbm_physical::transfer::TransferOptions,
) -> Result<TransferCompleteNotification> {
unimplemented!()
}
fn connect_remote(
&self,
instance_id: InstanceId,
metadata: Vec<SerializedLayout>,
) -> Result<ConnectRemoteResponse> {
self.inner.connect_remote(instance_id, metadata)
}
fn has_remote_metadata(&self, instance_id: InstanceId) -> bool {
self.inner.has_remote_metadata(instance_id)
}
#[expect(unused_variables)]
fn execute_remote_onboard_for_instance(
&self,
instance_id: InstanceId,
remote_logical_type: LogicalLayoutHandle,
src_block_ids: Vec<BlockId>,
dst: LogicalLayoutHandle,
dst_block_ids: Arc<[BlockId]>,
options: kvbm_physical::transfer::TransferOptions,
) -> Result<TransferCompleteNotification> {
unimplemented!()
}
}