Skip to main content

kcode_k1_txn_ordering/
lib.rs

1use kcode_k1_transaction::Transaction;
2pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
3pub use kcode_k1_transaction_store::TxId;
4use kcode_k1_transaction_store::{PutOutcome, StoreError, TransactionStore};
5use sha2::{Digest, Sha256};
6use std::cmp::Ordering;
7use std::collections::HashMap;
8use std::fmt;
9use std::fs::{self, File, OpenOptions};
10use std::io::{Read, Seek, SeekFrom, Write};
11use std::path::Path;
12use std::sync::{Arc, Mutex, MutexGuard};
13use std::time::{Duration, Instant};
14
15type Order = Vec<(TxId, SubsystemId)>;
16type Indexes = HashMap<TxId, usize>;
17
18pub trait Subsystem: Send + Sync + 'static {
19    fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
20    fn reorg(&self) -> Result<(), String>;
21}
22
23#[derive(Debug)]
24pub enum SubmitError {
25    MissingParent,
26    Other(String),
27}
28
29impl fmt::Display for SubmitError {
30    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
31        match self {
32            Self::MissingParent => formatter.write_str("missing parent"),
33            Self::Other(message) => formatter.write_str(message),
34        }
35    }
36}
37
38impl std::error::Error for SubmitError {}
39
40pub struct K1TxnOrdering {
41    state: Mutex<State>,
42    store: TransactionStore,
43}
44
45struct State {
46    file: File,
47    order: Order,
48    indexes: Indexes,
49    subscribers: HashMap<SubsystemId, SubscriberState>,
50    reopen_required: bool,
51}
52
53struct SubscriberState {
54    handler: Arc<dyn Subsystem>,
55    latest: Option<usize>,
56    in_commission: bool,
57}
58
59struct Candidate<'a> {
60    id: TxId,
61    parent: TxId,
62    creator: [u8; 32],
63    timestamp: u64,
64    subsystem: SubsystemId,
65    payload: &'a [u8],
66    bytes: &'a [u8],
67}
68
69#[derive(Debug, Eq, PartialEq)]
70enum ForkDecision {
71    Incoming,
72    Incumbent,
73    Duplicate,
74    Collision,
75}
76
77impl K1TxnOrdering {
78    pub fn open(root: &Path) -> Result<Self, String> {
79        let started = Instant::now();
80        let result = Self::open_inner(root);
81        let elapsed = started.elapsed();
82
83        if elapsed > Duration::from_millis(100) {
84            let outcome = if result.is_ok() { "ready" } else { "error" };
85            eprintln!(
86                "{{\"module\":\"kcode-k1-txn-ordering\",\"operation\":\"open\",\"elapsed_microseconds\":{},\"outcome\":\"{}\"}}",
87                elapsed.as_micros(),
88                outcome
89            );
90        }
91
92        result
93    }
94
95    pub fn register_subsystem(
96        &self,
97        subsystem: SubsystemId,
98        after: Option<TxId>,
99        handler: Arc<dyn Subsystem>,
100    ) -> Result<(), String> {
101        let mut state = self.lock_state();
102
103        if state.reopen_required {
104            return Err(reopen_required_message());
105        }
106
107        if state
108            .subscribers
109            .get(&subsystem)
110            .is_some_and(|subscriber| subscriber.in_commission)
111        {
112            return Err("subsystem is already registered and active".to_owned());
113        }
114
115        let (start, latest) = match after {
116            None => (0, None),
117            Some(id) => {
118                let index = state
119                    .indexes
120                    .get(&id)
121                    .copied()
122                    .ok_or_else(|| "registration checkpoint is not canonical".to_owned())?;
123
124                if state.order[index].1 != subsystem {
125                    return Err("registration checkpoint belongs to another subsystem".to_owned());
126                }
127
128                (index + 1, Some(index))
129            }
130        };
131
132        let mut subscriber = SubscriberState {
133            handler,
134            latest,
135            in_commission: false,
136        };
137
138        for index in start..state.order.len() {
139            let (id, entry_subsystem) = state.order[index];
140
141            if entry_subsystem != subsystem {
142                continue;
143            }
144
145            let bytes = match self.store_get(&mut state, id) {
146                Ok(Some(bytes)) => bytes,
147                Ok(None) => {
148                    state.subscribers.insert(subsystem, subscriber);
149                    return Err("canonical transaction bytes are missing".to_owned());
150                }
151                Err(message) => {
152                    state.subscribers.insert(subsystem, subscriber);
153                    return Err(message);
154                }
155            };
156
157            let transaction = match Transaction::parse(&bytes) {
158                Ok(transaction) => transaction,
159                Err(message) => {
160                    state.subscribers.insert(subsystem, subscriber);
161                    return Err(format!("canonical transaction is corrupt: {message}"));
162                }
163            };
164
165            if transaction.subsystem() != subsystem {
166                state.subscribers.insert(subsystem, subscriber);
167                return Err("canonical transaction has the wrong subsystem".to_owned());
168            }
169
170            if let Err(message) = subscriber.handler.submit_txn(id, transaction.payload()) {
171                state.subscribers.insert(subsystem, subscriber);
172                return Err(format!("subsystem replay callback failed: {message}"));
173            }
174
175            subscriber.latest = Some(index);
176        }
177
178        subscriber.in_commission = true;
179        state.subscribers.insert(subsystem, subscriber);
180        Ok(())
181    }
182
183    pub fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError> {
184        let parsed = Transaction::parse(transaction)
185            .map_err(|message| SubmitError::Other(format!("invalid transaction: {message}")))?;
186
187        let candidate = Candidate {
188            id: TxId::for_transaction(transaction),
189            parent: parsed.parent(),
190            creator: *parsed.creator(),
191            timestamp: parsed.timestamp(),
192            subsystem: parsed.subsystem(),
193            payload: parsed.payload(),
194            bytes: transaction,
195        };
196
197        if candidate.id == GENESIS_PARENT {
198            return Err(SubmitError::Other(
199                "transaction ID collides with the genesis sentinel".to_owned(),
200            ));
201        }
202
203        let mut state = self.lock_state();
204
205        if state.reopen_required {
206            return Err(SubmitError::Other(reopen_required_message()));
207        }
208
209        self.submit_candidate(&mut state, candidate)
210    }
211
212    pub fn contains(&self, id: TxId) -> bool {
213        id != GENESIS_PARENT && self.lock_state().indexes.contains_key(&id)
214    }
215
216    pub fn tip(&self) -> Option<TxId> {
217        self.lock_state().order.last().map(|entry| entry.0)
218    }
219
220    pub fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String> {
221        let state = self.lock_state();
222
223        let older_index = if older == GENESIS_PARENT {
224            -1_i128
225        } else {
226            state
227                .indexes
228                .get(&older)
229                .map(|index| *index as i128)
230                .ok_or_else(|| "older boundary is not canonical".to_owned())?
231        };
232
233        let newer_index = state
234            .indexes
235            .get(&newer)
236            .map(|index| *index as i128)
237            .ok_or_else(|| "newer boundary is not canonical".to_owned())?;
238
239        if older_index == newer_index {
240            return Ok(Vec::new());
241        }
242
243        if older_index > newer_index {
244            return Err("transaction boundaries are reversed".to_owned());
245        }
246
247        let interior = newer_index - older_index - 1;
248
249        if interior <= 128 {
250            return Ok(((older_index + 1)..newer_index)
251                .map(|index| state.order[index as usize].0)
252                .collect());
253        }
254
255        let distance = newer_index - older_index;
256
257        Ok((1_i128..=128)
258            .map(|k| {
259                let index = older_index + k * distance / 129;
260                state.order[index as usize].0
261            })
262            .collect())
263    }
264
265    pub fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String> {
266        {
267            let state = self.lock_state();
268
269            if id == GENESIS_PARENT || !state.indexes.contains_key(&id) {
270                return Ok(None);
271            }
272        }
273
274        match self.store.get(id) {
275            Ok(Some(bytes)) => Ok(Some(bytes)),
276            Ok(None) => Err("canonical transaction bytes are missing".to_owned()),
277            Err(error) => {
278                if store_requires_reopen(&error) {
279                    self.lock_state().reopen_required = true;
280                }
281                Err(store_error_message(error))
282            }
283        }
284    }
285
286    fn open_inner(root: &Path) -> Result<Self, String> {
287        prepare_root(root)?;
288
289        let ordering_path = root.join("ordering.dat");
290        let store_path = root.join("k1-transaction-store");
291        let ordering_type = path_type(&ordering_path)?;
292        let store_type = path_type(&store_path)?;
293
294        if ordering_type.is_some_and(|kind| !kind.is_file()) {
295            return Err("ordering.dat is not a regular file".to_owned());
296        }
297
298        if store_type.is_some_and(|kind| !kind.is_dir()) {
299            return Err("k1-transaction-store is not a directory".to_owned());
300        }
301
302        let (mut file, store) = match (ordering_type, store_type) {
303            (None, None) => {
304                let store = TransactionStore::create(&store_path).map_err(store_error_message)?;
305                let file = OpenOptions::new()
306                    .read(true)
307                    .write(true)
308                    .create_new(true)
309                    .open(&ordering_path)
310                    .map_err(|error| format!("cannot create ordering.dat: {error}"))?;
311                (file, store)
312            }
313            (Some(_), Some(_)) => {
314                let file = OpenOptions::new()
315                    .read(true)
316                    .write(true)
317                    .open(&ordering_path)
318                    .map_err(|error| format!("cannot open ordering.dat: {error}"))?;
319                let store = TransactionStore::open(&store_path).map_err(store_error_message)?;
320                (file, store)
321            }
322            _ => return Err("ordering root is incomplete".to_owned()),
323        };
324
325        let (order, indexes) = reconstruct_order(&mut file)?;
326
327        Ok(Self {
328            state: Mutex::new(State {
329                file,
330                order,
331                indexes,
332                subscribers: HashMap::new(),
333                reopen_required: false,
334            }),
335            store,
336        })
337    }
338
339    fn submit_candidate(
340        &self,
341        state: &mut State,
342        candidate: Candidate<'_>,
343    ) -> Result<(), SubmitError> {
344        if state.indexes.contains_key(&candidate.id) {
345            return match self.store_get(state, candidate.id) {
346                Ok(Some(bytes)) if bytes == candidate.bytes => Ok(()),
347                Ok(Some(_)) => Err(SubmitError::Other(
348                    "transaction ID collision with canonical bytes".to_owned(),
349                )),
350                Ok(None) => Err(SubmitError::Other(
351                    "canonical transaction bytes are missing".to_owned(),
352                )),
353                Err(message) => Err(SubmitError::Other(message)),
354            };
355        }
356
357        let shared_len = if candidate.parent == GENESIS_PARENT {
358            0
359        } else {
360            match state.indexes.get(&candidate.parent) {
361                Some(index) => index + 1,
362                None => return Err(SubmitError::MissingParent),
363            }
364        };
365
366        if shared_len == state.order.len() {
367            self.extend(state, candidate)
368        } else {
369            self.replace_fork(state, shared_len, candidate)
370        }
371    }
372
373    fn extend(&self, state: &mut State, candidate: Candidate<'_>) -> Result<(), SubmitError> {
374        self.persist_candidate(state, &candidate)?;
375        let index = state.order.len();
376        append_record(state, index, candidate.id, candidate.subsystem)
377            .map_err(SubmitError::Other)?;
378
379        state.indexes.insert(candidate.id, index);
380        state.order.push((candidate.id, candidate.subsystem));
381
382        if let Some(message) = deliver_live(
383            state,
384            candidate.subsystem,
385            candidate.id,
386            candidate.payload,
387            index,
388        ) {
389            return Err(SubmitError::Other(format!(
390                "transaction committed; {message}"
391            )));
392        }
393
394        Ok(())
395    }
396
397    fn replace_fork(
398        &self,
399        state: &mut State,
400        shared_len: usize,
401        candidate: Candidate<'_>,
402    ) -> Result<(), SubmitError> {
403        let (incumbent_id, incumbent_subsystem) = state.order[shared_len];
404        let incumbent_bytes = match self.store_get(state, incumbent_id) {
405            Ok(Some(bytes)) => bytes,
406            Ok(None) => {
407                return Err(SubmitError::Other(
408                    "canonical incumbent bytes are missing".to_owned(),
409                ));
410            }
411            Err(message) => return Err(SubmitError::Other(message)),
412        };
413
414        let incumbent = Transaction::parse(&incumbent_bytes).map_err(|message| {
415            SubmitError::Other(format!("canonical incumbent is corrupt: {message}"))
416        })?;
417
418        if incumbent.subsystem() != incumbent_subsystem {
419            return Err(SubmitError::Other(
420                "canonical incumbent has the wrong subsystem".to_owned(),
421            ));
422        }
423
424        if incumbent.parent() != candidate.parent {
425            return Err(SubmitError::Other(
426                "canonical incumbent has the wrong parent".to_owned(),
427            ));
428        }
429
430        match fork_decision(
431            &candidate.creator,
432            candidate.timestamp,
433            candidate.bytes,
434            incumbent.creator(),
435            incumbent.timestamp(),
436            &incumbent_bytes,
437        ) {
438            ForkDecision::Incumbent => {
439                return Err(SubmitError::Other(
440                    "fork loses canonical ordering".to_owned(),
441                ));
442            }
443            ForkDecision::Duplicate => return Ok(()),
444            ForkDecision::Collision => {
445                return Err(SubmitError::Other(
446                    "full transaction digest collision".to_owned(),
447                ));
448            }
449            ForkDecision::Incoming => {}
450        }
451
452        self.persist_candidate(state, &candidate)?;
453        truncate_records(state, shared_len).map_err(SubmitError::Other)?;
454        append_record(state, shared_len, candidate.id, candidate.subsystem)
455            .map_err(SubmitError::Other)?;
456
457        replace_memory(state, shared_len, candidate.id, candidate.subsystem);
458
459        let mut errors = notify_reorg(state, shared_len);
460
461        if let Some(message) = deliver_live(
462            state,
463            candidate.subsystem,
464            candidate.id,
465            candidate.payload,
466            shared_len,
467        ) {
468            errors.push(message);
469        }
470
471        if errors.is_empty() {
472            Ok(())
473        } else {
474            Err(SubmitError::Other(format!(
475                "transaction committed; {}",
476                errors.join("; ")
477            )))
478        }
479    }
480
481    fn persist_candidate(
482        &self,
483        state: &mut State,
484        candidate: &Candidate<'_>,
485    ) -> Result<(), SubmitError> {
486        match self.store.put(candidate.bytes) {
487            Ok(PutOutcome::Inserted(id)) | Ok(PutOutcome::Duplicate(id)) if id == candidate.id => {
488                Ok(())
489            }
490            Ok(_) => Err(SubmitError::Other(
491                "transaction store returned an unexpected ID".to_owned(),
492            )),
493            Err(error) => {
494                if store_requires_reopen(&error) {
495                    state.reopen_required = true;
496                }
497                Err(SubmitError::Other(store_error_message(error)))
498            }
499        }
500    }
501
502    fn store_get(&self, state: &mut State, id: TxId) -> Result<Option<Vec<u8>>, String> {
503        match self.store.get(id) {
504            Ok(bytes) => Ok(bytes),
505            Err(error) => {
506                if store_requires_reopen(&error) {
507                    state.reopen_required = true;
508                }
509                Err(store_error_message(error))
510            }
511        }
512    }
513
514    fn lock_state(&self) -> MutexGuard<'_, State> {
515        self.state
516            .lock()
517            .unwrap_or_else(|poisoned| poisoned.into_inner())
518    }
519}
520
521fn prepare_root(root: &Path) -> Result<(), String> {
522    match fs::symlink_metadata(root) {
523        Ok(metadata) if metadata.file_type().is_dir() => Ok(()),
524        Ok(_) => Err("ordering root is not a directory".to_owned()),
525        Err(error) if error.kind() == std::io::ErrorKind::NotFound => fs::create_dir_all(root)
526            .map_err(|error| format!("cannot create ordering root: {error}")),
527        Err(error) => Err(format!("cannot inspect ordering root: {error}")),
528    }
529}
530
531fn path_type(path: &Path) -> Result<Option<fs::FileType>, String> {
532    match fs::symlink_metadata(path) {
533        Ok(metadata) => Ok(Some(metadata.file_type())),
534        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
535        Err(error) => Err(format!("cannot inspect ordering component: {error}")),
536    }
537}
538
539fn reconstruct_order(file: &mut File) -> Result<(Order, Indexes), String> {
540    file.seek(SeekFrom::Start(0))
541        .map_err(|error| format!("cannot seek ordering.dat: {error}"))?;
542
543    let mut bytes = Vec::new();
544    file.read_to_end(&mut bytes)
545        .map_err(|error| format!("cannot read ordering.dat: {error}"))?;
546
547    if bytes.len() % 32 != 0 {
548        return Err("ordering.dat length is not a multiple of 32".to_owned());
549    }
550
551    let mut order = Vec::with_capacity(bytes.len() / 32);
552    let mut indexes = HashMap::with_capacity(bytes.len() / 32);
553
554    for chunk in bytes.chunks_exact(32) {
555        let id = TxId::from_bytes(chunk[..12].try_into().expect("fixed transaction ID range"));
556        let subsystem =
557            SubsystemId::from_bytes(chunk[12..].try_into().expect("fixed subsystem ID range"))
558                .map_err(|message| {
559                    format!("ordering.dat contains an invalid subsystem: {message}")
560                })?;
561
562        if id == GENESIS_PARENT {
563            return Err("ordering.dat contains the genesis sentinel".to_owned());
564        }
565
566        let index = order.len();
567
568        if indexes.insert(id, index).is_some() {
569            return Err("ordering.dat contains a duplicate transaction ID".to_owned());
570        }
571
572        order.push((id, subsystem));
573    }
574
575    Ok((order, indexes))
576}
577
578fn append_record(
579    state: &mut State,
580    record_index: usize,
581    id: TxId,
582    subsystem: SubsystemId,
583) -> Result<(), String> {
584    let expected_offset = record_offset(record_index)?;
585    let actual_offset = state
586        .file
587        .metadata()
588        .map_err(|error| format!("cannot inspect ordering.dat: {error}"))?
589        .len();
590
591    if actual_offset != expected_offset {
592        state.reopen_required = true;
593        return Err("ordering.dat changed outside this instance; reopen required".to_owned());
594    }
595
596    state
597        .file
598        .seek(SeekFrom::Start(expected_offset))
599        .map_err(|error| format!("cannot seek ordering.dat: {error}"))?;
600
601    let record = order_record(id, subsystem);
602
603    match state.file.write(&record) {
604        Ok(32) => {}
605        Ok(_) | Err(_) => {
606            state.reopen_required = true;
607            return Err("ordering append outcome is ambiguous; reopen required".to_owned());
608        }
609    }
610
611    if state.file.sync_data().is_err() {
612        state.reopen_required = true;
613        return Err("ordering append synchronization is ambiguous; reopen required".to_owned());
614    }
615
616    Ok(())
617}
618
619fn truncate_records(state: &mut State, records: usize) -> Result<(), String> {
620    let length = record_offset(records)?;
621
622    if state.file.set_len(length).is_err() {
623        state.reopen_required = true;
624        return Err("ordering truncation outcome is ambiguous; reopen required".to_owned());
625    }
626
627    if state.file.sync_data().is_err() {
628        state.reopen_required = true;
629        return Err("ordering truncation synchronization is ambiguous; reopen required".to_owned());
630    }
631
632    Ok(())
633}
634
635fn record_offset(records: usize) -> Result<u64, String> {
636    let records =
637        u64::try_from(records).map_err(|_| "ordering.dat offset exceeds u64".to_owned())?;
638
639    records
640        .checked_mul(32)
641        .ok_or_else(|| "ordering.dat offset exceeds u64".to_owned())
642}
643
644fn order_record(id: TxId, subsystem: SubsystemId) -> [u8; 32] {
645    let mut record = [0_u8; 32];
646    record[..12].copy_from_slice(id.as_bytes());
647    record[12..].copy_from_slice(subsystem.as_bytes());
648    record
649}
650
651fn replace_memory(state: &mut State, shared_len: usize, id: TxId, subsystem: SubsystemId) {
652    for (removed_id, _) in state.order.drain(shared_len..) {
653        state.indexes.remove(&removed_id);
654    }
655
656    state.indexes.insert(id, shared_len);
657    state.order.push((id, subsystem));
658}
659
660fn notify_reorg(state: &mut State, removed_start: usize) -> Vec<String> {
661    let mut errors = Vec::new();
662
663    for subscriber in state.subscribers.values_mut() {
664        let affected = subscriber.in_commission
665            && subscriber
666                .latest
667                .is_some_and(|latest| latest >= removed_start);
668
669        if !affected {
670            continue;
671        }
672
673        let result = subscriber.handler.reorg();
674        subscriber.in_commission = false;
675
676        if let Err(message) = result {
677            errors.push(format!("reorg callback failed: {message}"));
678        }
679    }
680
681    errors
682}
683
684fn deliver_live(
685    state: &mut State,
686    subsystem: SubsystemId,
687    id: TxId,
688    payload: &[u8],
689    index: usize,
690) -> Option<String> {
691    let subscriber = state.subscribers.get_mut(&subsystem)?;
692
693    if !subscriber.in_commission {
694        return None;
695    }
696
697    match subscriber.handler.submit_txn(id, payload) {
698        Ok(()) => {
699            subscriber.latest = Some(index);
700            None
701        }
702        Err(message) => {
703            subscriber.in_commission = false;
704            Some(format!("subsystem callback failed: {message}"))
705        }
706    }
707}
708
709fn fork_decision(
710    incoming_creator: &[u8; 32],
711    incoming_timestamp: u64,
712    incoming_bytes: &[u8],
713    incumbent_creator: &[u8; 32],
714    incumbent_timestamp: u64,
715    incumbent_bytes: &[u8],
716) -> ForkDecision {
717    match incoming_creator.cmp(incumbent_creator) {
718        Ordering::Less => return ForkDecision::Incoming,
719        Ordering::Greater => return ForkDecision::Incumbent,
720        Ordering::Equal => {}
721    }
722
723    match incoming_timestamp.cmp(&incumbent_timestamp) {
724        Ordering::Less => return ForkDecision::Incoming,
725        Ordering::Greater => return ForkDecision::Incumbent,
726        Ordering::Equal => {}
727    }
728
729    let incoming_digest: [u8; 32] = Sha256::digest(incoming_bytes).into();
730    let incumbent_digest: [u8; 32] = Sha256::digest(incumbent_bytes).into();
731
732    digest_decision(
733        incoming_digest,
734        incumbent_digest,
735        incoming_bytes == incumbent_bytes,
736    )
737}
738
739fn digest_decision(incoming: [u8; 32], incumbent: [u8; 32], equal_bytes: bool) -> ForkDecision {
740    match incoming.cmp(&incumbent) {
741        Ordering::Less => ForkDecision::Incoming,
742        Ordering::Greater => ForkDecision::Incumbent,
743        Ordering::Equal if equal_bytes => ForkDecision::Duplicate,
744        Ordering::Equal => ForkDecision::Collision,
745    }
746}
747
748fn store_requires_reopen(error: &StoreError) -> bool {
749    matches!(
750        error,
751        StoreError::OutcomeUnknown(_) | StoreError::ReopenRequired
752    )
753}
754
755fn store_error_message(error: StoreError) -> String {
756    format!("transaction store error: {error}")
757}
758
759fn reopen_required_message() -> String {
760    "instance requires reopening before further mutation".to_owned()
761}
762
763#[cfg(test)]
764mod tests {
765    use super::*;
766    use std::path::PathBuf;
767    use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering as AtomicOrdering};
768
769    static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
770
771    struct TempRoot {
772        path: PathBuf,
773    }
774
775    impl TempRoot {
776        fn new(label: &str) -> Self {
777            let sequence = NEXT_ROOT.fetch_add(1, AtomicOrdering::Relaxed);
778            let path = std::env::temp_dir().join(format!(
779                "kcode-k1-txn-ordering-{}-{}-{}",
780                std::process::id(),
781                sequence,
782                label
783            ));
784            let _ = fs::remove_dir_all(&path);
785            Self { path }
786        }
787    }
788
789    impl Drop for TempRoot {
790        fn drop(&mut self) {
791            let _ = fs::remove_dir_all(&self.path);
792        }
793    }
794
795    struct RecordingSubsystem {
796        submissions: Mutex<Vec<(TxId, Vec<u8>)>>,
797        reorgs: AtomicUsize,
798        fail_submit: AtomicBool,
799        fail_reorg: AtomicBool,
800    }
801
802    impl RecordingSubsystem {
803        fn new() -> Self {
804            Self {
805                submissions: Mutex::new(Vec::new()),
806                reorgs: AtomicUsize::new(0),
807                fail_submit: AtomicBool::new(false),
808                fail_reorg: AtomicBool::new(false),
809            }
810        }
811
812        fn submission_count(&self) -> usize {
813            self.submissions.lock().unwrap().len()
814        }
815    }
816
817    impl Subsystem for RecordingSubsystem {
818        fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
819            self.submissions
820                .lock()
821                .unwrap()
822                .push((id, payload.to_vec()));
823
824            if self.fail_submit.load(AtomicOrdering::Relaxed) {
825                Err("submit failure".to_owned())
826            } else {
827                Ok(())
828            }
829        }
830
831        fn reorg(&self) -> Result<(), String> {
832            self.reorgs.fetch_add(1, AtomicOrdering::Relaxed);
833
834            if self.fail_reorg.load(AtomicOrdering::Relaxed) {
835                Err("reorg failure".to_owned())
836            } else {
837                Ok(())
838            }
839        }
840    }
841
842    fn id(value: u64) -> TxId {
843        let mut bytes = [0_u8; 12];
844        bytes[..4].copy_from_slice(b"KTO!");
845        bytes[4..].copy_from_slice(&value.to_be_bytes());
846        TxId::from_bytes(bytes)
847    }
848
849    fn subsystem(value: u8) -> SubsystemId {
850        SubsystemId::from_bytes([value; 20]).unwrap()
851    }
852
853    fn transaction(
854        parent: TxId,
855        creator: [u8; 32],
856        timestamp: u64,
857        subsystem: SubsystemId,
858        payload: &[u8],
859        signature: u8,
860    ) -> Vec<u8> {
861        let mut bytes = Vec::with_capacity(136 + payload.len());
862        bytes.extend_from_slice(parent.as_bytes());
863        bytes.extend_from_slice(&timestamp.to_le_bytes());
864        bytes.extend_from_slice(&creator);
865        bytes.extend_from_slice(subsystem.as_bytes());
866        bytes.extend_from_slice(payload);
867        bytes.extend_from_slice(&[signature; 64]);
868        bytes
869    }
870
871    fn write_order(root: &Path, entries: &[(TxId, SubsystemId)]) {
872        let mut file = OpenOptions::new()
873            .write(true)
874            .truncate(true)
875            .open(root.join("ordering.dat"))
876            .unwrap();
877
878        for (id, subsystem) in entries {
879            file.write_all(&order_record(*id, *subsystem)).unwrap();
880        }
881
882        file.sync_data().unwrap();
883    }
884
885    #[test]
886    fn creates_reopens_and_validates_ordering_roots() {
887        let root = TempRoot::new("root");
888        let ordering = K1TxnOrdering::open(&root.path).unwrap();
889
890        assert!(root.path.join("ordering.dat").is_file());
891        assert!(root.path.join("k1-transaction-store").is_dir());
892
893        drop(ordering);
894        assert!(K1TxnOrdering::open(&root.path).is_ok());
895
896        let mixed = TempRoot::new("mixed");
897        fs::create_dir_all(&mixed.path).unwrap();
898        File::create(mixed.path.join("ordering.dat")).unwrap();
899        assert!(K1TxnOrdering::open(&mixed.path).is_err());
900
901        let mut malformed = OpenOptions::new()
902            .write(true)
903            .truncate(true)
904            .open(root.path.join("ordering.dat"))
905            .unwrap();
906        malformed.write_all(&[0_u8; 31]).unwrap();
907        malformed.sync_data().unwrap();
908        drop(malformed);
909        assert!(K1TxnOrdering::open(&root.path).is_err());
910
911        write_order(
912            &root.path,
913            &[(id(1), subsystem(b'a')), (id(1), subsystem(b'b'))],
914        );
915        assert!(K1TxnOrdering::open(&root.path).is_err());
916
917        let mut invalid_subsystem = [0_u8; 32];
918        invalid_subsystem[..12].copy_from_slice(id(2).as_bytes());
919        invalid_subsystem[12..].fill(0xff);
920        fs::write(root.path.join("ordering.dat"), invalid_subsystem).unwrap();
921        assert!(K1TxnOrdering::open(&root.path).is_err());
922    }
923
924    #[test]
925    fn provides_canonical_queries_and_even_sampling() {
926        let root = TempRoot::new("queries");
927        drop(K1TxnOrdering::open(&root.path).unwrap());
928
929        let entries: Vec<_> = (0_u64..260)
930            .map(|value| (id(value), subsystem(b'q')))
931            .collect();
932        write_order(&root.path, &entries);
933
934        let ordering = K1TxnOrdering::open(&root.path).unwrap();
935
936        assert!(ordering.contains(id(0)));
937        assert!(!ordering.contains(GENESIS_PARENT));
938        assert_eq!(ordering.tip(), Some(id(259)));
939        assert_eq!(ordering.get_txn(id(9999)).unwrap(), None);
940        assert!(ordering.get_txn(id(0)).is_err());
941
942        let short = ordering.between_txids(id(10), id(20)).unwrap();
943        assert_eq!(short, (11_u64..20).map(id).collect::<Vec<_>>());
944        assert!(ordering.between_txids(id(10), id(10)).unwrap().is_empty());
945        assert!(ordering.between_txids(id(20), id(10)).is_err());
946        assert!(ordering.between_txids(id(9999), id(10)).is_err());
947
948        let sampled = ordering.between_txids(GENESIS_PARENT, id(200)).unwrap();
949        assert_eq!(sampled.len(), 128);
950
951        for (offset, actual) in sampled.iter().enumerate() {
952            let k = offset as i128 + 1;
953            let expected = -1_i128 + k * 201 / 129;
954            assert_eq!(*actual, id(expected as u64));
955        }
956    }
957
958    #[test]
959    fn submits_replays_reorganizes_and_retains_removed_bytes() {
960        let root = TempRoot::new("workflow");
961        let ordering = K1TxnOrdering::open(&root.path).unwrap();
962        let subsystem_a = subsystem(b'a');
963        let subsystem_b = subsystem(b'b');
964        let handler_a = Arc::new(RecordingSubsystem::new());
965        let handler_b = Arc::new(RecordingSubsystem::new());
966
967        ordering
968            .register_subsystem(subsystem_a, None, handler_a.clone())
969            .unwrap();
970
971        let first = transaction(GENESIS_PARENT, [20; 32], 10, subsystem_a, b"first", 1);
972        let first_id = TxId::for_transaction(&first);
973        ordering.submit_txn(&first).unwrap();
974        ordering.submit_txn(&first).unwrap();
975        assert_eq!(handler_a.submission_count(), 1);
976
977        let missing = transaction(id(9999), [1; 32], 1, subsystem_a, b"missing", 2);
978        assert!(matches!(
979            ordering.submit_txn(&missing),
980            Err(SubmitError::MissingParent)
981        ));
982        assert!(!ordering.contains(TxId::for_transaction(&missing)));
983
984        let second = transaction(first_id, [50; 32], 20, subsystem_b, b"second", 3);
985        let second_id = TxId::for_transaction(&second);
986        ordering.submit_txn(&second).unwrap();
987
988        ordering
989            .register_subsystem(subsystem_b, None, handler_b.clone())
990            .unwrap();
991        assert_eq!(handler_b.submission_count(), 1);
992
993        let third = transaction(second_id, [50; 32], 30, subsystem_a, b"third", 4);
994        let third_id = TxId::for_transaction(&third);
995        ordering.submit_txn(&third).unwrap();
996        assert_eq!(handler_a.submission_count(), 2);
997
998        let replacement = transaction(first_id, [1; 32], 100, subsystem_a, b"replacement", 5);
999        let replacement_id = TxId::for_transaction(&replacement);
1000        ordering.submit_txn(&replacement).unwrap();
1001
1002        assert!(ordering.contains(first_id));
1003        assert!(ordering.contains(replacement_id));
1004        assert!(!ordering.contains(second_id));
1005        assert!(!ordering.contains(third_id));
1006        assert_eq!(handler_a.reorgs.load(AtomicOrdering::Relaxed), 1);
1007        assert_eq!(handler_b.reorgs.load(AtomicOrdering::Relaxed), 1);
1008        assert_eq!(handler_a.submission_count(), 2);
1009
1010        ordering
1011            .register_subsystem(subsystem_a, Some(first_id), handler_a.clone())
1012            .unwrap();
1013        assert_eq!(handler_a.submission_count(), 3);
1014
1015        ordering
1016            .register_subsystem(subsystem_b, None, handler_b.clone())
1017            .unwrap();
1018
1019        let extension = transaction(replacement_id, [2; 32], 200, subsystem_b, b"extension", 6);
1020        let extension_id = TxId::for_transaction(&extension);
1021        ordering.submit_txn(&extension).unwrap();
1022        assert!(ordering.contains(extension_id));
1023        assert_eq!(handler_b.submission_count(), 2);
1024
1025        let loser = transaction(first_id, [250; 32], 1, subsystem_b, b"loser", 7);
1026        let loser_id = TxId::for_transaction(&loser);
1027        assert!(matches!(
1028            ordering.submit_txn(&loser),
1029            Err(SubmitError::Other(_))
1030        ));
1031        assert!(!ordering.contains(loser_id));
1032
1033        assert_eq!(
1034            ordering.get_txn(replacement_id).unwrap().unwrap(),
1035            replacement
1036        );
1037        assert_eq!(ordering.get_txn(second_id).unwrap(), None);
1038
1039        drop(ordering);
1040
1041        let store = TransactionStore::open(&root.path.join("k1-transaction-store")).unwrap();
1042        assert!(store.contains(second_id));
1043        assert!(store.contains(third_id));
1044        assert!(!store.contains(loser_id));
1045    }
1046
1047    #[test]
1048    fn callback_failure_blocks_until_reregistration() {
1049        let root = TempRoot::new("callback");
1050        let ordering = K1TxnOrdering::open(&root.path).unwrap();
1051        let subsystem_c = subsystem(b'c');
1052        let handler = Arc::new(RecordingSubsystem::new());
1053
1054        handler.fail_submit.store(true, AtomicOrdering::Relaxed);
1055        ordering
1056            .register_subsystem(subsystem_c, None, handler.clone())
1057            .unwrap();
1058
1059        let first = transaction(GENESIS_PARENT, [1; 32], 1, subsystem_c, b"first", 1);
1060        let first_id = TxId::for_transaction(&first);
1061
1062        assert!(matches!(
1063            ordering.submit_txn(&first),
1064            Err(SubmitError::Other(_))
1065        ));
1066        assert!(ordering.contains(first_id));
1067
1068        let second = transaction(first_id, [1; 32], 2, subsystem_c, b"second", 2);
1069        ordering.submit_txn(&second).unwrap();
1070        assert_eq!(handler.submission_count(), 1);
1071
1072        handler.fail_submit.store(false, AtomicOrdering::Relaxed);
1073        ordering
1074            .register_subsystem(subsystem_c, None, handler.clone())
1075            .unwrap();
1076
1077        assert_eq!(handler.submission_count(), 3);
1078        assert!(
1079            ordering
1080                .register_subsystem(subsystem_c, None, handler.clone())
1081                .is_err()
1082        );
1083
1084        let wrong = subsystem(b'd');
1085        assert!(
1086            ordering
1087                .register_subsystem(wrong, Some(first_id), Arc::new(RecordingSubsystem::new()))
1088                .is_err()
1089        );
1090    }
1091
1092    #[test]
1093    fn fork_ranking_uses_creator_timestamp_then_digest() {
1094        let low = [1_u8; 32];
1095        let high = [2_u8; 32];
1096
1097        assert_eq!(
1098            fork_decision(&low, 10, b"x", &high, 1, b"y"),
1099            ForkDecision::Incoming
1100        );
1101        assert_eq!(
1102            fork_decision(&high, 1, b"x", &low, 10, b"y"),
1103            ForkDecision::Incumbent
1104        );
1105        assert_eq!(
1106            fork_decision(&low, 1, b"x", &low, 2, b"y"),
1107            ForkDecision::Incoming
1108        );
1109        assert_eq!(
1110            fork_decision(&low, 2, b"x", &low, 1, b"y"),
1111            ForkDecision::Incumbent
1112        );
1113        assert_eq!(
1114            fork_decision(&low, 1, b"same", &low, 1, b"same"),
1115            ForkDecision::Duplicate
1116        );
1117
1118        let left: [u8; 32] = Sha256::digest(b"left").into();
1119        let right: [u8; 32] = Sha256::digest(b"right").into();
1120        let expected = if left < right {
1121            ForkDecision::Incoming
1122        } else {
1123            ForkDecision::Incumbent
1124        };
1125
1126        assert_eq!(fork_decision(&low, 1, b"left", &low, 1, b"right"), expected);
1127        assert_eq!(
1128            digest_decision([4; 32], [4; 32], false),
1129            ForkDecision::Collision
1130        );
1131    }
1132
1133    #[test]
1134    fn ambiguous_order_change_leaves_an_orphan_and_requires_reopen() {
1135        let root = TempRoot::new("ambiguous");
1136        let ordering = K1TxnOrdering::open(&root.path).unwrap();
1137        let bytes = transaction(GENESIS_PARENT, [1; 32], 1, subsystem(b'e'), b"orphan", 1);
1138        let transaction_id = TxId::for_transaction(&bytes);
1139
1140        let mut external = OpenOptions::new()
1141            .append(true)
1142            .open(root.path.join("ordering.dat"))
1143            .unwrap();
1144        external.write_all(&[0]).unwrap();
1145        external.sync_data().unwrap();
1146
1147        assert!(matches!(
1148            ordering.submit_txn(&bytes),
1149            Err(SubmitError::Other(_))
1150        ));
1151        assert!(!ordering.contains(transaction_id));
1152        assert!(
1153            ordering
1154                .register_subsystem(subsystem(b'e'), None, Arc::new(RecordingSubsystem::new()))
1155                .is_err()
1156        );
1157
1158        drop(ordering);
1159
1160        let store = TransactionStore::open(&root.path.join("k1-transaction-store")).unwrap();
1161        assert!(store.contains(transaction_id));
1162    }
1163
1164    #[test]
1165    fn million_entry_open_and_scan_fixture() {
1166        let root = TempRoot::new("million");
1167        drop(K1TxnOrdering::open(&root.path).unwrap());
1168
1169        let count = 1_000_000_u64;
1170        let entry_subsystem = subsystem(b'm');
1171        let mut bytes = Vec::with_capacity(count as usize * 32);
1172
1173        for value in 0..count {
1174            bytes.extend_from_slice(id(value).as_bytes());
1175            bytes.extend_from_slice(entry_subsystem.as_bytes());
1176        }
1177
1178        let mut file = OpenOptions::new()
1179            .write(true)
1180            .truncate(true)
1181            .open(root.path.join("ordering.dat"))
1182            .unwrap();
1183        file.write_all(&bytes).unwrap();
1184        file.sync_data().unwrap();
1185        drop(file);
1186
1187        let started = Instant::now();
1188        let ordering = K1TxnOrdering::open(&root.path).unwrap();
1189        assert!(started.elapsed() < Duration::from_secs(5));
1190
1191        let scan_started = Instant::now();
1192        let matching = ordering
1193            .lock_state()
1194            .order
1195            .iter()
1196            .filter(|entry| entry.1 == entry_subsystem)
1197            .count();
1198        assert_eq!(matching, count as usize);
1199        assert!(scan_started.elapsed() < Duration::from_secs(1));
1200
1201        assert!(ordering.contains(id(0)));
1202        assert!(ordering.contains(id(count - 1)));
1203        assert_eq!(
1204            ordering
1205                .between_txids(GENESIS_PARENT, id(count - 1))
1206                .unwrap()
1207                .len(),
1208            128
1209        );
1210    }
1211}