1use std::sync::Arc;
33
34use anyhow::Result;
35use dashmap::DashMap;
36use tokio::sync::mpsc;
37use tokio::task::JoinHandle;
38use uuid::Uuid;
39
40use crate::leader::InstanceLeader;
41use crate::object::ObjectBlockOps;
42use crate::worker::RemoteDescriptor;
43use crate::{BlockId, G1, G2, G3, SequenceHash};
44use kvbm_common::LogicalLayoutHandle;
45use kvbm_logical::blocks::{BlockMetadata, BlockRegistry, WeakBlock};
46use kvbm_logical::manager::BlockManager;
47use kvbm_physical::transfer::{PhysicalLayout, TransferOptions};
48
49use super::handle::{TransferHandle, TransferId, TransferState};
50use super::pipeline::{
51 ChainOutput, ChainOutputRx, ObjectPipeline, ObjectPipelineConfig, Pipeline, PipelineConfig,
52 PipelineInput,
53};
54use super::queue::CancellableQueue;
55use super::settlement::{
56 PipelineLane, PipelineSettlementTracker, SettlementError, SettlementTarget, SettlementToken,
57 SettlementWaiter, wait_for_settlement,
58};
59use super::source::SourceBlocks;
60
61#[allow(dead_code)]
79pub struct OffloadEngine {
80 engine_id: Uuid,
82 leader: Arc<InstanceLeader>,
84 registry: Arc<BlockRegistry>,
86 g1_to_g2: Option<Pipeline<G1, G2>>,
88 g2_to_g3: Option<Pipeline<G2, G3>>,
90 g2_to_g4: Option<ObjectPipeline<G2>>,
92 transfers: Arc<DashMap<TransferId, Arc<std::sync::Mutex<TransferState>>>>,
94 _chain_router_handle: Option<JoinHandle<()>>,
96 _remote_g4_offload_handle: Option<JoinHandle<()>>,
98}
99
100impl OffloadEngine {
101 pub fn builder(leader: Arc<InstanceLeader>) -> OffloadEngineBuilder {
103 OffloadEngineBuilder::new(leader)
104 }
105
106 pub fn enqueue_g1_to_g2(&self, blocks: impl Into<SourceBlocks<G1>>) -> Result<TransferHandle> {
110 let pipeline = self
111 .g1_to_g2
112 .as_ref()
113 .ok_or_else(|| anyhow::anyhow!("G1→G2 pipeline not configured"))?;
114
115 self.enqueue_to_pipeline(pipeline, blocks.into())
116 }
117
118 pub fn enqueue_g1_to_g2_with_precondition(
125 &self,
126 blocks: impl Into<SourceBlocks<G1>>,
127 precondition: Option<velo::EventHandle>,
128 ) -> Result<TransferHandle> {
129 let pipeline = self
130 .g1_to_g2
131 .as_ref()
132 .ok_or_else(|| anyhow::anyhow!("G1→G2 pipeline not configured"))?;
133
134 self.enqueue_to_pipeline_with_precondition(pipeline, blocks.into(), precondition)
135 }
136
137 pub fn enqueue_g2_to_g3(&self, blocks: impl Into<SourceBlocks<G2>>) -> Result<TransferHandle> {
141 let pipeline = self
142 .g2_to_g3
143 .as_ref()
144 .ok_or_else(|| anyhow::anyhow!("G2→G3 pipeline not configured"))?;
145
146 self.enqueue_to_pipeline(pipeline, blocks.into())
147 }
148
149 pub fn enqueue_g2_to_g4(&self, blocks: impl Into<SourceBlocks<G2>>) -> Result<TransferHandle> {
153 let pipeline = self
154 .g2_to_g4
155 .as_ref()
156 .ok_or_else(|| anyhow::anyhow!("G2→G4 pipeline not configured"))?;
157
158 self.enqueue_to_object_pipeline(pipeline, blocks.into())
159 }
160
161 fn create_transfer<T: BlockMetadata>(
163 &self,
164 source: &SourceBlocks<T>,
165 ) -> (
166 TransferId,
167 Arc<std::sync::Mutex<TransferState>>,
168 TransferHandle,
169 ) {
170 let input_block_ids = self.extract_block_ids(source);
171 let transfer_id = TransferId::new();
172 let (state, handle) = TransferState::new(transfer_id, input_block_ids);
173 let state = Arc::new(std::sync::Mutex::new(state));
174 self.transfers.insert(transfer_id, state.clone());
175 (transfer_id, state, handle)
176 }
177
178 fn enqueue_to_pipeline<Src: BlockMetadata, Dst: BlockMetadata>(
180 &self,
181 pipeline: &Pipeline<Src, Dst>,
182 source: SourceBlocks<Src>,
183 ) -> Result<TransferHandle> {
184 let (transfer_id, state, handle) = self.create_transfer(&source);
185 if !pipeline.enqueue(transfer_id, source, state) {
186 tracing::warn!("Transfer {} was cancelled before enqueueing", transfer_id);
187 }
188 Ok(handle)
189 }
190
191 fn enqueue_to_pipeline_with_precondition<Src: BlockMetadata, Dst: BlockMetadata>(
193 &self,
194 pipeline: &Pipeline<Src, Dst>,
195 source: SourceBlocks<Src>,
196 precondition: Option<velo::EventHandle>,
197 ) -> Result<TransferHandle> {
198 let (transfer_id, state, handle) = self.create_transfer(&source);
199 state.lock().unwrap().precondition = precondition;
200 if !pipeline.enqueue(transfer_id, source, state) {
201 tracing::warn!("Transfer {} was cancelled before enqueueing", transfer_id);
202 }
203 Ok(handle)
204 }
205
206 fn enqueue_to_object_pipeline(
208 &self,
209 pipeline: &ObjectPipeline<G2>,
210 source: SourceBlocks<G2>,
211 ) -> Result<TransferHandle> {
212 let (transfer_id, state, handle) = self.create_transfer(&source);
213 if !pipeline.enqueue(transfer_id, source, state) {
214 tracing::warn!("Transfer {} was cancelled before enqueueing", transfer_id);
215 }
216 Ok(handle)
217 }
218
219 fn extract_block_ids<T: BlockMetadata>(&self, source: &SourceBlocks<T>) -> Vec<BlockId> {
224 match source {
225 SourceBlocks::External(blocks) => blocks.iter().map(|b| b.block_id).collect(),
226 SourceBlocks::Strong(blocks) => blocks.iter().map(|b| b.block_id()).collect(),
227 SourceBlocks::Weak(_) => Vec::new(), }
229 }
230
231 pub fn release_transfer(&self, transfer_id: TransferId) {
236 self.transfers.remove(&transfer_id);
237 }
238
239 pub fn active_transfer_count(&self) -> usize {
241 self.transfers.len()
242 }
243
244 pub fn has_g1_to_g2(&self) -> bool {
246 self.g1_to_g2.is_some()
247 }
248
249 pub fn has_g2_to_g3(&self) -> bool {
251 self.g2_to_g3.is_some()
252 }
253
254 pub fn has_g2_to_g4(&self) -> bool {
256 self.g2_to_g4.is_some()
257 }
258
259 pub fn settlement_token(&self) -> SettlementToken {
261 let mut checkpoints = [None; 3];
262 for lane in PipelineLane::ALL {
263 if let Some((tracker, _)) = self.settlement_tracker(lane) {
264 checkpoints[lane.index()] = Some(tracker.snapshot().checkpoint());
265 }
266 }
267 SettlementToken {
268 engine_id: self.engine_id,
269 checkpoints,
270 }
271 }
272
273 pub async fn settle_after(
279 &self,
280 token: SettlementToken,
281 target: SettlementTarget,
282 ) -> Result<(), SettlementError> {
283 token.validate_engine(self.engine_id)?;
284 if target.is_empty() {
285 return Ok(());
286 }
287
288 let mut waiters = Vec::new();
289 for lane in PipelineLane::ALL {
290 let delta = target.completed_batches(lane);
291 if delta == 0 {
292 continue;
293 }
294
295 let (tracker, auto_chain) = self
296 .settlement_tracker(lane)
297 .ok_or(SettlementError::LaneUnavailable { lane })?;
298 if auto_chain {
299 return Err(SettlementError::UnsupportedAutoChain { lane });
300 }
301 let checkpoint = token.checkpoint(lane)?;
302
303 waiters.push(SettlementWaiter::new(lane, tracker, checkpoint, delta)?);
305 }
306 wait_for_settlement(waiters).await
307 }
308
309 fn settlement_tracker(&self, lane: PipelineLane) -> Option<(&PipelineSettlementTracker, bool)> {
310 match lane {
311 PipelineLane::G1ToG2 => self
312 .g1_to_g2
313 .as_ref()
314 .map(|pipeline| (&pipeline.settlement, pipeline.auto_chain())),
315 PipelineLane::G2ToG3 => self
316 .g2_to_g3
317 .as_ref()
318 .map(|pipeline| (&pipeline.settlement, pipeline.auto_chain())),
319 PipelineLane::G2ToG4 => self
320 .g2_to_g4
321 .as_ref()
322 .map(|pipeline| (&pipeline.settlement, false)),
323 }
324 }
325}
326
327pub struct OffloadEngineBuilder {
329 leader: Arc<InstanceLeader>,
330 registry: Option<Arc<BlockRegistry>>,
331 g1_manager: Option<Arc<BlockManager<G1>>>,
332 g2_manager: Option<Arc<BlockManager<G2>>>,
333 g3_manager: Option<Arc<BlockManager<G3>>>,
334 object_ops: Option<Arc<dyn ObjectBlockOps>>,
336 g2_physical_layout: Option<PhysicalLayout>,
338 g1_to_g2_config: Option<PipelineConfig<G1, G2>>,
339 g2_to_g3_config: Option<PipelineConfig<G2, G3>>,
340 g2_to_g4_config: Option<ObjectPipelineConfig<G2>>,
342 runtime: Option<tokio::runtime::Handle>,
344 enable_remote_g4: bool,
346}
347
348impl OffloadEngineBuilder {
349 pub fn new(leader: Arc<InstanceLeader>) -> Self {
351 Self {
352 leader,
353 registry: None,
354 g1_manager: None,
355 g2_manager: None,
356 g3_manager: None,
357 object_ops: None,
358 g2_physical_layout: None,
359 g1_to_g2_config: None,
360 g2_to_g3_config: None,
361 g2_to_g4_config: None,
362 runtime: None,
363 enable_remote_g4: false,
364 }
365 }
366
367 pub fn with_runtime(mut self, runtime: tokio::runtime::Handle) -> Self {
372 self.runtime = Some(runtime);
373 self
374 }
375
376 pub fn with_registry(mut self, registry: Arc<BlockRegistry>) -> Self {
378 self.registry = Some(registry);
379 self
380 }
381
382 pub fn with_g1_manager(mut self, manager: Arc<BlockManager<G1>>) -> Self {
384 self.g1_manager = Some(manager);
385 self
386 }
387
388 pub fn with_g2_manager(mut self, manager: Arc<BlockManager<G2>>) -> Self {
390 self.g2_manager = Some(manager);
391 self
392 }
393
394 pub fn with_g3_manager(mut self, manager: Arc<BlockManager<G3>>) -> Self {
396 self.g3_manager = Some(manager);
397 self
398 }
399
400 pub fn with_object_ops(mut self, object_ops: Arc<dyn ObjectBlockOps>) -> Self {
405 self.object_ops = Some(object_ops);
406 self
407 }
408
409 pub fn with_g2_physical_layout(mut self, layout: PhysicalLayout) -> Self {
414 self.g2_physical_layout = Some(layout);
415 self
416 }
417
418 pub fn with_g1_to_g2_pipeline(mut self, config: PipelineConfig<G1, G2>) -> Self {
420 self.g1_to_g2_config = Some(config);
421 self
422 }
423
424 pub fn with_g2_to_g3_pipeline(mut self, config: PipelineConfig<G2, G3>) -> Self {
426 self.g2_to_g3_config = Some(config);
427 self
428 }
429
430 pub fn with_g2_to_g4_pipeline(mut self, config: ObjectPipelineConfig<G2>) -> Self {
438 self.g2_to_g4_config = Some(config);
439 self
440 }
441
442 pub fn with_enable_remote_g4(mut self, enable: bool) -> Self {
453 self.enable_remote_g4 = enable;
454 self
455 }
456
457 pub fn build(self) -> Result<OffloadEngine> {
459 let registry = self
460 .registry
461 .ok_or_else(|| anyhow::anyhow!("Block registry required"))?;
462
463 let runtime = self.runtime.unwrap_or_else(|| self.leader.runtime());
466
467 let mut g1_to_g2 = if let Some(config) = self.g1_to_g2_config {
471 let g2_manager = self
472 .g2_manager
473 .clone()
474 .ok_or_else(|| anyhow::anyhow!("G2 manager required for G1→G2 pipeline"))?;
475
476 Some(Pipeline::new(
477 config,
478 registry.clone(),
479 g2_manager,
480 self.leader.clone(),
481 LogicalLayoutHandle::G1,
482 LogicalLayoutHandle::G2,
483 runtime.clone(),
484 ))
485 } else {
486 None
487 };
488
489 let g2_to_g3 = if let Some(config) = self.g2_to_g3_config {
491 let g3_manager = self
492 .g3_manager
493 .ok_or_else(|| anyhow::anyhow!("G3 manager required for G2→G3 pipeline"))?;
494
495 Some(Pipeline::new(
496 config,
497 registry.clone(),
498 g3_manager,
499 self.leader.clone(),
500 LogicalLayoutHandle::G2,
501 LogicalLayoutHandle::G3,
502 runtime.clone(),
503 ))
504 } else {
505 None
506 };
507
508 let g2_to_g4 = if let Some(config) = self.g2_to_g4_config {
511 let object_ops = self
512 .object_ops
513 .ok_or_else(|| anyhow::anyhow!("ObjectBlockOps required for G2→G4 pipeline"))?;
514
515 Some(ObjectPipeline::new(
518 config,
519 object_ops,
520 LogicalLayoutHandle::G2,
521 self.leader.clone(),
522 runtime.clone(),
523 ))
524 } else {
525 None
526 };
527
528 let (remote_g4_tx, remote_g4_rx) = if self.enable_remote_g4 {
530 let (tx, rx) = mpsc::channel::<RemoteG4OffloadRequest>(64);
531 (Some(tx), Some(rx))
532 } else {
533 (None, None)
534 };
535
536 let chain_router_handle = if let Some(ref mut g1_to_g2_pipeline) = g1_to_g2 {
538 if g1_to_g2_pipeline.auto_chain() {
539 if let Some(chain_rx) = g1_to_g2_pipeline.take_chain_rx() {
540 let g2_to_g3_queue = g2_to_g3.as_ref().map(|p| p.eval_queue.clone());
542 let g2_to_g4_queue = g2_to_g4.as_ref().map(|p| p.eval_queue.clone());
543
544 let has_g2_to_g4_local = g2_to_g4_queue.is_some();
546 let has_g2_to_g4_remote = remote_g4_tx.is_some();
547
548 if g2_to_g3_queue.is_some() || has_g2_to_g4_local || has_g2_to_g4_remote {
550 tracing::debug!(
551 has_g2_to_g3 = g2_to_g3_queue.is_some(),
552 has_g2_to_g4_local,
553 has_g2_to_g4_remote,
554 "Spawning chain router for G1→G2 auto-chaining"
555 );
556 Some(runtime.spawn(chain_router_task(
557 chain_rx,
558 g2_to_g3_queue,
559 g2_to_g4_queue,
560 remote_g4_tx,
561 )))
562 } else {
563 tracing::debug!(
564 "G1→G2 auto_chain enabled but no downstream pipelines configured"
565 );
566 None
567 }
568 } else {
569 None
570 }
571 } else {
572 None
573 }
574 } else {
575 None
576 };
577
578 let remote_g4_offload_handle = if let Some(rx) = remote_g4_rx {
580 tracing::info!("Enabling remote G4 offload via workers' ObjectBlockOps");
581 Some(runtime.spawn(remote_g4_offload_task(rx, self.leader.clone())))
582 } else {
583 None
584 };
585
586 Ok(OffloadEngine {
587 engine_id: Uuid::new_v4(),
588 leader: self.leader,
589 registry,
590 g1_to_g2,
591 g2_to_g3,
592 g2_to_g4,
593 transfers: Arc::new(DashMap::new()),
594 _chain_router_handle: chain_router_handle,
595 _remote_g4_offload_handle: remote_g4_offload_handle,
596 })
597 }
598}
599
600struct RemoteG4OffloadRequest {
604 transfer_id: TransferId,
606 keys: Vec<SequenceHash>,
608 block_ids: Vec<BlockId>,
610}
611
612async fn chain_router_task(
618 mut chain_rx: ChainOutputRx<G2>,
619 g2_to_g3_queue: Option<Arc<CancellableQueue<PipelineInput<G2>>>>,
620 g2_to_g4_queue: Option<Arc<CancellableQueue<PipelineInput<G2>>>>,
621 remote_g4_tx: Option<mpsc::Sender<RemoteG4OffloadRequest>>,
622) {
623 while let Some(output) = chain_rx.recv().await {
624 let ChainOutput {
625 transfer_id,
626 blocks,
627 state,
628 } = output;
629
630 if blocks.is_empty() {
631 continue;
632 }
633
634 let weak_blocks: Vec<WeakBlock<G2>> =
637 blocks.iter().map(|block| block.downgrade()).collect();
638
639 let remote_g4_data: Option<(Vec<SequenceHash>, Vec<BlockId>)> = if remote_g4_tx.is_some() {
641 Some((
642 blocks.iter().map(|b| b.sequence_hash()).collect(),
643 blocks.iter().map(|b| b.block_id()).collect(),
644 ))
645 } else {
646 None
647 };
648
649 drop(blocks);
651
652 tracing::debug!(
653 %transfer_id,
654 num_blocks = weak_blocks.len(),
655 "Routing chain output to downstream pipelines as WeakBlocks"
656 );
657
658 if let Some(ref queue) = g2_to_g3_queue {
660 let input = PipelineInput {
661 transfer_id,
662 source: SourceBlocks::Weak(weak_blocks.clone()),
663 state: state.clone(),
664 };
665 if !queue.push(transfer_id, input) {
666 tracing::debug!(%transfer_id, "G2→G3 chain enqueue skipped (cancelled)");
667 }
668 }
669
670 if let Some(ref queue) = g2_to_g4_queue {
672 let input = PipelineInput {
673 transfer_id,
674 source: SourceBlocks::Weak(weak_blocks.clone()),
675 state: state.clone(),
676 };
677 if !queue.push(transfer_id, input) {
678 tracing::debug!(%transfer_id, "G2→G4 chain enqueue skipped (cancelled)");
679 }
680 }
681
682 if let (Some(tx), Some((keys, block_ids))) = (&remote_g4_tx, remote_g4_data) {
684 let request = RemoteG4OffloadRequest {
685 transfer_id,
686 keys,
687 block_ids,
688 };
689 if tx.send(request).await.is_err() {
690 tracing::debug!(%transfer_id, "Remote G4 offload channel closed");
691 }
692 }
693 }
694
695 tracing::debug!("Chain router task shutting down");
696}
697
698async fn remote_g4_offload_task(
705 mut rx: mpsc::Receiver<RemoteG4OffloadRequest>,
706 leader: Arc<InstanceLeader>,
707) {
708 tracing::info!("Remote G4 offload task started");
709
710 while let Some(request) = rx.recv().await {
711 let num_blocks = request.keys.len();
712 tracing::debug!(
713 %request.transfer_id,
714 num_blocks,
715 "Processing remote G4 offload request"
716 );
717
718 let result = leader.execute_remote_offload(
721 LogicalLayoutHandle::G2, RemoteDescriptor::Object {
723 keys: request.keys.clone(),
724 },
725 request.block_ids.clone(),
726 TransferOptions::default(),
727 );
728
729 match result {
730 Ok(notification) => {
731 match notification.await {
733 Ok(()) => {
734 tracing::info!(
735 %request.transfer_id,
736 num_blocks,
737 "Remote G4 offload completed successfully"
738 );
739 }
740 Err(e) => {
741 tracing::warn!(
742 %request.transfer_id,
743 num_blocks,
744 error = %e,
745 "Remote G4 offload failed"
746 );
747 }
748 }
749 }
750 Err(e) => {
751 tracing::warn!(
752 %request.transfer_id,
753 num_blocks,
754 error = %e,
755 "Failed to initiate remote G4 offload"
756 );
757 }
758 }
759 }
760
761 tracing::info!("Remote G4 offload task shutting down");
762}
763
764#[cfg(test)]
765mod tests {
766 use super::*;
767
768 #[test]
772 fn test_transfer_id_generation() {
773 let id1 = TransferId::new();
774 let id2 = TransferId::new();
775 assert_ne!(id1, id2);
776 }
777}