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, 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<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_queue_before_callback(harness: &dyn OrderingHarness) -> Result<(), String> {
239    run_scenario("queue before callback", || {
240        let root = TempRoot::new("queue");
241        let ordering = harness.open(&root.0).unwrap();
242        let owner = subsystem(b'a');
243        let queued = Arc::new(AtomicBool::new(false));
244        let handler = Arc::new(Recording::new(None, Some(queued.clone())));
245        ordering
246            .register_subsystem(owner, None, handler.clone())
247            .unwrap();
248        let first_queued = queued.clone();
249        let first = ordering.submit_local_txn(
250            1,
251            [1; 32],
252            owner,
253            b"first",
254            Box::new(|_| Ok([1; 64])),
255            Box::new(move |bytes| {
256                assert_eq!(Transaction::parse(bytes).unwrap().payload(), b"first");
257                first_queued.store(true, Ordering::Release);
258                Err("queue unavailable".to_owned())
259            }),
260        );
261        let first_id = ordering.tip().unwrap();
262        assert_committed(&first.unwrap_err(), first_id);
263        assert!(handler.queue_observed.load(Ordering::Acquire));
264        queued.store(false, Ordering::Release);
265        let second_queued = queued.clone();
266        let second = ordering.submit_local_txn(
267            2,
268            [1; 32],
269            owner,
270            b"second",
271            Box::new(|_| Ok([2; 64])),
272            Box::new(move |_| {
273                second_queued.store(true, Ordering::Release);
274                panic!("queue panic")
275            }),
276        );
277        let second_id = ordering.tip().unwrap();
278        assert_committed(&second.unwrap_err(), second_id);
279        assert!(handler.queue_observed.load(Ordering::Acquire));
280        ordering
281            .submit_local_txn(
282                3,
283                [1; 32],
284                owner,
285                b"third",
286                Box::new(|_| Ok([3; 64])),
287                Box::new(|_| Ok(())),
288            )
289            .unwrap();
290        assert_eq!(
291            handler.payloads(),
292            vec![b"first".to_vec(), b"second".to_vec(), b"third".to_vec()]
293        );
294    })
295}
296
297pub fn verify_independent_subsystem_lanes(harness: &dyn OrderingHarness) -> Result<(), String> {
298    run_scenario("independent subsystem lanes", || {
299        let root = TempRoot::new("lanes");
300        let ordering = harness.open(&root.0).unwrap();
301        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
302        let gate = Arc::new(Gate::new());
303        let _release = gate.release_on_drop();
304        let handler_a = Arc::new(Recording::new(Some(gate.clone()), None));
305        let handler_b = Recording::plain();
306        ordering
307            .register_subsystem(a, None, handler_a.clone())
308            .unwrap();
309        ordering
310            .register_subsystem(b, None, handler_b.clone())
311            .unwrap();
312        let (first_tx, first_rx) = mpsc::channel();
313        let first_ordering = ordering.clone();
314        let first = thread::spawn(move || {
315            first_tx
316                .send(first_ordering.submit_local_txn(
317                    1,
318                    [1; 32],
319                    a,
320                    b"a1",
321                    Box::new(|_| Ok([1; 64])),
322                    Box::new(|_| Ok(())),
323                ))
324                .unwrap();
325        });
326        gate.wait_for(1);
327        let (query_tx, query_rx) = mpsc::channel();
328        let query_ordering = ordering.clone();
329        let query = thread::spawn(move || {
330            let id = query_ordering.tip().unwrap();
331            query_tx
332                .send((
333                    id,
334                    query_ordering.contains(id),
335                    query_ordering.get_txn(id),
336                    query_ordering.between_txids(GENESIS_PARENT, id),
337                ))
338                .unwrap();
339        });
340        let (_id, contains, stored, between) = receive(&query_rx, "queries blocked by A");
341        assert!(contains);
342        assert!(!stored.unwrap().unwrap().is_empty());
343        assert!(between.unwrap().is_empty());
344        query.join().unwrap();
345        let (queued_tx, queued_rx) = mpsc::channel();
346        let (a_tx, a_rx) = mpsc::channel();
347        let second_ordering = ordering.clone();
348        let second = thread::spawn(move || {
349            a_tx.send(second_ordering.submit_local_txn(
350                2,
351                [1; 32],
352                a,
353                b"a2",
354                Box::new(|_| Ok([2; 64])),
355                Box::new(move |_| {
356                    queued_tx.send(()).unwrap();
357                    Ok(())
358                }),
359            ))
360            .unwrap();
361        });
362        receive(&queued_rx, "second A queue blocked");
363        assert_eq!(gate.count(), 1);
364        assert_pending(&a_rx, "second A caller");
365        let (b_tx, b_rx) = mpsc::channel();
366        let b_ordering = ordering.clone();
367        let b_work = thread::spawn(move || {
368            b_tx.send(b_ordering.submit_local_txn(
369                3,
370                [3; 32],
371                b,
372                b"b1",
373                Box::new(|_| Ok([3; 64])),
374                Box::new(|_| Ok(())),
375            ))
376            .unwrap();
377        });
378        receive(&b_rx, "B submission blocked by A").unwrap();
379        b_work.join().unwrap();
380        gate.release();
381        receive(&first_rx, "first A did not finish").unwrap();
382        first.join().unwrap();
383        receive(&a_rx, "second A did not finish").unwrap();
384        second.join().unwrap();
385        assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
386        assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
387    })
388}
389
390pub fn verify_signing_commit_exclusion(harness: &dyn OrderingHarness) -> Result<(), String> {
391    run_scenario("signing commit exclusion", || {
392        let root = TempRoot::new("signer");
393        let ordering = harness.open(&root.0).unwrap();
394        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
395        ordering
396            .register_subsystem(a, None, Recording::plain())
397            .unwrap();
398        ordering
399            .register_subsystem(b, None, Recording::plain())
400            .unwrap();
401        let gate = Arc::new(Gate::new());
402        let _release = gate.release_on_drop();
403        let first_gate = gate.clone();
404        let (first_tx, first_rx) = mpsc::channel();
405        let first_ordering = ordering.clone();
406        let first = thread::spawn(move || {
407            first_tx
408                .send(first_ordering.submit_local_txn(
409                    1,
410                    [1; 32],
411                    a,
412                    b"a",
413                    Box::new(move |_| {
414                        first_gate.enter();
415                        Ok([1; 64])
416                    }),
417                    Box::new(|_| Ok(())),
418                ))
419                .unwrap();
420        });
421        gate.wait_for(1);
422        let (submit_ready_tx, submit_ready_rx) = mpsc::channel();
423        let (submit_tx, submit_rx) = mpsc::channel();
424        let submit_ordering = ordering.clone();
425        let submit = thread::spawn(move || {
426            submit_ready_tx.send(()).unwrap();
427            submit_tx
428                .send(submit_ordering.submit_local_txn(
429                    2,
430                    [2; 32],
431                    b,
432                    b"b",
433                    Box::new(|_| Ok([2; 64])),
434                    Box::new(|_| Ok(())),
435                ))
436                .unwrap();
437        });
438        let (query_ready_tx, query_ready_rx) = mpsc::channel();
439        let (query_tx, query_rx) = mpsc::channel();
440        let query_ordering = ordering.clone();
441        let query = thread::spawn(move || {
442            query_ready_tx.send(()).unwrap();
443            query_tx.send(query_ordering.tip()).unwrap();
444        });
445        receive(&submit_ready_rx, "second writer did not start");
446        receive(&query_ready_rx, "query did not start");
447        assert_pending(&submit_rx, "second writer");
448        assert_pending(&query_rx, "tip query");
449        gate.release();
450        receive(&first_rx, "first writer did not finish").unwrap();
451        first.join().unwrap();
452        receive(&submit_rx, "second writer did not finish").unwrap();
453        assert!(receive(&query_rx, "tip query did not finish").is_some());
454        submit.join().unwrap();
455        query.join().unwrap();
456    })
457}
458
459pub fn verify_unrelated_callback_reentry(harness: &dyn OrderingHarness) -> Result<(), String> {
460    run_scenario("unrelated callback reentry", || {
461        let root = TempRoot::new("reentry");
462        let ordering = harness.open(&root.0).unwrap();
463        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
464        let handler_b = Recording::plain();
465        ordering
466            .register_subsystem(b, None, handler_b.clone())
467            .unwrap();
468        ordering
469            .register_subsystem(
470                a,
471                None,
472                Arc::new(Reentrant {
473                    ordering: ordering.clone(),
474                    target: b,
475                }),
476            )
477            .unwrap();
478        let (outer_tx, outer_rx) = mpsc::channel();
479        let outer_ordering = ordering.clone();
480        let outer = thread::spawn(move || {
481            outer_tx
482                .send(outer_ordering.submit_local_txn(
483                    1,
484                    [1; 32],
485                    a,
486                    b"outer",
487                    Box::new(|_| Ok([1; 64])),
488                    Box::new(|_| Ok(())),
489                ))
490                .unwrap();
491        });
492        receive(&outer_rx, "reentrant callback deadlocked").unwrap();
493        outer.join().unwrap();
494        assert_eq!(handler_b.payloads(), vec![b"reentered".to_vec()]);
495    })
496}