kvbm_engine/worker/velo/
service.rs1use 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#[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 pub fn worker(&self) -> &Arc<DirectWorker> {
46 &self.worker
47 }
48
49 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 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 let handler = Handler::unary_handler_async("kvbm.worker.local_transfer", move |ctx| {
70 let worker = worker.clone();
71
72 async move {
73 let message: LocalTransferMessage = serde_json::from_slice(&ctx.payload)?;
75
76 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 notification.await?;
93
94 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 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 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 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 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 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 let metadata: Vec<SerializedLayout> = message
210 .metadata
211 .into_iter()
212 .map(SerializedLayout::from_bytes)
213 .collect();
214
215 worker.connect_remote(message.instance_id, metadata)?;
217
218 Ok(Some(Bytes::new()))
220 })
221 .build();
222
223 self.messenger.register_handler(handler)?;
224 Ok(())
225 }
226
227 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 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 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 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 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 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 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 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}