use base64::{Engine as _, engine::general_purpose::URL_SAFE_NO_PAD};
use hmac::{Hmac, Mac};
use kcode_k1_invite_projection::{ApplyOutcome, InviteAction, InviteProjection};
use kcode_k1_peering::K1Peering;
use kcode_k1_transaction::SubsystemId;
use kcode_k1_transaction_id::TxId;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem};
use sha2::Sha256;
use std::collections::{HashMap, HashSet};
use std::path::Path;
use std::str::FromStr;
use std::sync::{Arc, Condvar, Mutex, MutexGuard};
use zeroize::Zeroize;
const SUBSYSTEM_NAME: &str = "k1-invites-subsystem";
const ISSUE_LENGTH: usize = 34;
const CONSUME_HEADER_LENGTH: usize = 78;
const REOPEN_REQUIRED: &str = "k1-invites is unavailable; reopen required";
type HmacSha256 = Hmac<Sha256>;
pub struct InviteCode {
bytes: [u8; 6],
}
impl InviteCode {
pub fn expose(&self) -> String {
URL_SAFE_NO_PAD.encode(self.bytes.as_slice())
}
}
impl FromStr for InviteCode {
type Err = String;
fn from_str(value: &str) -> Result<Self, Self::Err> {
let mut code = InviteCode { bytes: [0; 6] };
if value.len() != 8 {
return Err("invalid invite code".to_owned());
}
let written = match URL_SAFE_NO_PAD.decode_slice(value, &mut code.bytes) {
Ok(written) => written,
Err(_) => return Err("invalid invite code".to_owned()),
};
if written != code.bytes.len() {
return Err("invalid invite code".to_owned());
}
let mut canonical = [0_u8; 8];
let encoded = match URL_SAFE_NO_PAD.encode_slice(&code.bytes, &mut canonical) {
Ok(encoded) => encoded,
Err(_) => {
canonical.zeroize();
return Err("invalid invite code".to_owned());
}
};
if encoded != canonical.len() || canonical.as_slice() != value.as_bytes() {
canonical.zeroize();
return Err("invalid invite code".to_owned());
}
canonical.zeroize();
Ok(code)
}
}
impl Drop for InviteCode {
fn drop(&mut self) {
self.bytes.zeroize();
}
}
pub struct InviteVerifierKey {
bytes: [u8; 32],
}
impl InviteVerifierKey {
pub fn from_bytes(bytes: [u8; 32]) -> Self {
Self { bytes }
}
}
impl Drop for InviteVerifierKey {
fn drop(&mut self) {
self.bytes.zeroize();
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum InviteStatus {
Unknown,
Unused,
Consumed,
}
#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
pub struct UserId(TxId);
impl UserId {
pub const fn from_tx_id(txid: TxId) -> Self {
Self(txid)
}
pub const fn as_tx_id(self) -> TxId {
self.0
}
}
#[derive(Clone, Copy, Eq, Hash, PartialEq)]
pub struct RegistrationKey([u8; 32]);
impl RegistrationKey {
pub fn from_bytes(bytes: [u8; 32]) -> Self {
Self(bytes)
}
pub fn as_bytes(&self) -> &[u8; 32] {
&self.0
}
}
#[derive(Clone, Eq, PartialEq)]
pub struct Registration {
user_id: UserId,
registration_key: RegistrationKey,
data: Arc<Vec<u8>>,
}
impl Registration {
pub fn user_id(&self) -> UserId {
self.user_id
}
pub fn registration_key(&self) -> RegistrationKey {
self.registration_key
}
pub fn data(&self) -> &[u8] {
self.data.as_slice()
}
}
pub struct K1Invites {
subsystem: Arc<InviteSubsystem>,
_ordering: Arc<K1TxnOrdering>,
peering: Arc<K1Peering>,
verifier_key: InviteVerifierKey,
}
impl K1Invites {
pub fn open(
root: &Path,
ordering: Arc<K1TxnOrdering>,
peering: Arc<K1Peering>,
verifier_key: InviteVerifierKey,
) -> Result<Self, String> {
let projection = InviteProjection::open(root)?;
let snapshot = projection.snapshot()?;
let checkpoint = snapshot.checkpoint;
let mut state = FacadeState::new();
state
.issues
.try_reserve(snapshot.issues.len())
.map_err(|_| allocation_error())?;
state
.issue_commitments
.try_reserve(snapshot.issues.len())
.map_err(|_| allocation_error())?;
state
.registrations
.try_reserve(snapshot.registrations.len())
.map_err(|_| allocation_error())?;
state
.registration_keys
.try_reserve(snapshot.registrations.len())
.map_err(|_| allocation_error())?;
for issue in snapshot.issues {
if state.issues.contains_key(&issue.commitment)
|| state.issue_commitments.contains_key(&issue.issue_id)
{
return Err("invite projection snapshot is inconsistent".to_owned());
}
state.issues.insert(issue.commitment, issue.issue_id);
state
.issue_commitments
.insert(issue.issue_id, issue.commitment);
}
for accepted in snapshot.registrations {
let commitment = accepted.commitment;
if state.issue_commitments.get(&accepted.issue_id) != Some(&commitment) {
return Err("invite projection snapshot is inconsistent".to_owned());
}
let registration_key = RegistrationKey::from_bytes(accepted.registration_key);
if state.registrations.contains_key(&commitment)
|| state.registration_keys.contains_key(®istration_key)
{
return Err("invite projection snapshot is inconsistent".to_owned());
}
state.registrations.insert(
commitment,
Registration {
user_id: UserId(accepted.user_id),
registration_key,
data: Arc::new(accepted.data),
},
);
state.registration_keys.insert(registration_key, commitment);
}
let subsystem = Arc::new(InviteSubsystem {
projection: Mutex::new(projection),
state: Mutex::new(state),
});
let callback: Arc<dyn Subsystem> = subsystem.clone();
ordering.register_subsystem(subsystem_id()?, checkpoint, callback)?;
subsystem.ensure_available()?;
Ok(Self {
subsystem,
_ordering: ordering,
peering,
verifier_key,
})
}
pub fn status(&self, code: &InviteCode) -> Result<InviteStatus, String> {
let commitment = invite_commitment(&self.verifier_key, &code.bytes)?;
let state = lock_unpoison(&self.subsystem.state);
if !state.available {
return Err(REOPEN_REQUIRED.to_owned());
}
Ok(status_from_state(&state, commitment))
}
pub fn create(&self) -> Result<(TxId, InviteCode), String> {
enum Resolution {
Return(Result<TxId, String>),
Retry,
}
loop {
self.subsystem.ensure_available()?;
let subsystem_id = subsystem_id()?;
let mut code = InviteCode { bytes: [0_u8; 6] };
getrandom::fill(&mut code.bytes).map_err(|error| error.to_string())?;
let commitment = invite_commitment(&self.verifier_key, &code.bytes)?;
let reserved = {
let mut state = lock_unpoison(&self.subsystem.state);
if !state.available {
return Err(REOPEN_REQUIRED.to_owned());
}
if state.issues.contains_key(&commitment)
|| state.pending_issues.contains(&commitment)
{
false
} else {
state
.pending_issues
.try_reserve(1)
.map_err(|_| allocation_error())?;
state.pending_issues.insert(commitment);
true
}
};
if !reserved {
continue;
}
let payload = issue_payload(commitment);
let submission = self.peering.submit_txn(subsystem_id, &payload);
let resolution = {
let mut state = lock_unpoison(&self.subsystem.state);
let pending_was_present = state.pending_issues.remove(&commitment);
if !pending_was_present {
state.available = false;
Resolution::Return(Err(REOPEN_REQUIRED.to_owned()))
} else if !state.available {
Resolution::Return(Err(REOPEN_REQUIRED.to_owned()))
} else {
match submission {
Ok(returned_id) => match state.issues.get(&commitment).copied() {
Some(accepted_id) if accepted_id == returned_id => {
Resolution::Return(Ok(returned_id))
}
Some(_) => Resolution::Retry,
None => Resolution::Return(Err(
"issue transaction was not the winning Issue".to_owned(),
)),
},
Err(error) => match state.issues.get(&commitment).copied() {
Some(accepted_id) => Resolution::Return(Ok(accepted_id)),
None => Resolution::Return(Err(error)),
},
}
}
};
match resolution {
Resolution::Return(result) => return result.map(|id| (id, code)),
Resolution::Retry => continue,
}
}
}
pub fn consume_with_data(
&self,
code: &InviteCode,
registration_key: RegistrationKey,
data: &[u8],
) -> Result<UserId, String> {
self.subsystem.ensure_available()?;
let commitment = invite_commitment(&self.verifier_key, &code.bytes)?;
let shared_data = Arc::new(fallible_copy(data)?);
let candidate_cell = Arc::new(PendingResult::new());
let role = {
let mut state = lock_unpoison(&self.subsystem.state);
if !state.available {
return Err(REOPEN_REQUIRED.to_owned());
}
if let Some(accepted) = state.registrations.get(&commitment) {
ConsumeRole::Accepted(accepted.clone())
} else if let Some(pending) = state.pending_consumes.get(&commitment) {
ConsumeRole::Pending {
registration_key: pending.registration_key,
data: pending.data.clone(),
cell: pending.cell.clone(),
}
} else {
let issue_id = state
.issues
.get(&commitment)
.copied()
.ok_or_else(|| "unknown invite code".to_owned())?;
if state.registration_keys.contains_key(®istration_key)
|| state
.pending_registration_keys
.contains_key(®istration_key)
{
return Err("registration key is already used by another invite".to_owned());
}
state
.pending_consumes
.try_reserve(1)
.map_err(|_| allocation_error())?;
state
.pending_registration_keys
.try_reserve(1)
.map_err(|_| allocation_error())?;
state.pending_consumes.insert(
commitment,
PendingConsume {
registration_key,
data: shared_data.clone(),
cell: candidate_cell.clone(),
},
);
state
.pending_registration_keys
.insert(registration_key, commitment);
ConsumeRole::Leader {
issue_id,
data: shared_data.clone(),
cell: candidate_cell.clone(),
}
}
};
match role {
ConsumeRole::Accepted(accepted) => {
if accepted.registration_key == registration_key && accepted.data.as_slice() == data
{
Ok(accepted.user_id)
} else {
Err("invite was already consumed with different registration data".to_owned())
}
}
ConsumeRole::Pending {
registration_key: pending_key,
data: pending_data,
cell,
} => {
if pending_key != registration_key || pending_data.as_slice() != data {
return Err("invite has a conflicting consume request in progress".to_owned());
}
wait_for_pending(&cell)
}
ConsumeRole::Leader {
issue_id,
data,
cell,
} => {
let payload = match consume_payload(
issue_id,
commitment,
registration_key,
data.as_slice(),
) {
Ok(payload) => payload,
Err(error) => {
return self.abandon_pending(commitment, registration_key, &cell, error);
}
};
let subsystem_id = match subsystem_id() {
Ok(subsystem_id) => subsystem_id,
Err(error) => {
return self.abandon_pending(commitment, registration_key, &cell, error);
}
};
if let Err(error) = self.subsystem.ensure_available() {
return self.abandon_pending(commitment, registration_key, &cell, error);
}
let submission = self.peering.submit_txn(subsystem_id, &payload);
self.complete_pending(
commitment,
registration_key,
data.as_slice(),
&cell,
submission,
)
}
}
}
pub fn registrations(&self) -> Result<Vec<Registration>, String> {
self.subsystem.ensure_available()?;
let snapshot = {
let projection = lock_unpoison(&self.subsystem.projection);
projection.snapshot()
};
let snapshot = match snapshot {
Ok(snapshot) => snapshot,
Err(error) => {
self.subsystem.mark_unavailable();
return Err(error);
}
};
self.subsystem.ensure_available()?;
let mut registrations = Vec::new();
registrations
.try_reserve_exact(snapshot.registrations.len())
.map_err(|_| allocation_error())?;
for accepted in snapshot.registrations {
registrations.push(Registration {
user_id: UserId(accepted.user_id),
registration_key: RegistrationKey::from_bytes(accepted.registration_key),
data: Arc::new(accepted.data),
});
}
Ok(registrations)
}
fn abandon_pending(
&self,
commitment: [u8; 32],
registration_key: RegistrationKey,
cell: &Arc<PendingResult>,
error: String,
) -> Result<UserId, String> {
let consistent = self.remove_pending(commitment, registration_key, cell);
let result = if consistent {
Err(error)
} else {
Err(REOPEN_REQUIRED.to_owned())
};
publish_pending(cell, &result);
result
}
fn complete_pending(
&self,
commitment: [u8; 32],
registration_key: RegistrationKey,
data: &[u8],
cell: &Arc<PendingResult>,
submission: Result<TxId, String>,
) -> Result<UserId, String> {
let (available, accepted, consistent) = {
let mut state = lock_unpoison(&self.subsystem.state);
let available = state.available;
let accepted = state.registrations.get(&commitment).cloned();
let consistent = pending_matches(&state, commitment, registration_key, cell);
state.pending_consumes.remove(&commitment);
if state.pending_registration_keys.get(®istration_key) == Some(&commitment) {
state.pending_registration_keys.remove(®istration_key);
}
if !consistent {
state.available = false;
}
(available, accepted, consistent)
};
let result = if !available || !consistent {
Err(REOPEN_REQUIRED.to_owned())
} else {
match submission {
Err(error) => match accepted {
Some(accepted)
if accepted.registration_key == registration_key
&& accepted.data.as_slice() == data =>
{
Ok(accepted.user_id)
}
_ => Err(error),
},
Ok(returned_id) => {
if let Some(accepted) = accepted {
if accepted.user_id.as_tx_id() == returned_id
&& accepted.registration_key == registration_key
&& accepted.data.as_slice() == data
{
Ok(accepted.user_id)
} else {
Err("consume transaction was not the accepted registration".to_owned())
}
} else {
Err("consume transaction was a semantic loser".to_owned())
}
}
}
};
publish_pending(cell, &result);
result
}
fn remove_pending(
&self,
commitment: [u8; 32],
registration_key: RegistrationKey,
cell: &Arc<PendingResult>,
) -> bool {
let mut state = lock_unpoison(&self.subsystem.state);
let consistent = pending_matches(&state, commitment, registration_key, cell);
state.pending_consumes.remove(&commitment);
if state.pending_registration_keys.get(®istration_key) == Some(&commitment) {
state.pending_registration_keys.remove(®istration_key);
}
if !consistent {
state.available = false;
}
consistent
}
}
struct InviteSubsystem {
projection: Mutex<InviteProjection>,
state: Mutex<FacadeState>,
}
impl InviteSubsystem {
fn ensure_available(&self) -> Result<(), String> {
if lock_unpoison(&self.state).available {
Ok(())
} else {
Err(REOPEN_REQUIRED.to_owned())
}
}
fn apply_if_available(&self, action: InviteAction) -> Result<ApplyOutcome, String> {
let projection = lock_unpoison(&self.projection);
self.ensure_available()?;
projection.apply(action)
}
fn mark_unavailable(&self) {
lock_unpoison(&self.state).available = false;
}
fn fault(&self, error: String) -> Result<(), String> {
self.mark_unavailable();
Err(error)
}
fn accept_issue(&self, id: TxId, commitment: [u8; 32]) -> Result<(), String> {
let mut state = lock_unpoison(&self.state);
if !state.available {
return Err(REOPEN_REQUIRED.to_owned());
}
if state.issues.contains_key(&commitment) || state.issue_commitments.contains_key(&id) {
state.available = false;
return Err("accepted Issue contradicted invite indexes".to_owned());
}
if state.issues.try_reserve(1).is_err() || state.issue_commitments.try_reserve(1).is_err() {
state.available = false;
return Err(allocation_error());
}
state.issues.insert(commitment, id);
state.issue_commitments.insert(id, commitment);
Ok(())
}
fn accept_registration(
&self,
id: TxId,
issue_id: TxId,
commitment: [u8; 32],
registration_key: [u8; 32],
data: &[u8],
) -> Result<(), String> {
let copied_data = match fallible_copy(data) {
Ok(data) => Arc::new(data),
Err(error) => return self.fault(error),
};
let registration_key = RegistrationKey::from_bytes(registration_key);
let mut state = lock_unpoison(&self.state);
if !state.available {
return Err(REOPEN_REQUIRED.to_owned());
}
if state.issues.get(&commitment) != Some(&issue_id)
|| state.issue_commitments.get(&issue_id) != Some(&commitment)
|| state.registrations.contains_key(&commitment)
|| state.registration_keys.contains_key(®istration_key)
{
state.available = false;
return Err("accepted registration contradicted invite indexes".to_owned());
}
if state.registrations.try_reserve(1).is_err()
|| state.registration_keys.try_reserve(1).is_err()
{
state.available = false;
return Err(allocation_error());
}
state.registrations.insert(
commitment,
Registration {
user_id: UserId(id),
registration_key,
data: copied_data,
},
);
state.registration_keys.insert(registration_key, commitment);
Ok(())
}
}
impl Subsystem for InviteSubsystem {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
self.ensure_available()?;
let parsed = match parse_payload(id, payload) {
Ok(parsed) => parsed,
Err(error) => return self.fault(error),
};
let action = match projection_action(parsed) {
Ok(action) => action,
Err(error) => return self.fault(error),
};
let outcome = match self.apply_if_available(action) {
Ok(outcome) => outcome,
Err(error) => return self.fault(error),
};
match outcome {
ApplyOutcome::IssueAccepted(issue) => match parsed {
ParsedAction::Issue { id, commitment }
if issue.issue_id == id && issue.commitment == commitment =>
{
self.accept_issue(issue.issue_id, issue.commitment)
}
_ => self.fault("projection Issue outcome contradicted action".to_owned()),
},
ApplyOutcome::RegistrationAccepted(accepted) => match parsed {
ParsedAction::Consume {
id,
issue_id,
commitment,
registration_key,
data,
} if accepted.user_id == id
&& accepted.issue_id == issue_id
&& accepted.commitment == commitment
&& accepted.registration_key == registration_key
&& accepted.data.as_slice() == data =>
{
self.accept_registration(
accepted.user_id,
accepted.issue_id,
accepted.commitment,
accepted.registration_key,
accepted.data.as_slice(),
)
}
_ => self.fault("projection registration outcome contradicted action".to_owned()),
},
ApplyOutcome::Noop => Ok(()),
}
}
fn reorg(&self) -> Result<(), String> {
self.mark_unavailable();
let projection = lock_unpoison(&self.projection);
projection.discard()
}
}
struct FacadeState {
available: bool,
issues: HashMap<[u8; 32], TxId>,
issue_commitments: HashMap<TxId, [u8; 32]>,
registrations: HashMap<[u8; 32], Registration>,
registration_keys: HashMap<RegistrationKey, [u8; 32]>,
pending_issues: HashSet<[u8; 32]>,
pending_consumes: HashMap<[u8; 32], PendingConsume>,
pending_registration_keys: HashMap<RegistrationKey, [u8; 32]>,
}
impl FacadeState {
fn new() -> Self {
Self {
available: true,
issues: HashMap::new(),
issue_commitments: HashMap::new(),
registrations: HashMap::new(),
registration_keys: HashMap::new(),
pending_issues: HashSet::new(),
pending_consumes: HashMap::new(),
pending_registration_keys: HashMap::new(),
}
}
}
struct PendingConsume {
registration_key: RegistrationKey,
data: Arc<Vec<u8>>,
cell: Arc<PendingResult>,
}
struct PendingResult {
outcome: Mutex<Option<Result<UserId, String>>>,
ready: Condvar,
}
impl PendingResult {
fn new() -> Self {
Self {
outcome: Mutex::new(None),
ready: Condvar::new(),
}
}
}
enum ConsumeRole {
Accepted(Registration),
Pending {
registration_key: RegistrationKey,
data: Arc<Vec<u8>>,
cell: Arc<PendingResult>,
},
Leader {
issue_id: TxId,
data: Arc<Vec<u8>>,
cell: Arc<PendingResult>,
},
}
#[derive(Clone, Copy)]
enum ParsedAction<'a> {
Issue {
id: TxId,
commitment: [u8; 32],
},
Consume {
id: TxId,
issue_id: TxId,
commitment: [u8; 32],
registration_key: [u8; 32],
data: &'a [u8],
},
}
fn parse_payload<'a>(id: TxId, payload: &'a [u8]) -> Result<ParsedAction<'a>, String> {
if payload.len() < 2 || payload[0] != 2 {
return Err("malformed k1-invites payload".to_owned());
}
match payload[1] {
1 if payload.len() == ISSUE_LENGTH => Ok(ParsedAction::Issue {
id,
commitment: array_32(&payload[2..34]),
}),
2 if payload.len() >= CONSUME_HEADER_LENGTH => Ok(ParsedAction::Consume {
id,
issue_id: tx_id_from_slice(&payload[2..14]),
commitment: array_32(&payload[14..46]),
registration_key: array_32(&payload[46..78]),
data: &payload[78..],
}),
_ => Err("malformed k1-invites payload".to_owned()),
}
}
fn projection_action(parsed: ParsedAction<'_>) -> Result<InviteAction, String> {
match parsed {
ParsedAction::Issue { id, commitment } => Ok(InviteAction::Issue { id, commitment }),
ParsedAction::Consume {
id,
issue_id,
commitment,
registration_key,
data,
} => Ok(InviteAction::Consume {
id,
issue_id,
commitment,
registration_key,
data: fallible_copy(data)?,
}),
}
}
fn invite_commitment(verifier_key: &InviteVerifierKey, code: &[u8; 6]) -> Result<[u8; 32], String> {
let mut mac = HmacSha256::new_from_slice(&verifier_key.bytes)
.map_err(|_| "invalid invite verifier key".to_owned())?;
mac.update(b"k1-invite-v2");
mac.update(code);
let output = mac.finalize().into_bytes();
let mut commitment = [0_u8; 32];
commitment.copy_from_slice(&output);
Ok(commitment)
}
fn status_from_state(state: &FacadeState, commitment: [u8; 32]) -> InviteStatus {
if state.registrations.contains_key(&commitment) {
InviteStatus::Consumed
} else if state.issues.contains_key(&commitment) {
InviteStatus::Unused
} else {
InviteStatus::Unknown
}
}
fn issue_payload(commitment: [u8; 32]) -> [u8; ISSUE_LENGTH] {
let mut payload = [0_u8; ISSUE_LENGTH];
payload[0] = 2;
payload[1] = 1;
payload[2..].copy_from_slice(&commitment);
payload
}
fn consume_payload(
issue_id: TxId,
commitment: [u8; 32],
registration_key: RegistrationKey,
data: &[u8],
) -> Result<Vec<u8>, String> {
let total = CONSUME_HEADER_LENGTH
.checked_add(data.len())
.ok_or_else(|| "consume payload length overflow".to_owned())?;
let mut payload = Vec::new();
payload
.try_reserve_exact(total)
.map_err(|_| allocation_error())?;
payload.extend_from_slice(&[2, 2]);
payload.extend_from_slice(issue_id.as_bytes());
payload.extend_from_slice(&commitment);
payload.extend_from_slice(registration_key.as_bytes());
payload.extend_from_slice(data);
Ok(payload)
}
fn subsystem_id() -> Result<SubsystemId, String> {
SubsystemId::from_str(SUBSYSTEM_NAME).map_err(|error| error.to_string())
}
fn tx_id_from_slice(bytes: &[u8]) -> TxId {
let mut raw = [0_u8; 12];
raw.copy_from_slice(bytes);
TxId::from_bytes(raw)
}
fn array_32(bytes: &[u8]) -> [u8; 32] {
let mut array = [0_u8; 32];
array.copy_from_slice(bytes);
array
}
fn fallible_copy(bytes: &[u8]) -> Result<Vec<u8>, String> {
let mut copied = Vec::new();
copied
.try_reserve_exact(bytes.len())
.map_err(|_| allocation_error())?;
copied.extend_from_slice(bytes);
Ok(copied)
}
fn allocation_error() -> String {
"memory allocation failed".to_owned()
}
fn pending_matches(
state: &FacadeState,
commitment: [u8; 32],
registration_key: RegistrationKey,
cell: &Arc<PendingResult>,
) -> bool {
state
.pending_consumes
.get(&commitment)
.is_some_and(|pending| {
pending.registration_key == registration_key && Arc::ptr_eq(&pending.cell, cell)
})
&& state.pending_registration_keys.get(®istration_key) == Some(&commitment)
}
fn publish_pending(cell: &PendingResult, result: &Result<UserId, String>) {
{
let mut outcome = lock_unpoison(&cell.outcome);
*outcome = Some(result.clone());
}
cell.ready.notify_all();
}
fn wait_for_pending(cell: &PendingResult) -> Result<UserId, String> {
let mut outcome = lock_unpoison(&cell.outcome);
loop {
if let Some(result) = outcome.as_ref() {
return result.clone();
}
outcome = match cell.ready.wait(outcome) {
Ok(outcome) => outcome,
Err(poisoned) => poisoned.into_inner(),
};
}
}
fn lock_unpoison<T>(mutex: &Mutex<T>) -> MutexGuard<'_, T> {
match mutex.lock() {
Ok(guard) => guard,
Err(poisoned) => poisoned.into_inner(),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn tx_id(byte: u8) -> TxId {
TxId::from_bytes([byte; 12])
}
#[test]
fn user_id_converts_from_and_to_tx_id() {
let txid = tx_id(6);
assert_eq!(UserId::from_tx_id(txid).as_tx_id(), txid);
}
#[test]
fn invite_status_uses_only_accepted_state() {
let commitment = [4; 32];
let registration_key = RegistrationKey::from_bytes([5; 32]);
let mut state = FacadeState::new();
assert_eq!(status_from_state(&state, commitment), InviteStatus::Unknown);
state.issues.insert(commitment, tx_id(1));
assert_eq!(status_from_state(&state, commitment), InviteStatus::Unused);
state.pending_consumes.insert(
commitment,
PendingConsume {
registration_key,
data: Arc::new(b"pending".to_vec()),
cell: Arc::new(PendingResult::new()),
},
);
assert_eq!(status_from_state(&state, commitment), InviteStatus::Unused);
state.registrations.insert(
commitment,
Registration {
user_id: UserId(tx_id(2)),
registration_key,
data: Arc::new(b"accepted".to_vec()),
},
);
assert_eq!(
status_from_state(&state, commitment),
InviteStatus::Consumed
);
let copied = InviteStatus::Consumed;
assert_eq!(copied, InviteStatus::Consumed);
assert_eq!(format!("{copied:?}"), "Consumed");
}
#[test]
fn public_status_survives_restart() {
use std::sync::atomic::{AtomicU64, Ordering};
static NEXT_TEMP_ROOT: AtomicU64 = AtomicU64::new(0);
let unique = NEXT_TEMP_ROOT.fetch_add(1, Ordering::Relaxed);
let root = std::env::temp_dir().join(format!(
"kcode-k1-invites-status-{}-{unique}",
std::process::id()
));
let ordering_root = root.join("ordering");
let peering_root = root.join("peering");
let invites_root = root.join("invites");
let ordering =
Arc::new(K1TxnOrdering::open(&ordering_root).expect("open temporary ordering"));
let peering = Arc::new(
K1Peering::open(&peering_root, Arc::clone(&ordering)).expect("open temporary peering"),
);
let invites = K1Invites::open(
&invites_root,
Arc::clone(&ordering),
Arc::clone(&peering),
InviteVerifierKey::from_bytes([7; 32]),
)
.expect("open temporary invites");
let unknown = InviteCode::from_str("AAAAAAAA").expect("canonical unknown code");
assert_eq!(
invites.status(&unknown).expect("unknown status"),
InviteStatus::Unknown
);
let (_, code) = invites.create().expect("create invite");
let retained_code = code.expose();
assert_eq!(
invites.status(&code).expect("unused status"),
InviteStatus::Unused
);
invites
.consume_with_data(&code, RegistrationKey::from_bytes([8; 32]), b"account")
.expect("consume invite");
assert_eq!(
invites.status(&code).expect("consumed status"),
InviteStatus::Consumed
);
drop(invites);
drop(peering);
drop(ordering);
let ordering =
Arc::new(K1TxnOrdering::open(&ordering_root).expect("reopen temporary ordering"));
let peering = Arc::new(
K1Peering::open(&peering_root, Arc::clone(&ordering))
.expect("reopen temporary peering"),
);
let invites = K1Invites::open(
&invites_root,
Arc::clone(&ordering),
Arc::clone(&peering),
InviteVerifierKey::from_bytes([7; 32]),
)
.expect("reopen temporary invites");
let code = InviteCode::from_str(&retained_code).expect("parse retained code");
assert_eq!(
invites.status(&code).expect("restarted status"),
InviteStatus::Consumed
);
drop(invites);
drop(peering);
drop(ordering);
std::fs::remove_dir_all(&root).expect("remove temporary root");
}
#[test]
fn invite_code_is_strict_and_round_trips() {
let code = InviteCode {
bytes: [0, 1, 2, 253, 254, 255],
};
let exposed = code.expose();
assert_eq!(exposed.len(), 8);
let parsed = InviteCode::from_str(&exposed).expect("canonical code");
assert_eq!(parsed.expose(), exposed);
for invalid in ["", "AAAAAAA", "AAAAAAAA=", "AAAAAAAAA", "AAAAAA+/", "åååå"] {
assert!(InviteCode::from_str(invalid).is_err());
}
}
#[test]
fn commitment_uses_the_exact_domain_and_code_bytes() {
let key = InviteVerifierKey::from_bytes([7; 32]);
let code = [1, 2, 3, 4, 5, 6];
let actual = invite_commitment(&key, &code).expect("commitment");
let mut mac = HmacSha256::new_from_slice(&[7; 32]).expect("key");
mac.update(b"k1-invite-v2");
mac.update(&code);
let expected = mac.finalize().into_bytes();
assert_eq!(&actual[..], &expected[..]);
}
#[test]
fn version_two_wire_round_trips() {
let callback_id = tx_id(9);
let issue_id = tx_id(4);
let commitment = [3; 32];
let key = RegistrationKey::from_bytes([5; 32]);
let issue = issue_payload(commitment);
match parse_payload(callback_id, &issue).expect("Issue") {
ParsedAction::Issue {
id,
commitment: parsed,
} => {
assert_eq!(id, callback_id);
assert_eq!(parsed, commitment);
}
ParsedAction::Consume { .. } => panic!("wrong action"),
}
let consume =
consume_payload(issue_id, commitment, key, b"opaque\0bytes").expect("Consume payload");
assert_eq!(consume.len(), CONSUME_HEADER_LENGTH + 12);
match parse_payload(callback_id, &consume).expect("Consume") {
ParsedAction::Consume {
id,
issue_id: parsed_issue,
commitment: parsed_commitment,
registration_key,
data,
} => {
assert_eq!(id, callback_id);
assert_eq!(parsed_issue, issue_id);
assert_eq!(parsed_commitment, commitment);
assert_eq!(registration_key, [5; 32]);
assert_eq!(data, b"opaque\0bytes");
}
ParsedAction::Issue { .. } => panic!("wrong action"),
}
}
#[test]
fn malformed_payloads_are_rejected() {
let id = tx_id(1);
let mut extended_issue = vec![0; ISSUE_LENGTH + 1];
extended_issue[0] = 2;
extended_issue[1] = 1;
let mut short_consume = vec![0; CONSUME_HEADER_LENGTH - 1];
short_consume[0] = 2;
short_consume[1] = 2;
let malformed = vec![
Vec::new(),
vec![2],
vec![1, 1],
vec![2, 3],
vec![2, 1],
extended_issue,
short_consume,
];
for payload in malformed {
assert!(parse_payload(id, &payload).is_err());
}
}
#[test]
fn apply_waiting_behind_reorg_linearization_does_not_reach_projection() {
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::mpsc;
use std::thread;
static NEXT_TEMP_ROOT: AtomicU64 = AtomicU64::new(0);
let unique = NEXT_TEMP_ROOT.fetch_add(1, Ordering::Relaxed);
let root = std::env::temp_dir().join(format!(
"kcode-k1-invites-reorg-{}-{unique}",
std::process::id()
));
let projection = InviteProjection::open(&root).expect("open temporary projection");
let subsystem = Arc::new(InviteSubsystem {
projection: Mutex::new(projection),
state: Mutex::new(FacadeState::new()),
});
let projection_guard = lock_unpoison(&subsystem.projection);
let worker_subsystem = Arc::clone(&subsystem);
let (ready_sender, ready_receiver) = mpsc::channel();
let worker = thread::spawn(move || {
ready_sender.send(()).expect("signal helper readiness");
worker_subsystem.apply_if_available(InviteAction::Issue {
id: tx_id(8),
commitment: [9; 32],
})
});
ready_receiver.recv().expect("receive helper readiness");
subsystem.mark_unavailable();
drop(projection_guard);
match worker.join().expect("join helper thread") {
Err(error) => assert_eq!(error, REOPEN_REQUIRED),
Ok(_) => panic!("unavailable helper unexpectedly applied the Issue"),
}
drop(subsystem);
let projection = InviteProjection::open(&root).expect("reopen temporary projection");
let snapshot = projection
.snapshot()
.expect("snapshot temporary projection");
assert!(snapshot.checkpoint.is_none());
assert!(snapshot.issues.is_empty());
assert!(snapshot.registrations.is_empty());
drop(projection);
std::fs::remove_dir_all(&root).expect("remove temporary projection");
}
}