kvbm_engine/leader/session/
mod.rs1mod blocks;
26mod endpoint;
27mod handle;
28mod server_session;
29mod staging;
30mod state;
31
32mod initiator;
34mod messages;
35mod responder;
36pub mod transport;
37
38pub use blocks::BlockHolder;
44
45pub use endpoint::{SessionEndpoint, SessionMessageTx, session_message_channel};
47
48pub use server_session::{
50 ServerSession, ServerSessionCommand, ServerSessionHandle, ServerSessionOptions,
51 create_server_session,
52};
53
54pub use server_session::ServerSessionCommand as EndpointSessionCommand;
56pub use server_session::ServerSessionHandle as EndpointSessionHandle;
57
58pub use handle::{SessionHandle, SessionHandleStateTx, session_handle_state_channel};
60
61pub use state::{AttachmentState, ControlRole, SessionPhase};
63
64pub use messages::{BlockInfo, SessionMessage, SessionStateSnapshot};
66
67pub use initiator::InitiatorSession;
73pub use responder::ResponderSession;
74
75pub use server_session::ServerSessionOptions as ControllableSessionOptions;
77
78#[derive(Debug, Clone)]
80pub struct ControllableSessionResult {
81 pub session_id: super::SessionId,
83 pub local_g2_count: usize,
85 pub local_g3_count: usize,
87}
88
89pub use messages::{BlockMatch, OnboardMessage};
91
92pub use transport::{LocalTransport, MessageTransport, VeloTransport};
94
95use anyhow::Result;
96use dashmap::DashMap;
97use tokio::sync::mpsc;
98
99pub type SessionId = uuid::Uuid;
100pub type OnboardSessionTx = mpsc::Sender<OnboardMessage>;
101
102pub async fn dispatch_onboard_message(
108 sessions: &DashMap<SessionId, OnboardSessionTx>,
109 message: OnboardMessage,
110) -> Result<()> {
111 let session_id = message.session_id();
112
113 let sender = sessions.get(&session_id).map(|entry| entry.value().clone());
114 if let Some(sender) = sender {
115 sender
116 .send(message)
117 .await
118 .map_err(|e| anyhow::anyhow!("failed to send to session {session_id}: {e}"))?;
119 return Ok(());
120 }
121
122 anyhow::bail!("no session task registered for session {session_id}");
123}
124
125pub async fn dispatch_session_message(
130 sessions: &DashMap<SessionId, SessionMessageTx>,
131 message: SessionMessage,
132) -> Result<()> {
133 let session_id = message.session_id();
134
135 let sender = sessions.get(&session_id).map(|entry| entry.value().clone());
136 if let Some(sender) = sender {
137 sender
138 .send(message)
139 .await
140 .map_err(|e| anyhow::anyhow!("failed to send to session {session_id}: {e}"))?;
141 return Ok(());
142 }
143
144 anyhow::bail!("no session registered for session {session_id}");
145}
146
147#[cfg(test)]
148mod tests {
149 use super::*;
150
151 #[tokio::test]
152 async fn test_dispatch_onboard_message() {
153 let sessions: DashMap<SessionId, OnboardSessionTx> = DashMap::new();
154 let session_id = SessionId::new_v4();
155 let (tx, mut rx) = mpsc::channel(16);
156 sessions.insert(session_id, tx);
157
158 let msg = OnboardMessage::CloseSession {
159 requester: crate::InstanceId::new_v4(),
160 session_id,
161 };
162
163 dispatch_onboard_message(&sessions, msg).await.unwrap();
164
165 let received = rx.recv().await.unwrap();
166 assert_eq!(received.session_id(), session_id);
167 }
168
169 #[tokio::test]
170 async fn test_dispatch_session_message() {
171 let sessions: DashMap<SessionId, SessionMessageTx> = DashMap::new();
172 let session_id = SessionId::new_v4();
173 let (tx, mut rx) = mpsc::channel(16);
174 sessions.insert(session_id, tx);
175
176 let msg = SessionMessage::Close { session_id };
177
178 dispatch_session_message(&sessions, msg).await.unwrap();
179
180 let received = rx.recv().await.unwrap();
181 assert_eq!(received.session_id(), session_id);
182 }
183
184 #[tokio::test]
185 async fn test_dispatch_missing_onboard_session() {
186 let sessions: DashMap<SessionId, OnboardSessionTx> = DashMap::new();
187 let session_id = SessionId::new_v4();
188
189 let msg = OnboardMessage::CloseSession {
190 requester: crate::InstanceId::new_v4(),
191 session_id,
192 };
193
194 let result = dispatch_onboard_message(&sessions, msg).await;
195 assert!(result.is_err());
196 }
197
198 #[tokio::test]
199 async fn test_dispatch_missing_session_message() {
200 let sessions: DashMap<SessionId, SessionMessageTx> = DashMap::new();
201 let session_id = SessionId::new_v4();
202
203 let msg = SessionMessage::Close { session_id };
204
205 let result = dispatch_session_message(&sessions, msg).await;
206 assert!(result.is_err());
207 }
208}