Skip to main content

kcode_k1_txn_ordering/
lib.rs

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(&registration.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(|(&registered, 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(&registration.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}