pub use kcode_k1_canonical_chain::{SubmitError, TxId};
use kcode_k1_transaction::Transaction;
pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
use std::any::Any;
use std::fs;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex, mpsc};
use std::thread;
use std::time::Duration;
pub type TestSigner = Box<dyn FnOnce(&[u8]) -> Result<[u8; 64], String> + Send + 'static>;
pub type TestQueuePropagation = Box<dyn FnOnce(&[u8]) -> Result<(), String> + Send + 'static>;
pub trait TestSubsystem: Send + Sync + 'static {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
fn reorg(&self) -> Result<(), String>;
}
pub trait OrderingCandidate: Send + Sync + 'static {
fn register_subsystem(
&self,
subsystem: SubsystemId,
after: Option<TxId>,
handler: Arc<dyn TestSubsystem>,
) -> Result<(), String>;
fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError>;
fn submit_local_txn(
&self,
timestamp: u64,
creator: [u8; 32],
subsystem: SubsystemId,
payload: &[u8],
signer: TestSigner,
queue_propagation: TestQueuePropagation,
) -> Result<Vec<u8>, String>;
fn contains(&self, id: TxId) -> bool;
fn tip(&self) -> Option<TxId>;
fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String>;
fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String>;
}
pub trait OrderingHarness: Send + Sync {
fn open(&self, root: &Path) -> Result<Arc<dyn OrderingCandidate>, String>;
}
const TIMEOUT: Duration = Duration::from_secs(2);
const QUIET: Duration = Duration::from_millis(50);
static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
struct TempRoot(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-live-testkit-{}-{number}-{label}",
std::process::id()
));
let _ = fs::remove_dir_all(&path);
Self(path)
}
}
impl Drop for TempRoot {
fn drop(&mut self) {
if !thread::panicking() {
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();
let state = self.changed.wait_while(state, |state| !state.1).unwrap();
assert!(state.1);
}
fn wait_for(&self, count: usize) {
let state = self.state.lock().unwrap();
let (state, _) = self
.changed
.wait_timeout_while(state, TIMEOUT, |state| state.0 < count)
.unwrap();
assert!(state.0 >= count, "gate entry timed out");
}
fn count(&self) -> usize {
self.state.lock().unwrap().0
}
fn release(&self) {
self.state.lock().unwrap().1 = true;
self.changed.notify_all();
}
fn release_on_drop(self: &Arc<Self>) -> GateRelease {
GateRelease(self.clone())
}
}
struct GateRelease(Arc<Gate>);
impl Drop for GateRelease {
fn drop(&mut self) {
self.0.release();
}
}
struct Recording {
payloads: Mutex<Vec<Vec<u8>>>,
submit_gate: Option<Arc<Gate>>,
queue_flag: Option<Arc<AtomicBool>>,
queue_observed: AtomicBool,
}
impl Recording {
fn new(submit_gate: Option<Arc<Gate>>, queue_flag: Option<Arc<AtomicBool>>) -> Self {
Self {
payloads: Mutex::new(Vec::new()),
submit_gate,
queue_flag,
queue_observed: AtomicBool::new(false),
}
}
fn plain() -> Arc<Self> {
Arc::new(Self::new(None, None))
}
fn payloads(&self) -> Vec<Vec<u8>> {
self.payloads.lock().unwrap().clone()
}
}
impl TestSubsystem 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();
}
Ok(())
}
fn reorg(&self) -> Result<(), String> {
Ok(())
}
}
struct Reentrant {
ordering: Arc<dyn OrderingCandidate>,
target: SubsystemId,
}
impl TestSubsystem for Reentrant {
fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
self.ordering
.submit_local_txn(
99,
[9; 32],
self.target,
b"reentered",
Box::new(|_| Ok([9; 64])),
Box::new(|_| Ok(())),
)
.map(|_| ())
}
fn reorg(&self) -> Result<(), String> {
Ok(())
}
}
fn subsystem(value: u8) -> SubsystemId {
SubsystemId::from_bytes([value; 20]).unwrap()
}
fn assert_committed(message: &str, id: TxId) {
assert!(message.contains("TxId"));
assert!(message.contains("was committed"));
assert!(message.contains(&format!("{id:?}")));
}
fn receive<T>(receiver: &mpsc::Receiver<T>, label: &str) -> T {
receiver
.recv_timeout(TIMEOUT)
.unwrap_or_else(|error| panic!("{label}: {error}"))
}
fn assert_pending<T>(receiver: &mpsc::Receiver<T>, label: &str) {
match receiver.recv_timeout(QUIET) {
Err(mpsc::RecvTimeoutError::Timeout) => {}
Err(mpsc::RecvTimeoutError::Disconnected) => panic!("{label} disconnected"),
Ok(_) => panic!("{label} completed while it should have been blocked"),
}
}
fn panic_message(payload: Box<dyn Any + Send>) -> String {
if let Some(message) = payload.downcast_ref::<String>() {
message.clone()
} else if let Some(message) = payload.downcast_ref::<&str>() {
(*message).to_owned()
} else {
"non-string panic".to_owned()
}
}
fn run_scenario<F>(name: &str, scenario: F) -> Result<(), String>
where
F: FnOnce(),
{
match catch_unwind(AssertUnwindSafe(scenario)) {
Ok(()) => Ok(()),
Err(payload) => Err(format!("{name}: {}", panic_message(payload))),
}
}
pub fn verify_queue_before_callback(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("queue before callback", || {
let root = TempRoot::new("queue");
let ordering = harness.open(&root.0).unwrap();
let owner = subsystem(b'a');
let queued = Arc::new(AtomicBool::new(false));
let handler = Arc::new(Recording::new(None, Some(queued.clone())));
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
let first_queued = queued.clone();
let first = ordering.submit_local_txn(
1,
[1; 32],
owner,
b"first",
Box::new(|_| Ok([1; 64])),
Box::new(move |bytes| {
assert_eq!(Transaction::parse(bytes).unwrap().payload(), b"first");
first_queued.store(true, Ordering::Release);
Err("queue unavailable".to_owned())
}),
);
let first_id = ordering.tip().unwrap();
assert_committed(&first.unwrap_err(), first_id);
assert!(handler.queue_observed.load(Ordering::Acquire));
queued.store(false, Ordering::Release);
let second_queued = queued.clone();
let second = ordering.submit_local_txn(
2,
[1; 32],
owner,
b"second",
Box::new(|_| Ok([2; 64])),
Box::new(move |_| {
second_queued.store(true, Ordering::Release);
panic!("queue panic")
}),
);
let second_id = ordering.tip().unwrap();
assert_committed(&second.unwrap_err(), second_id);
assert!(handler.queue_observed.load(Ordering::Acquire));
ordering
.submit_local_txn(
3,
[1; 32],
owner,
b"third",
Box::new(|_| Ok([3; 64])),
Box::new(|_| Ok(())),
)
.unwrap();
assert_eq!(
handler.payloads(),
vec![b"first".to_vec(), b"second".to_vec(), b"third".to_vec()]
);
})
}
pub fn verify_independent_subsystem_lanes(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("independent subsystem lanes", || {
let root = TempRoot::new("lanes");
let ordering = harness.open(&root.0).unwrap();
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let gate = Arc::new(Gate::new());
let _release = gate.release_on_drop();
let handler_a = Arc::new(Recording::new(Some(gate.clone()), 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_tx, first_rx) = mpsc::channel();
let first_ordering = ordering.clone();
let first = thread::spawn(move || {
first_tx
.send(first_ordering.submit_local_txn(
1,
[1; 32],
a,
b"a1",
Box::new(|_| Ok([1; 64])),
Box::new(|_| Ok(())),
))
.unwrap();
});
gate.wait_for(1);
let (query_tx, query_rx) = mpsc::channel();
let query_ordering = ordering.clone();
let query = thread::spawn(move || {
let id = query_ordering.tip().unwrap();
query_tx
.send((
id,
query_ordering.contains(id),
query_ordering.get_txn(id),
query_ordering.between_txids(GENESIS_PARENT, id),
))
.unwrap();
});
let (_id, contains, stored, between) = receive(&query_rx, "queries blocked by A");
assert!(contains);
assert!(!stored.unwrap().unwrap().is_empty());
assert!(between.unwrap().is_empty());
query.join().unwrap();
let (queued_tx, queued_rx) = mpsc::channel();
let (a_tx, a_rx) = mpsc::channel();
let second_ordering = ordering.clone();
let second = thread::spawn(move || {
a_tx.send(second_ordering.submit_local_txn(
2,
[1; 32],
a,
b"a2",
Box::new(|_| Ok([2; 64])),
Box::new(move |_| {
queued_tx.send(()).unwrap();
Ok(())
}),
))
.unwrap();
});
receive(&queued_rx, "second A queue blocked");
assert_eq!(gate.count(), 1);
assert_pending(&a_rx, "second A caller");
let (b_tx, b_rx) = mpsc::channel();
let b_ordering = ordering.clone();
let b_work = thread::spawn(move || {
b_tx.send(b_ordering.submit_local_txn(
3,
[3; 32],
b,
b"b1",
Box::new(|_| Ok([3; 64])),
Box::new(|_| Ok(())),
))
.unwrap();
});
receive(&b_rx, "B submission blocked by A").unwrap();
b_work.join().unwrap();
gate.release();
receive(&first_rx, "first A did not finish").unwrap();
first.join().unwrap();
receive(&a_rx, "second A did not finish").unwrap();
second.join().unwrap();
assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
})
}
pub fn verify_signing_commit_exclusion(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("signing commit exclusion", || {
let root = TempRoot::new("signer");
let ordering = harness.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 gate = Arc::new(Gate::new());
let _release = gate.release_on_drop();
let first_gate = gate.clone();
let (first_tx, first_rx) = mpsc::channel();
let first_ordering = ordering.clone();
let first = thread::spawn(move || {
first_tx
.send(first_ordering.submit_local_txn(
1,
[1; 32],
a,
b"a",
Box::new(move |_| {
first_gate.enter();
Ok([1; 64])
}),
Box::new(|_| Ok(())),
))
.unwrap();
});
gate.wait_for(1);
let (submit_ready_tx, submit_ready_rx) = mpsc::channel();
let (submit_tx, submit_rx) = mpsc::channel();
let submit_ordering = ordering.clone();
let submit = thread::spawn(move || {
submit_ready_tx.send(()).unwrap();
submit_tx
.send(submit_ordering.submit_local_txn(
2,
[2; 32],
b,
b"b",
Box::new(|_| Ok([2; 64])),
Box::new(|_| Ok(())),
))
.unwrap();
});
let (query_ready_tx, query_ready_rx) = mpsc::channel();
let (query_tx, query_rx) = mpsc::channel();
let query_ordering = ordering.clone();
let query = thread::spawn(move || {
query_ready_tx.send(()).unwrap();
query_tx.send(query_ordering.tip()).unwrap();
});
receive(&submit_ready_rx, "second writer did not start");
receive(&query_ready_rx, "query did not start");
assert_pending(&submit_rx, "second writer");
assert_pending(&query_rx, "tip query");
gate.release();
receive(&first_rx, "first writer did not finish").unwrap();
first.join().unwrap();
receive(&submit_rx, "second writer did not finish").unwrap();
assert!(receive(&query_rx, "tip query did not finish").is_some());
submit.join().unwrap();
query.join().unwrap();
})
}
pub fn verify_unrelated_callback_reentry(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("unrelated callback reentry", || {
let root = TempRoot::new("reentry");
let ordering = harness.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();
let (outer_tx, outer_rx) = mpsc::channel();
let outer_ordering = ordering.clone();
let outer = thread::spawn(move || {
outer_tx
.send(outer_ordering.submit_local_txn(
1,
[1; 32],
a,
b"outer",
Box::new(|_| Ok([1; 64])),
Box::new(|_| Ok(())),
))
.unwrap();
});
receive(&outer_rx, "reentrant callback deadlocked").unwrap();
outer.join().unwrap();
assert_eq!(handler_b.payloads(), vec![b"reentered".to_vec()]);
})
}