use std::{
collections::HashMap,
path::Path,
sync::{Arc, Mutex, MutexGuard},
};
use kcode_k1_peering::K1Peering;
use kcode_k1_persons_projection::{ApplyOutcome, Projection};
pub use kcode_k1_persons_wire::PersonId;
use kcode_k1_persons_wire::{PersonsWire, decode};
pub use kcode_k1_txn_ordering::TxId;
use kcode_k1_txn_ordering::{K1TxnOrdering, Subsystem, SubsystemId};
const SUBSYSTEM_NAME: &str = "k1-persons-subsystem";
const UNAVAILABLE: &str = "persons facade is unavailable";
pub struct K1Persons {
_ordering: Arc<K1TxnOrdering>,
peering: Arc<K1Peering>,
projection: Arc<Projection>,
state: Arc<Mutex<State>>,
subsystem: SubsystemId,
}
struct PersonSubsystem {
projection: Arc<Projection>,
state: Arc<Mutex<State>>,
}
struct State {
available: bool,
pending: HashMap<Vec<u8>, usize>,
outcomes: HashMap<TxId, CapturedOutcome>,
}
struct CapturedOutcome {
payload: Vec<u8>,
outcome: ApplyOutcome,
}
impl State {
fn new() -> Self {
Self {
available: true,
pending: HashMap::new(),
outcomes: HashMap::new(),
}
}
fn fault(&mut self) {
self.available = false;
self.pending.clear();
self.outcomes.clear();
}
}
impl K1Persons {
pub fn open(
root: &Path,
ordering: Arc<K1TxnOrdering>,
peering: Arc<K1Peering>,
) -> Result<Self, String> {
let (projection, checkpoint) = Projection::open(root, &ordering)?;
let projection = Arc::new(projection);
let state = Arc::new(Mutex::new(State::new()));
let subsystem = SubsystemId::from_str(SUBSYSTEM_NAME)?;
let handler = Arc::new(PersonSubsystem {
projection: projection.clone(),
state: state.clone(),
});
if let Err(error) = ordering.register_subsystem(subsystem, checkpoint, handler) {
mark_fault(&state);
return Err(error);
}
Ok(Self {
_ordering: ordering,
peering,
projection,
state,
subsystem,
})
}
pub fn create(&self, name: String) -> Result<PersonId, String> {
let wire = PersonsWire::create(name)?;
let (transaction, outcome) = self.submit(wire.as_bytes())?;
match outcome {
ApplyOutcome::Applied(person) if person.as_tx_id() == transaction => Ok(person),
ApplyOutcome::Rejected(reason) => Err(reason),
ApplyOutcome::Applied(_) | ApplyOutcome::Unchanged(_) => {
mark_fault(&self.state);
Err(UNAVAILABLE.to_owned())
}
}
}
pub fn update(&self, person: PersonId, name: String) -> Result<(), String> {
let wire = PersonsWire::update(person, name)?;
let (_, outcome) = self.submit(wire.as_bytes())?;
match outcome {
ApplyOutcome::Applied(_) | ApplyOutcome::Unchanged(_) => Ok(()),
ApplyOutcome::Rejected(reason) => Err(reason),
}
}
pub fn resolve(&self, canonical: PersonId, alias: PersonId) -> Result<PersonId, String> {
let wire = PersonsWire::resolve(canonical, alias);
let (_, outcome) = self.submit(wire.as_bytes())?;
match outcome {
ApplyOutcome::Applied(root) | ApplyOutcome::Unchanged(root) => Ok(root),
ApplyOutcome::Rejected(reason) => Err(reason),
}
}
pub fn get(&self, person: PersonId) -> Result<Option<String>, String> {
ensure_available(&self.state)?;
let result = self.projection.get(person);
ensure_available(&self.state)?;
result
}
fn submit(&self, payload: &[u8]) -> Result<(TxId, ApplyOutcome), String> {
self.reserve(payload)?;
match self.peering.submit_txn(self.subsystem, payload) {
Ok(transaction) => self
.claim(payload, transaction)
.map(|outcome| (transaction, outcome)),
Err(error) => Err(self.finish_submission_error(payload, error)),
}
}
fn reserve(&self, payload: &[u8]) -> Result<(), String> {
let mut state = lock_state(&self.state)?;
if !state.available {
return Err(UNAVAILABLE.to_owned());
}
if let Some(count) = state.pending.get_mut(payload) {
if let Some(next) = count.checked_add(1) {
*count = next;
return Ok(());
}
state.fault();
return Err(UNAVAILABLE.to_owned());
}
if state.pending.try_reserve(1).is_err() {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
let key = match copy_payload(payload) {
Ok(key) => key,
Err(()) => {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
};
state.pending.insert(key, 1);
Ok(())
}
fn claim(&self, payload: &[u8], transaction: TxId) -> Result<ApplyOutcome, String> {
let mut state = lock_state(&self.state)?;
if !state.available {
return Err(UNAVAILABLE.to_owned());
}
let captured = match state.outcomes.remove(&transaction) {
Some(captured) => captured,
None => {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
};
if captured.payload != payload || !release_pending(&mut state, payload) {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
Ok(captured.outcome)
}
fn finish_submission_error(&self, payload: &[u8], error: String) -> String {
let mut state = match lock_state(&self.state) {
Ok(state) => state,
Err(error) => return error,
};
if !state.available {
return UNAVAILABLE.to_owned();
}
if !release_pending(&mut state, payload) {
state.fault();
return UNAVAILABLE.to_owned();
}
error
}
}
impl Subsystem for PersonSubsystem {
fn submit_txn(&self, transaction: TxId, payload: &[u8]) -> Result<(), String> {
ensure_available(&self.state)?;
let action = match decode(payload) {
Ok(action) => action,
Err(error) => {
mark_fault(&self.state);
return Err(error);
}
};
let outcome = match self.projection.apply(transaction, action) {
Ok(outcome) => outcome,
Err(error) => {
mark_fault(&self.state);
return Err(error);
}
};
self.capture(transaction, payload, outcome)
}
fn reorg(&self) -> Result<(), String> {
let lock_error = match self.state.lock() {
Ok(mut state) => {
state.fault();
None
}
Err(error) => {
error.into_inner().fault();
Some(UNAVAILABLE.to_owned())
}
};
self.projection.clear()?;
match lock_error {
Some(error) => Err(error),
None => Ok(()),
}
}
}
impl PersonSubsystem {
fn capture(
&self,
transaction: TxId,
payload: &[u8],
outcome: ApplyOutcome,
) -> Result<(), String> {
let mut state = lock_state(&self.state)?;
if !state.available {
return Err(UNAVAILABLE.to_owned());
}
match state.pending.get(payload).copied() {
None => return Ok(()),
Some(0) => {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
Some(_) => {}
}
if state.outcomes.contains_key(&transaction) || state.outcomes.try_reserve(1).is_err() {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
let exact_payload = match copy_payload(payload) {
Ok(payload) => payload,
Err(()) => {
state.fault();
return Err(UNAVAILABLE.to_owned());
}
};
if state
.outcomes
.insert(
transaction,
CapturedOutcome {
payload: exact_payload,
outcome,
},
)
.is_some()
{
state.fault();
return Err(UNAVAILABLE.to_owned());
}
Ok(())
}
}
fn copy_payload(payload: &[u8]) -> Result<Vec<u8>, ()> {
let mut copy = Vec::new();
copy.try_reserve_exact(payload.len()).map_err(|_| ())?;
copy.extend_from_slice(payload);
Ok(copy)
}
fn release_pending(state: &mut State, payload: &[u8]) -> bool {
let remove = match state.pending.get_mut(payload) {
Some(count) if *count > 1 => {
*count -= 1;
false
}
Some(count) if *count == 1 => true,
_ => return false,
};
if remove {
state.pending.remove(payload);
state
.outcomes
.retain(|_, captured| captured.payload.as_slice() != payload);
}
true
}
fn ensure_available(state: &Mutex<State>) -> Result<(), String> {
let state = lock_state(state)?;
if state.available {
Ok(())
} else {
Err(UNAVAILABLE.to_owned())
}
}
fn lock_state(state: &Mutex<State>) -> Result<MutexGuard<'_, State>, String> {
match state.lock() {
Ok(state) => Ok(state),
Err(error) => {
error.into_inner().fault();
Err(UNAVAILABLE.to_owned())
}
}
}
fn mark_fault(state: &Mutex<State>) {
match state.lock() {
Ok(mut state) => state.fault(),
Err(error) => error.into_inner().fault(),
}
}