use std::ops::Range;
use std::sync::Arc;
use uuid::Uuid;
use web_time::Duration;
use crate::dbs::{Capabilities, DurableSession, Session};
use crate::key::root::se::{Se, SePrefix};
use crate::kvs::TransactionType::{Read, Write};
use crate::kvs::{Datastore, LockType};
async fn mem_ds() -> Arc<Datastore> {
Arc::new(
Datastore::builder()
.with_capabilities(Capabilities::all())
.build_with_path("memory")
.await
.unwrap(),
)
}
async fn count_range(ds: &Datastore, range: Range<Vec<u8>>) -> usize {
let tx = ds.transaction(Read, LockType::Optimistic).await.unwrap();
let res = tx.getr(range, None).await.unwrap();
let _ = tx.cancel().await;
res.len()
}
async fn durable_session_count(ds: &Datastore) -> usize {
let range = SePrefix {}.encode_range().unwrap();
count_range(ds, range).await
}
async fn durable_session_entry(ds: &Datastore, id: Uuid) -> Option<DurableSession> {
let tx = ds.transaction(Read, LockType::Optimistic).await.unwrap();
let res = tx.get(&Se::new(id), None).await.unwrap();
let _ = tx.cancel().await;
res
}
async fn put_entry_with_expiry(ds: &Datastore, id: Uuid, session: &Session, expires_at: u64) {
let value = DurableSession::from_session(session, expires_at);
let tx = ds.transaction(Write, LockType::Optimistic).await.unwrap();
tx.set(&Se::new(id), &value).await.unwrap();
tx.commit().await.unwrap();
}
fn sample_session(id: Uuid) -> Session {
let mut session = Session::owner().with_ns("app").with_db("app").with_rt(true);
session.id = Some(id);
session
}
#[tokio::test]
async fn persist_and_load_round_trip() {
let ds = mem_ds().await;
let id = Uuid::new_v4();
let session = sample_session(id);
ds.persist_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap();
assert_eq!(durable_session_count(&ds).await, 1);
let (loaded, _expires_at) =
ds.load_rpc_session(id).await.unwrap().expect("session should load");
assert_eq!(loaded, session);
assert!(ds.load_rpc_session(Uuid::new_v4()).await.unwrap().is_none());
assert_eq!(durable_session_count(&ds).await, 1);
}
#[tokio::test]
async fn expired_entry_is_deleted_on_load() {
let ds = mem_ds().await;
let id = Uuid::new_v4();
let session = sample_session(id);
put_entry_with_expiry(&ds, id, &session, 1).await;
assert!(ds.load_rpc_session(id).await.unwrap().is_none());
assert_eq!(durable_session_count(&ds).await, 0);
}
#[tokio::test]
async fn load_does_not_modify_the_expiry() {
let ds = mem_ds().await;
let id = Uuid::new_v4();
let session = sample_session(id);
ds.persist_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap();
let before = durable_session_entry(&ds, id).await.unwrap().expires_at;
let (loaded, reported) = ds.load_rpc_session(id).await.unwrap().expect("session should load");
assert_eq!(loaded, session);
assert_eq!(reported, before, "load should report the stored expiry it read");
assert_eq!(
durable_session_entry(&ds, id).await.unwrap().expires_at,
before,
"load must not touch the stored expiry"
);
}
#[tokio::test]
async fn delete_removes_the_entry_and_is_idempotent() {
let ds = mem_ds().await;
let id = Uuid::new_v4();
let session = sample_session(id);
ds.persist_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap();
ds.delete_rpc_session(id).await.unwrap();
assert_eq!(durable_session_count(&ds).await, 0);
assert!(ds.load_rpc_session(id).await.unwrap().is_none());
ds.delete_rpc_session(id).await.unwrap();
}
#[tokio::test]
async fn create_is_conditional_on_absence() {
let ds = mem_ds().await;
let id = Uuid::new_v4();
let session = sample_session(id);
assert!(ds.create_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap());
assert!(!ds.create_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap());
assert_eq!(durable_session_count(&ds).await, 1);
}
#[tokio::test]
async fn update_only_refreshes_a_present_entry_and_never_resurrects() {
let ds = mem_ds().await;
let id = Uuid::new_v4();
let session = sample_session(id);
assert!(!ds.update_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap());
assert_eq!(durable_session_count(&ds).await, 0);
ds.create_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap();
assert!(ds.update_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap());
ds.delete_rpc_session(id).await.unwrap();
assert!(!ds.update_rpc_session(id, &session, Duration::from_secs(60)).await.unwrap());
assert_eq!(durable_session_count(&ds).await, 0, "update must not resurrect a deleted session");
}
#[tokio::test]
async fn purge_deletes_only_expired_entries() {
let ds = mem_ds().await;
let expired_a = Uuid::new_v4();
let expired_b = Uuid::new_v4();
let live = Uuid::new_v4();
put_entry_with_expiry(&ds, expired_a, &sample_session(expired_a), 1).await;
put_entry_with_expiry(&ds, expired_b, &sample_session(expired_b), 2).await;
ds.persist_rpc_session(live, &sample_session(live), Duration::from_secs(3600)).await.unwrap();
assert_eq!(durable_session_count(&ds).await, 3);
ds.purge_expired_rpc_sessions(&Duration::from_secs(60)).await.unwrap();
assert_eq!(durable_session_count(&ds).await, 1);
assert!(durable_session_entry(&ds, live).await.is_some());
assert!(durable_session_entry(&ds, expired_a).await.is_none());
assert!(durable_session_entry(&ds, expired_b).await.is_none());
}