use kcode_k1_canonical_chain::{CanonicalChain, CommitOutcome};
pub use kcode_k1_canonical_chain::{SubmitError, TxId};
use kcode_k1_transaction::Transaction;
pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
use std::collections::HashMap;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Condvar, Mutex, MutexGuard};
use std::thread;
use std::time::{Duration, Instant};
pub trait Subsystem: Send + Sync + 'static {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
fn reorg(&self) -> Result<(), String>;
}
pub struct K1TxnOrdering {
writer: Mutex<WriterState>,
chain: Mutex<CanonicalChain>,
}
struct WriterState {
registrations: HashMap<SubsystemId, Registration>,
}
struct Registration {
lane: Arc<Lane>,
latest: Option<TxId>,
next_ticket: u64,
mode: RegistrationMode,
}
#[derive(Clone, Copy, Eq, PartialEq)]
enum RegistrationMode {
Replaying,
Active,
OutOfCommission,
}
struct Lane {
handler: Arc<dyn Subsystem>,
progress: Mutex<LaneProgress>,
changed: Condvar,
running: AtomicBool,
faulted: AtomicBool,
}
struct LaneProgress {
serving: u64,
fault: Option<String>,
}
struct Reservation {
lane: Arc<Lane>,
ticket: u64,
}
enum CommittedWork {
Duplicate,
Extension {
id: TxId,
delivery: Option<Reservation>,
},
Reorganization {
id: TxId,
reorgs: Vec<Reservation>,
replacement: Option<Reservation>,
},
}
impl K1TxnOrdering {
pub fn open(root: &Path) -> Result<Self, String> {
let started = Instant::now();
let result = CanonicalChain::open(root);
let elapsed = started.elapsed();
if elapsed > Duration::from_millis(100) {
eprintln!(
"{{\"module\":\"kcode-k1-txn-ordering\",\"operation\":\"open\",\"elapsed_microseconds\":{},\"outcome\":\"{}\"}}",
elapsed.as_micros(),
if result.is_ok() { "ready" } else { "error" }
);
}
result.map(|chain| Self {
writer: Mutex::new(WriterState {
registrations: HashMap::new(),
}),
chain: Mutex::new(chain),
})
}
pub fn register_subsystem(
&self,
subsystem: SubsystemId,
after: Option<TxId>,
handler: Arc<dyn Subsystem>,
) -> Result<(), String> {
let (mut cursor, lane) = {
let mut writer = self.lock_writer();
if let Some(registration) = writer.registrations.get(&subsystem) {
ensure_replaceable(registration)?;
}
let chain = self.lock_chain();
let cursor = chain.replay_cursor(subsystem, after)?;
let lane = Arc::new(Lane::new(handler));
writer.registrations.insert(
subsystem,
Registration {
lane: lane.clone(),
latest: after,
next_ticket: 0,
mode: RegistrationMode::Replaying,
},
);
(cursor, lane)
};
let mut latest = after;
loop {
let mut next = match self.lock_chain().replay_next(&mut cursor) {
Ok(next) => next,
Err(message) => return Err(fault_cursor(&lane, message)),
};
if next.is_none() {
let final_result = {
let mut writer = self.lock_writer();
let chain = self.lock_chain();
match chain.replay_next(&mut cursor) {
Ok(None) => {
let registration = writer
.registrations
.get_mut(&subsystem)
.expect("replaying registration exists");
if !Arc::ptr_eq(®istration.lane, &lane) {
return Err("registration changed during replay".to_owned());
}
registration.latest = latest;
registration.mode = RegistrationMode::Active;
return Ok(());
}
result => result,
}
};
next = match final_result {
Ok(next) => next,
Err(message) => return Err(fault_cursor(&lane, message)),
};
}
let replayed = next.expect("replay result contains a transaction");
let transaction = match Transaction::parse(&replayed.bytes) {
Ok(transaction) => transaction,
Err(message) => {
let failure = committed_error(
replayed.id,
vec![format!(
"canonical replay transaction was invalid: {message}"
)],
);
lane.record_fault(failure.clone());
return Err(failure);
}
};
if let Err(message) = lane.run_replay(replayed.id, transaction.payload()) {
return Err(committed_error(replayed.id, vec![message]));
}
latest = Some(replayed.id);
}
}
pub fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError> {
let parsed = Transaction::parse(transaction)
.map_err(|message| SubmitError::Other(format!("invalid transaction: {message}")))?;
let payload = parsed.payload();
let work = {
let mut writer = self.lock_writer();
let mut chain = self.lock_chain();
match chain.submit_validated(transaction)? {
CommitOutcome::Duplicate => CommittedWork::Duplicate,
CommitOutcome::Extension { id, subsystem } => CommittedWork::Extension {
id,
delivery: reserve_active_delivery(&mut writer, subsystem, id),
},
CommitOutcome::Reorganization { id, subsystem } => {
let affected: Vec<SubsystemId> = writer
.registrations
.iter()
.filter_map(|(®istered, registration)| {
(registration.is_active()
&& registration
.latest
.is_some_and(|latest| !chain.contains(latest)))
.then_some(registered)
})
.collect();
let mut reorgs = Vec::with_capacity(affected.len());
for affected_subsystem in affected {
let registration = writer
.registrations
.get_mut(&affected_subsystem)
.expect("affected registration exists");
reorgs.push(registration.reserve_reorg());
registration.mode = RegistrationMode::OutOfCommission;
}
let replacement = reserve_active_delivery(&mut writer, subsystem, id);
CommittedWork::Reorganization {
id,
reorgs,
replacement,
}
}
}
};
match work {
CommittedWork::Duplicate => Ok(()),
CommittedWork::Extension { id, delivery } => {
let Some(delivery) = delivery else {
return Ok(());
};
delivery
.run_delivery(id, payload)
.map_err(|message| SubmitError::Other(committed_error(id, vec![message])))
}
CommittedWork::Reorganization {
id,
reorgs,
replacement,
} => {
let failures = run_reorganization(id, payload, reorgs, replacement);
if failures.is_empty() {
Ok(())
} else {
Err(SubmitError::Other(committed_error(id, failures)))
}
}
}
}
pub fn submit_local_txn<F, Q>(
&self,
timestamp: u64,
creator: [u8; 32],
subsystem: SubsystemId,
payload: &[u8],
signer: F,
queue_propagation: Q,
) -> Result<Vec<u8>, String>
where
F: FnOnce(&[u8]) -> Result<[u8; 64], String>,
Q: FnOnce(&[u8]) -> Result<(), String>,
{
let (bytes, id, delivery) = {
let mut writer = self.lock_writer();
let lane = writer
.registrations
.get(&subsystem)
.filter(|registration| registration.is_active())
.map(|registration| registration.lane.clone())
.ok_or_else(|| "target subsystem is not registered and active".to_owned())?;
let mut chain = self.lock_chain();
let bytes =
chain.submit_local(timestamp, creator, subsystem, payload, move |prefix| {
match catch_unwind(AssertUnwindSafe(|| signer(prefix))) {
Ok(result) => result,
Err(_) => Err("signer panicked".to_owned()),
}
})?;
let id = TxId::for_transaction(&bytes);
let registration = writer
.registrations
.get_mut(&subsystem)
.expect("prechecked registration exists");
if !Arc::ptr_eq(®istration.lane, &lane) {
return Err(committed_error(
id,
vec!["target registration changed during local commit".to_owned()],
));
}
let delivery = registration.reserve_delivery(id);
(bytes, id, delivery)
};
let mut failures = Vec::new();
match catch_unwind(AssertUnwindSafe(|| queue_propagation(&bytes))) {
Ok(Ok(())) => {}
Ok(Err(message)) => failures.push(format!("queue propagation failed: {message}")),
Err(_) => failures.push("queue propagation panicked".to_owned()),
}
if let Err(message) = delivery.run_delivery(id, payload) {
failures.push(message);
}
if failures.is_empty() {
Ok(bytes)
} else {
Err(committed_error(id, failures))
}
}
pub fn contains(&self, id: TxId) -> bool {
self.lock_chain().contains(id)
}
pub fn tip(&self) -> Option<TxId> {
self.lock_chain().tip()
}
pub fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String> {
self.lock_chain().between_txids(older, newer)
}
pub fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String> {
self.lock_chain().get_txn(id)
}
fn lock_writer(&self) -> MutexGuard<'_, WriterState> {
self.writer.lock().expect("KTO writer mutex poisoned")
}
fn lock_chain(&self) -> MutexGuard<'_, CanonicalChain> {
self.chain.lock().expect("KTO chain mutex poisoned")
}
}
impl Registration {
fn is_active(&self) -> bool {
self.mode == RegistrationMode::Active && !self.lane.faulted.load(Ordering::Acquire)
}
fn reserve_delivery(&mut self, id: TxId) -> Reservation {
let reservation = Reservation {
lane: self.lane.clone(),
ticket: self.next_ticket,
};
self.next_ticket += 1;
self.latest = Some(id);
reservation
}
fn reserve_reorg(&mut self) -> Reservation {
let reservation = Reservation {
lane: self.lane.clone(),
ticket: self.next_ticket,
};
self.next_ticket += 1;
reservation
}
}
impl Lane {
fn new(handler: Arc<dyn Subsystem>) -> Self {
Self {
handler,
progress: Mutex::new(LaneProgress {
serving: 0,
fault: None,
}),
changed: Condvar::new(),
running: AtomicBool::new(false),
faulted: AtomicBool::new(false),
}
}
fn run_replay(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
self.running.store(true, Ordering::Release);
let result = invoke_callback("replay submit", || self.handler.submit_txn(id, payload));
if let Err(message) = &result {
self.record_fault(message.clone());
}
self.running.store(false, Ordering::Release);
result
}
fn run_ticket<F>(&self, ticket: u64, operation: &str, callback: F) -> Result<(), String>
where
F: FnOnce() -> Result<(), String>,
{
let mut progress = self.progress.lock().expect("subsystem lane mutex poisoned");
while progress.serving < ticket {
progress = self
.changed
.wait(progress)
.expect("subsystem lane mutex poisoned");
}
if progress.serving > ticket {
return Err(format!("{operation} ticket was already completed"));
}
if let Some(reason) = progress.fault.clone() {
progress.serving += 1;
self.changed.notify_all();
return Err(format!(
"{operation} callback was not run because the registration faulted: {reason}"
));
}
self.running.store(true, Ordering::Release);
drop(progress);
let result = invoke_callback(operation, callback);
let mut progress = self.progress.lock().expect("subsystem lane mutex poisoned");
if let Err(message) = &result
&& progress.fault.is_none()
{
progress.fault = Some(message.clone());
self.faulted.store(true, Ordering::Release);
}
progress.serving += 1;
self.running.store(false, Ordering::Release);
self.changed.notify_all();
result
}
fn record_fault(&self, message: String) {
let mut progress = self.progress.lock().expect("subsystem lane mutex poisoned");
if progress.fault.is_none() {
progress.fault = Some(message);
self.faulted.store(true, Ordering::Release);
}
}
fn is_quiescent(&self, next_ticket: u64) -> bool {
let progress = self.progress.lock().expect("subsystem lane mutex poisoned");
progress.serving == next_ticket && !self.running.load(Ordering::Acquire)
}
}
impl Reservation {
fn run_delivery(self, id: TxId, payload: &[u8]) -> Result<(), String> {
self.lane.run_ticket(self.ticket, "submit", || {
self.lane.handler.submit_txn(id, payload)
})
}
fn run_reorg(self) -> Result<(), String> {
self.lane
.run_ticket(self.ticket, "reorg", || self.lane.handler.reorg())
}
}
fn invoke_callback<F>(operation: &str, callback: F) -> Result<(), String>
where
F: FnOnce() -> Result<(), String>,
{
match catch_unwind(AssertUnwindSafe(callback)) {
Ok(Ok(())) => Ok(()),
Ok(Err(message)) => Err(format!("{operation} callback failed: {message}")),
Err(_) => Err(format!("{operation} callback panicked")),
}
}
fn ensure_replaceable(registration: &Registration) -> Result<(), String> {
if registration.mode == RegistrationMode::Replaying
&& !registration.lane.faulted.load(Ordering::Acquire)
{
return Err("subsystem registration is replaying".to_owned());
}
if registration.mode == RegistrationMode::Active
&& !registration.lane.faulted.load(Ordering::Acquire)
{
return Err("subsystem is already registered and active".to_owned());
}
if !registration.lane.is_quiescent(registration.next_ticket) {
return Err("previous subsystem lane work has not quiesced".to_owned());
}
Ok(())
}
fn reserve_active_delivery(
writer: &mut WriterState,
subsystem: SubsystemId,
id: TxId,
) -> Option<Reservation> {
writer
.registrations
.get_mut(&subsystem)
.filter(|registration| registration.is_active())
.map(|registration| registration.reserve_delivery(id))
}
fn run_reorganization(
id: TxId,
payload: &[u8],
reorgs: Vec<Reservation>,
replacement: Option<Reservation>,
) -> Vec<String> {
thread::scope(|scope| {
let mut jobs = Vec::with_capacity(reorgs.len() + usize::from(replacement.is_some()));
for reservation in reorgs {
jobs.push(scope.spawn(move || reservation.run_reorg()));
}
if let Some(reservation) = replacement {
jobs.push(scope.spawn(move || reservation.run_delivery(id, payload)));
}
let mut failures = Vec::new();
for job in jobs {
match job.join() {
Ok(Ok(())) => {}
Ok(Err(message)) => failures.push(message),
Err(_) => failures.push("callback lane panicked".to_owned()),
}
}
failures
})
}
fn fault_cursor(lane: &Lane, message: String) -> String {
let failure = format!("subsystem replay cursor failed: {message}");
lane.record_fault(failure.clone());
failure
}
fn committed_error(id: TxId, failures: Vec<String>) -> String {
format!("TxId {id:?} was committed; {}", failures.join("; "))
}
#[cfg(test)]
mod tests {
use super::*;
use kcode_k1_transaction::build_signed_transaction;
use std::fs;
use std::sync::atomic::{AtomicU64, AtomicUsize};
use std::sync::mpsc;
use std::time::Duration;
static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
struct TempRoot(std::path::PathBuf);
impl TempRoot {
fn new(label: &str) -> Self {
let number = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"kcode-k1-txn-ordering-{}-{number}-{label}",
std::process::id()
));
let _ = fs::remove_dir_all(&path);
Self(path)
}
}
impl Drop for TempRoot {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.0);
}
}
struct Gate {
state: Mutex<(usize, bool)>,
changed: Condvar,
}
impl Gate {
fn new() -> Self {
Self {
state: Mutex::new((0, false)),
changed: Condvar::new(),
}
}
fn enter(&self) {
let mut state = self.state.lock().unwrap();
state.0 += 1;
self.changed.notify_all();
while !state.1 {
state = self.changed.wait(state).unwrap();
}
}
fn wait_for(&self, count: usize) {
let mut state = self.state.lock().unwrap();
while state.0 < count {
state = self.changed.wait(state).unwrap();
}
}
fn count(&self) -> usize {
self.state.lock().unwrap().0
}
fn release(&self) {
self.state.lock().unwrap().1 = true;
self.changed.notify_all();
}
}
struct Recording {
payloads: Mutex<Vec<Vec<u8>>>,
reorgs: AtomicUsize,
submit_gate: Option<Arc<Gate>>,
reorg_gate: Option<Arc<Gate>>,
queue_flag: Option<Arc<AtomicBool>>,
queue_observed: AtomicBool,
fail_submit: AtomicBool,
panic_submit: AtomicBool,
}
impl Recording {
fn new(
submit_gate: Option<Arc<Gate>>,
reorg_gate: Option<Arc<Gate>>,
queue_flag: Option<Arc<AtomicBool>>,
) -> Self {
Self {
payloads: Mutex::new(Vec::new()),
reorgs: AtomicUsize::new(0),
submit_gate,
reorg_gate,
queue_flag,
queue_observed: AtomicBool::new(false),
fail_submit: AtomicBool::new(false),
panic_submit: AtomicBool::new(false),
}
}
fn plain() -> Arc<Self> {
Arc::new(Self::new(None, None, None))
}
fn payloads(&self) -> Vec<Vec<u8>> {
self.payloads.lock().unwrap().clone()
}
}
impl Subsystem for Recording {
fn submit_txn(&self, _id: TxId, payload: &[u8]) -> Result<(), String> {
self.payloads.lock().unwrap().push(payload.to_vec());
if let Some(flag) = &self.queue_flag {
self.queue_observed
.store(flag.load(Ordering::Acquire), Ordering::Release);
}
if let Some(gate) = &self.submit_gate {
gate.enter();
}
if self.panic_submit.load(Ordering::Acquire) {
panic!("submit panic");
}
if self.fail_submit.load(Ordering::Acquire) {
Err("submit failure".to_owned())
} else {
Ok(())
}
}
fn reorg(&self) -> Result<(), String> {
self.reorgs.fetch_add(1, Ordering::Relaxed);
if let Some(gate) = &self.reorg_gate {
gate.enter();
}
Ok(())
}
}
struct Reentrant {
ordering: Arc<K1TxnOrdering>,
target: SubsystemId,
}
impl Subsystem for Reentrant {
fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
self.ordering
.submit_local_txn(
99,
[9; 32],
self.target,
b"reentered",
|_| Ok([9; 64]),
|_| Ok(()),
)
.map(|_| ())
}
fn reorg(&self) -> Result<(), String> {
Ok(())
}
}
fn subsystem(value: u8) -> SubsystemId {
SubsystemId::from_bytes([value; 20]).unwrap()
}
fn transaction(
parent: TxId,
creator: u8,
timestamp: u64,
subsystem: SubsystemId,
payload: &[u8],
) -> Vec<u8> {
build_signed_transaction(parent, timestamp, [creator; 32], subsystem, payload, |_| {
Ok([creator; 64])
})
.unwrap()
}
fn assert_committed(message: &str, id: TxId) {
assert!(message.contains("TxId"));
assert!(message.contains("was committed"));
assert!(message.contains(&format!("{id:?}")));
}
#[test]
fn queue_precedes_callback_and_queue_failures_still_integrate() {
let root = TempRoot::new("queue");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let owner = subsystem(b'a');
let queued = Arc::new(AtomicBool::new(false));
let handler = Arc::new(Recording::new(None, None, Some(queued.clone())));
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
let result = ordering.submit_local_txn(
1,
[1; 32],
owner,
b"first",
|_| Ok([1; 64]),
|bytes| {
assert_eq!(Transaction::parse(bytes).unwrap().payload(), b"first");
queued.store(true, Ordering::Release);
Err("queue unavailable".to_owned())
},
);
let first_id = ordering.tip().unwrap();
assert_committed(&result.unwrap_err(), first_id);
assert!(handler.queue_observed.load(Ordering::Acquire));
queued.store(false, Ordering::Release);
let result = ordering.submit_local_txn(
2,
[1; 32],
owner,
b"second",
|_| Ok([2; 64]),
|_| {
queued.store(true, Ordering::Release);
panic!("queue panic")
},
);
let second_id = ordering.tip().unwrap();
assert_committed(&result.unwrap_err(), second_id);
assert_eq!(
handler.payloads(),
vec![b"first".to_vec(), b"second".to_vec()]
);
assert!(
ordering
.submit_local_txn(3, [1; 32], owner, b"third", |_| Ok([3; 64]), |_| Ok(()))
.is_ok()
);
}
#[test]
fn subsystem_lanes_order_a_without_blocking_b_or_queries() {
let root = TempRoot::new("lanes");
let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let gate = Arc::new(Gate::new());
let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
let handler_b = Recording::plain();
ordering
.register_subsystem(a, None, handler_a.clone())
.unwrap();
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
let first_ordering = ordering.clone();
let first = thread::spawn(move || {
first_ordering.submit_local_txn(1, [1; 32], a, b"a1", |_| Ok([1; 64]), |_| Ok(()))
});
gate.wait_for(1);
let first_id = ordering.tip().unwrap();
assert!(ordering.contains(first_id));
assert!(!ordering.get_txn(first_id).unwrap().unwrap().is_empty());
let (queued_tx, queued_rx) = mpsc::channel();
let (a_tx, a_rx) = mpsc::channel();
let second_ordering = ordering.clone();
thread::spawn(move || {
let result = second_ordering.submit_local_txn(
2,
[1; 32],
a,
b"a2",
|_| Ok([2; 64]),
|_| {
queued_tx.send(()).unwrap();
Ok(())
},
);
a_tx.send(result).unwrap();
});
queued_rx.recv_timeout(Duration::from_secs(2)).unwrap();
assert_eq!(gate.count(), 1);
assert!(matches!(
a_rx.recv_timeout(Duration::from_millis(50)),
Err(mpsc::RecvTimeoutError::Timeout)
));
let (b_tx, b_rx) = mpsc::channel();
let b_ordering = ordering.clone();
thread::spawn(move || {
b_tx.send(b_ordering.submit_local_txn(
3,
[3; 32],
b,
b"b1",
|_| Ok([3; 64]),
|_| Ok(()),
))
.unwrap();
});
assert!(b_rx.recv_timeout(Duration::from_secs(2)).unwrap().is_ok());
gate.release();
assert!(first.join().unwrap().is_ok());
assert!(a_rx.recv_timeout(Duration::from_secs(2)).unwrap().is_ok());
assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
}
#[test]
fn signer_holds_writer_and_chain_until_local_commit() {
let root = TempRoot::new("signer");
let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
ordering
.register_subsystem(a, None, Recording::plain())
.unwrap();
ordering
.register_subsystem(b, None, Recording::plain())
.unwrap();
let signer_gate = Arc::new(Gate::new());
let first_ordering = ordering.clone();
let first_gate = signer_gate.clone();
let first = thread::spawn(move || {
first_ordering.submit_local_txn(
1,
[1; 32],
a,
b"a",
|_| {
first_gate.enter();
Ok([1; 64])
},
|_| Ok(()),
)
});
signer_gate.wait_for(1);
let (submit_tx, submit_rx) = mpsc::channel();
let submit_ordering = ordering.clone();
thread::spawn(move || {
submit_tx
.send(submit_ordering.submit_local_txn(
2,
[2; 32],
b,
b"b",
|_| Ok([2; 64]),
|_| Ok(()),
))
.unwrap();
});
let (query_tx, query_rx) = mpsc::channel();
let query_ordering = ordering.clone();
thread::spawn(move || query_tx.send(query_ordering.tip()).unwrap());
assert!(matches!(
submit_rx.recv_timeout(Duration::from_millis(50)),
Err(mpsc::RecvTimeoutError::Timeout)
));
assert!(matches!(
query_rx.recv_timeout(Duration::from_millis(50)),
Err(mpsc::RecvTimeoutError::Timeout)
));
signer_gate.release();
assert!(first.join().unwrap().is_ok());
assert!(
submit_rx
.recv_timeout(Duration::from_secs(2))
.unwrap()
.is_ok()
);
assert!(
query_rx
.recv_timeout(Duration::from_secs(2))
.unwrap()
.is_some()
);
}
#[test]
fn replay_isolated_by_subsystem_and_catches_one_concurrent_commit() {
let root = TempRoot::new("replay");
let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let first = transaction(GENESIS_PARENT, 8, 1, a, b"a1");
ordering.submit_txn(&first).unwrap();
let gate = Arc::new(Gate::new());
let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
let register_ordering = ordering.clone();
let register_handler = handler_a.clone();
let (register_tx, register_rx) = mpsc::channel();
thread::spawn(move || {
register_tx
.send(register_ordering.register_subsystem(a, None, register_handler))
.unwrap();
});
gate.wait_for(1);
let handler_b = Recording::plain();
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
let b_bytes = ordering
.submit_local_txn(2, [2; 32], b, b"b1", |_| Ok([2; 64]), |_| Ok(()))
.unwrap();
let second = transaction(TxId::for_transaction(&b_bytes), 8, 3, a, b"a2");
ordering.submit_txn(&second).unwrap();
assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
gate.release();
assert!(
register_rx
.recv_timeout(Duration::from_secs(2))
.unwrap()
.is_ok()
);
assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
}
#[test]
fn callback_error_and_panic_report_committed_ids_and_isolate_faults() {
let root = TempRoot::new("faults");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let failing = Recording::plain();
failing.fail_submit.store(true, Ordering::Release);
let handler_b = Recording::plain();
ordering
.register_subsystem(a, None, failing.clone())
.unwrap();
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
let message = ordering
.submit_local_txn(1, [1; 32], a, b"failure", |_| Ok([1; 64]), |_| Ok(()))
.unwrap_err();
let failed_id = ordering.tip().unwrap();
assert_committed(&message, failed_id);
ordering
.submit_local_txn(2, [2; 32], b, b"b1", |_| Ok([2; 64]), |_| Ok(()))
.unwrap();
let panicking = Recording::plain();
panicking.panic_submit.store(true, Ordering::Release);
ordering
.register_subsystem(a, Some(failed_id), panicking.clone())
.unwrap();
let peer = transaction(ordering.tip().unwrap(), 3, 3, a, b"panic");
let peer_id = TxId::for_transaction(&peer);
let message = match ordering.submit_txn(&peer) {
Err(SubmitError::Other(message)) => message,
result => panic!("unexpected peer result: {result:?}"),
};
assert_committed(&message, peer_id);
assert!(ordering.contains(peer_id));
ordering
.submit_local_txn(4, [4; 32], b, b"b2", |_| Ok([4; 64]), |_| Ok(()))
.unwrap();
assert_eq!(handler_b.payloads(), vec![b"b1".to_vec(), b"b2".to_vec()]);
}
#[test]
fn reorg_fanout_releases_global_locks_and_unaffected_delivery() {
let root = TempRoot::new("reorg");
let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let reorg_gate = Arc::new(Gate::new());
let handler_a = Arc::new(Recording::new(None, Some(reorg_gate.clone()), None));
let b_gate = Arc::new(Gate::new());
b_gate.release();
let handler_b = Arc::new(Recording::new(Some(b_gate.clone()), None, None));
ordering
.register_subsystem(a, None, handler_a.clone())
.unwrap();
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
let first = transaction(GENESIS_PARENT, 20, 1, a, b"first");
let first_id = TxId::for_transaction(&first);
ordering.submit_txn(&first).unwrap();
let incumbent = transaction(first_id, 50, 2, a, b"incumbent");
ordering.submit_txn(&incumbent).unwrap();
let replacement = transaction(first_id, 1, 3, b, b"replacement");
let replacement_id = TxId::for_transaction(&replacement);
let replacement_ordering = ordering.clone();
let (replacement_tx, replacement_rx) = mpsc::channel();
thread::spawn(move || {
replacement_tx
.send(replacement_ordering.submit_txn(&replacement))
.unwrap();
});
reorg_gate.wait_for(1);
b_gate.wait_for(1);
assert!(ordering.contains(replacement_id));
let (local_tx, local_rx) = mpsc::channel();
let local_ordering = ordering.clone();
thread::spawn(move || {
local_tx
.send(local_ordering.submit_local_txn(
4,
[4; 32],
b,
b"after",
|_| Ok([4; 64]),
|_| Ok(()),
))
.unwrap();
});
assert!(
local_rx
.recv_timeout(Duration::from_secs(2))
.unwrap()
.is_ok()
);
reorg_gate.release();
assert!(
replacement_rx
.recv_timeout(Duration::from_secs(2))
.unwrap()
.is_ok()
);
assert_eq!(handler_a.reorgs.load(Ordering::Acquire), 1);
assert_eq!(
handler_b.payloads(),
vec![b"replacement".to_vec(), b"after".to_vec()]
);
}
#[test]
fn callback_can_reenter_an_unrelated_subsystem() {
let root = TempRoot::new("reentry");
let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let handler_b = Recording::plain();
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
ordering
.register_subsystem(
a,
None,
Arc::new(Reentrant {
ordering: ordering.clone(),
target: b,
}),
)
.unwrap();
ordering
.submit_local_txn(1, [1; 32], a, b"outer", |_| Ok([1; 64]), |_| Ok(()))
.unwrap();
assert_eq!(handler_b.payloads(), vec![b"reentered".to_vec()]);
}
#[test]
fn duplicate_restart_queries_and_root_behavior_are_preserved() {
let root = TempRoot::new("restart");
let owner = subsystem(b'z');
let first = transaction(GENESIS_PARENT, 1, 1, owner, b"first");
let first_id = TxId::for_transaction(&first);
let second = transaction(first_id, 1, 2, owner, b"second");
let second_id = TxId::for_transaction(&second);
{
let ordering = K1TxnOrdering::open(&root.0).unwrap();
ordering.submit_txn(&first).unwrap();
ordering.submit_txn(&first).unwrap();
ordering.submit_txn(&second).unwrap();
assert_eq!(ordering.tip(), Some(second_id));
assert_eq!(
ordering.between_txids(GENESIS_PARENT, second_id).unwrap(),
vec![first_id]
);
assert_eq!(ordering.get_txn(first_id).unwrap(), Some(first.clone()));
}
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let handler = Recording::plain();
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
ordering.submit_txn(&second).unwrap();
assert_eq!(
handler.payloads(),
vec![b"first".to_vec(), b"second".to_vec()]
);
assert!(ordering.contains(first_id));
assert_eq!(ordering.tip(), Some(second_id));
}
}