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}