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::path::Path;
7use std::sync::{Arc, Mutex, MutexGuard};
8use std::time::{Duration, Instant};
9
10pub trait Subsystem: Send + Sync + 'static {
11    fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
12    fn reorg(&self) -> Result<(), String>;
13}
14
15pub struct K1TxnOrdering {
16    writer: Mutex<WriterState>,
17    chain: Mutex<CanonicalChain>,
18}
19
20struct WriterState {
21    subscribers: HashMap<SubsystemId, SubscriberState>,
22}
23
24struct SubscriberState {
25    handler: Arc<dyn Subsystem>,
26    latest: Option<TxId>,
27    in_commission: bool,
28}
29
30impl K1TxnOrdering {
31    pub fn open(root: &Path) -> Result<Self, String> {
32        let started = Instant::now();
33        let result = CanonicalChain::open(root);
34        let elapsed = started.elapsed();
35
36        if elapsed > Duration::from_millis(100) {
37            eprintln!(
38                "{{\"module\":\"kcode-k1-txn-ordering\",\"operation\":\"open\",\"elapsed_microseconds\":{},\"outcome\":\"{}\"}}",
39                elapsed.as_micros(),
40                if result.is_ok() { "ready" } else { "error" }
41            );
42        }
43
44        result.map(|chain| Self {
45            writer: Mutex::new(WriterState {
46                subscribers: HashMap::new(),
47            }),
48            chain: Mutex::new(chain),
49        })
50    }
51
52    pub fn register_subsystem(
53        &self,
54        subsystem: SubsystemId,
55        after: Option<TxId>,
56        handler: Arc<dyn Subsystem>,
57    ) -> Result<(), String> {
58        let mut writer = self.lock_writer();
59
60        if writer
61            .subscribers
62            .get(&subsystem)
63            .is_some_and(|subscriber| subscriber.in_commission)
64        {
65            return Err("subsystem is already registered and active".to_owned());
66        }
67
68        let mut cursor = self.lock_chain().replay_cursor(subsystem, after)?;
69        let mut subscriber = SubscriberState {
70            handler,
71            latest: after,
72            in_commission: false,
73        };
74
75        loop {
76            let next = match self.lock_chain().replay_next(&mut cursor) {
77                Ok(next) => next,
78                Err(message) => {
79                    writer.subscribers.insert(subsystem, subscriber);
80                    return Err(message);
81                }
82            };
83            let Some(replayed) = next else {
84                break;
85            };
86            let transaction = Transaction::parse(&replayed.bytes)
87                .expect("canonical chain validates replay transaction bytes");
88
89            if let Err(message) = subscriber
90                .handler
91                .submit_txn(replayed.id, transaction.payload())
92            {
93                writer.subscribers.insert(subsystem, subscriber);
94                return Err(format!("subsystem replay callback failed: {message}"));
95            }
96            subscriber.latest = Some(replayed.id);
97        }
98
99        subscriber.in_commission = true;
100        writer.subscribers.insert(subsystem, subscriber);
101        Ok(())
102    }
103
104    pub fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError> {
105        let parsed = Transaction::parse(transaction)
106            .map_err(|message| SubmitError::Other(format!("invalid transaction: {message}")))?;
107        let payload = parsed.payload();
108        let mut writer = self.lock_writer();
109
110        let (outcome, affected) = {
111            let mut chain = self.lock_chain();
112            let outcome = chain.submit_validated(transaction)?;
113            let affected = if matches!(outcome, CommitOutcome::Reorganization { .. }) {
114                writer
115                    .subscribers
116                    .iter()
117                    .filter_map(|(&subsystem, subscriber)| {
118                        (subscriber.in_commission
119                            && subscriber.latest.is_some_and(|id| !chain.contains(id)))
120                        .then_some(subsystem)
121                    })
122                    .collect()
123            } else {
124                Vec::new()
125            };
126            (outcome, affected)
127        };
128
129        match outcome {
130            CommitOutcome::Duplicate => {}
131            CommitOutcome::Extension { id, subsystem } => {
132                deliver_live(&mut writer, subsystem, id, payload);
133            }
134            CommitOutcome::Reorganization { id, subsystem } => {
135                notify_reorg(&mut writer, &affected);
136                deliver_live(&mut writer, subsystem, id, payload);
137            }
138        }
139        Ok(())
140    }
141
142    pub fn submit_local_txn<F>(
143        &self,
144        timestamp: u64,
145        creator: [u8; 32],
146        subsystem: SubsystemId,
147        payload: &[u8],
148        signer: F,
149    ) -> Result<Vec<u8>, String>
150    where
151        F: FnOnce(&[u8]) -> Result<[u8; 64], String>,
152    {
153        let mut writer = self.lock_writer();
154        let bytes = self
155            .lock_chain()
156            .submit_local(timestamp, creator, subsystem, payload, signer)?;
157        let id = TxId::for_transaction(&bytes);
158        deliver_live(&mut writer, subsystem, id, payload);
159        Ok(bytes)
160    }
161
162    pub fn contains(&self, id: TxId) -> bool {
163        self.lock_chain().contains(id)
164    }
165
166    pub fn tip(&self) -> Option<TxId> {
167        self.lock_chain().tip()
168    }
169
170    pub fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String> {
171        self.lock_chain().between_txids(older, newer)
172    }
173
174    pub fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String> {
175        self.lock_chain().get_txn(id)
176    }
177
178    fn lock_writer(&self) -> MutexGuard<'_, WriterState> {
179        self.writer.lock().expect("KTO writer mutex poisoned")
180    }
181
182    fn lock_chain(&self) -> MutexGuard<'_, CanonicalChain> {
183        self.chain.lock().expect("KTO chain mutex poisoned")
184    }
185}
186
187fn deliver_live(writer: &mut WriterState, subsystem: SubsystemId, id: TxId, payload: &[u8]) {
188    let Some(subscriber) = writer.subscribers.get_mut(&subsystem) else {
189        return;
190    };
191    if !subscriber.in_commission {
192        return;
193    }
194
195    match subscriber.handler.submit_txn(id, payload) {
196        Ok(()) => subscriber.latest = Some(id),
197        Err(_) => subscriber.in_commission = false,
198    }
199}
200
201fn notify_reorg(writer: &mut WriterState, affected: &[SubsystemId]) {
202    for subsystem in affected {
203        let subscriber = writer
204            .subscribers
205            .get_mut(subsystem)
206            .expect("affected registration still exists");
207        let _ = subscriber.handler.reorg();
208        subscriber.in_commission = false;
209    }
210}
211
212#[cfg(test)]
213mod tests {
214    use super::*;
215    use kcode_k1_transaction::build_signed_transaction;
216    use std::fs;
217    use std::panic::{AssertUnwindSafe, catch_unwind};
218    use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
219    use std::sync::{Condvar, mpsc};
220    use std::thread;
221    use std::time::Duration;
222
223    static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
224
225    struct TempRoot(std::path::PathBuf);
226
227    impl TempRoot {
228        fn new(label: &str) -> Self {
229            let number = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
230            let path = std::env::temp_dir().join(format!(
231                "kcode-k1-txn-ordering-{}-{number}-{label}",
232                std::process::id()
233            ));
234            let _ = fs::remove_dir_all(&path);
235            Self(path)
236        }
237    }
238
239    impl Drop for TempRoot {
240        fn drop(&mut self) {
241            let _ = fs::remove_dir_all(&self.0);
242        }
243    }
244
245    struct Recorder {
246        submissions: Mutex<Vec<(TxId, Vec<u8>)>>,
247        reorgs: AtomicUsize,
248        fail_submit: AtomicBool,
249        fail_reorg: AtomicBool,
250        panic_submit: AtomicBool,
251        panic_reorg: AtomicBool,
252    }
253
254    impl Recorder {
255        fn new() -> Self {
256            Self {
257                submissions: Mutex::new(Vec::new()),
258                reorgs: AtomicUsize::new(0),
259                fail_submit: AtomicBool::new(false),
260                fail_reorg: AtomicBool::new(false),
261                panic_submit: AtomicBool::new(false),
262                panic_reorg: AtomicBool::new(false),
263            }
264        }
265
266        fn submissions(&self) -> usize {
267            self.submissions.lock().unwrap().len()
268        }
269    }
270
271    impl Subsystem for Recorder {
272        fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
273            if self.panic_submit.load(Ordering::Relaxed) {
274                panic!("submit callback panic");
275            }
276            self.submissions
277                .lock()
278                .unwrap()
279                .push((id, payload.to_vec()));
280            if self.fail_submit.load(Ordering::Relaxed) {
281                Err("submit failed".to_owned())
282            } else {
283                Ok(())
284            }
285        }
286
287        fn reorg(&self) -> Result<(), String> {
288            self.reorgs.fetch_add(1, Ordering::Relaxed);
289            if self.panic_reorg.load(Ordering::Relaxed) {
290                panic!("reorg callback panic");
291            }
292            if self.fail_reorg.load(Ordering::Relaxed) {
293                Err("reorg failed".to_owned())
294            } else {
295                Ok(())
296            }
297        }
298    }
299
300    struct Blocker {
301        entered: AtomicBool,
302        state: Mutex<bool>,
303        changed: Condvar,
304    }
305
306    impl Blocker {
307        fn new() -> Self {
308            Self {
309                entered: AtomicBool::new(false),
310                state: Mutex::new(false),
311                changed: Condvar::new(),
312            }
313        }
314
315        fn wait_until_entered(&self) {
316            let mut released = self.state.lock().unwrap();
317            while !self.entered.load(Ordering::Acquire) {
318                released = self.changed.wait(released).unwrap();
319            }
320        }
321
322        fn release(&self) {
323            *self.state.lock().unwrap() = true;
324            self.changed.notify_all();
325        }
326    }
327
328    impl Subsystem for Blocker {
329        fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
330            let mut released = self.state.lock().unwrap();
331            self.entered.store(true, Ordering::Release);
332            self.changed.notify_all();
333            while !*released {
334                released = self.changed.wait(released).unwrap();
335            }
336            Ok(())
337        }
338
339        fn reorg(&self) -> Result<(), String> {
340            Ok(())
341        }
342    }
343
344    fn subsystem(value: u8) -> SubsystemId {
345        SubsystemId::from_bytes([value; 20]).unwrap()
346    }
347
348    fn transaction(
349        parent: TxId,
350        creator: u8,
351        timestamp: u64,
352        subsystem: SubsystemId,
353        payload: &[u8],
354    ) -> Vec<u8> {
355        build_signed_transaction(parent, timestamp, [creator; 32], subsystem, payload, |_| {
356            Ok([creator; 64])
357        })
358        .unwrap()
359    }
360
361    #[test]
362    fn local_submission_parents_returns_bytes_and_ignores_callback_error() {
363        let root = TempRoot::new("local");
364        let ordering = K1TxnOrdering::open(&root.0).unwrap();
365        let owner = subsystem(b'a');
366        let handler = Arc::new(Recorder::new());
367        ordering
368            .register_subsystem(owner, None, handler.clone())
369            .unwrap();
370
371        let first = ordering
372            .submit_local_txn(1, [1; 32], owner, b"first", |prefix| {
373                assert_eq!(&prefix[..12], GENESIS_PARENT.as_bytes());
374                Ok([2; 64])
375            })
376            .unwrap();
377        let first_id = TxId::for_transaction(&first);
378        assert_eq!(Transaction::parse(&first).unwrap().parent(), GENESIS_PARENT);
379        assert_eq!(handler.submissions(), 1);
380
381        handler.fail_submit.store(true, Ordering::Relaxed);
382        let second = ordering
383            .submit_local_txn(2, [1; 32], owner, b"second", |_| Ok([3; 64]))
384            .unwrap();
385        let second_id = TxId::for_transaction(&second);
386        assert_eq!(Transaction::parse(&second).unwrap().parent(), first_id);
387        assert!(ordering.contains(second_id));
388        assert_eq!(handler.submissions(), 2);
389
390        handler.fail_submit.store(false, Ordering::Relaxed);
391        ordering
392            .submit_local_txn(3, [1; 32], owner, b"third", |_| Ok([4; 64]))
393            .unwrap();
394        assert_eq!(handler.submissions(), 2);
395        ordering
396            .register_subsystem(owner, Some(first_id), handler.clone())
397            .unwrap();
398        assert_eq!(handler.submissions(), 4);
399    }
400
401    #[test]
402    fn remote_submission_reorganizes_and_callback_errors_are_non_authoritative() {
403        let root = TempRoot::new("remote");
404        let ordering = K1TxnOrdering::open(&root.0).unwrap();
405        let (a, b) = (subsystem(b'a'), subsystem(b'b'));
406        let (handler_a, handler_b) = (Arc::new(Recorder::new()), Arc::new(Recorder::new()));
407        ordering
408            .register_subsystem(a, None, handler_a.clone())
409            .unwrap();
410        ordering
411            .register_subsystem(b, None, handler_b.clone())
412            .unwrap();
413
414        let first = transaction(GENESIS_PARENT, 20, 1, a, b"first");
415        let first_id = TxId::for_transaction(&first);
416        ordering.submit_txn(&first).unwrap();
417        ordering.submit_txn(&first).unwrap();
418        assert_eq!(handler_a.submissions(), 1);
419
420        let second = transaction(first_id, 50, 2, b, b"second");
421        let second_id = TxId::for_transaction(&second);
422        ordering.submit_txn(&second).unwrap();
423        let third = transaction(second_id, 50, 3, a, b"third");
424        let third_id = TxId::for_transaction(&third);
425        ordering.submit_txn(&third).unwrap();
426
427        handler_a.fail_reorg.store(true, Ordering::Relaxed);
428        let replacement = transaction(first_id, 1, 99, a, b"replacement");
429        let replacement_id = TxId::for_transaction(&replacement);
430        ordering.submit_txn(&replacement).unwrap();
431
432        assert!(ordering.contains(replacement_id));
433        assert!(!ordering.contains(second_id));
434        assert!(!ordering.contains(third_id));
435        assert_eq!(handler_a.reorgs.load(Ordering::Relaxed), 1);
436        assert_eq!(handler_b.reorgs.load(Ordering::Relaxed), 1);
437        assert_eq!(handler_a.submissions(), 2);
438
439        ordering
440            .register_subsystem(a, Some(first_id), handler_a.clone())
441            .unwrap();
442        assert_eq!(handler_a.submissions(), 3);
443
444        handler_b.fail_submit.store(true, Ordering::Relaxed);
445        ordering
446            .register_subsystem(b, None, handler_b.clone())
447            .unwrap();
448        let extension = transaction(replacement_id, 1, 100, b, b"extension");
449        assert!(ordering.submit_txn(&extension).is_ok());
450        assert!(ordering.contains(TxId::for_transaction(&extension)));
451    }
452
453    #[test]
454    fn callback_blocks_writers_but_not_post_commit_queries() {
455        let root = TempRoot::new("blocking");
456        let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
457        let owner = subsystem(b'c');
458        let blocker = Arc::new(Blocker::new());
459        ordering
460            .register_subsystem(owner, None, blocker.clone())
461            .unwrap();
462
463        let first_ordering = ordering.clone();
464        let first = thread::spawn(move || {
465            first_ordering
466                .submit_local_txn(1, [1; 32], owner, b"first", |_| Ok([1; 64]))
467                .unwrap()
468        });
469        blocker.wait_until_entered();
470
471        let tip = ordering.tip().unwrap();
472        assert!(ordering.contains(tip));
473
474        let (sent, received) = mpsc::channel();
475        let second_ordering = ordering.clone();
476        thread::spawn(move || {
477            let result =
478                second_ordering.submit_local_txn(2, [1; 32], owner, b"second", |_| Ok([2; 64]));
479            sent.send(result).unwrap();
480        });
481        assert!(matches!(
482            received.recv_timeout(Duration::from_millis(50)),
483            Err(mpsc::RecvTimeoutError::Timeout)
484        ));
485
486        blocker.release();
487        first.join().unwrap();
488        assert!(
489            received
490                .recv_timeout(Duration::from_secs(5))
491                .unwrap()
492                .is_ok()
493        );
494    }
495
496    #[test]
497    fn registration_replay_failure_is_retryable_and_queries_delegate() {
498        let root = TempRoot::new("replay");
499        let ordering = K1TxnOrdering::open(&root.0).unwrap();
500        let owner = subsystem(b'd');
501        let first = transaction(GENESIS_PARENT, 1, 1, owner, b"first");
502        let first_id = TxId::for_transaction(&first);
503        ordering.submit_txn(&first).unwrap();
504
505        let handler = Arc::new(Recorder::new());
506        handler.fail_submit.store(true, Ordering::Relaxed);
507        assert!(
508            ordering
509                .register_subsystem(owner, None, handler.clone())
510                .is_err()
511        );
512        handler.fail_submit.store(false, Ordering::Relaxed);
513        ordering
514            .register_subsystem(owner, None, handler.clone())
515            .unwrap();
516
517        assert_eq!(handler.submissions(), 2);
518        assert_eq!(ordering.tip(), Some(first_id));
519        assert_eq!(ordering.get_txn(first_id).unwrap(), Some(first));
520        assert!(
521            ordering
522                .between_txids(GENESIS_PARENT, first_id)
523                .unwrap()
524                .is_empty()
525        );
526        assert!(ordering.register_subsystem(owner, None, handler).is_err());
527    }
528
529    #[test]
530    fn callback_panics_poison_writers_without_hiding_committed_chain() {
531        let root = TempRoot::new("submit-panic");
532        let ordering = K1TxnOrdering::open(&root.0).unwrap();
533        let owner = subsystem(b'e');
534        let handler = Arc::new(Recorder::new());
535        ordering
536            .register_subsystem(owner, None, handler.clone())
537            .unwrap();
538        handler.panic_submit.store(true, Ordering::Relaxed);
539        let panicking_txn = transaction(GENESIS_PARENT, 1, 1, owner, b"panic");
540        let id = TxId::for_transaction(&panicking_txn);
541        assert!(catch_unwind(AssertUnwindSafe(|| ordering.submit_txn(&panicking_txn))).is_err());
542        assert!(ordering.contains(id));
543        assert!(catch_unwind(AssertUnwindSafe(|| ordering.submit_txn(&panicking_txn))).is_err());
544
545        let root = TempRoot::new("reorg-panic");
546        let ordering = K1TxnOrdering::open(&root.0).unwrap();
547        let handler = Arc::new(Recorder::new());
548        ordering
549            .register_subsystem(owner, None, handler.clone())
550            .unwrap();
551        let first = transaction(GENESIS_PARENT, 20, 1, owner, b"first");
552        let first_id = TxId::for_transaction(&first);
553        ordering.submit_txn(&first).unwrap();
554        let incumbent = transaction(first_id, 50, 2, owner, b"incumbent");
555        ordering.submit_txn(&incumbent).unwrap();
556        handler.panic_reorg.store(true, Ordering::Relaxed);
557        let replacement = transaction(first_id, 1, 3, owner, b"replacement");
558        let replacement_id = TxId::for_transaction(&replacement);
559        assert!(catch_unwind(AssertUnwindSafe(|| ordering.submit_txn(&replacement))).is_err());
560        assert!(ordering.contains(replacement_id));
561        assert!(
562            catch_unwind(AssertUnwindSafe(|| {
563                ordering.register_subsystem(owner, Some(first_id), handler)
564            }))
565            .is_err()
566        );
567    }
568}