use ::velo::{Handler, Messenger};
use anyhow::Result;
use bytes::Bytes;
use dashmap::DashMap;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use crate::leader::session::{
OnboardMessage, OnboardSessionTx, SessionId, SessionMessage, SessionMessageTx,
dispatch_onboard_message, dispatch_session_message,
};
use kvbm_physical::manager::SerializedLayout;
pub type ExportMetadataCallback = Arc<
dyn Fn() -> Pin<Box<dyn Future<Output = Result<Vec<SerializedLayout>>> + Send>> + Send + Sync,
>;
pub struct VeloLeaderService {
messenger: Arc<Messenger>,
sessions: Arc<DashMap<SessionId, OnboardSessionTx>>,
spawn_responder: Option<Arc<dyn Fn(OnboardMessage) -> Result<()> + Send + Sync>>,
session_sessions: Option<Arc<DashMap<SessionId, SessionMessageTx>>>,
export_metadata: Option<ExportMetadataCallback>,
}
impl VeloLeaderService {
pub fn new(
messenger: Arc<Messenger>,
sessions: Arc<DashMap<SessionId, OnboardSessionTx>>,
) -> Self {
Self {
messenger,
sessions,
spawn_responder: None,
session_sessions: None,
export_metadata: None,
}
}
pub fn with_spawn_responder<F>(mut self, f: F) -> Self
where
F: Fn(OnboardMessage) -> Result<()> + Send + Sync + 'static,
{
self.spawn_responder = Some(Arc::new(f));
self
}
pub fn with_session_sessions(
mut self,
sessions: Arc<DashMap<SessionId, SessionMessageTx>>,
) -> Self {
self.session_sessions = Some(sessions);
self
}
pub fn with_export_metadata(mut self, callback: ExportMetadataCallback) -> Self {
self.export_metadata = Some(callback);
self
}
pub fn register_handlers(self) -> Result<()> {
self.register_onboard_handler()?;
if self.session_sessions.is_some() {
self.register_session_handler()?;
}
if self.export_metadata.is_some() {
self.register_export_metadata_handler()?;
}
Ok(())
}
fn register_onboard_handler(&self) -> Result<()> {
let sessions = self.sessions.clone();
let spawn_responder = self.spawn_responder.clone();
let handler = Handler::am_handler_async("kvbm.leader.onboard", move |ctx| {
let sessions = sessions.clone();
let spawn_responder = spawn_responder.clone();
async move {
let message: OnboardMessage = serde_json::from_slice(&ctx.payload)
.map_err(|e| anyhow::anyhow!("failed to deserialize OnboardMessage: {e}"))?;
let session_id = message.session_id();
tracing::debug!(
variant = message.variant_name(),
%session_id,
"Received onboard message"
);
if matches!(message, OnboardMessage::CreateSession { .. })
&& !sessions.contains_key(&session_id)
{
tracing::debug!(%session_id, "Spawning new ResponderSession");
if let Some(ref spawner) = spawn_responder {
spawner(message.clone()).ok(); }
}
tracing::debug!(%session_id, "Dispatching message to session");
dispatch_onboard_message(&sessions, message).await?;
Ok(())
}
})
.build();
self.messenger.register_handler(handler)?;
Ok(())
}
fn register_session_handler(&self) -> Result<()> {
let session_sessions = self
.session_sessions
.clone()
.expect("session_sessions required for handler registration");
let handler = Handler::am_handler_async("kvbm.leader.session", move |ctx| {
let session_sessions = session_sessions.clone();
async move {
let message: SessionMessage = serde_json::from_slice(&ctx.payload)
.map_err(|e| anyhow::anyhow!("failed to deserialize SessionMessage: {e}"))?;
let session_id = message.session_id();
tracing::debug!(
variant = message.variant_name(),
%session_id,
"Received session message"
);
dispatch_session_message(&session_sessions, message).await?;
Ok(())
}
})
.build();
self.messenger.register_handler(handler)?;
Ok(())
}
fn register_export_metadata_handler(&self) -> Result<()> {
let export_metadata = self
.export_metadata
.clone()
.expect("export_metadata callback required for handler registration");
let handler = Handler::unary_handler_async("kvbm.leader.export_metadata", move |_ctx| {
let export_metadata = export_metadata.clone();
async move {
tracing::debug!("Received export_metadata request");
let metadata_vec = export_metadata().await?;
let serialized = serde_json::to_vec(&metadata_vec)?;
tracing::debug!(
count = metadata_vec.len(),
"Returning worker metadata entries"
);
Ok(Some(Bytes::from(serialized)))
}
})
.build();
self.messenger.register_handler(handler)?;
Ok(())
}
}