Skip to main content

kcode_k1_txn_ordering_testkit/
lib.rs

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(&register_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}