1pub use kcode_k1_canonical_chain::{SubmitError, TxId};
2pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
3use kcode_k1_transaction::{Transaction, build_signed_transaction};
4use std::any::Any;
5use std::fs;
6use std::panic::{AssertUnwindSafe, catch_unwind};
7use std::path::{Path, PathBuf};
8use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
9use std::sync::{Arc, Condvar, Mutex, mpsc};
10use std::thread;
11use std::time::Duration;
12
13pub type TestSigner = Box<dyn FnOnce(&[u8]) -> Result<[u8; 64], String> + Send + 'static>;
14pub type TestQueuePropagation = Box<dyn FnOnce(&[u8]) -> Result<(), String> + Send + 'static>;
15
16pub trait TestSubsystem: Send + Sync + 'static {
17 fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
18 fn reorg(&self) -> Result<(), String>;
19}
20
21pub trait OrderingCandidate: Send + Sync + 'static {
22 fn register_subsystem(
23 &self,
24 subsystem: SubsystemId,
25 after: Option<TxId>,
26 handler: Arc<dyn TestSubsystem>,
27 ) -> Result<(), String>;
28 fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError>;
29 fn submit_local_txn(
30 &self,
31 timestamp: u64,
32 creator: [u8; 32],
33 subsystem: SubsystemId,
34 payload: &[u8],
35 signer: TestSigner,
36 queue_propagation: TestQueuePropagation,
37 ) -> Result<Vec<u8>, String>;
38 fn contains(&self, id: TxId) -> bool;
39 fn tip(&self) -> Option<TxId>;
40 fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String>;
41 fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String>;
42}
43
44pub trait OrderingHarness: Send + Sync {
45 fn open(&self, root: &Path) -> Result<Arc<dyn OrderingCandidate>, String>;
46}
47
48const TIMEOUT: Duration = Duration::from_secs(2);
49const QUIET: Duration = Duration::from_millis(50);
50static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
51
52struct TempRoot(PathBuf);
53
54impl TempRoot {
55 fn new(label: &str) -> Self {
56 let number = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
57 let path = std::env::temp_dir().join(format!(
58 "kcode-k1-txn-ordering-testkit-{}-{number}-{label}",
59 std::process::id()
60 ));
61 let _ = fs::remove_dir_all(&path);
62 Self(path)
63 }
64}
65
66impl Drop for TempRoot {
67 fn drop(&mut self) {
68 if !thread::panicking() {
69 let _ = fs::remove_dir_all(&self.0);
70 }
71 }
72}
73
74struct Gate {
75 state: Mutex<(usize, bool)>,
76 changed: Condvar,
77}
78
79impl Gate {
80 fn new() -> Self {
81 Self {
82 state: Mutex::new((0, false)),
83 changed: Condvar::new(),
84 }
85 }
86
87 fn enter(&self) {
88 let mut state = self.state.lock().unwrap();
89 state.0 += 1;
90 self.changed.notify_all();
91 let state = self.changed.wait_while(state, |state| !state.1).unwrap();
92 assert!(state.1);
93 }
94
95 fn wait_for(&self, count: usize) {
96 let state = self.state.lock().unwrap();
97 let (state, _) = self
98 .changed
99 .wait_timeout_while(state, TIMEOUT, |state| state.0 < count)
100 .unwrap();
101 assert!(state.0 >= count, "gate entry timed out");
102 }
103
104 fn count(&self) -> usize {
105 self.state.lock().unwrap().0
106 }
107
108 fn release(&self) {
109 self.state.lock().unwrap().1 = true;
110 self.changed.notify_all();
111 }
112
113 fn release_on_drop(self: &Arc<Self>) -> GateRelease {
114 GateRelease(self.clone())
115 }
116}
117
118struct GateRelease(Arc<Gate>);
119
120impl Drop for GateRelease {
121 fn drop(&mut self) {
122 self.0.release();
123 }
124}
125
126struct Recording {
127 payloads: Mutex<Vec<Vec<u8>>>,
128 reorgs: AtomicUsize,
129 submit_gate: Option<Arc<Gate>>,
130 reorg_gate: Option<Arc<Gate>>,
131 queue_flag: Option<Arc<AtomicBool>>,
132 queue_observed: AtomicBool,
133 fail_submit: AtomicBool,
134 panic_submit: AtomicBool,
135}
136
137impl Recording {
138 fn new(
139 submit_gate: Option<Arc<Gate>>,
140 reorg_gate: Option<Arc<Gate>>,
141 queue_flag: Option<Arc<AtomicBool>>,
142 ) -> Self {
143 Self {
144 payloads: Mutex::new(Vec::new()),
145 reorgs: AtomicUsize::new(0),
146 submit_gate,
147 reorg_gate,
148 queue_flag,
149 queue_observed: AtomicBool::new(false),
150 fail_submit: AtomicBool::new(false),
151 panic_submit: AtomicBool::new(false),
152 }
153 }
154
155 fn plain() -> Arc<Self> {
156 Arc::new(Self::new(None, None, None))
157 }
158
159 fn payloads(&self) -> Vec<Vec<u8>> {
160 self.payloads.lock().unwrap().clone()
161 }
162}
163
164impl TestSubsystem for Recording {
165 fn submit_txn(&self, _id: TxId, payload: &[u8]) -> Result<(), String> {
166 self.payloads.lock().unwrap().push(payload.to_vec());
167 if let Some(flag) = &self.queue_flag {
168 self.queue_observed
169 .store(flag.load(Ordering::Acquire), Ordering::Release);
170 }
171 if let Some(gate) = &self.submit_gate {
172 gate.enter();
173 }
174 if self.panic_submit.load(Ordering::Acquire) {
175 panic!("submit panic");
176 }
177 if self.fail_submit.load(Ordering::Acquire) {
178 Err("submit failure".to_owned())
179 } else {
180 Ok(())
181 }
182 }
183
184 fn reorg(&self) -> Result<(), String> {
185 self.reorgs.fetch_add(1, Ordering::Relaxed);
186 if let Some(gate) = &self.reorg_gate {
187 gate.enter();
188 }
189 Ok(())
190 }
191}
192
193struct Reentrant {
194 ordering: Arc<dyn OrderingCandidate>,
195 target: SubsystemId,
196}
197
198impl TestSubsystem for Reentrant {
199 fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
200 self.ordering
201 .submit_local_txn(
202 99,
203 [9; 32],
204 self.target,
205 b"reentered",
206 Box::new(|_| Ok([9; 64])),
207 Box::new(|_| Ok(())),
208 )
209 .map(|_| ())
210 }
211
212 fn reorg(&self) -> Result<(), String> {
213 Ok(())
214 }
215}
216
217fn subsystem(value: u8) -> SubsystemId {
218 SubsystemId::from_bytes([value; 20]).unwrap()
219}
220
221fn transaction(
222 parent: TxId,
223 creator: u8,
224 timestamp: u64,
225 subsystem: SubsystemId,
226 payload: &[u8],
227) -> Vec<u8> {
228 build_signed_transaction(parent, timestamp, [creator; 32], subsystem, payload, |_| {
229 Ok([creator; 64])
230 })
231 .unwrap()
232}
233
234fn assert_committed(message: &str, id: TxId) {
235 assert!(message.contains("TxId"));
236 assert!(message.contains("was committed"));
237 assert!(message.contains(&format!("{id:?}")));
238}
239
240fn receive<T>(receiver: &mpsc::Receiver<T>, label: &str) -> T {
241 receiver
242 .recv_timeout(TIMEOUT)
243 .unwrap_or_else(|error| panic!("{label}: {error}"))
244}
245
246fn assert_pending<T>(receiver: &mpsc::Receiver<T>, label: &str) {
247 match receiver.recv_timeout(QUIET) {
248 Err(mpsc::RecvTimeoutError::Timeout) => {}
249 Err(mpsc::RecvTimeoutError::Disconnected) => panic!("{label} disconnected"),
250 Ok(_) => panic!("{label} completed while it should have been blocked"),
251 }
252}
253
254fn panic_message(payload: Box<dyn Any + Send>) -> String {
255 if let Some(message) = payload.downcast_ref::<String>() {
256 message.clone()
257 } else if let Some(message) = payload.downcast_ref::<&str>() {
258 (*message).to_owned()
259 } else {
260 "non-string panic".to_owned()
261 }
262}
263
264fn run_scenario<F>(name: &str, scenario: F) -> Result<(), String>
265where
266 F: FnOnce(),
267{
268 match catch_unwind(AssertUnwindSafe(scenario)) {
269 Ok(()) => Ok(()),
270 Err(payload) => Err(format!("{name}: {}", panic_message(payload))),
271 }
272}
273
274pub fn verify_queue_before_callback(harness: &dyn OrderingHarness) -> Result<(), String> {
275 run_scenario("queue before callback", || {
276 let root = TempRoot::new("queue");
277 let ordering = harness.open(&root.0).unwrap();
278 let owner = subsystem(b'a');
279 let queued = Arc::new(AtomicBool::new(false));
280 let handler = Arc::new(Recording::new(None, None, Some(queued.clone())));
281 ordering
282 .register_subsystem(owner, None, handler.clone())
283 .unwrap();
284 let first_queued = queued.clone();
285 let first = ordering.submit_local_txn(
286 1,
287 [1; 32],
288 owner,
289 b"first",
290 Box::new(|_| Ok([1; 64])),
291 Box::new(move |bytes| {
292 assert_eq!(Transaction::parse(bytes).unwrap().payload(), b"first");
293 first_queued.store(true, Ordering::Release);
294 Err("queue unavailable".to_owned())
295 }),
296 );
297 let first_id = ordering.tip().unwrap();
298 assert_committed(&first.unwrap_err(), first_id);
299 assert!(handler.queue_observed.load(Ordering::Acquire));
300 queued.store(false, Ordering::Release);
301 let second_queued = queued.clone();
302 let second = ordering.submit_local_txn(
303 2,
304 [1; 32],
305 owner,
306 b"second",
307 Box::new(|_| Ok([2; 64])),
308 Box::new(move |_| {
309 second_queued.store(true, Ordering::Release);
310 panic!("queue panic")
311 }),
312 );
313 let second_id = ordering.tip().unwrap();
314 assert_committed(&second.unwrap_err(), second_id);
315 assert!(handler.queue_observed.load(Ordering::Acquire));
316 ordering
317 .submit_local_txn(
318 3,
319 [1; 32],
320 owner,
321 b"third",
322 Box::new(|_| Ok([3; 64])),
323 Box::new(|_| Ok(())),
324 )
325 .unwrap();
326 assert_eq!(
327 handler.payloads(),
328 vec![b"first".to_vec(), b"second".to_vec(), b"third".to_vec()]
329 );
330 })
331}
332
333pub fn verify_independent_subsystem_lanes(harness: &dyn OrderingHarness) -> Result<(), String> {
334 run_scenario("independent subsystem lanes", || {
335 let root = TempRoot::new("lanes");
336 let ordering = harness.open(&root.0).unwrap();
337 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
338 let gate = Arc::new(Gate::new());
339 let _release = gate.release_on_drop();
340 let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
341 let handler_b = Recording::plain();
342 ordering
343 .register_subsystem(a, None, handler_a.clone())
344 .unwrap();
345 ordering
346 .register_subsystem(b, None, handler_b.clone())
347 .unwrap();
348 let (first_tx, first_rx) = mpsc::channel();
349 let first_ordering = ordering.clone();
350 let first = thread::spawn(move || {
351 first_tx
352 .send(first_ordering.submit_local_txn(
353 1,
354 [1; 32],
355 a,
356 b"a1",
357 Box::new(|_| Ok([1; 64])),
358 Box::new(|_| Ok(())),
359 ))
360 .unwrap();
361 });
362 gate.wait_for(1);
363 let (query_tx, query_rx) = mpsc::channel();
364 let query_ordering = ordering.clone();
365 let query = thread::spawn(move || {
366 let id = query_ordering.tip().unwrap();
367 query_tx
368 .send((
369 id,
370 query_ordering.contains(id),
371 query_ordering.get_txn(id),
372 query_ordering.between_txids(GENESIS_PARENT, id),
373 ))
374 .unwrap();
375 });
376 let (_id, contains, stored, between) = receive(&query_rx, "queries blocked by A");
377 assert!(contains);
378 assert!(!stored.unwrap().unwrap().is_empty());
379 assert!(between.unwrap().is_empty());
380 query.join().unwrap();
381 let (queued_tx, queued_rx) = mpsc::channel();
382 let (a_tx, a_rx) = mpsc::channel();
383 let second_ordering = ordering.clone();
384 let second = thread::spawn(move || {
385 a_tx.send(second_ordering.submit_local_txn(
386 2,
387 [1; 32],
388 a,
389 b"a2",
390 Box::new(|_| Ok([2; 64])),
391 Box::new(move |_| {
392 queued_tx.send(()).unwrap();
393 Ok(())
394 }),
395 ))
396 .unwrap();
397 });
398 receive(&queued_rx, "second A queue blocked");
399 assert_eq!(gate.count(), 1);
400 assert_pending(&a_rx, "second A caller");
401 let (b_tx, b_rx) = mpsc::channel();
402 let b_ordering = ordering.clone();
403 let b_work = thread::spawn(move || {
404 b_tx.send(b_ordering.submit_local_txn(
405 3,
406 [3; 32],
407 b,
408 b"b1",
409 Box::new(|_| Ok([3; 64])),
410 Box::new(|_| Ok(())),
411 ))
412 .unwrap();
413 });
414 receive(&b_rx, "B submission blocked by A").unwrap();
415 b_work.join().unwrap();
416 gate.release();
417 receive(&first_rx, "first A did not finish").unwrap();
418 first.join().unwrap();
419 receive(&a_rx, "second A did not finish").unwrap();
420 second.join().unwrap();
421 assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
422 assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
423 })
424}
425
426pub fn verify_signing_commit_exclusion(harness: &dyn OrderingHarness) -> Result<(), String> {
427 run_scenario("signing commit exclusion", || {
428 let root = TempRoot::new("signer");
429 let ordering = harness.open(&root.0).unwrap();
430 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
431 ordering
432 .register_subsystem(a, None, Recording::plain())
433 .unwrap();
434 ordering
435 .register_subsystem(b, None, Recording::plain())
436 .unwrap();
437 let gate = Arc::new(Gate::new());
438 let _release = gate.release_on_drop();
439 let first_gate = gate.clone();
440 let (first_tx, first_rx) = mpsc::channel();
441 let first_ordering = ordering.clone();
442 let first = thread::spawn(move || {
443 first_tx
444 .send(first_ordering.submit_local_txn(
445 1,
446 [1; 32],
447 a,
448 b"a",
449 Box::new(move |_| {
450 first_gate.enter();
451 Ok([1; 64])
452 }),
453 Box::new(|_| Ok(())),
454 ))
455 .unwrap();
456 });
457 gate.wait_for(1);
458 let (submit_ready_tx, submit_ready_rx) = mpsc::channel();
459 let (submit_tx, submit_rx) = mpsc::channel();
460 let submit_ordering = ordering.clone();
461 let submit = thread::spawn(move || {
462 submit_ready_tx.send(()).unwrap();
463 submit_tx
464 .send(submit_ordering.submit_local_txn(
465 2,
466 [2; 32],
467 b,
468 b"b",
469 Box::new(|_| Ok([2; 64])),
470 Box::new(|_| Ok(())),
471 ))
472 .unwrap();
473 });
474 let (query_ready_tx, query_ready_rx) = mpsc::channel();
475 let (query_tx, query_rx) = mpsc::channel();
476 let query_ordering = ordering.clone();
477 let query = thread::spawn(move || {
478 query_ready_tx.send(()).unwrap();
479 query_tx.send(query_ordering.tip()).unwrap();
480 });
481 receive(&submit_ready_rx, "second writer did not start");
482 receive(&query_ready_rx, "query did not start");
483 assert_pending(&submit_rx, "second writer");
484 assert_pending(&query_rx, "tip query");
485 gate.release();
486 receive(&first_rx, "first writer did not finish").unwrap();
487 first.join().unwrap();
488 receive(&submit_rx, "second writer did not finish").unwrap();
489 assert!(receive(&query_rx, "tip query did not finish").is_some());
490 submit.join().unwrap();
491 query.join().unwrap();
492 })
493}
494
495pub fn verify_replay_live_handoff(harness: &dyn OrderingHarness) -> Result<(), String> {
496 run_scenario("replay live handoff", || {
497 let root = TempRoot::new("replay");
498 let ordering = harness.open(&root.0).unwrap();
499 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
500 let first = transaction(GENESIS_PARENT, 8, 1, a, b"a1");
501 ordering.submit_txn(&first).unwrap();
502 let gate = Arc::new(Gate::new());
503 let _release = gate.release_on_drop();
504 let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
505 let (register_tx, register_rx) = mpsc::channel();
506 let register_ordering = ordering.clone();
507 let register_handler = handler_a.clone();
508 let register = thread::spawn(move || {
509 register_tx
510 .send(register_ordering.register_subsystem(a, None, register_handler))
511 .unwrap();
512 });
513 gate.wait_for(1);
514 let handler_b = Recording::plain();
515 let (b_tx, b_rx) = mpsc::channel();
516 let b_ordering = ordering.clone();
517 let b_handler = handler_b.clone();
518 let b_work = thread::spawn(move || {
519 let result = b_ordering
520 .register_subsystem(b, None, b_handler)
521 .and_then(|()| {
522 b_ordering.submit_local_txn(
523 2,
524 [2; 32],
525 b,
526 b"b1",
527 Box::new(|_| Ok([2; 64])),
528 Box::new(|_| Ok(())),
529 )
530 });
531 b_tx.send(result).unwrap();
532 });
533 let b_bytes = receive(&b_rx, "B work blocked by replay").unwrap();
534 b_work.join().unwrap();
535 let second = transaction(TxId::for_transaction(&b_bytes), 8, 3, a, b"a2");
536 let (peer_tx, peer_rx) = mpsc::channel();
537 let peer_ordering = ordering.clone();
538 let peer = thread::spawn(move || peer_tx.send(peer_ordering.submit_txn(&second)).unwrap());
539 receive(&peer_rx, "A commit blocked by replay").unwrap();
540 peer.join().unwrap();
541 assert_eq!(gate.count(), 1);
542 assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
543 gate.release();
544 receive(®ister_rx, "A registration did not finish").unwrap();
545 register.join().unwrap();
546 assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
547 })
548}
549
550pub fn verify_callback_failure_isolation(harness: &dyn OrderingHarness) -> Result<(), String> {
551 run_scenario("callback failure isolation", || {
552 let root = TempRoot::new("faults");
553 let ordering = harness.open(&root.0).unwrap();
554 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
555 let gate = Arc::new(Gate::new());
556 let _release = gate.release_on_drop();
557 let failing = Arc::new(Recording::new(Some(gate.clone()), None, None));
558 failing.fail_submit.store(true, Ordering::Release);
559 let handler_b = Recording::plain();
560 ordering
561 .register_subsystem(a, None, failing.clone())
562 .unwrap();
563 ordering
564 .register_subsystem(b, None, handler_b.clone())
565 .unwrap();
566 let (first_id_tx, first_id_rx) = mpsc::channel();
567 let (first_tx, first_rx) = mpsc::channel();
568 let first_ordering = ordering.clone();
569 let first = thread::spawn(move || {
570 first_tx
571 .send(first_ordering.submit_local_txn(
572 1,
573 [1; 32],
574 a,
575 b"failure",
576 Box::new(|_| Ok([1; 64])),
577 Box::new(move |bytes| {
578 first_id_tx.send(TxId::for_transaction(bytes)).unwrap();
579 Ok(())
580 }),
581 ))
582 .unwrap();
583 });
584 let first_id = receive(&first_id_rx, "first queue did not run");
585 gate.wait_for(1);
586 let (second_id_tx, second_id_rx) = mpsc::channel();
587 let (second_tx, second_rx) = mpsc::channel();
588 let second_ordering = ordering.clone();
589 let second = thread::spawn(move || {
590 second_tx
591 .send(second_ordering.submit_local_txn(
592 2,
593 [1; 32],
594 a,
595 b"skipped",
596 Box::new(|_| Ok([2; 64])),
597 Box::new(move |bytes| {
598 second_id_tx.send(TxId::for_transaction(bytes)).unwrap();
599 Ok(())
600 }),
601 ))
602 .unwrap();
603 });
604 let second_id = receive(&second_id_rx, "second queue did not run");
605 assert_eq!(gate.count(), 1);
606 assert_pending(&second_rx, "second faulted-lane caller");
607 gate.release();
608 let first_message = receive(&first_rx, "first fault did not return").unwrap_err();
609 let second_message = receive(&second_rx, "second fault did not return").unwrap_err();
610 assert_committed(&first_message, first_id);
611 assert_committed(&second_message, second_id);
612 first.join().unwrap();
613 second.join().unwrap();
614 assert_eq!(failing.payloads(), vec![b"failure".to_vec()]);
615 let (b1_tx, b1_rx) = mpsc::channel();
616 let b1_ordering = ordering.clone();
617 let b1 = thread::spawn(move || {
618 b1_tx
619 .send(b1_ordering.submit_local_txn(
620 3,
621 [3; 32],
622 b,
623 b"b1",
624 Box::new(|_| Ok([3; 64])),
625 Box::new(|_| Ok(())),
626 ))
627 .unwrap();
628 });
629 receive(&b1_rx, "B blocked by A fault").unwrap();
630 b1.join().unwrap();
631 let panicking = Recording::plain();
632 panicking.panic_submit.store(true, Ordering::Release);
633 ordering
634 .register_subsystem(a, Some(second_id), panicking)
635 .unwrap();
636 let peer_bytes = transaction(ordering.tip().unwrap(), 4, 4, a, b"panic");
637 let peer_id = TxId::for_transaction(&peer_bytes);
638 let (peer_tx, peer_rx) = mpsc::channel();
639 let peer_ordering = ordering.clone();
640 let peer =
641 thread::spawn(move || peer_tx.send(peer_ordering.submit_txn(&peer_bytes)).unwrap());
642 let message = match receive(&peer_rx, "panicking callback did not return") {
643 Err(SubmitError::Other(message)) => message,
644 result => panic!("unexpected peer result: {result:?}"),
645 };
646 peer.join().unwrap();
647 assert_committed(&message, peer_id);
648 assert!(ordering.contains(peer_id));
649 ordering
650 .submit_local_txn(
651 5,
652 [5; 32],
653 b,
654 b"b2",
655 Box::new(|_| Ok([5; 64])),
656 Box::new(|_| Ok(())),
657 )
658 .unwrap();
659 assert_eq!(handler_b.payloads(), vec![b"b1".to_vec(), b"b2".to_vec()]);
660 })
661}
662
663pub fn verify_reorganization_isolation(harness: &dyn OrderingHarness) -> Result<(), String> {
664 run_scenario("reorganization isolation", || {
665 let root = TempRoot::new("reorg");
666 let ordering = harness.open(&root.0).unwrap();
667 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
668 let reorg_gate = Arc::new(Gate::new());
669 let _reorg_release = reorg_gate.release_on_drop();
670 let handler_a = Arc::new(Recording::new(None, Some(reorg_gate.clone()), None));
671 let b_gate = Arc::new(Gate::new());
672 b_gate.release();
673 let handler_b = Arc::new(Recording::new(Some(b_gate.clone()), None, None));
674 ordering
675 .register_subsystem(a, None, handler_a.clone())
676 .unwrap();
677 ordering
678 .register_subsystem(b, None, handler_b.clone())
679 .unwrap();
680 let first = transaction(GENESIS_PARENT, 20, 1, a, b"first");
681 let first_id = TxId::for_transaction(&first);
682 ordering.submit_txn(&first).unwrap();
683 let incumbent = transaction(first_id, 50, 2, a, b"incumbent");
684 let incumbent_id = TxId::for_transaction(&incumbent);
685 ordering.submit_txn(&incumbent).unwrap();
686 let replacement = transaction(first_id, 1, 3, b, b"replacement");
687 let replacement_id = TxId::for_transaction(&replacement);
688 let (replacement_tx, replacement_rx) = mpsc::channel();
689 let replacement_ordering = ordering.clone();
690 let replacement_work = thread::spawn(move || {
691 replacement_tx
692 .send(replacement_ordering.submit_txn(&replacement))
693 .unwrap();
694 });
695 reorg_gate.wait_for(1);
696 b_gate.wait_for(1);
697 let (unrelated_tx, unrelated_rx) = mpsc::channel();
698 let unrelated_ordering = ordering.clone();
699 let unrelated = thread::spawn(move || {
700 let state = (
701 unrelated_ordering.tip(),
702 unrelated_ordering.contains(replacement_id),
703 unrelated_ordering.contains(incumbent_id),
704 );
705 let result = unrelated_ordering.submit_local_txn(
706 4,
707 [4; 32],
708 b,
709 b"after",
710 Box::new(|_| Ok([4; 64])),
711 Box::new(|_| Ok(())),
712 );
713 unrelated_tx.send((state, result)).unwrap();
714 });
715 let ((tip, replacement_visible, incumbent_visible), result) =
716 receive(&unrelated_rx, "unrelated work blocked by reorg");
717 assert_eq!(tip, Some(replacement_id));
718 assert!(replacement_visible);
719 assert!(!incumbent_visible);
720 result.unwrap();
721 unrelated.join().unwrap();
722 reorg_gate.release();
723 receive(&replacement_rx, "reorganization did not finish").unwrap();
724 replacement_work.join().unwrap();
725 assert_eq!(handler_a.reorgs.load(Ordering::Acquire), 1);
726 assert_eq!(
727 handler_a.payloads(),
728 vec![b"first".to_vec(), b"incumbent".to_vec()]
729 );
730 assert_eq!(
731 handler_b.payloads(),
732 vec![b"replacement".to_vec(), b"after".to_vec()]
733 );
734 let restored = Recording::plain();
735 ordering
736 .register_subsystem(a, Some(first_id), restored.clone())
737 .unwrap();
738 ordering
739 .submit_local_txn(
740 5,
741 [5; 32],
742 a,
743 b"restored",
744 Box::new(|_| Ok([5; 64])),
745 Box::new(|_| Ok(())),
746 )
747 .unwrap();
748 assert_eq!(restored.payloads(), vec![b"restored".to_vec()]);
749 assert_eq!(
750 handler_a.payloads(),
751 vec![b"first".to_vec(), b"incumbent".to_vec()]
752 );
753 })
754}
755
756pub fn verify_unrelated_callback_reentry(harness: &dyn OrderingHarness) -> Result<(), String> {
757 run_scenario("unrelated callback reentry", || {
758 let root = TempRoot::new("reentry");
759 let ordering = harness.open(&root.0).unwrap();
760 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
761 let handler_b = Recording::plain();
762 ordering
763 .register_subsystem(b, None, handler_b.clone())
764 .unwrap();
765 ordering
766 .register_subsystem(
767 a,
768 None,
769 Arc::new(Reentrant {
770 ordering: ordering.clone(),
771 target: b,
772 }),
773 )
774 .unwrap();
775 let (outer_tx, outer_rx) = mpsc::channel();
776 let outer_ordering = ordering.clone();
777 let outer = thread::spawn(move || {
778 outer_tx
779 .send(outer_ordering.submit_local_txn(
780 1,
781 [1; 32],
782 a,
783 b"outer",
784 Box::new(|_| Ok([1; 64])),
785 Box::new(|_| Ok(())),
786 ))
787 .unwrap();
788 });
789 receive(&outer_rx, "reentrant callback deadlocked").unwrap();
790 outer.join().unwrap();
791 assert_eq!(handler_b.payloads(), vec![b"reentered".to_vec()]);
792 })
793}
794
795pub fn verify_restart_duplicates_and_queries(harness: &dyn OrderingHarness) -> Result<(), String> {
796 run_scenario("restart duplicates and queries", || {
797 let root = TempRoot::new("restart");
798 let owner = subsystem(b'z');
799 let first = transaction(GENESIS_PARENT, 1, 1, owner, b"first");
800 let first_id = TxId::for_transaction(&first);
801 let second = transaction(first_id, 1, 2, owner, b"second");
802 let second_id = TxId::for_transaction(&second);
803 let ordering = harness.open(&root.0).unwrap();
804 let gate = Arc::new(Gate::new());
805 let _release = gate.release_on_drop();
806 let first_handler = Arc::new(Recording::new(Some(gate.clone()), None, None));
807 ordering
808 .register_subsystem(owner, None, first_handler.clone())
809 .unwrap();
810 let (original_tx, original_rx) = mpsc::channel();
811 let original_ordering = ordering.clone();
812 let original_bytes = first.clone();
813 let original = thread::spawn(move || {
814 original_tx
815 .send(original_ordering.submit_txn(&original_bytes))
816 .unwrap();
817 });
818 gate.wait_for(1);
819 let (duplicate_tx, duplicate_rx) = mpsc::channel();
820 let duplicate_ordering = ordering.clone();
821 let duplicate_bytes = first.clone();
822 let duplicate = thread::spawn(move || {
823 duplicate_tx
824 .send(duplicate_ordering.submit_txn(&duplicate_bytes))
825 .unwrap();
826 });
827 receive(&duplicate_rx, "duplicate blocked by original callback").unwrap();
828 duplicate.join().unwrap();
829 assert_pending(&original_rx, "original callback");
830 assert_eq!(first_handler.payloads(), vec![b"first".to_vec()]);
831 gate.release();
832 receive(&original_rx, "original callback did not finish").unwrap();
833 original.join().unwrap();
834 ordering.submit_txn(&second).unwrap();
835 ordering.submit_txn(&second).unwrap();
836 assert_eq!(
837 first_handler.payloads(),
838 vec![b"first".to_vec(), b"second".to_vec()]
839 );
840 assert!(ordering.contains(first_id));
841 assert_eq!(ordering.tip(), Some(second_id));
842 assert_eq!(
843 ordering.between_txids(GENESIS_PARENT, second_id).unwrap(),
844 vec![first_id]
845 );
846 assert_eq!(ordering.get_txn(first_id).unwrap(), Some(first.clone()));
847 drop(ordering);
848 let reopened = harness.open(&root.0).unwrap();
849 let replayed = Recording::plain();
850 reopened
851 .register_subsystem(owner, None, replayed.clone())
852 .unwrap();
853 reopened.submit_txn(&second).unwrap();
854 assert_eq!(
855 replayed.payloads(),
856 vec![b"first".to_vec(), b"second".to_vec()]
857 );
858 assert!(reopened.contains(first_id));
859 assert_eq!(reopened.tip(), Some(second_id));
860 let empty_root = TempRoot::new("empty");
861 fs::create_dir_all(&empty_root.0).unwrap();
862 let empty = harness.open(&empty_root.0).unwrap();
863 assert_eq!(empty.tip(), None);
864 })
865}