use super::*;
use crate::object::ObjectBlockOps;
use futures::future::BoxFuture;
use parking_lot::RwLock;
use std::collections::HashSet;
use std::sync::OnceLock;
#[derive(Clone)]
pub struct VeloWorkerClient {
messenger: Arc<Messenger>,
remote: InstanceId,
g1_handle: Arc<OnceLock<LayoutHandle>>,
g2_handle: Arc<OnceLock<LayoutHandle>>,
g3_handle: Arc<OnceLock<LayoutHandle>>,
connected_instances: Arc<RwLock<HashSet<InstanceId>>>,
}
impl WorkerTransfers for VeloWorkerClient {
fn execute_local_transfer(
&self,
src: LogicalLayoutHandle,
dst: LogicalLayoutHandle,
src_block_ids: Arc<[BlockId]>,
dst_block_ids: Arc<[BlockId]>,
options: TransferOptions,
) -> Result<TransferCompleteNotification> {
let event = self.messenger.events().new_event()?;
let awaiter = self.messenger.events().awaiter(event.handle())?;
let options = SerializableTransferOptions {
layer_range: options.layer_range,
nixl_write_notification: options.nixl_write_notification,
bounce_buffer_handle: None,
bounce_buffer_block_ids: None,
};
let message = LocalTransferMessage {
src,
dst,
src_block_ids: src_block_ids.to_vec(),
dst_block_ids: dst_block_ids.to_vec(),
options,
};
let bytes = Bytes::from(serde_json::to_vec(&message)?);
let nova = self.messenger.clone();
let remote_instance = self.remote;
self.messenger.tracker().spawn_on(
async move {
let result = nova
.unary("kvbm.worker.local_transfer")?
.raw_payload(bytes)
.instance(remote_instance)
.send()
.await;
match result {
Ok(_) => event.trigger(),
Err(e) => event.poison(e.to_string()),
}
},
self.messenger.runtime(),
);
Ok(TransferCompleteNotification::from_awaiter(awaiter))
}
fn execute_remote_onboard(
&self,
src: RemoteDescriptor,
dst: LogicalLayoutHandle,
dst_block_ids: Arc<[BlockId]>,
options: TransferOptions,
) -> Result<TransferCompleteNotification> {
let event = self.messenger.events().new_event()?;
let awaiter = self.messenger.events().awaiter(event.handle())?;
let options = SerializableTransferOptions {
layer_range: options.layer_range,
nixl_write_notification: options.nixl_write_notification,
bounce_buffer_handle: None,
bounce_buffer_block_ids: None,
};
let message = RemoteOnboardMessage {
src,
dst,
dst_block_ids: dst_block_ids.to_vec(),
options,
};
let bytes = Bytes::from(serde_json::to_vec(&message)?);
let nova = self.messenger.clone();
let remote_instance = self.remote;
self.messenger.tracker().spawn_on(
async move {
let result = nova
.unary("kvbm.worker.remote_onboard")?
.raw_payload(bytes)
.instance(remote_instance)
.send()
.await;
match result {
Ok(_) => event.trigger(),
Err(e) => event.poison(e.to_string()),
}
},
self.messenger.runtime(),
);
Ok(TransferCompleteNotification::from_awaiter(awaiter))
}
fn execute_remote_offload(
&self,
src: LogicalLayoutHandle,
src_block_ids: Arc<[BlockId]>,
dst: RemoteDescriptor,
options: TransferOptions,
) -> Result<TransferCompleteNotification> {
let event = self.messenger.events().new_event()?;
let awaiter = self.messenger.events().awaiter(event.handle())?;
let options = SerializableTransferOptions {
layer_range: options.layer_range,
nixl_write_notification: options.nixl_write_notification,
bounce_buffer_handle: None,
bounce_buffer_block_ids: None,
};
let message = RemoteOffloadMessage {
src,
dst,
src_block_ids: src_block_ids.to_vec(),
options,
};
let bytes = Bytes::from(serde_json::to_vec(&message)?);
let nova = self.messenger.clone();
let remote_instance = self.remote;
self.messenger.tracker().spawn_on(
async move {
let result = nova
.unary("kvbm.worker.remote_offload")?
.raw_payload(bytes)
.instance(remote_instance)
.send()
.await;
match result {
Ok(_) => event.trigger(),
Err(e) => event.poison(e.to_string()),
}
},
self.messenger.runtime(),
);
Ok(TransferCompleteNotification::from_awaiter(awaiter))
}
fn connect_remote(
&self,
instance_id: InstanceId,
metadata: Vec<SerializedLayout>,
) -> Result<ConnectRemoteResponse> {
let serialized_metadata: Vec<Vec<u8>> =
metadata.iter().map(|m| m.as_bytes().to_vec()).collect();
let message = ConnectRemoteMessage {
instance_id,
metadata: serialized_metadata,
};
let bytes = Bytes::from(serde_json::to_vec(&message)?);
let event = self.messenger.events().new_event()?;
let awaiter = self.messenger.events().awaiter(event.handle())?;
let nova = self.messenger.clone();
let remote_instance = self.remote;
let connected = self.connected_instances.clone();
let target_instance = instance_id;
self.messenger.tracker().spawn_on(
async move {
let result = nova
.unary("kvbm.worker.connect_remote")?
.raw_payload(bytes)
.instance(remote_instance)
.send()
.await;
match result {
Ok(_) => {
connected.write().insert(target_instance);
event.trigger()
}
Err(e) => event.poison(e.to_string()),
}
},
self.messenger.runtime(),
);
Ok(ConnectRemoteResponse::from_awaiter(awaiter))
}
fn has_remote_metadata(&self, instance_id: InstanceId) -> bool {
self.connected_instances.read().contains(&instance_id)
}
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: TransferOptions,
) -> Result<TransferCompleteNotification> {
let message = ExecuteRemoteOnboardForInstanceMessage {
instance_id,
remote_logical_type,
src_block_ids,
dst,
dst_block_ids: dst_block_ids.to_vec(),
options: SerializableTransferOptions::from(options),
};
let bytes = Bytes::from(serde_json::to_vec(&message)?);
let event = self.messenger.events().new_event()?;
let awaiter = self.messenger.events().awaiter(event.handle())?;
let nova = self.messenger.clone();
let remote_instance = self.remote;
self.messenger.tracker().spawn_on(
async move {
let result = nova
.unary("kvbm.worker.remote_onboard_for_instance")?
.raw_payload(bytes)
.instance(remote_instance)
.send()
.await;
match result {
Ok(_) => event.trigger(),
Err(e) => event.poison(e.to_string()),
}
},
self.messenger.runtime(),
);
Ok(TransferCompleteNotification::from_awaiter(awaiter))
}
}
impl Worker for VeloWorkerClient {
fn g1_handle(&self) -> Option<LayoutHandle> {
self.g1_handle.get().copied()
}
fn g2_handle(&self) -> Option<LayoutHandle> {
self.g2_handle.get().copied()
}
fn g3_handle(&self) -> Option<LayoutHandle> {
self.g3_handle.get().copied()
}
fn export_metadata(&self) -> Result<SerializedLayoutResponse> {
let unary_result = self
.messenger
.unary("kvbm.worker.export_metadata")?
.instance(self.remote)
.send();
let future = async move {
let bytes = unary_result.await?;
Ok(SerializedLayout::from_bytes(bytes.to_vec()))
};
Ok(SerializedLayoutResponse::from_boxed(Box::pin(future)))
}
fn import_metadata(&self, metadata: SerializedLayout) -> Result<ImportMetadataResponse> {
let unary_result = self
.messenger
.unary("kvbm.worker.import_metadata")?
.raw_payload(Bytes::from(metadata.as_bytes().to_vec()))
.instance(self.remote)
.send();
let future = async move {
let bytes = unary_result.await?;
serde_json::from_slice(&bytes).map_err(|e| {
anyhow::anyhow!("Failed to deserialize import_metadata response: {}", e)
})
};
Ok(ImportMetadataResponse::from_boxed(Box::pin(future)))
}
}
impl VeloWorkerClient {
pub fn new(messenger: Arc<Messenger>, remote: InstanceId) -> Self {
Self {
messenger,
remote,
g1_handle: Arc::new(OnceLock::new()),
g2_handle: Arc::new(OnceLock::new()),
g3_handle: Arc::new(OnceLock::new()),
connected_instances: Arc::new(RwLock::new(HashSet::new())),
}
}
pub fn configure_layout_handles(&self, metadata: &SerializedLayout) -> Result<()> {
let unpacked = metadata.unpack()?;
for desc in &unpacked.layouts {
match desc.logical_type {
LogicalLayoutHandle::G1 => {
self.g1_handle.set(desc.handle).ok();
}
LogicalLayoutHandle::G2 => {
self.g2_handle.set(desc.handle).ok();
}
LogicalLayoutHandle::G3 => {
self.g3_handle.set(desc.handle).ok();
}
_ => {}
}
}
Ok(())
}
pub fn get_layout_config(&self) -> Result<::velo::TypedUnaryResult<LayoutConfig>> {
let instance = self.remote;
let awaiter = self
.messenger
.typed_unary::<LayoutConfig>("kvbm.worker.get_layout_config")?
.instance(instance)
.send();
Ok(awaiter)
}
pub fn configure_layouts(
&self,
config: LeaderLayoutConfig,
) -> Result<::velo::TypedUnaryResult<WorkerLayoutResponse>> {
let instance = self.remote;
let awaiter = self
.messenger
.typed_unary::<WorkerLayoutResponse>("kvbm.worker.configure_layouts")?
.payload(config)?
.instance(instance)
.send();
Ok(awaiter)
}
}
impl ObjectBlockOps for VeloWorkerClient {
fn has_blocks(
&self,
keys: Vec<SequenceHash>,
) -> BoxFuture<'static, Vec<(SequenceHash, Option<usize>)>> {
let message = ObjectHasBlocksMessage { keys: keys.clone() };
let bytes = match serde_json::to_vec(&message) {
Ok(b) => Bytes::from(b),
Err(_) => {
return Box::pin(async move { keys.into_iter().map(|k| (k, None)).collect() });
}
};
let nova = self.messenger.clone();
let remote = self.remote;
Box::pin(async move {
let result = nova
.unary("kvbm.worker.object_has_blocks")
.ok()
.map(|u| u.raw_payload(bytes).instance(remote).send());
match result {
Some(unary_result) => match unary_result.await {
Ok(response_bytes) => {
match serde_json::from_slice::<ObjectHasBlocksResponse>(&response_bytes) {
Ok(response) => response.results,
Err(_) => keys.into_iter().map(|k| (k, None)).collect(),
}
}
Err(_) => keys.into_iter().map(|k| (k, None)).collect(),
},
None => keys.into_iter().map(|k| (k, None)).collect(),
}
})
}
fn put_blocks(
&self,
keys: Vec<SequenceHash>,
src_layout: LogicalLayoutHandle,
block_ids: Vec<BlockId>,
) -> BoxFuture<'static, Vec<Result<SequenceHash, SequenceHash>>> {
let message = ObjectPutBlocksMessage {
keys: keys.clone(),
layout: src_layout,
block_ids,
};
let bytes = match serde_json::to_vec(&message) {
Ok(b) => Bytes::from(b),
Err(_) => return Box::pin(async move { keys.into_iter().map(Err).collect() }),
};
let nova = self.messenger.clone();
let remote = self.remote;
Box::pin(async move {
let result = nova
.unary("kvbm.worker.object_put_blocks")
.ok()
.map(|u| u.raw_payload(bytes).instance(remote).send());
match result {
Some(unary_result) => match unary_result.await {
Ok(response_bytes) => {
match serde_json::from_slice::<ObjectPutGetBlocksResponse>(&response_bytes)
{
Ok(response) => response.into_results(),
Err(_) => keys.into_iter().map(Err).collect(),
}
}
Err(_) => keys.into_iter().map(Err).collect(),
},
None => keys.into_iter().map(Err).collect(),
}
})
}
fn get_blocks(
&self,
keys: Vec<SequenceHash>,
dst_layout: LogicalLayoutHandle,
block_ids: Vec<BlockId>,
) -> BoxFuture<'static, Vec<Result<SequenceHash, SequenceHash>>> {
let message = ObjectGetBlocksMessage {
keys: keys.clone(),
layout: dst_layout,
block_ids,
};
let bytes = match serde_json::to_vec(&message) {
Ok(b) => Bytes::from(b),
Err(_) => return Box::pin(async move { keys.into_iter().map(Err).collect() }),
};
let nova = self.messenger.clone();
let remote = self.remote;
Box::pin(async move {
let result = nova
.unary("kvbm.worker.object_get_blocks")
.ok()
.map(|u| u.raw_payload(bytes).instance(remote).send());
match result {
Some(unary_result) => match unary_result.await {
Ok(response_bytes) => {
match serde_json::from_slice::<ObjectPutGetBlocksResponse>(&response_bytes)
{
Ok(response) => response.into_results(),
Err(_) => keys.into_iter().map(Err).collect(),
}
}
Err(_) => keys.into_iter().map(Err).collect(),
},
None => keys.into_iter().map(Err).collect(),
}
})
}
}