use std::sync::Arc;
use boatramp_core::time::now_unix_ms;
use boatramp_handlers::SessionController;
use crate::session_store::{SessionStore, StoreError};
fn to_handler_err(err: StoreError) -> boatramp_handlers::SessionError {
use boatramp_core::session::SessionError as Core;
use boatramp_handlers::SessionError as H;
match err {
StoreError::Session(Core::FrameTooLarge) => H::FrameTooLarge,
StoreError::Session(Core::BufferFull) => H::BufferFull,
StoreError::Session(Core::Closed) => H::Closed,
StoreError::NotFound => H::Closed,
StoreError::PrincipalMismatch => H::AccessDenied,
StoreError::ProjectSessionsFull => H::Other("project session limit reached".into()),
StoreError::Kv(m) => H::Other(format!("session store: {m}")),
StoreError::Corrupt(m) => H::Other(format!("session record: {m}")),
}
}
pub(crate) struct ServerSessionController {
store: SessionStore,
project: String,
id: String,
}
impl ServerSessionController {
pub(crate) fn new(
store: SessionStore,
project: impl Into<String>,
id: impl Into<String>,
) -> Self {
Self {
store,
project: project.into(),
id: id.into(),
}
}
}
#[async_trait::async_trait]
impl SessionController for ServerSessionController {
async fn send(&self, payload: Vec<u8>) -> Result<u64, boatramp_handlers::SessionError> {
self.store
.send(&self.project, &self.id, payload, now_unix_ms())
.await
.map_err(to_handler_err)
}
async fn checkpoint(&self, snapshot: Vec<u8>) -> Result<(), boatramp_handlers::SessionError> {
self.store
.checkpoint(&self.project, &self.id, snapshot, now_unix_ms())
.await
.map_err(to_handler_err)
}
async fn close(&self, reason: String) -> Result<(), boatramp_handlers::SessionError> {
self.store
.close(&self.project, &self.id, &reason, now_unix_ms())
.await
.map_err(to_handler_err)
}
}
pub(crate) fn controller(
store: SessionStore,
project: &str,
id: &str,
) -> Arc<dyn SessionController> {
Arc::new(ServerSessionController::new(store, project, id))
}
#[cfg(test)]
mod tests {
use super::*;
use boatramp_core::kv::MemoryKv;
use boatramp_core::session::SessionLimits;
fn store() -> SessionStore {
SessionStore::new(Arc::new(MemoryKv::new()), SessionLimits::default())
}
#[tokio::test]
async fn controller_routes_send_checkpoint_close_to_its_session() {
let s = store();
s.open("acme", "sess", "GET /a", None, now_unix_ms())
.await
.unwrap();
let ctl = ServerSessionController::new(s.clone(), "acme", "sess");
assert_eq!(ctl.send(b"hello".to_vec()).await.unwrap(), 1);
ctl.checkpoint(b"cp".to_vec()).await.unwrap();
assert_eq!(s.frames_since("acme", "sess", 0).await.unwrap().len(), 1);
assert_eq!(
s.resumed("acme", "sess").await.unwrap(),
Some(b"cp".to_vec())
);
ctl.close("done".into()).await.unwrap();
assert!(matches!(
ctl.send(b"x".to_vec()).await,
Err(boatramp_handlers::SessionError::Closed)
));
}
#[tokio::test]
async fn a_reaped_session_surfaces_as_closed() {
let s = store();
let ctl = ServerSessionController::new(s.clone(), "acme", "ghost");
assert!(matches!(
ctl.send(b"x".to_vec()).await,
Err(boatramp_handlers::SessionError::Closed)
));
}
}