pub use kcode_k1_canonical_chain::{SubmitError, TxId};
pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
use kcode_k1_transaction::{Transaction, build_signed_transaction};
use std::any::Any;
use std::fs;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, 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-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>>>,
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 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();
}
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<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 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:?}")));
}
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, 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, 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_replay_live_handoff(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("replay live handoff", || {
let root = TempRoot::new("replay");
let ordering = harness.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 _release = gate.release_on_drop();
let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
let (register_tx, register_rx) = mpsc::channel();
let register_ordering = ordering.clone();
let register_handler = handler_a.clone();
let register = thread::spawn(move || {
register_tx
.send(register_ordering.register_subsystem(a, None, register_handler))
.unwrap();
});
gate.wait_for(1);
let handler_b = Recording::plain();
let (b_tx, b_rx) = mpsc::channel();
let b_ordering = ordering.clone();
let b_handler = handler_b.clone();
let b_work = thread::spawn(move || {
let result = b_ordering
.register_subsystem(b, None, b_handler)
.and_then(|()| {
b_ordering.submit_local_txn(
2,
[2; 32],
b,
b"b1",
Box::new(|_| Ok([2; 64])),
Box::new(|_| Ok(())),
)
});
b_tx.send(result).unwrap();
});
let b_bytes = receive(&b_rx, "B work blocked by replay").unwrap();
b_work.join().unwrap();
let second = transaction(TxId::for_transaction(&b_bytes), 8, 3, a, b"a2");
let (peer_tx, peer_rx) = mpsc::channel();
let peer_ordering = ordering.clone();
let peer = thread::spawn(move || peer_tx.send(peer_ordering.submit_txn(&second)).unwrap());
receive(&peer_rx, "A commit blocked by replay").unwrap();
peer.join().unwrap();
assert_eq!(gate.count(), 1);
assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
gate.release();
receive(®ister_rx, "A registration did not finish").unwrap();
register.join().unwrap();
assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
})
}
pub fn verify_callback_failure_isolation(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("callback failure isolation", || {
let root = TempRoot::new("faults");
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 failing = Arc::new(Recording::new(Some(gate.clone()), None, None));
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 (first_id_tx, first_id_rx) = mpsc::channel();
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"failure",
Box::new(|_| Ok([1; 64])),
Box::new(move |bytes| {
first_id_tx.send(TxId::for_transaction(bytes)).unwrap();
Ok(())
}),
))
.unwrap();
});
let first_id = receive(&first_id_rx, "first queue did not run");
gate.wait_for(1);
let (second_id_tx, second_id_rx) = mpsc::channel();
let (second_tx, second_rx) = mpsc::channel();
let second_ordering = ordering.clone();
let second = thread::spawn(move || {
second_tx
.send(second_ordering.submit_local_txn(
2,
[1; 32],
a,
b"skipped",
Box::new(|_| Ok([2; 64])),
Box::new(move |bytes| {
second_id_tx.send(TxId::for_transaction(bytes)).unwrap();
Ok(())
}),
))
.unwrap();
});
let second_id = receive(&second_id_rx, "second queue did not run");
assert_eq!(gate.count(), 1);
assert_pending(&second_rx, "second faulted-lane caller");
gate.release();
let first_message = receive(&first_rx, "first fault did not return").unwrap_err();
let second_message = receive(&second_rx, "second fault did not return").unwrap_err();
assert_committed(&first_message, first_id);
assert_committed(&second_message, second_id);
first.join().unwrap();
second.join().unwrap();
assert_eq!(failing.payloads(), vec![b"failure".to_vec()]);
let (b1_tx, b1_rx) = mpsc::channel();
let b1_ordering = ordering.clone();
let b1 = thread::spawn(move || {
b1_tx
.send(b1_ordering.submit_local_txn(
3,
[3; 32],
b,
b"b1",
Box::new(|_| Ok([3; 64])),
Box::new(|_| Ok(())),
))
.unwrap();
});
receive(&b1_rx, "B blocked by A fault").unwrap();
b1.join().unwrap();
let panicking = Recording::plain();
panicking.panic_submit.store(true, Ordering::Release);
ordering
.register_subsystem(a, Some(second_id), panicking)
.unwrap();
let peer_bytes = transaction(ordering.tip().unwrap(), 4, 4, a, b"panic");
let peer_id = TxId::for_transaction(&peer_bytes);
let (peer_tx, peer_rx) = mpsc::channel();
let peer_ordering = ordering.clone();
let peer =
thread::spawn(move || peer_tx.send(peer_ordering.submit_txn(&peer_bytes)).unwrap());
let message = match receive(&peer_rx, "panicking callback did not return") {
Err(SubmitError::Other(message)) => message,
result => panic!("unexpected peer result: {result:?}"),
};
peer.join().unwrap();
assert_committed(&message, peer_id);
assert!(ordering.contains(peer_id));
ordering
.submit_local_txn(
5,
[5; 32],
b,
b"b2",
Box::new(|_| Ok([5; 64])),
Box::new(|_| Ok(())),
)
.unwrap();
assert_eq!(handler_b.payloads(), vec![b"b1".to_vec(), b"b2".to_vec()]);
})
}
pub fn verify_reorganization_isolation(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("reorganization isolation", || {
let root = TempRoot::new("reorg");
let ordering = harness.open(&root.0).unwrap();
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let reorg_gate = Arc::new(Gate::new());
let _reorg_release = reorg_gate.release_on_drop();
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");
let incumbent_id = TxId::for_transaction(&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_tx, replacement_rx) = mpsc::channel();
let replacement_ordering = ordering.clone();
let replacement_work = thread::spawn(move || {
replacement_tx
.send(replacement_ordering.submit_txn(&replacement))
.unwrap();
});
reorg_gate.wait_for(1);
b_gate.wait_for(1);
let (unrelated_tx, unrelated_rx) = mpsc::channel();
let unrelated_ordering = ordering.clone();
let unrelated = thread::spawn(move || {
let state = (
unrelated_ordering.tip(),
unrelated_ordering.contains(replacement_id),
unrelated_ordering.contains(incumbent_id),
);
let result = unrelated_ordering.submit_local_txn(
4,
[4; 32],
b,
b"after",
Box::new(|_| Ok([4; 64])),
Box::new(|_| Ok(())),
);
unrelated_tx.send((state, result)).unwrap();
});
let ((tip, replacement_visible, incumbent_visible), result) =
receive(&unrelated_rx, "unrelated work blocked by reorg");
assert_eq!(tip, Some(replacement_id));
assert!(replacement_visible);
assert!(!incumbent_visible);
result.unwrap();
unrelated.join().unwrap();
reorg_gate.release();
receive(&replacement_rx, "reorganization did not finish").unwrap();
replacement_work.join().unwrap();
assert_eq!(handler_a.reorgs.load(Ordering::Acquire), 1);
assert_eq!(
handler_a.payloads(),
vec![b"first".to_vec(), b"incumbent".to_vec()]
);
assert_eq!(
handler_b.payloads(),
vec![b"replacement".to_vec(), b"after".to_vec()]
);
let restored = Recording::plain();
ordering
.register_subsystem(a, Some(first_id), restored.clone())
.unwrap();
ordering
.submit_local_txn(
5,
[5; 32],
a,
b"restored",
Box::new(|_| Ok([5; 64])),
Box::new(|_| Ok(())),
)
.unwrap();
assert_eq!(restored.payloads(), vec![b"restored".to_vec()]);
assert_eq!(
handler_a.payloads(),
vec![b"first".to_vec(), b"incumbent".to_vec()]
);
})
}
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()]);
})
}
pub fn verify_restart_duplicates_and_queries(harness: &dyn OrderingHarness) -> Result<(), String> {
run_scenario("restart duplicates and queries", || {
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 = harness.open(&root.0).unwrap();
let gate = Arc::new(Gate::new());
let _release = gate.release_on_drop();
let first_handler = Arc::new(Recording::new(Some(gate.clone()), None, None));
ordering
.register_subsystem(owner, None, first_handler.clone())
.unwrap();
let (original_tx, original_rx) = mpsc::channel();
let original_ordering = ordering.clone();
let original_bytes = first.clone();
let original = thread::spawn(move || {
original_tx
.send(original_ordering.submit_txn(&original_bytes))
.unwrap();
});
gate.wait_for(1);
let (duplicate_tx, duplicate_rx) = mpsc::channel();
let duplicate_ordering = ordering.clone();
let duplicate_bytes = first.clone();
let duplicate = thread::spawn(move || {
duplicate_tx
.send(duplicate_ordering.submit_txn(&duplicate_bytes))
.unwrap();
});
receive(&duplicate_rx, "duplicate blocked by original callback").unwrap();
duplicate.join().unwrap();
assert_pending(&original_rx, "original callback");
assert_eq!(first_handler.payloads(), vec![b"first".to_vec()]);
gate.release();
receive(&original_rx, "original callback did not finish").unwrap();
original.join().unwrap();
ordering.submit_txn(&second).unwrap();
ordering.submit_txn(&second).unwrap();
assert_eq!(
first_handler.payloads(),
vec![b"first".to_vec(), b"second".to_vec()]
);
assert!(ordering.contains(first_id));
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()));
drop(ordering);
let reopened = harness.open(&root.0).unwrap();
let replayed = Recording::plain();
reopened
.register_subsystem(owner, None, replayed.clone())
.unwrap();
reopened.submit_txn(&second).unwrap();
assert_eq!(
replayed.payloads(),
vec![b"first".to_vec(), b"second".to_vec()]
);
assert!(reopened.contains(first_id));
assert_eq!(reopened.tip(), Some(second_id));
let empty_root = TempRoot::new("empty");
fs::create_dir_all(&empty_root.0).unwrap();
let empty = harness.open(&empty_root.0).unwrap();
assert_eq!(empty.tip(), None);
})
}