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}