use std::net::SocketAddr;
use crate::types::{DatabaseId, Lsn, TenantId};
use super::store::SessionStore;
impl SessionStore {
pub fn note_own_write(
&self,
addr: &SocketAddr,
database_id: DatabaseId,
tenant_id: TenantId,
collection: &str,
version: Lsn,
) {
if version == Lsn::ZERO {
return;
}
self.write_session(addr, |session| {
let slot = session
.own_write_versions
.entry((database_id, tenant_id, collection.to_string()))
.or_insert(Lsn::ZERO);
if version > *slot {
*slot = version;
}
});
}
pub fn own_write_version(
&self,
addr: &SocketAddr,
database_id: DatabaseId,
tenant_id: TenantId,
collection: &str,
) -> Lsn {
self.read_session(addr, |session| {
session
.own_write_versions
.get(&(database_id, tenant_id, collection.to_string()))
.copied()
.unwrap_or(Lsn::ZERO)
})
.unwrap_or(Lsn::ZERO)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn addr() -> SocketAddr {
"127.0.0.1:5601".parse().expect("test addr")
}
fn store_with_session() -> (SessionStore, SocketAddr) {
let sessions = SessionStore::new();
let a = addr();
sessions.ensure_session(a);
(sessions, a)
}
#[test]
fn absent_collection_returns_zero() {
let (sessions, a) = store_with_session();
assert_eq!(
sessions.own_write_version(&a, DatabaseId::DEFAULT, TenantId::new(1), "bread"),
Lsn::ZERO
);
}
#[test]
fn records_and_returns_own_write_version() {
let (sessions, a) = store_with_session();
sessions.note_own_write(
&a,
DatabaseId::DEFAULT,
TenantId::new(1),
"bread",
Lsn::new(2),
);
assert_eq!(
sessions.own_write_version(&a, DatabaseId::DEFAULT, TenantId::new(1), "bread"),
Lsn::new(2)
);
}
#[test]
fn keeps_the_maximum_version() {
let (sessions, a) = store_with_session();
sessions.note_own_write(
&a,
DatabaseId::DEFAULT,
TenantId::new(1),
"bread",
Lsn::new(5),
);
sessions.note_own_write(
&a,
DatabaseId::DEFAULT,
TenantId::new(1),
"bread",
Lsn::new(3),
);
assert_eq!(
sessions.own_write_version(&a, DatabaseId::DEFAULT, TenantId::new(1), "bread"),
Lsn::new(5)
);
}
#[test]
fn zero_version_is_ignored() {
let (sessions, a) = store_with_session();
sessions.note_own_write(
&a,
DatabaseId::DEFAULT,
TenantId::new(1),
"bread",
Lsn::ZERO,
);
assert_eq!(
sessions.own_write_version(&a, DatabaseId::DEFAULT, TenantId::new(1), "bread"),
Lsn::ZERO
);
}
#[test]
fn scoped_by_database_tenant_and_collection() {
let (sessions, a) = store_with_session();
sessions.note_own_write(
&a,
DatabaseId::DEFAULT,
TenantId::new(1),
"bread",
Lsn::new(7),
);
assert_eq!(
sessions.own_write_version(&a, DatabaseId::DEFAULT, TenantId::new(1), "milk"),
Lsn::ZERO
);
assert_eq!(
sessions.own_write_version(&a, DatabaseId::DEFAULT, TenantId::new(2), "bread"),
Lsn::ZERO
);
}
}