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}