Skip to main content

kcode_k1_txn_ordering_live_testkit/
lib.rs

1pub use kcode_k1_canonical_chain::{SubmitError, TxId};
2use kcode_k1_transaction::Transaction;
3pub use kcode_k1_transaction::{GENESIS_PARENT, REGISTER_AT_TIP, SubsystemId};
4use std::any::Any;
5use std::fs;
6use std::panic::{AssertUnwindSafe, catch_unwind};
7use std::path::{Path, PathBuf};
8use std::sync::atomic::{AtomicBool, AtomicU64, 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<(TxId, 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-live-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    submit_gate: Option<Arc<Gate>>,
129    queue_flag: Option<Arc<AtomicBool>>,
130    queue_observed: AtomicBool,
131}
132
133impl Recording {
134    fn new(submit_gate: Option<Arc<Gate>>, queue_flag: Option<Arc<AtomicBool>>) -> Self {
135        Self {
136            payloads: Mutex::new(Vec::new()),
137            submit_gate,
138            queue_flag,
139            queue_observed: AtomicBool::new(false),
140        }
141    }
142
143    fn plain() -> Arc<Self> {
144        Arc::new(Self::new(None, None))
145    }
146
147    fn payloads(&self) -> Vec<Vec<u8>> {
148        self.payloads.lock().unwrap().clone()
149    }
150}
151
152impl TestSubsystem for Recording {
153    fn submit_txn(&self, _id: TxId, payload: &[u8]) -> Result<(), String> {
154        self.payloads.lock().unwrap().push(payload.to_vec());
155        if let Some(flag) = &self.queue_flag {
156            self.queue_observed
157                .store(flag.load(Ordering::Acquire), Ordering::Release);
158        }
159        if let Some(gate) = &self.submit_gate {
160            gate.enter();
161        }
162        Ok(())
163    }
164
165    fn reorg(&self) -> Result<(), String> {
166        Ok(())
167    }
168}
169
170struct Reentrant {
171    ordering: Arc<dyn OrderingCandidate>,
172    target: SubsystemId,
173}
174
175impl TestSubsystem for Reentrant {
176    fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
177        self.ordering
178            .submit_local_txn(
179                99,
180                [9; 32],
181                self.target,
182                b"reentered",
183                Box::new(|_| Ok([9; 64])),
184                Box::new(|_| Ok(())),
185            )
186            .map(|_| ())
187    }
188
189    fn reorg(&self) -> Result<(), String> {
190        Ok(())
191    }
192}
193
194fn subsystem(value: u8) -> SubsystemId {
195    SubsystemId::from_bytes([value; 20]).unwrap()
196}
197
198fn assert_committed(message: &str, id: TxId) {
199    assert!(message.contains("TxId"));
200    assert!(message.contains("was committed"));
201    assert!(message.contains(&format!("{id:?}")));
202}
203
204fn receive<T>(receiver: &mpsc::Receiver<T>, label: &str) -> T {
205    receiver
206        .recv_timeout(TIMEOUT)
207        .unwrap_or_else(|error| panic!("{label}: {error}"))
208}
209
210fn assert_pending<T>(receiver: &mpsc::Receiver<T>, label: &str) {
211    match receiver.recv_timeout(QUIET) {
212        Err(mpsc::RecvTimeoutError::Timeout) => {}
213        Err(mpsc::RecvTimeoutError::Disconnected) => panic!("{label} disconnected"),
214        Ok(_) => panic!("{label} completed while it should have been blocked"),
215    }
216}
217
218fn panic_message(payload: Box<dyn Any + Send>) -> String {
219    if let Some(message) = payload.downcast_ref::<String>() {
220        message.clone()
221    } else if let Some(message) = payload.downcast_ref::<&str>() {
222        (*message).to_owned()
223    } else {
224        "non-string panic".to_owned()
225    }
226}
227
228fn run_scenario<F>(name: &str, scenario: F) -> Result<(), String>
229where
230    F: FnOnce(),
231{
232    match catch_unwind(AssertUnwindSafe(scenario)) {
233        Ok(()) => Ok(()),
234        Err(payload) => Err(format!("{name}: {}", panic_message(payload))),
235    }
236}
237
238pub fn verify_registration_sentinels(harness: &dyn OrderingHarness) -> Result<(), String> {
239    run_scenario("registration sentinels", || {
240        let replay_root = TempRoot::new("sentinel-genesis");
241        let owner = subsystem(b's');
242        {
243            let ordering = harness.open(&replay_root.0).unwrap();
244            ordering
245                .register_subsystem(owner, None, Recording::plain())
246                .unwrap();
247            ordering
248                .submit_local_txn(
249                    1,
250                    [1; 32],
251                    owner,
252                    b"historical",
253                    Box::new(|_| Ok([1; 64])),
254                    Box::new(|_| Ok(())),
255                )
256                .unwrap();
257        }
258        let replayed = harness.open(&replay_root.0).unwrap();
259        let replay_handler = Recording::plain();
260        replayed
261            .register_subsystem(owner, Some(GENESIS_PARENT), replay_handler.clone())
262            .unwrap();
263        assert_eq!(replay_handler.payloads(), vec![b"historical".to_vec()]);
264        drop(replayed);
265
266        let tip_root = TempRoot::new("sentinel-tip");
267        {
268            let ordering = harness.open(&tip_root.0).unwrap();
269            ordering
270                .register_subsystem(owner, None, Recording::plain())
271                .unwrap();
272            ordering
273                .submit_local_txn(
274                    1,
275                    [1; 32],
276                    owner,
277                    b"ignored",
278                    Box::new(|_| Ok([1; 64])),
279                    Box::new(|_| Ok(())),
280                )
281                .unwrap();
282        }
283        let at_tip = harness.open(&tip_root.0).unwrap();
284        let tip_handler = Recording::plain();
285        at_tip
286            .register_subsystem(owner, Some(REGISTER_AT_TIP), tip_handler.clone())
287            .unwrap();
288        assert!(tip_handler.payloads().is_empty());
289        let (delivered_id, delivered_bytes) = at_tip
290            .submit_local_txn(
291                2,
292                [1; 32],
293                owner,
294                b"delivered",
295                Box::new(|_| Ok([2; 64])),
296                Box::new(|_| Ok(())),
297            )
298            .unwrap();
299        assert_eq!(delivered_id, TxId::for_transaction(&delivered_bytes));
300        assert_eq!(at_tip.tip(), Some(delivered_id));
301        assert_eq!(tip_handler.payloads(), vec![b"delivered".to_vec()]);
302
303        let empty_root = TempRoot::new("sentinel-empty");
304        let empty = harness.open(&empty_root.0).unwrap();
305        let empty_handler = Recording::plain();
306        empty
307            .register_subsystem(owner, Some(REGISTER_AT_TIP), empty_handler.clone())
308            .unwrap();
309        assert!(empty_handler.payloads().is_empty());
310        empty
311            .submit_local_txn(
312                1,
313                [2; 32],
314                owner,
315                b"first",
316                Box::new(|_| Ok([3; 64])),
317                Box::new(|_| Ok(())),
318            )
319            .unwrap();
320        assert_eq!(empty_handler.payloads(), vec![b"first".to_vec()]);
321    })
322}
323
324pub fn verify_queue_before_callback(harness: &dyn OrderingHarness) -> Result<(), String> {
325    run_scenario("queue before callback", || {
326        let root = TempRoot::new("queue");
327        let ordering = harness.open(&root.0).unwrap();
328        let owner = subsystem(b'a');
329        let queued = Arc::new(AtomicBool::new(false));
330        let handler = Arc::new(Recording::new(None, Some(queued.clone())));
331        ordering
332            .register_subsystem(owner, None, handler.clone())
333            .unwrap();
334        let first_queued = queued.clone();
335        let first = ordering.submit_local_txn(
336            1,
337            [1; 32],
338            owner,
339            b"first",
340            Box::new(|_| Ok([1; 64])),
341            Box::new(move |bytes| {
342                assert_eq!(Transaction::parse(bytes).unwrap().payload(), b"first");
343                first_queued.store(true, Ordering::Release);
344                Err("queue unavailable".to_owned())
345            }),
346        );
347        let first_id = ordering.tip().unwrap();
348        assert_committed(&first.unwrap_err(), first_id);
349        assert!(handler.queue_observed.load(Ordering::Acquire));
350        queued.store(false, Ordering::Release);
351        let second_queued = queued.clone();
352        let second = ordering.submit_local_txn(
353            2,
354            [1; 32],
355            owner,
356            b"second",
357            Box::new(|_| Ok([2; 64])),
358            Box::new(move |_| {
359                second_queued.store(true, Ordering::Release);
360                panic!("queue panic")
361            }),
362        );
363        let second_id = ordering.tip().unwrap();
364        assert_committed(&second.unwrap_err(), second_id);
365        assert!(handler.queue_observed.load(Ordering::Acquire));
366        ordering
367            .submit_local_txn(
368                3,
369                [1; 32],
370                owner,
371                b"third",
372                Box::new(|_| Ok([3; 64])),
373                Box::new(|_| Ok(())),
374            )
375            .unwrap();
376        assert_eq!(
377            handler.payloads(),
378            vec![b"first".to_vec(), b"second".to_vec(), b"third".to_vec()]
379        );
380    })
381}
382
383pub fn verify_independent_subsystem_lanes(harness: &dyn OrderingHarness) -> Result<(), String> {
384    run_scenario("independent subsystem lanes", || {
385        let root = TempRoot::new("lanes");
386        let ordering = harness.open(&root.0).unwrap();
387        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
388        let gate = Arc::new(Gate::new());
389        let _release = gate.release_on_drop();
390        let handler_a = Arc::new(Recording::new(Some(gate.clone()), None));
391        let handler_b = Recording::plain();
392        ordering
393            .register_subsystem(a, None, handler_a.clone())
394            .unwrap();
395        ordering
396            .register_subsystem(b, None, handler_b.clone())
397            .unwrap();
398        let (first_tx, first_rx) = mpsc::channel();
399        let first_ordering = ordering.clone();
400        let first = thread::spawn(move || {
401            first_tx
402                .send(first_ordering.submit_local_txn(
403                    1,
404                    [1; 32],
405                    a,
406                    b"a1",
407                    Box::new(|_| Ok([1; 64])),
408                    Box::new(|_| Ok(())),
409                ))
410                .unwrap();
411        });
412        gate.wait_for(1);
413        let (query_tx, query_rx) = mpsc::channel();
414        let query_ordering = ordering.clone();
415        let query = thread::spawn(move || {
416            let id = query_ordering.tip().unwrap();
417            query_tx
418                .send((
419                    id,
420                    query_ordering.contains(id),
421                    query_ordering.get_txn(id),
422                    query_ordering.between_txids(GENESIS_PARENT, id),
423                ))
424                .unwrap();
425        });
426        let (_id, contains, stored, between) = receive(&query_rx, "queries blocked by A");
427        assert!(contains);
428        assert!(!stored.unwrap().unwrap().is_empty());
429        assert!(between.unwrap().is_empty());
430        query.join().unwrap();
431        let (queued_tx, queued_rx) = mpsc::channel();
432        let (a_tx, a_rx) = mpsc::channel();
433        let second_ordering = ordering.clone();
434        let second = thread::spawn(move || {
435            a_tx.send(second_ordering.submit_local_txn(
436                2,
437                [1; 32],
438                a,
439                b"a2",
440                Box::new(|_| Ok([2; 64])),
441                Box::new(move |_| {
442                    queued_tx.send(()).unwrap();
443                    Ok(())
444                }),
445            ))
446            .unwrap();
447        });
448        receive(&queued_rx, "second A queue blocked");
449        assert_eq!(gate.count(), 1);
450        assert_pending(&a_rx, "second A caller");
451        let (b_tx, b_rx) = mpsc::channel();
452        let b_ordering = ordering.clone();
453        let b_work = thread::spawn(move || {
454            b_tx.send(b_ordering.submit_local_txn(
455                3,
456                [3; 32],
457                b,
458                b"b1",
459                Box::new(|_| Ok([3; 64])),
460                Box::new(|_| Ok(())),
461            ))
462            .unwrap();
463        });
464        receive(&b_rx, "B submission blocked by A").unwrap();
465        b_work.join().unwrap();
466        gate.release();
467        receive(&first_rx, "first A did not finish").unwrap();
468        first.join().unwrap();
469        receive(&a_rx, "second A did not finish").unwrap();
470        second.join().unwrap();
471        assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
472        assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
473    })
474}
475
476pub fn verify_signing_commit_exclusion(harness: &dyn OrderingHarness) -> Result<(), String> {
477    run_scenario("signing commit exclusion", || {
478        let root = TempRoot::new("signer");
479        let ordering = harness.open(&root.0).unwrap();
480        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
481        ordering
482            .register_subsystem(a, None, Recording::plain())
483            .unwrap();
484        ordering
485            .register_subsystem(b, None, Recording::plain())
486            .unwrap();
487        let gate = Arc::new(Gate::new());
488        let _release = gate.release_on_drop();
489        let first_gate = gate.clone();
490        let (first_tx, first_rx) = mpsc::channel();
491        let first_ordering = ordering.clone();
492        let first = thread::spawn(move || {
493            first_tx
494                .send(first_ordering.submit_local_txn(
495                    1,
496                    [1; 32],
497                    a,
498                    b"a",
499                    Box::new(move |_| {
500                        first_gate.enter();
501                        Ok([1; 64])
502                    }),
503                    Box::new(|_| Ok(())),
504                ))
505                .unwrap();
506        });
507        gate.wait_for(1);
508        let (submit_ready_tx, submit_ready_rx) = mpsc::channel();
509        let (submit_tx, submit_rx) = mpsc::channel();
510        let submit_ordering = ordering.clone();
511        let submit = thread::spawn(move || {
512            submit_ready_tx.send(()).unwrap();
513            submit_tx
514                .send(submit_ordering.submit_local_txn(
515                    2,
516                    [2; 32],
517                    b,
518                    b"b",
519                    Box::new(|_| Ok([2; 64])),
520                    Box::new(|_| Ok(())),
521                ))
522                .unwrap();
523        });
524        let (query_ready_tx, query_ready_rx) = mpsc::channel();
525        let (query_tx, query_rx) = mpsc::channel();
526        let query_ordering = ordering.clone();
527        let query = thread::spawn(move || {
528            query_ready_tx.send(()).unwrap();
529            query_tx.send(query_ordering.tip()).unwrap();
530        });
531        receive(&submit_ready_rx, "second writer did not start");
532        receive(&query_ready_rx, "query did not start");
533        assert_pending(&submit_rx, "second writer");
534        assert_pending(&query_rx, "tip query");
535        gate.release();
536        receive(&first_rx, "first writer did not finish").unwrap();
537        first.join().unwrap();
538        receive(&submit_rx, "second writer did not finish").unwrap();
539        assert!(receive(&query_rx, "tip query did not finish").is_some());
540        submit.join().unwrap();
541        query.join().unwrap();
542    })
543}
544
545pub fn verify_unrelated_callback_reentry(harness: &dyn OrderingHarness) -> Result<(), String> {
546    run_scenario("unrelated callback reentry", || {
547        let root = TempRoot::new("reentry");
548        let ordering = harness.open(&root.0).unwrap();
549        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
550        let handler_b = Recording::plain();
551        ordering
552            .register_subsystem(b, None, handler_b.clone())
553            .unwrap();
554        ordering
555            .register_subsystem(
556                a,
557                None,
558                Arc::new(Reentrant {
559                    ordering: ordering.clone(),
560                    target: b,
561                }),
562            )
563            .unwrap();
564        let (outer_tx, outer_rx) = mpsc::channel();
565        let outer_ordering = ordering.clone();
566        let outer = thread::spawn(move || {
567            outer_tx
568                .send(outer_ordering.submit_local_txn(
569                    1,
570                    [1; 32],
571                    a,
572                    b"outer",
573                    Box::new(|_| Ok([1; 64])),
574                    Box::new(|_| Ok(())),
575                ))
576                .unwrap();
577        });
578        receive(&outer_rx, "reentrant callback deadlocked").unwrap();
579        outer.join().unwrap();
580        assert_eq!(handler_b.payloads(), vec![b"reentered".to_vec()]);
581    })
582}