use std::sync::Arc;
use anyhow::Result;
use dashmap::DashMap;
use tokio::sync::mpsc;
use tokio::task::JoinHandle;
use uuid::Uuid;
use crate::leader::InstanceLeader;
use crate::object::ObjectBlockOps;
use crate::worker::RemoteDescriptor;
use crate::{BlockId, G1, G2, G3, SequenceHash};
use kvbm_common::LogicalLayoutHandle;
use kvbm_logical::blocks::{BlockMetadata, BlockRegistry, WeakBlock};
use kvbm_logical::manager::BlockManager;
use kvbm_physical::transfer::{PhysicalLayout, TransferOptions};
use super::handle::{TransferHandle, TransferId, TransferState};
use super::pipeline::{
ChainOutput, ChainOutputRx, ObjectPipeline, ObjectPipelineConfig, Pipeline, PipelineConfig,
PipelineInput,
};
use super::queue::CancellableQueue;
use super::settlement::{
PipelineLane, PipelineSettlementTracker, SettlementError, SettlementTarget, SettlementToken,
SettlementWaiter, wait_for_settlement,
};
use super::source::SourceBlocks;
#[allow(dead_code)]
pub struct OffloadEngine {
engine_id: Uuid,
leader: Arc<InstanceLeader>,
registry: Arc<BlockRegistry>,
g1_to_g2: Option<Pipeline<G1, G2>>,
g2_to_g3: Option<Pipeline<G2, G3>>,
g2_to_g4: Option<ObjectPipeline<G2>>,
transfers: Arc<DashMap<TransferId, Arc<std::sync::Mutex<TransferState>>>>,
_chain_router_handle: Option<JoinHandle<()>>,
_remote_g4_offload_handle: Option<JoinHandle<()>>,
}
impl OffloadEngine {
pub fn builder(leader: Arc<InstanceLeader>) -> OffloadEngineBuilder {
OffloadEngineBuilder::new(leader)
}
pub fn enqueue_g1_to_g2(&self, blocks: impl Into<SourceBlocks<G1>>) -> Result<TransferHandle> {
let pipeline = self
.g1_to_g2
.as_ref()
.ok_or_else(|| anyhow::anyhow!("G1→G2 pipeline not configured"))?;
self.enqueue_to_pipeline(pipeline, blocks.into())
}
pub fn enqueue_g1_to_g2_with_precondition(
&self,
blocks: impl Into<SourceBlocks<G1>>,
precondition: Option<velo::EventHandle>,
) -> Result<TransferHandle> {
let pipeline = self
.g1_to_g2
.as_ref()
.ok_or_else(|| anyhow::anyhow!("G1→G2 pipeline not configured"))?;
self.enqueue_to_pipeline_with_precondition(pipeline, blocks.into(), precondition)
}
pub fn enqueue_g2_to_g3(&self, blocks: impl Into<SourceBlocks<G2>>) -> Result<TransferHandle> {
let pipeline = self
.g2_to_g3
.as_ref()
.ok_or_else(|| anyhow::anyhow!("G2→G3 pipeline not configured"))?;
self.enqueue_to_pipeline(pipeline, blocks.into())
}
pub fn enqueue_g2_to_g4(&self, blocks: impl Into<SourceBlocks<G2>>) -> Result<TransferHandle> {
let pipeline = self
.g2_to_g4
.as_ref()
.ok_or_else(|| anyhow::anyhow!("G2→G4 pipeline not configured"))?;
self.enqueue_to_object_pipeline(pipeline, blocks.into())
}
fn create_transfer<T: BlockMetadata>(
&self,
source: &SourceBlocks<T>,
) -> (
TransferId,
Arc<std::sync::Mutex<TransferState>>,
TransferHandle,
) {
let input_block_ids = self.extract_block_ids(source);
let transfer_id = TransferId::new();
let (state, handle) = TransferState::new(transfer_id, input_block_ids);
let state = Arc::new(std::sync::Mutex::new(state));
self.transfers.insert(transfer_id, state.clone());
(transfer_id, state, handle)
}
fn enqueue_to_pipeline<Src: BlockMetadata, Dst: BlockMetadata>(
&self,
pipeline: &Pipeline<Src, Dst>,
source: SourceBlocks<Src>,
) -> Result<TransferHandle> {
let (transfer_id, state, handle) = self.create_transfer(&source);
if !pipeline.enqueue(transfer_id, source, state) {
tracing::warn!("Transfer {} was cancelled before enqueueing", transfer_id);
}
Ok(handle)
}
fn enqueue_to_pipeline_with_precondition<Src: BlockMetadata, Dst: BlockMetadata>(
&self,
pipeline: &Pipeline<Src, Dst>,
source: SourceBlocks<Src>,
precondition: Option<velo::EventHandle>,
) -> Result<TransferHandle> {
let (transfer_id, state, handle) = self.create_transfer(&source);
state.lock().unwrap().precondition = precondition;
if !pipeline.enqueue(transfer_id, source, state) {
tracing::warn!("Transfer {} was cancelled before enqueueing", transfer_id);
}
Ok(handle)
}
fn enqueue_to_object_pipeline(
&self,
pipeline: &ObjectPipeline<G2>,
source: SourceBlocks<G2>,
) -> Result<TransferHandle> {
let (transfer_id, state, handle) = self.create_transfer(&source);
if !pipeline.enqueue(transfer_id, source, state) {
tracing::warn!("Transfer {} was cancelled before enqueueing", transfer_id);
}
Ok(handle)
}
fn extract_block_ids<T: BlockMetadata>(&self, source: &SourceBlocks<T>) -> Vec<BlockId> {
match source {
SourceBlocks::External(blocks) => blocks.iter().map(|b| b.block_id).collect(),
SourceBlocks::Strong(blocks) => blocks.iter().map(|b| b.block_id()).collect(),
SourceBlocks::Weak(_) => Vec::new(), }
}
pub fn release_transfer(&self, transfer_id: TransferId) {
self.transfers.remove(&transfer_id);
}
pub fn active_transfer_count(&self) -> usize {
self.transfers.len()
}
pub fn has_g1_to_g2(&self) -> bool {
self.g1_to_g2.is_some()
}
pub fn has_g2_to_g3(&self) -> bool {
self.g2_to_g3.is_some()
}
pub fn has_g2_to_g4(&self) -> bool {
self.g2_to_g4.is_some()
}
pub fn settlement_token(&self) -> SettlementToken {
let mut checkpoints = [None; 3];
for lane in PipelineLane::ALL {
if let Some((tracker, _)) = self.settlement_tracker(lane) {
checkpoints[lane.index()] = Some(tracker.snapshot().checkpoint());
}
}
SettlementToken {
engine_id: self.engine_id,
checkpoints,
}
}
pub async fn settle_after(
&self,
token: SettlementToken,
target: SettlementTarget,
) -> Result<(), SettlementError> {
token.validate_engine(self.engine_id)?;
if target.is_empty() {
return Ok(());
}
let mut waiters = Vec::new();
for lane in PipelineLane::ALL {
let delta = target.completed_batches(lane);
if delta == 0 {
continue;
}
let (tracker, auto_chain) = self
.settlement_tracker(lane)
.ok_or(SettlementError::LaneUnavailable { lane })?;
if auto_chain {
return Err(SettlementError::UnsupportedAutoChain { lane });
}
let checkpoint = token.checkpoint(lane)?;
waiters.push(SettlementWaiter::new(lane, tracker, checkpoint, delta)?);
}
wait_for_settlement(waiters).await
}
fn settlement_tracker(&self, lane: PipelineLane) -> Option<(&PipelineSettlementTracker, bool)> {
match lane {
PipelineLane::G1ToG2 => self
.g1_to_g2
.as_ref()
.map(|pipeline| (&pipeline.settlement, pipeline.auto_chain())),
PipelineLane::G2ToG3 => self
.g2_to_g3
.as_ref()
.map(|pipeline| (&pipeline.settlement, pipeline.auto_chain())),
PipelineLane::G2ToG4 => self
.g2_to_g4
.as_ref()
.map(|pipeline| (&pipeline.settlement, false)),
}
}
}
pub struct OffloadEngineBuilder {
leader: Arc<InstanceLeader>,
registry: Option<Arc<BlockRegistry>>,
g1_manager: Option<Arc<BlockManager<G1>>>,
g2_manager: Option<Arc<BlockManager<G2>>>,
g3_manager: Option<Arc<BlockManager<G3>>>,
object_ops: Option<Arc<dyn ObjectBlockOps>>,
g2_physical_layout: Option<PhysicalLayout>,
g1_to_g2_config: Option<PipelineConfig<G1, G2>>,
g2_to_g3_config: Option<PipelineConfig<G2, G3>>,
g2_to_g4_config: Option<ObjectPipelineConfig<G2>>,
runtime: Option<tokio::runtime::Handle>,
enable_remote_g4: bool,
}
impl OffloadEngineBuilder {
pub fn new(leader: Arc<InstanceLeader>) -> Self {
Self {
leader,
registry: None,
g1_manager: None,
g2_manager: None,
g3_manager: None,
object_ops: None,
g2_physical_layout: None,
g1_to_g2_config: None,
g2_to_g3_config: None,
g2_to_g4_config: None,
runtime: None,
enable_remote_g4: false,
}
}
pub fn with_runtime(mut self, runtime: tokio::runtime::Handle) -> Self {
self.runtime = Some(runtime);
self
}
pub fn with_registry(mut self, registry: Arc<BlockRegistry>) -> Self {
self.registry = Some(registry);
self
}
pub fn with_g1_manager(mut self, manager: Arc<BlockManager<G1>>) -> Self {
self.g1_manager = Some(manager);
self
}
pub fn with_g2_manager(mut self, manager: Arc<BlockManager<G2>>) -> Self {
self.g2_manager = Some(manager);
self
}
pub fn with_g3_manager(mut self, manager: Arc<BlockManager<G3>>) -> Self {
self.g3_manager = Some(manager);
self
}
pub fn with_object_ops(mut self, object_ops: Arc<dyn ObjectBlockOps>) -> Self {
self.object_ops = Some(object_ops);
self
}
pub fn with_g2_physical_layout(mut self, layout: PhysicalLayout) -> Self {
self.g2_physical_layout = Some(layout);
self
}
pub fn with_g1_to_g2_pipeline(mut self, config: PipelineConfig<G1, G2>) -> Self {
self.g1_to_g2_config = Some(config);
self
}
pub fn with_g2_to_g3_pipeline(mut self, config: PipelineConfig<G2, G3>) -> Self {
self.g2_to_g3_config = Some(config);
self
}
pub fn with_g2_to_g4_pipeline(mut self, config: ObjectPipelineConfig<G2>) -> Self {
self.g2_to_g4_config = Some(config);
self
}
pub fn with_enable_remote_g4(mut self, enable: bool) -> Self {
self.enable_remote_g4 = enable;
self
}
pub fn build(self) -> Result<OffloadEngine> {
let registry = self
.registry
.ok_or_else(|| anyhow::anyhow!("Block registry required"))?;
let runtime = self.runtime.unwrap_or_else(|| self.leader.runtime());
let mut g1_to_g2 = if let Some(config) = self.g1_to_g2_config {
let g2_manager = self
.g2_manager
.clone()
.ok_or_else(|| anyhow::anyhow!("G2 manager required for G1→G2 pipeline"))?;
Some(Pipeline::new(
config,
registry.clone(),
g2_manager,
self.leader.clone(),
LogicalLayoutHandle::G1,
LogicalLayoutHandle::G2,
runtime.clone(),
))
} else {
None
};
let g2_to_g3 = if let Some(config) = self.g2_to_g3_config {
let g3_manager = self
.g3_manager
.ok_or_else(|| anyhow::anyhow!("G3 manager required for G2→G3 pipeline"))?;
Some(Pipeline::new(
config,
registry.clone(),
g3_manager,
self.leader.clone(),
LogicalLayoutHandle::G2,
LogicalLayoutHandle::G3,
runtime.clone(),
))
} else {
None
};
let g2_to_g4 = if let Some(config) = self.g2_to_g4_config {
let object_ops = self
.object_ops
.ok_or_else(|| anyhow::anyhow!("ObjectBlockOps required for G2→G4 pipeline"))?;
Some(ObjectPipeline::new(
config,
object_ops,
LogicalLayoutHandle::G2,
self.leader.clone(),
runtime.clone(),
))
} else {
None
};
let (remote_g4_tx, remote_g4_rx) = if self.enable_remote_g4 {
let (tx, rx) = mpsc::channel::<RemoteG4OffloadRequest>(64);
(Some(tx), Some(rx))
} else {
(None, None)
};
let chain_router_handle = if let Some(ref mut g1_to_g2_pipeline) = g1_to_g2 {
if g1_to_g2_pipeline.auto_chain() {
if let Some(chain_rx) = g1_to_g2_pipeline.take_chain_rx() {
let g2_to_g3_queue = g2_to_g3.as_ref().map(|p| p.eval_queue.clone());
let g2_to_g4_queue = g2_to_g4.as_ref().map(|p| p.eval_queue.clone());
let has_g2_to_g4_local = g2_to_g4_queue.is_some();
let has_g2_to_g4_remote = remote_g4_tx.is_some();
if g2_to_g3_queue.is_some() || has_g2_to_g4_local || has_g2_to_g4_remote {
tracing::debug!(
has_g2_to_g3 = g2_to_g3_queue.is_some(),
has_g2_to_g4_local,
has_g2_to_g4_remote,
"Spawning chain router for G1→G2 auto-chaining"
);
Some(runtime.spawn(chain_router_task(
chain_rx,
g2_to_g3_queue,
g2_to_g4_queue,
remote_g4_tx,
)))
} else {
tracing::debug!(
"G1→G2 auto_chain enabled but no downstream pipelines configured"
);
None
}
} else {
None
}
} else {
None
}
} else {
None
};
let remote_g4_offload_handle = if let Some(rx) = remote_g4_rx {
tracing::info!("Enabling remote G4 offload via workers' ObjectBlockOps");
Some(runtime.spawn(remote_g4_offload_task(rx, self.leader.clone())))
} else {
None
};
Ok(OffloadEngine {
engine_id: Uuid::new_v4(),
leader: self.leader,
registry,
g1_to_g2,
g2_to_g3,
g2_to_g4,
transfers: Arc::new(DashMap::new()),
_chain_router_handle: chain_router_handle,
_remote_g4_offload_handle: remote_g4_offload_handle,
})
}
}
struct RemoteG4OffloadRequest {
transfer_id: TransferId,
keys: Vec<SequenceHash>,
block_ids: Vec<BlockId>,
}
async fn chain_router_task(
mut chain_rx: ChainOutputRx<G2>,
g2_to_g3_queue: Option<Arc<CancellableQueue<PipelineInput<G2>>>>,
g2_to_g4_queue: Option<Arc<CancellableQueue<PipelineInput<G2>>>>,
remote_g4_tx: Option<mpsc::Sender<RemoteG4OffloadRequest>>,
) {
while let Some(output) = chain_rx.recv().await {
let ChainOutput {
transfer_id,
blocks,
state,
} = output;
if blocks.is_empty() {
continue;
}
let weak_blocks: Vec<WeakBlock<G2>> =
blocks.iter().map(|block| block.downgrade()).collect();
let remote_g4_data: Option<(Vec<SequenceHash>, Vec<BlockId>)> = if remote_g4_tx.is_some() {
Some((
blocks.iter().map(|b| b.sequence_hash()).collect(),
blocks.iter().map(|b| b.block_id()).collect(),
))
} else {
None
};
drop(blocks);
tracing::debug!(
%transfer_id,
num_blocks = weak_blocks.len(),
"Routing chain output to downstream pipelines as WeakBlocks"
);
if let Some(ref queue) = g2_to_g3_queue {
let input = PipelineInput {
transfer_id,
source: SourceBlocks::Weak(weak_blocks.clone()),
state: state.clone(),
};
if !queue.push(transfer_id, input) {
tracing::debug!(%transfer_id, "G2→G3 chain enqueue skipped (cancelled)");
}
}
if let Some(ref queue) = g2_to_g4_queue {
let input = PipelineInput {
transfer_id,
source: SourceBlocks::Weak(weak_blocks.clone()),
state: state.clone(),
};
if !queue.push(transfer_id, input) {
tracing::debug!(%transfer_id, "G2→G4 chain enqueue skipped (cancelled)");
}
}
if let (Some(tx), Some((keys, block_ids))) = (&remote_g4_tx, remote_g4_data) {
let request = RemoteG4OffloadRequest {
transfer_id,
keys,
block_ids,
};
if tx.send(request).await.is_err() {
tracing::debug!(%transfer_id, "Remote G4 offload channel closed");
}
}
}
tracing::debug!("Chain router task shutting down");
}
async fn remote_g4_offload_task(
mut rx: mpsc::Receiver<RemoteG4OffloadRequest>,
leader: Arc<InstanceLeader>,
) {
tracing::info!("Remote G4 offload task started");
while let Some(request) = rx.recv().await {
let num_blocks = request.keys.len();
tracing::debug!(
%request.transfer_id,
num_blocks,
"Processing remote G4 offload request"
);
let result = leader.execute_remote_offload(
LogicalLayoutHandle::G2, RemoteDescriptor::Object {
keys: request.keys.clone(),
},
request.block_ids.clone(),
TransferOptions::default(),
);
match result {
Ok(notification) => {
match notification.await {
Ok(()) => {
tracing::info!(
%request.transfer_id,
num_blocks,
"Remote G4 offload completed successfully"
);
}
Err(e) => {
tracing::warn!(
%request.transfer_id,
num_blocks,
error = %e,
"Remote G4 offload failed"
);
}
}
}
Err(e) => {
tracing::warn!(
%request.transfer_id,
num_blocks,
error = %e,
"Failed to initiate remote G4 offload"
);
}
}
}
tracing::info!("Remote G4 offload task shutting down");
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_transfer_id_generation() {
let id1 = TransferId::new();
let id2 = TransferId::new();
assert_ne!(id1, id2);
}
}