1use kcode_k1_canonical_chain::{CanonicalChain, CommitOutcome};
2pub use kcode_k1_canonical_chain::{SubmitError, TxId};
3use kcode_k1_transaction::Transaction;
4pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
5use std::collections::HashMap;
6use std::panic::{AssertUnwindSafe, catch_unwind};
7use std::path::Path;
8use std::sync::atomic::{AtomicBool, Ordering};
9use std::sync::{Arc, Condvar, Mutex, MutexGuard};
10use std::thread;
11use std::time::{Duration, Instant};
12
13pub trait Subsystem: Send + Sync + 'static {
14 fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
15 fn reorg(&self) -> Result<(), String>;
16}
17
18pub struct K1TxnOrdering {
19 writer: Mutex<WriterState>,
20 chain: Mutex<CanonicalChain>,
21}
22
23struct WriterState {
24 registrations: HashMap<SubsystemId, Registration>,
25}
26
27struct Registration {
28 lane: Arc<Lane>,
29 latest: Option<TxId>,
30 next_ticket: u64,
31 mode: RegistrationMode,
32}
33
34#[derive(Clone, Copy, Eq, PartialEq)]
35enum RegistrationMode {
36 Replaying,
37 Active,
38 OutOfCommission,
39}
40
41struct Lane {
42 handler: Arc<dyn Subsystem>,
43 progress: Mutex<LaneProgress>,
44 changed: Condvar,
45 running: AtomicBool,
46 faulted: AtomicBool,
47}
48
49struct LaneProgress {
50 serving: u64,
51 fault: Option<String>,
52}
53
54struct Reservation {
55 lane: Arc<Lane>,
56 ticket: u64,
57}
58
59enum CommittedWork {
60 Duplicate,
61 Extension {
62 id: TxId,
63 delivery: Option<Reservation>,
64 },
65 Reorganization {
66 id: TxId,
67 reorgs: Vec<Reservation>,
68 replacement: Option<Reservation>,
69 },
70}
71
72impl K1TxnOrdering {
73 pub fn open(root: &Path) -> Result<Self, String> {
74 let started = Instant::now();
75 let result = CanonicalChain::open(root);
76 let elapsed = started.elapsed();
77 if elapsed > Duration::from_millis(100) {
78 eprintln!(
79 "{{\"module\":\"kcode-k1-txn-ordering\",\"operation\":\"open\",\"elapsed_microseconds\":{},\"outcome\":\"{}\"}}",
80 elapsed.as_micros(),
81 if result.is_ok() { "ready" } else { "error" }
82 );
83 }
84 result.map(|chain| Self {
85 writer: Mutex::new(WriterState {
86 registrations: HashMap::new(),
87 }),
88 chain: Mutex::new(chain),
89 })
90 }
91
92 pub fn register_subsystem(
93 &self,
94 subsystem: SubsystemId,
95 after: Option<TxId>,
96 handler: Arc<dyn Subsystem>,
97 ) -> Result<(), String> {
98 let (mut cursor, lane) = {
99 let mut writer = self.lock_writer();
100 if let Some(registration) = writer.registrations.get(&subsystem) {
101 ensure_replaceable(registration)?;
102 }
103 let chain = self.lock_chain();
104 let cursor = chain.replay_cursor(subsystem, after)?;
105 let lane = Arc::new(Lane::new(handler));
106 writer.registrations.insert(
107 subsystem,
108 Registration {
109 lane: lane.clone(),
110 latest: after,
111 next_ticket: 0,
112 mode: RegistrationMode::Replaying,
113 },
114 );
115 (cursor, lane)
116 };
117 let mut latest = after;
118 loop {
119 let mut next = match self.lock_chain().replay_next(&mut cursor) {
120 Ok(next) => next,
121 Err(message) => return Err(fault_cursor(&lane, message)),
122 };
123 if next.is_none() {
124 let final_result = {
125 let mut writer = self.lock_writer();
126 let chain = self.lock_chain();
127 match chain.replay_next(&mut cursor) {
128 Ok(None) => {
129 let registration = writer
130 .registrations
131 .get_mut(&subsystem)
132 .expect("replaying registration exists");
133 if !Arc::ptr_eq(®istration.lane, &lane) {
134 return Err("registration changed during replay".to_owned());
135 }
136 registration.latest = latest;
137 registration.mode = RegistrationMode::Active;
138 return Ok(());
139 }
140 result => result,
141 }
142 };
143 next = match final_result {
144 Ok(next) => next,
145 Err(message) => return Err(fault_cursor(&lane, message)),
146 };
147 }
148 let replayed = next.expect("replay result contains a transaction");
149 let transaction = match Transaction::parse(&replayed.bytes) {
150 Ok(transaction) => transaction,
151 Err(message) => {
152 let failure = committed_error(
153 replayed.id,
154 vec![format!(
155 "canonical replay transaction was invalid: {message}"
156 )],
157 );
158 lane.record_fault(failure.clone());
159 return Err(failure);
160 }
161 };
162 if let Err(message) = lane.run_replay(replayed.id, transaction.payload()) {
163 return Err(committed_error(replayed.id, vec![message]));
164 }
165 latest = Some(replayed.id);
166 }
167 }
168
169 pub fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError> {
170 let parsed = Transaction::parse(transaction)
171 .map_err(|message| SubmitError::Other(format!("invalid transaction: {message}")))?;
172 let payload = parsed.payload();
173 let work = {
174 let mut writer = self.lock_writer();
175 let mut chain = self.lock_chain();
176 match chain.submit_validated(transaction)? {
177 CommitOutcome::Duplicate => CommittedWork::Duplicate,
178 CommitOutcome::Extension { id, subsystem } => CommittedWork::Extension {
179 id,
180 delivery: reserve_active_delivery(&mut writer, subsystem, id),
181 },
182 CommitOutcome::Reorganization { id, subsystem } => {
183 let affected: Vec<SubsystemId> = writer
184 .registrations
185 .iter()
186 .filter_map(|(®istered, registration)| {
187 (registration.is_active()
188 && registration
189 .latest
190 .is_some_and(|latest| !chain.contains(latest)))
191 .then_some(registered)
192 })
193 .collect();
194 let mut reorgs = Vec::with_capacity(affected.len());
195 for affected_subsystem in affected {
196 let registration = writer
197 .registrations
198 .get_mut(&affected_subsystem)
199 .expect("affected registration exists");
200 reorgs.push(registration.reserve_reorg());
201 registration.mode = RegistrationMode::OutOfCommission;
202 }
203 let replacement = reserve_active_delivery(&mut writer, subsystem, id);
204 CommittedWork::Reorganization {
205 id,
206 reorgs,
207 replacement,
208 }
209 }
210 }
211 };
212 match work {
213 CommittedWork::Duplicate => Ok(()),
214 CommittedWork::Extension { id, delivery } => {
215 let Some(delivery) = delivery else {
216 return Ok(());
217 };
218 delivery
219 .run_delivery(id, payload)
220 .map_err(|message| SubmitError::Other(committed_error(id, vec![message])))
221 }
222 CommittedWork::Reorganization {
223 id,
224 reorgs,
225 replacement,
226 } => {
227 let failures = run_reorganization(id, payload, reorgs, replacement);
228 if failures.is_empty() {
229 Ok(())
230 } else {
231 Err(SubmitError::Other(committed_error(id, failures)))
232 }
233 }
234 }
235 }
236
237 pub fn submit_local_txn<F, Q>(
238 &self,
239 timestamp: u64,
240 creator: [u8; 32],
241 subsystem: SubsystemId,
242 payload: &[u8],
243 signer: F,
244 queue_propagation: Q,
245 ) -> Result<Vec<u8>, String>
246 where
247 F: FnOnce(&[u8]) -> Result<[u8; 64], String>,
248 Q: FnOnce(&[u8]) -> Result<(), String>,
249 {
250 let (bytes, id, delivery) = {
251 let mut writer = self.lock_writer();
252 let lane = writer
253 .registrations
254 .get(&subsystem)
255 .filter(|registration| registration.is_active())
256 .map(|registration| registration.lane.clone())
257 .ok_or_else(|| "target subsystem is not registered and active".to_owned())?;
258 let mut chain = self.lock_chain();
259 let bytes =
260 chain.submit_local(timestamp, creator, subsystem, payload, move |prefix| {
261 match catch_unwind(AssertUnwindSafe(|| signer(prefix))) {
262 Ok(result) => result,
263 Err(_) => Err("signer panicked".to_owned()),
264 }
265 })?;
266 let id = TxId::for_transaction(&bytes);
267 let registration = writer
268 .registrations
269 .get_mut(&subsystem)
270 .expect("prechecked registration exists");
271 if !Arc::ptr_eq(®istration.lane, &lane) {
272 return Err(committed_error(
273 id,
274 vec!["target registration changed during local commit".to_owned()],
275 ));
276 }
277 let delivery = registration.reserve_delivery(id);
278 (bytes, id, delivery)
279 };
280 let mut failures = Vec::new();
281 match catch_unwind(AssertUnwindSafe(|| queue_propagation(&bytes))) {
282 Ok(Ok(())) => {}
283 Ok(Err(message)) => failures.push(format!("queue propagation failed: {message}")),
284 Err(_) => failures.push("queue propagation panicked".to_owned()),
285 }
286 if let Err(message) = delivery.run_delivery(id, payload) {
287 failures.push(message);
288 }
289 if failures.is_empty() {
290 Ok(bytes)
291 } else {
292 Err(committed_error(id, failures))
293 }
294 }
295
296 pub fn contains(&self, id: TxId) -> bool {
297 self.lock_chain().contains(id)
298 }
299
300 pub fn tip(&self) -> Option<TxId> {
301 self.lock_chain().tip()
302 }
303
304 pub fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String> {
305 self.lock_chain().between_txids(older, newer)
306 }
307
308 pub fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String> {
309 self.lock_chain().get_txn(id)
310 }
311
312 fn lock_writer(&self) -> MutexGuard<'_, WriterState> {
313 self.writer.lock().expect("KTO writer mutex poisoned")
314 }
315
316 fn lock_chain(&self) -> MutexGuard<'_, CanonicalChain> {
317 self.chain.lock().expect("KTO chain mutex poisoned")
318 }
319}
320
321impl Registration {
322 fn is_active(&self) -> bool {
323 self.mode == RegistrationMode::Active && !self.lane.faulted.load(Ordering::Acquire)
324 }
325
326 fn reserve_delivery(&mut self, id: TxId) -> Reservation {
327 let reservation = Reservation {
328 lane: self.lane.clone(),
329 ticket: self.next_ticket,
330 };
331 self.next_ticket += 1;
332 self.latest = Some(id);
333 reservation
334 }
335
336 fn reserve_reorg(&mut self) -> Reservation {
337 let reservation = Reservation {
338 lane: self.lane.clone(),
339 ticket: self.next_ticket,
340 };
341 self.next_ticket += 1;
342 reservation
343 }
344}
345
346impl Lane {
347 fn new(handler: Arc<dyn Subsystem>) -> Self {
348 Self {
349 handler,
350 progress: Mutex::new(LaneProgress {
351 serving: 0,
352 fault: None,
353 }),
354 changed: Condvar::new(),
355 running: AtomicBool::new(false),
356 faulted: AtomicBool::new(false),
357 }
358 }
359
360 fn run_replay(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
361 self.running.store(true, Ordering::Release);
362 let result = invoke_callback("replay submit", || self.handler.submit_txn(id, payload));
363 if let Err(message) = &result {
364 self.record_fault(message.clone());
365 }
366 self.running.store(false, Ordering::Release);
367 result
368 }
369
370 fn run_ticket<F>(&self, ticket: u64, operation: &str, callback: F) -> Result<(), String>
371 where
372 F: FnOnce() -> Result<(), String>,
373 {
374 let mut progress = self.progress.lock().expect("subsystem lane mutex poisoned");
375 while progress.serving < ticket {
376 progress = self
377 .changed
378 .wait(progress)
379 .expect("subsystem lane mutex poisoned");
380 }
381 if progress.serving > ticket {
382 return Err(format!("{operation} ticket was already completed"));
383 }
384 if let Some(reason) = progress.fault.clone() {
385 progress.serving += 1;
386 self.changed.notify_all();
387 return Err(format!(
388 "{operation} callback was not run because the registration faulted: {reason}"
389 ));
390 }
391 self.running.store(true, Ordering::Release);
392 drop(progress);
393 let result = invoke_callback(operation, callback);
394 let mut progress = self.progress.lock().expect("subsystem lane mutex poisoned");
395 if let Err(message) = &result
396 && progress.fault.is_none()
397 {
398 progress.fault = Some(message.clone());
399 self.faulted.store(true, Ordering::Release);
400 }
401 progress.serving += 1;
402 self.running.store(false, Ordering::Release);
403 self.changed.notify_all();
404 result
405 }
406
407 fn record_fault(&self, message: String) {
408 let mut progress = self.progress.lock().expect("subsystem lane mutex poisoned");
409 if progress.fault.is_none() {
410 progress.fault = Some(message);
411 self.faulted.store(true, Ordering::Release);
412 }
413 }
414
415 fn is_quiescent(&self, next_ticket: u64) -> bool {
416 let progress = self.progress.lock().expect("subsystem lane mutex poisoned");
417 progress.serving == next_ticket && !self.running.load(Ordering::Acquire)
418 }
419}
420
421impl Reservation {
422 fn run_delivery(self, id: TxId, payload: &[u8]) -> Result<(), String> {
423 self.lane.run_ticket(self.ticket, "submit", || {
424 self.lane.handler.submit_txn(id, payload)
425 })
426 }
427
428 fn run_reorg(self) -> Result<(), String> {
429 self.lane
430 .run_ticket(self.ticket, "reorg", || self.lane.handler.reorg())
431 }
432}
433
434fn invoke_callback<F>(operation: &str, callback: F) -> Result<(), String>
435where
436 F: FnOnce() -> Result<(), String>,
437{
438 match catch_unwind(AssertUnwindSafe(callback)) {
439 Ok(Ok(())) => Ok(()),
440 Ok(Err(message)) => Err(format!("{operation} callback failed: {message}")),
441 Err(_) => Err(format!("{operation} callback panicked")),
442 }
443}
444
445fn ensure_replaceable(registration: &Registration) -> Result<(), String> {
446 if registration.mode == RegistrationMode::Replaying
447 && !registration.lane.faulted.load(Ordering::Acquire)
448 {
449 return Err("subsystem registration is replaying".to_owned());
450 }
451 if registration.mode == RegistrationMode::Active
452 && !registration.lane.faulted.load(Ordering::Acquire)
453 {
454 return Err("subsystem is already registered and active".to_owned());
455 }
456 if !registration.lane.is_quiescent(registration.next_ticket) {
457 return Err("previous subsystem lane work has not quiesced".to_owned());
458 }
459 Ok(())
460}
461
462fn reserve_active_delivery(
463 writer: &mut WriterState,
464 subsystem: SubsystemId,
465 id: TxId,
466) -> Option<Reservation> {
467 writer
468 .registrations
469 .get_mut(&subsystem)
470 .filter(|registration| registration.is_active())
471 .map(|registration| registration.reserve_delivery(id))
472}
473
474fn run_reorganization(
475 id: TxId,
476 payload: &[u8],
477 reorgs: Vec<Reservation>,
478 replacement: Option<Reservation>,
479) -> Vec<String> {
480 thread::scope(|scope| {
481 let mut jobs = Vec::with_capacity(reorgs.len() + usize::from(replacement.is_some()));
482 for reservation in reorgs {
483 jobs.push(scope.spawn(move || reservation.run_reorg()));
484 }
485 if let Some(reservation) = replacement {
486 jobs.push(scope.spawn(move || reservation.run_delivery(id, payload)));
487 }
488 let mut failures = Vec::new();
489 for job in jobs {
490 match job.join() {
491 Ok(Ok(())) => {}
492 Ok(Err(message)) => failures.push(message),
493 Err(_) => failures.push("callback lane panicked".to_owned()),
494 }
495 }
496 failures
497 })
498}
499
500fn fault_cursor(lane: &Lane, message: String) -> String {
501 let failure = format!("subsystem replay cursor failed: {message}");
502 lane.record_fault(failure.clone());
503 failure
504}
505
506fn committed_error(id: TxId, failures: Vec<String>) -> String {
507 format!("TxId {id:?} was committed; {}", failures.join("; "))
508}
509
510#[cfg(test)]
511mod tests {
512 use super::*;
513 use kcode_k1_transaction::build_signed_transaction;
514 use std::fs;
515 use std::sync::atomic::{AtomicU64, AtomicUsize};
516 use std::sync::mpsc;
517 use std::time::Duration;
518
519 static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
520
521 struct TempRoot(std::path::PathBuf);
522
523 impl TempRoot {
524 fn new(label: &str) -> Self {
525 let number = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
526 let path = std::env::temp_dir().join(format!(
527 "kcode-k1-txn-ordering-{}-{number}-{label}",
528 std::process::id()
529 ));
530 let _ = fs::remove_dir_all(&path);
531 Self(path)
532 }
533 }
534
535 impl Drop for TempRoot {
536 fn drop(&mut self) {
537 let _ = fs::remove_dir_all(&self.0);
538 }
539 }
540
541 struct Gate {
542 state: Mutex<(usize, bool)>,
543 changed: Condvar,
544 }
545
546 impl Gate {
547 fn new() -> Self {
548 Self {
549 state: Mutex::new((0, false)),
550 changed: Condvar::new(),
551 }
552 }
553
554 fn enter(&self) {
555 let mut state = self.state.lock().unwrap();
556 state.0 += 1;
557 self.changed.notify_all();
558 while !state.1 {
559 state = self.changed.wait(state).unwrap();
560 }
561 }
562
563 fn wait_for(&self, count: usize) {
564 let mut state = self.state.lock().unwrap();
565 while state.0 < count {
566 state = self.changed.wait(state).unwrap();
567 }
568 }
569
570 fn count(&self) -> usize {
571 self.state.lock().unwrap().0
572 }
573
574 fn release(&self) {
575 self.state.lock().unwrap().1 = true;
576 self.changed.notify_all();
577 }
578 }
579
580 struct Recording {
581 payloads: Mutex<Vec<Vec<u8>>>,
582 reorgs: AtomicUsize,
583 submit_gate: Option<Arc<Gate>>,
584 reorg_gate: Option<Arc<Gate>>,
585 queue_flag: Option<Arc<AtomicBool>>,
586 queue_observed: AtomicBool,
587 fail_submit: AtomicBool,
588 panic_submit: AtomicBool,
589 }
590
591 impl Recording {
592 fn new(
593 submit_gate: Option<Arc<Gate>>,
594 reorg_gate: Option<Arc<Gate>>,
595 queue_flag: Option<Arc<AtomicBool>>,
596 ) -> Self {
597 Self {
598 payloads: Mutex::new(Vec::new()),
599 reorgs: AtomicUsize::new(0),
600 submit_gate,
601 reorg_gate,
602 queue_flag,
603 queue_observed: AtomicBool::new(false),
604 fail_submit: AtomicBool::new(false),
605 panic_submit: AtomicBool::new(false),
606 }
607 }
608
609 fn plain() -> Arc<Self> {
610 Arc::new(Self::new(None, None, None))
611 }
612
613 fn payloads(&self) -> Vec<Vec<u8>> {
614 self.payloads.lock().unwrap().clone()
615 }
616 }
617
618 impl Subsystem for Recording {
619 fn submit_txn(&self, _id: TxId, payload: &[u8]) -> Result<(), String> {
620 self.payloads.lock().unwrap().push(payload.to_vec());
621 if let Some(flag) = &self.queue_flag {
622 self.queue_observed
623 .store(flag.load(Ordering::Acquire), Ordering::Release);
624 }
625 if let Some(gate) = &self.submit_gate {
626 gate.enter();
627 }
628 if self.panic_submit.load(Ordering::Acquire) {
629 panic!("submit panic");
630 }
631 if self.fail_submit.load(Ordering::Acquire) {
632 Err("submit failure".to_owned())
633 } else {
634 Ok(())
635 }
636 }
637
638 fn reorg(&self) -> Result<(), String> {
639 self.reorgs.fetch_add(1, Ordering::Relaxed);
640 if let Some(gate) = &self.reorg_gate {
641 gate.enter();
642 }
643 Ok(())
644 }
645 }
646
647 struct Reentrant {
648 ordering: Arc<K1TxnOrdering>,
649 target: SubsystemId,
650 }
651
652 impl Subsystem for Reentrant {
653 fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
654 self.ordering
655 .submit_local_txn(
656 99,
657 [9; 32],
658 self.target,
659 b"reentered",
660 |_| Ok([9; 64]),
661 |_| Ok(()),
662 )
663 .map(|_| ())
664 }
665
666 fn reorg(&self) -> Result<(), String> {
667 Ok(())
668 }
669 }
670
671 fn subsystem(value: u8) -> SubsystemId {
672 SubsystemId::from_bytes([value; 20]).unwrap()
673 }
674
675 fn transaction(
676 parent: TxId,
677 creator: u8,
678 timestamp: u64,
679 subsystem: SubsystemId,
680 payload: &[u8],
681 ) -> Vec<u8> {
682 build_signed_transaction(parent, timestamp, [creator; 32], subsystem, payload, |_| {
683 Ok([creator; 64])
684 })
685 .unwrap()
686 }
687
688 fn assert_committed(message: &str, id: TxId) {
689 assert!(message.contains("TxId"));
690 assert!(message.contains("was committed"));
691 assert!(message.contains(&format!("{id:?}")));
692 }
693
694 #[test]
695 fn queue_precedes_callback_and_queue_failures_still_integrate() {
696 let root = TempRoot::new("queue");
697 let ordering = K1TxnOrdering::open(&root.0).unwrap();
698 let owner = subsystem(b'a');
699 let queued = Arc::new(AtomicBool::new(false));
700 let handler = Arc::new(Recording::new(None, None, Some(queued.clone())));
701 ordering
702 .register_subsystem(owner, None, handler.clone())
703 .unwrap();
704 let result = ordering.submit_local_txn(
705 1,
706 [1; 32],
707 owner,
708 b"first",
709 |_| Ok([1; 64]),
710 |bytes| {
711 assert_eq!(Transaction::parse(bytes).unwrap().payload(), b"first");
712 queued.store(true, Ordering::Release);
713 Err("queue unavailable".to_owned())
714 },
715 );
716 let first_id = ordering.tip().unwrap();
717 assert_committed(&result.unwrap_err(), first_id);
718 assert!(handler.queue_observed.load(Ordering::Acquire));
719 queued.store(false, Ordering::Release);
720 let result = ordering.submit_local_txn(
721 2,
722 [1; 32],
723 owner,
724 b"second",
725 |_| Ok([2; 64]),
726 |_| {
727 queued.store(true, Ordering::Release);
728 panic!("queue panic")
729 },
730 );
731 let second_id = ordering.tip().unwrap();
732 assert_committed(&result.unwrap_err(), second_id);
733 assert_eq!(
734 handler.payloads(),
735 vec![b"first".to_vec(), b"second".to_vec()]
736 );
737 assert!(
738 ordering
739 .submit_local_txn(3, [1; 32], owner, b"third", |_| Ok([3; 64]), |_| Ok(()))
740 .is_ok()
741 );
742 }
743
744 #[test]
745 fn subsystem_lanes_order_a_without_blocking_b_or_queries() {
746 let root = TempRoot::new("lanes");
747 let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
748 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
749 let gate = Arc::new(Gate::new());
750 let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
751 let handler_b = Recording::plain();
752 ordering
753 .register_subsystem(a, None, handler_a.clone())
754 .unwrap();
755 ordering
756 .register_subsystem(b, None, handler_b.clone())
757 .unwrap();
758 let first_ordering = ordering.clone();
759 let first = thread::spawn(move || {
760 first_ordering.submit_local_txn(1, [1; 32], a, b"a1", |_| Ok([1; 64]), |_| Ok(()))
761 });
762 gate.wait_for(1);
763 let first_id = ordering.tip().unwrap();
764 assert!(ordering.contains(first_id));
765 assert!(!ordering.get_txn(first_id).unwrap().unwrap().is_empty());
766 let (queued_tx, queued_rx) = mpsc::channel();
767 let (a_tx, a_rx) = mpsc::channel();
768 let second_ordering = ordering.clone();
769 thread::spawn(move || {
770 let result = second_ordering.submit_local_txn(
771 2,
772 [1; 32],
773 a,
774 b"a2",
775 |_| Ok([2; 64]),
776 |_| {
777 queued_tx.send(()).unwrap();
778 Ok(())
779 },
780 );
781 a_tx.send(result).unwrap();
782 });
783 queued_rx.recv_timeout(Duration::from_secs(2)).unwrap();
784 assert_eq!(gate.count(), 1);
785 assert!(matches!(
786 a_rx.recv_timeout(Duration::from_millis(50)),
787 Err(mpsc::RecvTimeoutError::Timeout)
788 ));
789 let (b_tx, b_rx) = mpsc::channel();
790 let b_ordering = ordering.clone();
791 thread::spawn(move || {
792 b_tx.send(b_ordering.submit_local_txn(
793 3,
794 [3; 32],
795 b,
796 b"b1",
797 |_| Ok([3; 64]),
798 |_| Ok(()),
799 ))
800 .unwrap();
801 });
802 assert!(b_rx.recv_timeout(Duration::from_secs(2)).unwrap().is_ok());
803 gate.release();
804 assert!(first.join().unwrap().is_ok());
805 assert!(a_rx.recv_timeout(Duration::from_secs(2)).unwrap().is_ok());
806 assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
807 assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
808 }
809
810 #[test]
811 fn signer_holds_writer_and_chain_until_local_commit() {
812 let root = TempRoot::new("signer");
813 let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
814 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
815 ordering
816 .register_subsystem(a, None, Recording::plain())
817 .unwrap();
818 ordering
819 .register_subsystem(b, None, Recording::plain())
820 .unwrap();
821 let signer_gate = Arc::new(Gate::new());
822 let first_ordering = ordering.clone();
823 let first_gate = signer_gate.clone();
824 let first = thread::spawn(move || {
825 first_ordering.submit_local_txn(
826 1,
827 [1; 32],
828 a,
829 b"a",
830 |_| {
831 first_gate.enter();
832 Ok([1; 64])
833 },
834 |_| Ok(()),
835 )
836 });
837 signer_gate.wait_for(1);
838 let (submit_tx, submit_rx) = mpsc::channel();
839 let submit_ordering = ordering.clone();
840 thread::spawn(move || {
841 submit_tx
842 .send(submit_ordering.submit_local_txn(
843 2,
844 [2; 32],
845 b,
846 b"b",
847 |_| Ok([2; 64]),
848 |_| Ok(()),
849 ))
850 .unwrap();
851 });
852 let (query_tx, query_rx) = mpsc::channel();
853 let query_ordering = ordering.clone();
854 thread::spawn(move || query_tx.send(query_ordering.tip()).unwrap());
855 assert!(matches!(
856 submit_rx.recv_timeout(Duration::from_millis(50)),
857 Err(mpsc::RecvTimeoutError::Timeout)
858 ));
859 assert!(matches!(
860 query_rx.recv_timeout(Duration::from_millis(50)),
861 Err(mpsc::RecvTimeoutError::Timeout)
862 ));
863 signer_gate.release();
864 assert!(first.join().unwrap().is_ok());
865 assert!(
866 submit_rx
867 .recv_timeout(Duration::from_secs(2))
868 .unwrap()
869 .is_ok()
870 );
871 assert!(
872 query_rx
873 .recv_timeout(Duration::from_secs(2))
874 .unwrap()
875 .is_some()
876 );
877 }
878
879 #[test]
880 fn replay_isolated_by_subsystem_and_catches_one_concurrent_commit() {
881 let root = TempRoot::new("replay");
882 let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
883 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
884 let first = transaction(GENESIS_PARENT, 8, 1, a, b"a1");
885 ordering.submit_txn(&first).unwrap();
886 let gate = Arc::new(Gate::new());
887 let handler_a = Arc::new(Recording::new(Some(gate.clone()), None, None));
888 let register_ordering = ordering.clone();
889 let register_handler = handler_a.clone();
890 let (register_tx, register_rx) = mpsc::channel();
891 thread::spawn(move || {
892 register_tx
893 .send(register_ordering.register_subsystem(a, None, register_handler))
894 .unwrap();
895 });
896 gate.wait_for(1);
897 let handler_b = Recording::plain();
898 ordering
899 .register_subsystem(b, None, handler_b.clone())
900 .unwrap();
901 let b_bytes = ordering
902 .submit_local_txn(2, [2; 32], b, b"b1", |_| Ok([2; 64]), |_| Ok(()))
903 .unwrap();
904 let second = transaction(TxId::for_transaction(&b_bytes), 8, 3, a, b"a2");
905 ordering.submit_txn(&second).unwrap();
906 assert_eq!(handler_b.payloads(), vec![b"b1".to_vec()]);
907 gate.release();
908 assert!(
909 register_rx
910 .recv_timeout(Duration::from_secs(2))
911 .unwrap()
912 .is_ok()
913 );
914 assert_eq!(handler_a.payloads(), vec![b"a1".to_vec(), b"a2".to_vec()]);
915 }
916
917 #[test]
918 fn callback_error_and_panic_report_committed_ids_and_isolate_faults() {
919 let root = TempRoot::new("faults");
920 let ordering = K1TxnOrdering::open(&root.0).unwrap();
921 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
922 let failing = Recording::plain();
923 failing.fail_submit.store(true, Ordering::Release);
924 let handler_b = Recording::plain();
925 ordering
926 .register_subsystem(a, None, failing.clone())
927 .unwrap();
928 ordering
929 .register_subsystem(b, None, handler_b.clone())
930 .unwrap();
931 let message = ordering
932 .submit_local_txn(1, [1; 32], a, b"failure", |_| Ok([1; 64]), |_| Ok(()))
933 .unwrap_err();
934 let failed_id = ordering.tip().unwrap();
935 assert_committed(&message, failed_id);
936 ordering
937 .submit_local_txn(2, [2; 32], b, b"b1", |_| Ok([2; 64]), |_| Ok(()))
938 .unwrap();
939 let panicking = Recording::plain();
940 panicking.panic_submit.store(true, Ordering::Release);
941 ordering
942 .register_subsystem(a, Some(failed_id), panicking.clone())
943 .unwrap();
944 let peer = transaction(ordering.tip().unwrap(), 3, 3, a, b"panic");
945 let peer_id = TxId::for_transaction(&peer);
946 let message = match ordering.submit_txn(&peer) {
947 Err(SubmitError::Other(message)) => message,
948 result => panic!("unexpected peer result: {result:?}"),
949 };
950 assert_committed(&message, peer_id);
951 assert!(ordering.contains(peer_id));
952 ordering
953 .submit_local_txn(4, [4; 32], b, b"b2", |_| Ok([4; 64]), |_| Ok(()))
954 .unwrap();
955 assert_eq!(handler_b.payloads(), vec![b"b1".to_vec(), b"b2".to_vec()]);
956 }
957
958 #[test]
959 fn reorg_fanout_releases_global_locks_and_unaffected_delivery() {
960 let root = TempRoot::new("reorg");
961 let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
962 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
963 let reorg_gate = Arc::new(Gate::new());
964 let handler_a = Arc::new(Recording::new(None, Some(reorg_gate.clone()), None));
965 let b_gate = Arc::new(Gate::new());
966 b_gate.release();
967 let handler_b = Arc::new(Recording::new(Some(b_gate.clone()), None, None));
968 ordering
969 .register_subsystem(a, None, handler_a.clone())
970 .unwrap();
971 ordering
972 .register_subsystem(b, None, handler_b.clone())
973 .unwrap();
974 let first = transaction(GENESIS_PARENT, 20, 1, a, b"first");
975 let first_id = TxId::for_transaction(&first);
976 ordering.submit_txn(&first).unwrap();
977 let incumbent = transaction(first_id, 50, 2, a, b"incumbent");
978 ordering.submit_txn(&incumbent).unwrap();
979 let replacement = transaction(first_id, 1, 3, b, b"replacement");
980 let replacement_id = TxId::for_transaction(&replacement);
981 let replacement_ordering = ordering.clone();
982 let (replacement_tx, replacement_rx) = mpsc::channel();
983 thread::spawn(move || {
984 replacement_tx
985 .send(replacement_ordering.submit_txn(&replacement))
986 .unwrap();
987 });
988 reorg_gate.wait_for(1);
989 b_gate.wait_for(1);
990 assert!(ordering.contains(replacement_id));
991 let (local_tx, local_rx) = mpsc::channel();
992 let local_ordering = ordering.clone();
993 thread::spawn(move || {
994 local_tx
995 .send(local_ordering.submit_local_txn(
996 4,
997 [4; 32],
998 b,
999 b"after",
1000 |_| Ok([4; 64]),
1001 |_| Ok(()),
1002 ))
1003 .unwrap();
1004 });
1005 assert!(
1006 local_rx
1007 .recv_timeout(Duration::from_secs(2))
1008 .unwrap()
1009 .is_ok()
1010 );
1011 reorg_gate.release();
1012 assert!(
1013 replacement_rx
1014 .recv_timeout(Duration::from_secs(2))
1015 .unwrap()
1016 .is_ok()
1017 );
1018 assert_eq!(handler_a.reorgs.load(Ordering::Acquire), 1);
1019 assert_eq!(
1020 handler_b.payloads(),
1021 vec![b"replacement".to_vec(), b"after".to_vec()]
1022 );
1023 }
1024
1025 #[test]
1026 fn callback_can_reenter_an_unrelated_subsystem() {
1027 let root = TempRoot::new("reentry");
1028 let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
1029 let (a, b) = (subsystem(b'a'), subsystem(b'b'));
1030 let handler_b = Recording::plain();
1031 ordering
1032 .register_subsystem(b, None, handler_b.clone())
1033 .unwrap();
1034 ordering
1035 .register_subsystem(
1036 a,
1037 None,
1038 Arc::new(Reentrant {
1039 ordering: ordering.clone(),
1040 target: b,
1041 }),
1042 )
1043 .unwrap();
1044 ordering
1045 .submit_local_txn(1, [1; 32], a, b"outer", |_| Ok([1; 64]), |_| Ok(()))
1046 .unwrap();
1047 assert_eq!(handler_b.payloads(), vec![b"reentered".to_vec()]);
1048 }
1049
1050 #[test]
1051 fn duplicate_restart_queries_and_root_behavior_are_preserved() {
1052 let root = TempRoot::new("restart");
1053 let owner = subsystem(b'z');
1054 let first = transaction(GENESIS_PARENT, 1, 1, owner, b"first");
1055 let first_id = TxId::for_transaction(&first);
1056 let second = transaction(first_id, 1, 2, owner, b"second");
1057 let second_id = TxId::for_transaction(&second);
1058 {
1059 let ordering = K1TxnOrdering::open(&root.0).unwrap();
1060 ordering.submit_txn(&first).unwrap();
1061 ordering.submit_txn(&first).unwrap();
1062 ordering.submit_txn(&second).unwrap();
1063 assert_eq!(ordering.tip(), Some(second_id));
1064 assert_eq!(
1065 ordering.between_txids(GENESIS_PARENT, second_id).unwrap(),
1066 vec![first_id]
1067 );
1068 assert_eq!(ordering.get_txn(first_id).unwrap(), Some(first.clone()));
1069 }
1070 let ordering = K1TxnOrdering::open(&root.0).unwrap();
1071 let handler = Recording::plain();
1072 ordering
1073 .register_subsystem(owner, None, handler.clone())
1074 .unwrap();
1075 ordering.submit_txn(&second).unwrap();
1076 assert_eq!(
1077 handler.payloads(),
1078 vec![b"first".to_vec(), b"second".to_vec()]
1079 );
1080 assert!(ordering.contains(first_id));
1081 assert_eq!(ordering.tip(), Some(second_id));
1082 }
1083}