kcode-k1-groups-projection 0.1.0

Durable derived projection and authorization semantics for K1 groups
Documentation
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()
    );
}