use std::ops::Range;
use std::sync::Arc;
use serde::{Deserialize, Serialize};
use crate::{BlockId, G2, InstanceId, SequenceHash};
use kvbm_logical::blocks::ImmutableBlock;
use kvbm_physical::manager::LayoutHandle;
use super::SessionId;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum OnboardMessage {
CreateSession {
requester: InstanceId,
session_id: SessionId,
sequence_hashes: Vec<SequenceHash>,
},
SearchComplete {
responder: InstanceId,
session_id: SessionId,
},
G2Results {
responder: InstanceId,
session_id: SessionId,
sequence_hashes: Vec<SequenceHash>,
block_ids: Vec<BlockId>,
},
G3Results {
responder: InstanceId,
session_id: SessionId,
sequence_hashes: Vec<SequenceHash>,
},
HoldBlocks {
requester: InstanceId,
session_id: SessionId,
hold_hashes: Vec<SequenceHash>,
drop_hashes: Vec<SequenceHash>,
},
StageBlocks {
requester: InstanceId,
session_id: SessionId,
stage_hashes: Vec<SequenceHash>,
},
BlocksReady {
responder: InstanceId,
session_id: SessionId,
sequence_hashes: Vec<SequenceHash>,
block_ids: Vec<BlockId>,
},
Acknowledged {
responder: InstanceId,
session_id: SessionId,
},
ReleaseBlocks {
requester: InstanceId,
session_id: SessionId,
release_hashes: Vec<SequenceHash>,
},
CloseSession {
requester: InstanceId,
session_id: SessionId,
},
G4Results {
session_id: SessionId,
found_hashes: Vec<(SequenceHash, usize)>,
},
G4LoadComplete {
session_id: SessionId,
success: Vec<SequenceHash>,
failures: Vec<(SequenceHash, String)>,
#[serde(skip)]
blocks: Arc<Vec<ImmutableBlock<G2>>>,
},
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BlockMatch {
pub sequence_hash: SequenceHash,
pub block_id: BlockId,
}
impl OnboardMessage {
pub fn session_id(&self) -> SessionId {
match self {
OnboardMessage::CreateSession { session_id, .. }
| OnboardMessage::SearchComplete { session_id, .. }
| OnboardMessage::G2Results { session_id, .. }
| OnboardMessage::G3Results { session_id, .. }
| OnboardMessage::HoldBlocks { session_id, .. }
| OnboardMessage::StageBlocks { session_id, .. }
| OnboardMessage::BlocksReady { session_id, .. }
| OnboardMessage::Acknowledged { session_id, .. }
| OnboardMessage::ReleaseBlocks { session_id, .. }
| OnboardMessage::CloseSession { session_id, .. }
| OnboardMessage::G4Results { session_id, .. }
| OnboardMessage::G4LoadComplete { session_id, .. } => *session_id,
}
}
pub fn instance_id(&self) -> InstanceId {
match self {
OnboardMessage::CreateSession { requester, .. }
| OnboardMessage::HoldBlocks { requester, .. }
| OnboardMessage::StageBlocks { requester, .. }
| OnboardMessage::ReleaseBlocks { requester, .. }
| OnboardMessage::CloseSession { requester, .. } => *requester,
OnboardMessage::SearchComplete { responder, .. }
| OnboardMessage::G2Results { responder, .. }
| OnboardMessage::G3Results { responder, .. }
| OnboardMessage::BlocksReady { responder, .. }
| OnboardMessage::Acknowledged { responder, .. } => *responder,
OnboardMessage::G4Results { .. } | OnboardMessage::G4LoadComplete { .. } => {
panic!("G4 messages are internal and do not have an instance ID")
}
}
}
pub fn variant_name(&self) -> &'static str {
match self {
OnboardMessage::CreateSession { .. } => "CreateSession",
OnboardMessage::SearchComplete { .. } => "SearchComplete",
OnboardMessage::G2Results { .. } => "G2Results",
OnboardMessage::G3Results { .. } => "G3Results",
OnboardMessage::HoldBlocks { .. } => "HoldBlocks",
OnboardMessage::StageBlocks { .. } => "StageBlocks",
OnboardMessage::BlocksReady { .. } => "BlocksReady",
OnboardMessage::Acknowledged { .. } => "Acknowledged",
OnboardMessage::ReleaseBlocks { .. } => "ReleaseBlocks",
OnboardMessage::CloseSession { .. } => "CloseSession",
OnboardMessage::G4Results { .. } => "G4Results",
OnboardMessage::G4LoadComplete { .. } => "G4LoadComplete",
}
}
}
use super::state::{ControlRole, SessionPhase};
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum SessionMessage {
Attach {
peer: InstanceId,
session_id: SessionId,
as_role: ControlRole,
},
Detach {
peer: InstanceId,
session_id: SessionId,
},
YieldControl {
peer: InstanceId,
session_id: SessionId,
},
AcquireControl {
peer: InstanceId,
session_id: SessionId,
},
TriggerStaging {
session_id: SessionId,
},
HoldBlocks {
session_id: SessionId,
hold_hashes: Vec<SequenceHash>,
},
ReleaseBlocks {
session_id: SessionId,
release_hashes: Vec<SequenceHash>,
},
BlocksPulled {
session_id: SessionId,
pulled_hashes: Vec<SequenceHash>,
},
StateResponse {
session_id: SessionId,
state: SessionStateSnapshot,
},
BlocksStaged {
session_id: SessionId,
staged_blocks: Vec<BlockInfo>,
remaining: usize,
layer_range: Option<Range<usize>>,
},
Close {
session_id: SessionId,
},
Error {
session_id: SessionId,
message: String,
},
}
impl SessionMessage {
pub fn session_id(&self) -> SessionId {
match self {
SessionMessage::Attach { session_id, .. }
| SessionMessage::Detach { session_id, .. }
| SessionMessage::YieldControl { session_id, .. }
| SessionMessage::AcquireControl { session_id, .. }
| SessionMessage::TriggerStaging { session_id, .. }
| SessionMessage::HoldBlocks { session_id, .. }
| SessionMessage::ReleaseBlocks { session_id, .. }
| SessionMessage::BlocksPulled { session_id, .. }
| SessionMessage::StateResponse { session_id, .. }
| SessionMessage::BlocksStaged { session_id, .. }
| SessionMessage::Close { session_id, .. }
| SessionMessage::Error { session_id, .. } => *session_id,
}
}
pub fn peer(&self) -> Option<InstanceId> {
match self {
SessionMessage::Attach { peer, .. }
| SessionMessage::Detach { peer, .. }
| SessionMessage::YieldControl { peer, .. }
| SessionMessage::AcquireControl { peer, .. } => Some(*peer),
_ => None,
}
}
pub fn is_control_command(&self) -> bool {
matches!(
self,
SessionMessage::TriggerStaging { .. }
| SessionMessage::HoldBlocks { .. }
| SessionMessage::ReleaseBlocks { .. }
| SessionMessage::BlocksPulled { .. }
)
}
pub fn is_state_response(&self) -> bool {
matches!(
self,
SessionMessage::StateResponse { .. } | SessionMessage::BlocksStaged { .. }
)
}
pub fn variant_name(&self) -> &'static str {
match self {
SessionMessage::Attach { .. } => "Attach",
SessionMessage::Detach { .. } => "Detach",
SessionMessage::YieldControl { .. } => "YieldControl",
SessionMessage::AcquireControl { .. } => "AcquireControl",
SessionMessage::TriggerStaging { .. } => "TriggerStaging",
SessionMessage::HoldBlocks { .. } => "HoldBlocks",
SessionMessage::ReleaseBlocks { .. } => "ReleaseBlocks",
SessionMessage::BlocksPulled { .. } => "BlocksPulled",
SessionMessage::StateResponse { .. } => "StateResponse",
SessionMessage::BlocksStaged { .. } => "BlocksStaged",
SessionMessage::Close { .. } => "Close",
SessionMessage::Error { .. } => "Error",
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct SessionStateSnapshot {
pub phase: SessionPhase,
pub control_role: ControlRole,
pub g2_blocks: Vec<BlockInfo>,
pub g3_pending: usize,
#[serde(default)]
pub ready_layer_range: Option<Range<usize>>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct BlockInfo {
pub block_id: BlockId,
pub sequence_hash: SequenceHash,
pub layout_handle: LayoutHandle,
}