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(×tamp.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}