Skip to main content

kvbm_engine/worker/velo/
service.rs

1// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4use kvbm_physical::manager::SerializedLayout;
5
6use super::{
7    Arc, ConnectRemoteMessage, DirectWorker, ExecuteRemoteOnboardForInstanceMessage,
8    LocalTransferMessage, ObjectGetBlocksMessage, ObjectHasBlocksMessage, ObjectHasBlocksResponse,
9    ObjectPutBlocksMessage, ObjectPutGetBlocksResponse, RemoteOffloadMessage, RemoteOnboardMessage,
10    Result, TransferOptions, WorkerTransfers,
11};
12use crate::object::ObjectBlockOps;
13
14use bytes::Bytes;
15use derive_builder::Builder;
16
17use ::velo::{Handler, Messenger};
18
19/// Builder for VeloWorkerService - provides flexibility in construction.
20///
21/// Use this builder when you need to:
22/// - Pass a pre-built DirectWorker (when caller manages layout registration)
23/// - Pass a pre-built TransferManager (service creates DirectWorker)
24/// - Have more control over worker configuration
25#[derive(Builder)]
26#[builder(pattern = "owned")]
27pub struct VeloWorkerService {
28    messenger: Arc<Messenger>,
29    worker: Arc<DirectWorker>,
30}
31
32impl VeloWorkerService {
33    pub fn new(messenger: Arc<Messenger>, worker: Arc<DirectWorker>) -> Result<Self> {
34        let service = Self { messenger, worker };
35        service.register_handlers()?;
36        Ok(service)
37    }
38
39    /// Access the underlying DirectWorker.
40    ///
41    /// This is useful for:
42    /// - Registering additional layouts after service creation
43    /// - Exporting metadata for handshake
44    /// - Accessing the TransferManager
45    pub fn worker(&self) -> &Arc<DirectWorker> {
46        &self.worker
47    }
48
49    /// Register all worker handlers with Nova
50    fn register_handlers(&self) -> Result<()> {
51        self.register_local_transfer_handler()?;
52        self.register_remote_onboard_handler()?;
53        self.register_remote_offload_handler()?;
54        self.register_import_metadata_handler()?;
55        self.register_export_metadata_handler()?;
56        self.register_connect_remote_handler()?;
57        self.register_execute_remote_onboard_for_instance_handler()?;
58        // Object storage handlers
59        self.register_object_has_blocks_handler()?;
60        self.register_object_put_blocks_handler()?;
61        self.register_object_get_blocks_handler()?;
62        Ok(())
63    }
64
65    fn register_local_transfer_handler(&self) -> Result<()> {
66        let worker = self.worker.clone();
67
68        // Use unary_handler_async for explicit response (client waits for transfer completion)
69        let handler = Handler::unary_handler_async("kvbm.worker.local_transfer", move |ctx| {
70            let worker = worker.clone();
71
72            async move {
73                // Deserialize the message
74                let message: LocalTransferMessage = serde_json::from_slice(&ctx.payload)?;
75
76                // Convert options and resolve bounce buffer if present
77                let bounce_buffer_parts = message.options.bounce_buffer_parts();
78                let mut options: TransferOptions = message.options.into();
79                if let Some((handle, block_ids)) = bounce_buffer_parts {
80                    options.bounce_buffer = Some(worker.create_bounce_buffer(handle, block_ids)?);
81                }
82
83                let notification = worker.execute_local_transfer(
84                    message.src,
85                    message.dst,
86                    Arc::from(message.src_block_ids),
87                    Arc::from(message.dst_block_ids),
88                    options,
89                )?;
90
91                // Await the transfer completion
92                notification.await?;
93
94                // Return empty response to signal success
95                Ok(Some(Bytes::new()))
96            }
97        })
98        .build();
99
100        self.messenger.register_handler(handler)?;
101        Ok(())
102    }
103
104    fn register_remote_onboard_handler(&self) -> Result<()> {
105        let worker = self.worker.clone();
106
107        // Use unary_handler_async for explicit response (works with unary client)
108        let handler = Handler::unary_handler_async("kvbm.worker.remote_onboard", move |ctx| {
109            let worker = worker.clone();
110
111            async move {
112                let message: RemoteOnboardMessage = serde_json::from_slice(&ctx.payload)?;
113
114                // Convert options and resolve bounce buffer if present
115                let bounce_buffer_parts = message.options.bounce_buffer_parts();
116                let mut options: TransferOptions = message.options.into();
117                if let Some((handle, block_ids)) = bounce_buffer_parts {
118                    options.bounce_buffer = Some(worker.create_bounce_buffer(handle, block_ids)?);
119                }
120
121                let notification = worker.execute_remote_onboard(
122                    message.src,
123                    message.dst,
124                    Arc::from(message.dst_block_ids),
125                    options,
126                )?;
127
128                notification.await?;
129
130                Ok(Some(Bytes::new()))
131            }
132        })
133        .build();
134
135        self.messenger.register_handler(handler)?;
136        Ok(())
137    }
138
139    fn register_remote_offload_handler(&self) -> Result<()> {
140        let worker = self.worker.clone();
141
142        // Use unary_handler_async for explicit response (works with unary client)
143        let handler = Handler::unary_handler_async("kvbm.worker.remote_offload", move |ctx| {
144            let worker = worker.clone();
145
146            async move {
147                let message: RemoteOffloadMessage = serde_json::from_slice(&ctx.payload)?;
148
149                // Convert options and resolve bounce buffer if present
150                let bounce_buffer_parts = message.options.bounce_buffer_parts();
151                let mut options: TransferOptions = message.options.into();
152                if let Some((handle, block_ids)) = bounce_buffer_parts {
153                    options.bounce_buffer = Some(worker.create_bounce_buffer(handle, block_ids)?);
154                }
155
156                let notification = worker.execute_remote_offload(
157                    message.src,
158                    Arc::from(message.src_block_ids),
159                    message.dst,
160                    options,
161                )?;
162
163                notification.await?;
164
165                Ok(Some(Bytes::new()))
166            }
167        })
168        .build();
169
170        self.messenger.register_handler(handler)?;
171        Ok(())
172    }
173
174    fn register_import_metadata_handler(&self) -> Result<()> {
175        let worker = self.worker.clone();
176
177        let handler = Handler::unary_handler("kvbm.worker.import_metadata", move |ctx| {
178            let metadata = SerializedLayout::from_bytes(ctx.payload.to_vec());
179            let handles = worker.import_metadata(metadata)?;
180            Ok(Some(Bytes::from(serde_json::to_vec(&handles)?)))
181        })
182        .build();
183
184        self.messenger.register_handler(handler)?;
185        Ok(())
186    }
187
188    fn register_export_metadata_handler(&self) -> Result<()> {
189        let worker = self.worker.clone();
190
191        let handler = Handler::unary_handler("kvbm.worker.export_metadata", move |_ctx| {
192            let response = worker.export_metadata()?;
193            Ok(Some(Bytes::from(response.as_bytes().to_vec())))
194        })
195        .build();
196
197        self.messenger.register_handler(handler)?;
198        Ok(())
199    }
200
201    /// Register handler for connect_remote - stores remote instance metadata in local worker
202    fn register_connect_remote_handler(&self) -> Result<()> {
203        let worker = self.worker.clone();
204
205        let handler = Handler::unary_handler("kvbm.worker.connect_remote", move |ctx| {
206            let message: ConnectRemoteMessage = serde_json::from_slice(&ctx.payload)?;
207
208            // Deserialize metadata (SerializedLayout stored as raw bytes)
209            let metadata: Vec<SerializedLayout> = message
210                .metadata
211                .into_iter()
212                .map(SerializedLayout::from_bytes)
213                .collect();
214
215            // Call DirectWorker.connect_remote()
216            worker.connect_remote(message.instance_id, metadata)?;
217
218            // Return empty response to signal success
219            Ok(Some(Bytes::new()))
220        })
221        .build();
222
223        self.messenger.register_handler(handler)?;
224        Ok(())
225    }
226
227    /// Register handler for execute_remote_onboard_for_instance - pulls from remote using instance ID
228    fn register_execute_remote_onboard_for_instance_handler(&self) -> Result<()> {
229        let worker = self.worker.clone();
230
231        let handler =
232            Handler::unary_handler_async("kvbm.worker.remote_onboard_for_instance", move |ctx| {
233                let worker = worker.clone();
234                async move {
235                    let message: ExecuteRemoteOnboardForInstanceMessage =
236                        serde_json::from_slice(&ctx.payload)?;
237
238                    // Convert options and resolve bounce buffer if present
239                    let bounce_buffer_parts = message.options.bounce_buffer_parts();
240                    let mut options: TransferOptions = message.options.into();
241                    if let Some((handle, block_ids)) = bounce_buffer_parts {
242                        options.bounce_buffer =
243                            Some(worker.create_bounce_buffer(handle, block_ids)?);
244                    }
245
246                    let notification = worker.execute_remote_onboard_for_instance(
247                        message.instance_id,
248                        message.remote_logical_type,
249                        message.src_block_ids,
250                        message.dst,
251                        Arc::from(message.dst_block_ids),
252                        options,
253                    )?;
254
255                    notification.await?;
256                    Ok(Some(Bytes::new()))
257                }
258            })
259            .build();
260
261        self.messenger.register_handler(handler)?;
262        Ok(())
263    }
264
265    // ========================================================================
266    // Object Storage Handlers
267    // ========================================================================
268
269    /// Register handler for object_has_blocks - check if blocks exist in object storage
270    fn register_object_has_blocks_handler(&self) -> Result<()> {
271        let worker = self.worker.clone();
272
273        let handler = Handler::unary_handler_async("kvbm.worker.object_has_blocks", move |ctx| {
274            let worker = worker.clone();
275
276            async move {
277                let message: ObjectHasBlocksMessage = serde_json::from_slice(&ctx.payload)?;
278
279                // Call DirectWorker's ObjectBlockOps implementation
280                let results = worker.has_blocks(message.keys).await;
281
282                let response = ObjectHasBlocksResponse { results };
283                Ok(Some(Bytes::from(serde_json::to_vec(&response)?)))
284            }
285        })
286        .build();
287
288        self.messenger.register_handler(handler)?;
289        Ok(())
290    }
291
292    /// Register handler for object_put_blocks - upload blocks to object storage
293    fn register_object_put_blocks_handler(&self) -> Result<()> {
294        let worker = self.worker.clone();
295
296        let handler = Handler::unary_handler_async("kvbm.worker.object_put_blocks", move |ctx| {
297            let worker = worker.clone();
298
299            async move {
300                let message: ObjectPutBlocksMessage = serde_json::from_slice(&ctx.payload)?;
301
302                // Call DirectWorker's ObjectBlockOps implementation
303                // DirectWorker resolves logical handle to physical layout internally
304                let results = worker
305                    .put_blocks(message.keys, message.layout, message.block_ids)
306                    .await;
307
308                let response = ObjectPutGetBlocksResponse::from_results(results);
309                Ok(Some(Bytes::from(serde_json::to_vec(&response)?)))
310            }
311        })
312        .build();
313
314        self.messenger.register_handler(handler)?;
315        Ok(())
316    }
317
318    /// Register handler for object_get_blocks - download blocks from object storage
319    fn register_object_get_blocks_handler(&self) -> Result<()> {
320        let worker = self.worker.clone();
321
322        let handler = Handler::unary_handler_async("kvbm.worker.object_get_blocks", move |ctx| {
323            let worker = worker.clone();
324
325            async move {
326                let message: ObjectGetBlocksMessage = serde_json::from_slice(&ctx.payload)?;
327
328                // Call DirectWorker's ObjectBlockOps implementation
329                // DirectWorker resolves logical handle to physical layout internally
330                let results = worker
331                    .get_blocks(message.keys, message.layout, message.block_ids)
332                    .await;
333
334                let response = ObjectPutGetBlocksResponse::from_results(results);
335                Ok(Some(Bytes::from(serde_json::to_vec(&response)?)))
336            }
337        })
338        .build();
339
340        self.messenger.register_handler(handler)?;
341        Ok(())
342    }
343}