use std::{
sync::atomic::{AtomicU64, Ordering},
time::{Duration, Instant},
};
use kcode_k1_groups_projection::{GroupId, ModelId, Projection, TxId, UserId};
use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId, build_signed_transaction};
use kcode_k1_txn_ordering::K1TxnOrdering;
use rusqlite::{Connection, TransactionBehavior, params};
use tempfile::TempDir;
const SUBSYSTEM: &str = "k1-groups-subsystem";
const GROUPS: usize = 10_000;
const HUMANS_PER_GROUP: usize = 10;
const MODELS_PER_GROUP: usize = 10;
const OPEN_LIMIT: Duration = Duration::from_secs(5);
const LOOKUP_LIMIT: Duration = Duration::from_millis(200);
static TIMESTAMP: AtomicU64 = AtomicU64::new(1);
fn callback(ordering: &K1TxnOrdering) -> TxId {
let timestamp = TIMESTAMP.fetch_add(1, Ordering::SeqCst);
let transaction = build_signed_transaction(
ordering.tip().unwrap_or(GENESIS_PARENT),
timestamp,
[7; 32],
SubsystemId::from_str(SUBSYSTEM).unwrap(),
&[],
|_| Ok([0; 64]),
)
.unwrap();
let txid = TxId::for_transaction(&transaction);
ordering.submit_txn(&transaction).unwrap();
txid
}
fn user(value: usize) -> UserId {
let mut bytes = [0; 12];
bytes[..8].copy_from_slice(&(value as u64).to_le_bytes());
bytes[8..].copy_from_slice(b"user");
UserId::from_tx_id(TxId::from_bytes(bytes))
}
fn model(value: usize) -> ModelId {
let mut bytes = [0; 32];
bytes[..8].copy_from_slice(&(value as u64).to_le_bytes());
bytes[8..13].copy_from_slice(b"model");
ModelId::from_bytes(bytes)
}
fn measure<T, F>(mut operation: F) -> (Duration, T)
where
F: FnMut() -> T,
{
let mut longest = Duration::ZERO;
let mut result = None;
for _ in 0..5 {
let started = Instant::now();
result = Some(operation());
let elapsed = started.elapsed();
if elapsed > longest {
longest = elapsed;
}
}
(longest, result.unwrap())
}
fn main() {
let temporary = TempDir::new().unwrap();
let ordering = K1TxnOrdering::open(&temporary.path().join("ordering")).unwrap();
let projection_root = temporary.path().join("projection");
let (projection, cursor) = Projection::open(&projection_root, &ordering).unwrap();
assert_eq!(cursor, None);
drop(projection);
let fixture_started = Instant::now();
let mut groups = Vec::with_capacity(GROUPS);
for _ in 0..GROUPS {
groups.push(callback(&ordering));
}
let mut connection = Connection::open(projection_root.join("groups.sqlite3")).unwrap();
connection
.execute_batch("PRAGMA synchronous = FULL; PRAGMA foreign_keys = ON;")
.unwrap();
let transaction = connection
.transaction_with_behavior(TransactionBehavior::Immediate)
.unwrap();
{
let mut insert_group = transaction
.prepare("INSERT INTO groups (group_id, revision) VALUES (?1, ?1)")
.unwrap();
let mut insert_user = transaction
.prepare("INSERT INTO user_memberships (group_id, user_id, role) VALUES (?1, ?2, ?3)")
.unwrap();
let mut insert_model = transaction
.prepare("INSERT INTO model_memberships (group_id, model_id) VALUES (?1, ?2)")
.unwrap();
for group in &groups {
insert_group
.execute(params![group.as_bytes().as_slice()])
.unwrap();
for index in 0..HUMANS_PER_GROUP {
let identity = user(index).as_tx_id();
let role = if index == 0 { 2_i64 } else { 0_i64 };
insert_user
.execute(params![
group.as_bytes().as_slice(),
identity.as_bytes().as_slice(),
role
])
.unwrap();
}
for index in 0..MODELS_PER_GROUP {
insert_model
.execute(params![
group.as_bytes().as_slice(),
model(index).as_bytes().as_slice()
])
.unwrap();
}
}
}
transaction
.execute(
"UPDATE metadata SET last_applied_txid = ?1 WHERE singleton = 1",
params![groups.last().unwrap().as_bytes().as_slice()],
)
.unwrap();
transaction.commit().unwrap();
drop(connection);
let fixture_elapsed = fixture_started.elapsed();
let mut valid_open_max = Duration::ZERO;
let mut samples = 0;
let (projection, cursor) = loop {
let started = Instant::now();
let opened = Projection::open(&projection_root, &ordering).unwrap();
let elapsed = started.elapsed();
if elapsed > valid_open_max {
valid_open_max = elapsed;
}
assert!(
elapsed < OPEN_LIMIT,
"valid open took {} microseconds",
elapsed.as_micros()
);
samples += 1;
if samples == 3 {
break opened;
}
drop(opened);
};
assert_eq!(cursor, groups.last().copied());
let selected_group = GroupId::new(groups[GROUPS / 2]);
let (get_max, group) = measure(|| projection.get(selected_group).unwrap().unwrap());
let (user_max, user_groups) = measure(|| projection.groups_for_user(user(1)).unwrap());
let (model_max, model_groups) = measure(|| projection.groups_for_model(model(1)).unwrap());
let (memberships_max, memberships) =
measure(|| projection.memberships(user(1), model(1)).unwrap());
for (name, elapsed) in [
("get", get_max),
("groups_for_user", user_max),
("groups_for_model", model_max),
("memberships", memberships_max),
] {
assert!(
elapsed < LOOKUP_LIMIT,
"{name} took {} microseconds",
elapsed.as_micros()
);
}
assert_eq!(group.id(), selected_group);
assert_eq!(group.users().len(), HUMANS_PER_GROUP);
assert_eq!(group.models().len(), MODELS_PER_GROUP);
assert_eq!(user_groups.len(), GROUPS);
assert_eq!(model_groups.len(), GROUPS);
assert_eq!(memberships.revision(), groups.last().copied());
assert_eq!(memberships.user_groups().len(), GROUPS);
assert_eq!(memberships.model_groups().len(), GROUPS);
assert_eq!(memberships.shared_groups().len(), GROUPS);
drop(projection);
let connection = Connection::open(projection_root.join("groups.sqlite3")).unwrap();
let last_user = user(HUMANS_PER_GROUP - 1).as_tx_id();
assert_eq!(
connection
.execute(
"UPDATE user_memberships SET role = 99 WHERE group_id = ?1 AND user_id = ?2",
params![
groups.last().unwrap().as_bytes().as_slice(),
last_user.as_bytes().as_slice()
],
)
.unwrap(),
1
);
drop(connection);
let rebuild_started = Instant::now();
let (rebuilt, cursor) = Projection::open(&projection_root, &ordering).unwrap();
let rebuild_elapsed = rebuild_started.elapsed();
assert!(
rebuild_elapsed < OPEN_LIMIT,
"rebuild open took {} microseconds",
rebuild_elapsed.as_micros()
);
assert_eq!(cursor, None);
assert!(rebuilt.get(selected_group).unwrap().is_none());
assert!(rebuilt.groups_for_user(user(1)).unwrap().is_empty());
println!(
"realistic_scale groups={GROUPS} human_memberships={} model_memberships={} fixture_us={} valid_open_max_us={} rebuild_open_us={} get_max_us={} user_lookup_max_us={} model_lookup_max_us={} memberships_lookup_max_us={}",
GROUPS * HUMANS_PER_GROUP,
GROUPS * MODELS_PER_GROUP,
fixture_elapsed.as_micros(),
valid_open_max.as_micros(),
rebuild_elapsed.as_micros(),
get_max.as_micros(),
user_max.as_micros(),
model_max.as_micros(),
memberships_max.as_micros()
);
}