use ::velo::Messenger;
use anyhow::Result;
use bytes::Bytes;
use dashmap::DashMap;
use std::sync::Arc;
use crate::InstanceId;
use kvbm_physical::manager::SerializedLayout;
use super::{
OnboardSessionTx, SessionId, SessionMessageTx, dispatch_onboard_message,
messages::{OnboardMessage, SessionMessage},
};
pub enum MessageTransport {
Velo(VeloTransport),
Local(LocalTransport),
}
impl MessageTransport {
pub fn velo(messenger: Arc<Messenger>) -> Self {
Self::Velo(VeloTransport::new(messenger))
}
pub fn local(
sessions: Arc<DashMap<SessionId, OnboardSessionTx>>,
session_sessions: Arc<DashMap<SessionId, SessionMessageTx>>,
) -> Self {
Self::Local(LocalTransport::new(sessions, session_sessions))
}
pub async fn send(&self, target: InstanceId, message: OnboardMessage) -> Result<()> {
match self {
MessageTransport::Velo(transport) => transport.send(target, message).await,
MessageTransport::Local(transport) => transport.send(target, message).await,
}
}
pub async fn request_metadata(&self, target: InstanceId) -> Result<Vec<SerializedLayout>> {
match self {
MessageTransport::Velo(transport) => transport.request_metadata(target).await,
MessageTransport::Local(_) => {
anyhow::bail!("request_metadata not supported for local transport")
}
}
}
pub async fn send_session(&self, target: InstanceId, message: SessionMessage) -> Result<()> {
match self {
MessageTransport::Velo(transport) => transport.send_session(target, message).await,
MessageTransport::Local(transport) => transport.send_session(target, message).await,
}
}
}
pub struct VeloTransport {
messenger: Arc<Messenger>,
}
impl VeloTransport {
pub fn new(messenger: Arc<Messenger>) -> Self {
Self { messenger }
}
pub async fn send(&self, target: InstanceId, message: OnboardMessage) -> Result<()> {
tracing::debug!(
msg = message.variant_name(),
target = %target,
"Sending message"
);
let bytes = Bytes::from(serde_json::to_vec(&message)?);
self.messenger
.am_send("kvbm.leader.onboard")?
.raw_payload(bytes)
.instance(target)
.send()
.await?;
tracing::debug!(target = %target, "Successfully sent");
Ok(())
}
pub async fn request_metadata(&self, target: InstanceId) -> Result<Vec<SerializedLayout>> {
tracing::debug!(target = %target, "Requesting metadata from instance");
let response: Bytes = self
.messenger
.unary("kvbm.leader.export_metadata")?
.instance(target)
.send()
.await?;
let metadata: Vec<SerializedLayout> = serde_json::from_slice(&response)?;
tracing::debug!(
count = metadata.len(),
target = %target,
"Received metadata entries"
);
Ok(metadata)
}
pub async fn send_session(&self, target: InstanceId, message: SessionMessage) -> Result<()> {
tracing::debug!(
msg = message.variant_name(),
target = %target,
"Sending Session"
);
let bytes = Bytes::from(serde_json::to_vec(&message)?);
self.messenger
.am_send("kvbm.leader.session")?
.raw_payload(bytes)
.instance(target)
.send()
.await?;
tracing::debug!(target = %target, "Successfully sent session msg");
Ok(())
}
}
pub struct LocalTransport {
sessions: Arc<DashMap<SessionId, OnboardSessionTx>>,
session_sessions: Arc<DashMap<SessionId, SessionMessageTx>>,
}
impl LocalTransport {
pub fn new(
sessions: Arc<DashMap<SessionId, OnboardSessionTx>>,
session_sessions: Arc<DashMap<SessionId, SessionMessageTx>>,
) -> Self {
Self {
sessions,
session_sessions,
}
}
pub async fn send(&self, _target: InstanceId, message: OnboardMessage) -> Result<()> {
dispatch_onboard_message(&self.sessions, message).await
}
pub async fn send_session(&self, _target: InstanceId, message: SessionMessage) -> Result<()> {
let session_id = message.session_id();
let sender = self
.session_sessions
.get(&session_id)
.map(|entry| entry.value().clone());
if let Some(sender) = sender {
sender
.send(message)
.await
.map_err(|e| anyhow::anyhow!("failed to send to session {session_id}: {e}"))?;
return Ok(());
}
anyhow::bail!("no session registered for session {session_id}");
}
}