use std::{
fs,
path::{Path, PathBuf},
sync::{
Arc, Barrier,
atomic::{AtomicU64, Ordering},
},
thread,
};
use super::*;
static NEXT_ROOT: AtomicU64 = AtomicU64::new(1);
fn user(value: u8) -> UserId {
UserId::from_tx_id(TxId::from_bytes([value; 12]))
}
fn model(value: u8) -> ModelId {
ModelId::from_bytes([value; 32])
}
fn test_root(name: &str) -> PathBuf {
let sequence = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
let root = std::env::temp_dir().join(format!(
"kcode-k1-groups-{name}-{}-{sequence}",
std::process::id()
));
let _ = fs::remove_dir_all(&root);
fs::create_dir_all(&root).unwrap();
root
}
fn open_stack(root: &Path) -> (Arc<K1TxnOrdering>, Arc<K1Peering>, K1Groups) {
let ordering = Arc::new(K1TxnOrdering::open(&root.join("ordering")).unwrap());
let peering = Arc::new(K1Peering::open(&root.join("peering"), ordering.clone()).unwrap());
let groups = K1Groups::open(&root.join("groups"), ordering.clone(), peering.clone()).unwrap();
(ordering, peering, groups)
}
#[test]
fn payload_round_trips_exact_version_one_shapes() {
let operation_id = [19_u8; 16];
let group = GroupId::new(TxId::from_bytes([7; 12]));
let actor = user(8);
let target = user(9);
let actions = vec![
Action::Create { owner: actor },
Action::SetUserRole {
group,
actor,
user: target,
role: None,
},
Action::SetUserRole {
group,
actor,
user: target,
role: Some(GroupRole::User),
},
Action::SetUserRole {
group,
actor,
user: target,
role: Some(GroupRole::Admin),
},
Action::SetUserRole {
group,
actor,
user: target,
role: Some(GroupRole::Owner),
},
Action::SetModelMembership {
group,
actor,
model: model(10),
present: false,
},
Action::SetModelMembership {
group,
actor,
model: model(10),
present: true,
},
];
for action in actions {
let payload = action.payload(operation_id);
let expected = match action {
Action::Create { .. } => 30,
Action::SetUserRole { .. } => 55,
Action::SetModelMembership { .. } => 75,
};
assert_eq!(payload.len(), expected);
assert_eq!(parse_payload(&payload).unwrap(), (operation_id, action));
}
}
#[test]
fn payload_parser_rejects_every_structural_class() {
assert!(parse_payload(&[]).is_err());
assert!(parse_payload(&[1]).is_err());
let operation_id = [23_u8; 16];
let group = GroupId::new(TxId::from_bytes([11; 12]));
let create = Action::Create { owner: user(1) }.payload(operation_id);
let role = Action::SetUserRole {
group,
actor: user(1),
user: user(2),
role: Some(GroupRole::Admin),
}
.payload(operation_id);
let membership = Action::SetModelMembership {
group,
actor: user(1),
model: model(3),
present: true,
}
.payload(operation_id);
for valid in [&create, &role, &membership] {
let mut short = valid.clone();
short.pop();
assert!(parse_payload(&short).is_err());
let mut trailing = valid.clone();
trailing.push(0);
assert!(parse_payload(&trailing).is_err());
}
let mut version = create.clone();
version[0] = 2;
assert!(parse_payload(&version).is_err());
let mut kind = create.clone();
kind[1] = 9;
assert!(parse_payload(&kind).is_err());
for invalid in [4, 255] {
let mut payload = role.clone();
payload[54] = invalid;
assert!(parse_payload(&payload).is_err());
}
for invalid in [2, 255] {
let mut payload = membership.clone();
payload[74] = invalid;
assert!(parse_payload(&payload).is_err());
}
}
#[test]
fn reconciliation_decisions_cover_all_submission_outcomes() {
let callback_txid = TxId::from_bytes([31; 12]);
let other_txid = TxId::from_bytes([32; 12]);
for kind in [
OutcomeKind::Applied,
OutcomeKind::Unchanged,
OutcomeKind::Rejected,
] {
let evidence = EvidenceSummary::One {
txid: callback_txid,
kind,
};
assert_eq!(
reconcile_decision(&Ok(callback_txid), evidence),
ReconcileDecision::UseCallback(kind)
);
assert_eq!(
reconcile_decision(&Err("committed error".to_string()), evidence),
ReconcileDecision::UseCallback(kind)
);
}
assert_eq!(
reconcile_decision(&Err("precommit".to_string()), EvidenceSummary::None),
ReconcileDecision::UseSubmissionError
);
assert!(matches!(
reconcile_decision(&Ok(callback_txid), EvidenceSummary::None),
ReconcileDecision::Fault(_)
));
assert!(matches!(
reconcile_decision(
&Ok(other_txid),
EvidenceSummary::One {
txid: callback_txid,
kind: OutcomeKind::Applied
}
),
ReconcileDecision::Fault(_)
));
assert!(matches!(
reconcile_decision(&Ok(callback_txid), EvidenceSummary::Ambiguous),
ReconcileDecision::Fault(_)
));
}
#[test]
fn full_stack_mutations_queries_restart_and_cursor() {
let root = test_root("full-stack");
let (ordering, peering, groups) = open_stack(&root);
let owner = user(1);
let admin = user(2);
let member = user(3);
let model = model(4);
let created = groups.create(owner).unwrap();
let group = created.group_id();
assert_eq!(group.txid(), created.txid());
let fetched = groups.get(group).unwrap().unwrap();
assert_eq!(fetched.id(), group);
assert_eq!(fetched.users().len(), 1);
assert_eq!(fetched.users()[0].user_id(), owner);
assert_eq!(fetched.users()[0].role(), GroupRole::Owner);
groups
.set_user_role(owner, group, admin, Some(GroupRole::Admin))
.unwrap();
groups
.set_user_role(admin, group, member, Some(GroupRole::User))
.unwrap();
assert!(
groups
.set_user_role(admin, group, member, Some(GroupRole::Admin))
.is_err()
);
assert!(
groups
.set_model_membership(admin, group, model, true)
.is_err()
);
groups
.set_model_membership(owner, group, model, true)
.unwrap();
let removed = groups
.set_model_membership(owner, group, model, false)
.unwrap();
let unchanged = groups
.set_model_membership(owner, group, model, false)
.unwrap();
assert_eq!(unchanged, removed);
groups
.set_model_membership(owner, group, model, true)
.unwrap();
groups.set_user_role(admin, group, member, None).unwrap();
groups
.set_user_role(owner, group, member, Some(GroupRole::User))
.unwrap();
let user_groups = groups.groups_for_user(member).unwrap();
assert_eq!(user_groups, vec![group]);
let model_groups = groups.groups_for_model(model).unwrap();
assert_eq!(model_groups, vec![group]);
let memberships = groups.memberships(member, model).unwrap();
assert_eq!(memberships.user_groups(), &[group]);
assert_eq!(memberships.model_groups(), &[group]);
assert_eq!(memberships.shared_groups(), &[group]);
let cursor_before_rejection = memberships.revision().unwrap();
assert!(groups.set_user_role(owner, group, owner, None).is_err());
assert!(groups.set_user_role(member, group, admin, None).is_err());
let cursor_after_rejection = groups
.memberships(member, model)
.unwrap()
.revision()
.unwrap();
assert_ne!(cursor_after_rejection, cursor_before_rejection);
drop(groups);
drop(peering);
drop(ordering);
let (ordering, peering, groups) = open_stack(&root);
assert_eq!(
groups.memberships(member, model).unwrap().revision(),
Some(cursor_after_rejection)
);
assert_eq!(
groups.get(group).unwrap().unwrap().revision().group_id(),
group
);
drop(groups);
drop(peering);
drop(ordering);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn many_to_many_reverse_and_intersection() {
let root = test_root("many-to-many");
let (ordering, peering, groups) = open_stack(&root);
let owner = user(1);
let member = user(2);
let first_model = model(1);
let second_model = model(2);
let first = groups.create(owner).unwrap().group_id();
let second = groups.create(owner).unwrap().group_id();
for group in [first, second] {
groups
.set_user_role(owner, group, member, Some(GroupRole::User))
.unwrap();
groups
.set_model_membership(owner, group, first_model, true)
.unwrap();
}
groups
.set_model_membership(owner, second, second_model, true)
.unwrap();
let user_groups = groups.groups_for_user(member).unwrap();
assert_eq!(user_groups.len(), 2);
assert!(user_groups.contains(&first));
assert!(user_groups.contains(&second));
let model_groups = groups.groups_for_model(first_model).unwrap();
assert_eq!(model_groups.len(), 2);
let memberships = groups.memberships(member, second_model).unwrap();
assert_eq!(memberships.user_groups().len(), 2);
assert_eq!(memberships.model_groups(), &[second]);
assert_eq!(memberships.shared_groups(), &[second]);
drop(groups);
drop(peering);
drop(ordering);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn concurrent_submissions_and_immediate_revocation() {
let root = test_root("concurrent");
let (ordering, peering, groups) = open_stack(&root);
let groups = Arc::new(groups);
let owner = user(1);
let revoked_admin = user(2);
let first = groups.create(owner).unwrap().group_id();
let second = groups.create(owner).unwrap().group_id();
groups
.set_user_role(owner, first, revoked_admin, Some(GroupRole::Admin))
.unwrap();
groups
.set_user_role(owner, first, revoked_admin, None)
.unwrap();
assert!(
groups
.set_user_role(revoked_admin, first, user(9), Some(GroupRole::User))
.is_err()
);
let barrier = Arc::new(Barrier::new(3));
let first_groups = groups.clone();
let first_barrier = barrier.clone();
let same_group = thread::spawn(move || {
first_barrier.wait();
first_groups.set_user_role(owner, first, user(3), Some(GroupRole::User))
});
let second_groups = groups.clone();
let second_barrier = barrier.clone();
let different_group = thread::spawn(move || {
second_barrier.wait();
second_groups.set_model_membership(owner, second, model(7), true)
});
barrier.wait();
assert!(same_group.join().unwrap().is_ok());
assert!(different_group.join().unwrap().is_ok());
let barrier = Arc::new(Barrier::new(3));
let first_groups = groups.clone();
let first_barrier = barrier.clone();
let same_first = thread::spawn(move || {
first_barrier.wait();
first_groups.set_user_role(owner, first, user(4), Some(GroupRole::User))
});
let second_groups = groups.clone();
let second_barrier = barrier.clone();
let same_second = thread::spawn(move || {
second_barrier.wait();
second_groups.set_user_role(owner, first, user(5), Some(GroupRole::Admin))
});
barrier.wait();
assert!(same_first.join().unwrap().is_ok());
assert!(same_second.join().unwrap().is_ok());
assert_eq!(groups.get(first).unwrap().unwrap().users().len(), 4);
assert_eq!(groups.get(second).unwrap().unwrap().models(), &[model(7)]);
drop(groups);
drop(peering);
drop(ordering);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn malformed_canonical_callback_faults_and_replay_fails_closed() {
let root = test_root("malformed");
let (ordering, peering, groups) = open_stack(&root);
let group = groups.create(user(1)).unwrap().group_id();
let error = peering
.submit_txn(SubsystemId::from_str(SUBSYSTEM_NAME).unwrap(), &[1, 9])
.unwrap_err();
assert!(error.contains("committed"));
assert!(groups.get(group).is_err());
drop(groups);
let reopened = K1Groups::open(&root.join("groups"), ordering.clone(), peering.clone());
assert!(reopened.is_err());
drop(peering);
drop(ordering);
fs::remove_dir_all(root).unwrap();
}
#[test]
fn reorg_callback_clears_projection_and_invalidates_facade() {
let root = test_root("reorg");
let (ordering, peering, groups) = open_stack(&root);
let group = groups.create(user(1)).unwrap().group_id();
let callback = GroupsSubsystem {
projection: groups.projection.clone(),
shared: groups.shared.clone(),
};
callback.reorg().unwrap();
assert!(groups.projection.get(group).unwrap().is_none());
assert!(groups.get(group).is_err());
drop(groups);
drop(peering);
drop(ordering);
fs::remove_dir_all(root).unwrap();
}