use kcode_k1_canonical_chain::{CanonicalChain, CommitOutcome};
pub use kcode_k1_canonical_chain::{SubmitError, TxId};
use kcode_k1_transaction::Transaction;
pub use kcode_k1_transaction::{GENESIS_PARENT, SubsystemId};
use std::collections::HashMap;
use std::path::Path;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::{Duration, Instant};
pub trait Subsystem: Send + Sync + 'static {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String>;
fn reorg(&self) -> Result<(), String>;
}
pub struct K1TxnOrdering {
writer: Mutex<WriterState>,
chain: Mutex<CanonicalChain>,
}
struct WriterState {
subscribers: HashMap<SubsystemId, SubscriberState>,
}
struct SubscriberState {
handler: Arc<dyn Subsystem>,
latest: Option<TxId>,
in_commission: bool,
}
impl K1TxnOrdering {
pub fn open(root: &Path) -> Result<Self, String> {
let started = Instant::now();
let result = CanonicalChain::open(root);
let elapsed = started.elapsed();
if elapsed > Duration::from_millis(100) {
eprintln!(
"{{\"module\":\"kcode-k1-txn-ordering\",\"operation\":\"open\",\"elapsed_microseconds\":{},\"outcome\":\"{}\"}}",
elapsed.as_micros(),
if result.is_ok() { "ready" } else { "error" }
);
}
result.map(|chain| Self {
writer: Mutex::new(WriterState {
subscribers: HashMap::new(),
}),
chain: Mutex::new(chain),
})
}
pub fn register_subsystem(
&self,
subsystem: SubsystemId,
after: Option<TxId>,
handler: Arc<dyn Subsystem>,
) -> Result<(), String> {
let mut writer = self.lock_writer();
if writer
.subscribers
.get(&subsystem)
.is_some_and(|subscriber| subscriber.in_commission)
{
return Err("subsystem is already registered and active".to_owned());
}
let mut cursor = self.lock_chain().replay_cursor(subsystem, after)?;
let mut subscriber = SubscriberState {
handler,
latest: after,
in_commission: false,
};
loop {
let next = match self.lock_chain().replay_next(&mut cursor) {
Ok(next) => next,
Err(message) => {
writer.subscribers.insert(subsystem, subscriber);
return Err(message);
}
};
let Some(replayed) = next else {
break;
};
let transaction = Transaction::parse(&replayed.bytes)
.expect("canonical chain validates replay transaction bytes");
if let Err(message) = subscriber
.handler
.submit_txn(replayed.id, transaction.payload())
{
writer.subscribers.insert(subsystem, subscriber);
return Err(format!("subsystem replay callback failed: {message}"));
}
subscriber.latest = Some(replayed.id);
}
subscriber.in_commission = true;
writer.subscribers.insert(subsystem, subscriber);
Ok(())
}
pub fn submit_txn(&self, transaction: &[u8]) -> Result<(), SubmitError> {
let parsed = Transaction::parse(transaction)
.map_err(|message| SubmitError::Other(format!("invalid transaction: {message}")))?;
let payload = parsed.payload();
let mut writer = self.lock_writer();
let (outcome, affected) = {
let mut chain = self.lock_chain();
let outcome = chain.submit_validated(transaction)?;
let affected = if matches!(outcome, CommitOutcome::Reorganization { .. }) {
writer
.subscribers
.iter()
.filter_map(|(&subsystem, subscriber)| {
(subscriber.in_commission
&& subscriber.latest.is_some_and(|id| !chain.contains(id)))
.then_some(subsystem)
})
.collect()
} else {
Vec::new()
};
(outcome, affected)
};
match outcome {
CommitOutcome::Duplicate => {}
CommitOutcome::Extension { id, subsystem } => {
deliver_live(&mut writer, subsystem, id, payload);
}
CommitOutcome::Reorganization { id, subsystem } => {
notify_reorg(&mut writer, &affected);
deliver_live(&mut writer, subsystem, id, payload);
}
}
Ok(())
}
pub fn submit_local_txn<F>(
&self,
timestamp: u64,
creator: [u8; 32],
subsystem: SubsystemId,
payload: &[u8],
signer: F,
) -> Result<Vec<u8>, String>
where
F: FnOnce(&[u8]) -> Result<[u8; 64], String>,
{
let mut writer = self.lock_writer();
let bytes = self
.lock_chain()
.submit_local(timestamp, creator, subsystem, payload, signer)?;
let id = TxId::for_transaction(&bytes);
deliver_live(&mut writer, subsystem, id, payload);
Ok(bytes)
}
pub fn contains(&self, id: TxId) -> bool {
self.lock_chain().contains(id)
}
pub fn tip(&self) -> Option<TxId> {
self.lock_chain().tip()
}
pub fn between_txids(&self, older: TxId, newer: TxId) -> Result<Vec<TxId>, String> {
self.lock_chain().between_txids(older, newer)
}
pub fn get_txn(&self, id: TxId) -> Result<Option<Vec<u8>>, String> {
self.lock_chain().get_txn(id)
}
fn lock_writer(&self) -> MutexGuard<'_, WriterState> {
self.writer.lock().expect("KTO writer mutex poisoned")
}
fn lock_chain(&self) -> MutexGuard<'_, CanonicalChain> {
self.chain.lock().expect("KTO chain mutex poisoned")
}
}
fn deliver_live(writer: &mut WriterState, subsystem: SubsystemId, id: TxId, payload: &[u8]) {
let Some(subscriber) = writer.subscribers.get_mut(&subsystem) else {
return;
};
if !subscriber.in_commission {
return;
}
match subscriber.handler.submit_txn(id, payload) {
Ok(()) => subscriber.latest = Some(id),
Err(_) => subscriber.in_commission = false,
}
}
fn notify_reorg(writer: &mut WriterState, affected: &[SubsystemId]) {
for subsystem in affected {
let subscriber = writer
.subscribers
.get_mut(subsystem)
.expect("affected registration still exists");
let _ = subscriber.handler.reorg();
subscriber.in_commission = false;
}
}
#[cfg(test)]
mod tests {
use super::*;
use kcode_k1_transaction::build_signed_transaction;
use std::fs;
use std::panic::{AssertUnwindSafe, catch_unwind};
use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering};
use std::sync::{Condvar, mpsc};
use std::thread;
use std::time::Duration;
static NEXT_ROOT: AtomicU64 = AtomicU64::new(0);
struct TempRoot(std::path::PathBuf);
impl TempRoot {
fn new(label: &str) -> Self {
let number = NEXT_ROOT.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"kcode-k1-txn-ordering-{}-{number}-{label}",
std::process::id()
));
let _ = fs::remove_dir_all(&path);
Self(path)
}
}
impl Drop for TempRoot {
fn drop(&mut self) {
let _ = fs::remove_dir_all(&self.0);
}
}
struct Recorder {
submissions: Mutex<Vec<(TxId, Vec<u8>)>>,
reorgs: AtomicUsize,
fail_submit: AtomicBool,
fail_reorg: AtomicBool,
panic_submit: AtomicBool,
panic_reorg: AtomicBool,
}
impl Recorder {
fn new() -> Self {
Self {
submissions: Mutex::new(Vec::new()),
reorgs: AtomicUsize::new(0),
fail_submit: AtomicBool::new(false),
fail_reorg: AtomicBool::new(false),
panic_submit: AtomicBool::new(false),
panic_reorg: AtomicBool::new(false),
}
}
fn submissions(&self) -> usize {
self.submissions.lock().unwrap().len()
}
}
impl Subsystem for Recorder {
fn submit_txn(&self, id: TxId, payload: &[u8]) -> Result<(), String> {
if self.panic_submit.load(Ordering::Relaxed) {
panic!("submit callback panic");
}
self.submissions
.lock()
.unwrap()
.push((id, payload.to_vec()));
if self.fail_submit.load(Ordering::Relaxed) {
Err("submit failed".to_owned())
} else {
Ok(())
}
}
fn reorg(&self) -> Result<(), String> {
self.reorgs.fetch_add(1, Ordering::Relaxed);
if self.panic_reorg.load(Ordering::Relaxed) {
panic!("reorg callback panic");
}
if self.fail_reorg.load(Ordering::Relaxed) {
Err("reorg failed".to_owned())
} else {
Ok(())
}
}
}
struct Blocker {
entered: AtomicBool,
state: Mutex<bool>,
changed: Condvar,
}
impl Blocker {
fn new() -> Self {
Self {
entered: AtomicBool::new(false),
state: Mutex::new(false),
changed: Condvar::new(),
}
}
fn wait_until_entered(&self) {
let mut released = self.state.lock().unwrap();
while !self.entered.load(Ordering::Acquire) {
released = self.changed.wait(released).unwrap();
}
}
fn release(&self) {
*self.state.lock().unwrap() = true;
self.changed.notify_all();
}
}
impl Subsystem for Blocker {
fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
let mut released = self.state.lock().unwrap();
self.entered.store(true, Ordering::Release);
self.changed.notify_all();
while !*released {
released = self.changed.wait(released).unwrap();
}
Ok(())
}
fn reorg(&self) -> Result<(), String> {
Ok(())
}
}
fn subsystem(value: u8) -> SubsystemId {
SubsystemId::from_bytes([value; 20]).unwrap()
}
fn transaction(
parent: TxId,
creator: u8,
timestamp: u64,
subsystem: SubsystemId,
payload: &[u8],
) -> Vec<u8> {
build_signed_transaction(parent, timestamp, [creator; 32], subsystem, payload, |_| {
Ok([creator; 64])
})
.unwrap()
}
#[test]
fn local_submission_parents_returns_bytes_and_ignores_callback_error() {
let root = TempRoot::new("local");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let owner = subsystem(b'a');
let handler = Arc::new(Recorder::new());
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
let first = ordering
.submit_local_txn(1, [1; 32], owner, b"first", |prefix| {
assert_eq!(&prefix[..12], GENESIS_PARENT.as_bytes());
Ok([2; 64])
})
.unwrap();
let first_id = TxId::for_transaction(&first);
assert_eq!(Transaction::parse(&first).unwrap().parent(), GENESIS_PARENT);
assert_eq!(handler.submissions(), 1);
handler.fail_submit.store(true, Ordering::Relaxed);
let second = ordering
.submit_local_txn(2, [1; 32], owner, b"second", |_| Ok([3; 64]))
.unwrap();
let second_id = TxId::for_transaction(&second);
assert_eq!(Transaction::parse(&second).unwrap().parent(), first_id);
assert!(ordering.contains(second_id));
assert_eq!(handler.submissions(), 2);
handler.fail_submit.store(false, Ordering::Relaxed);
ordering
.submit_local_txn(3, [1; 32], owner, b"third", |_| Ok([4; 64]))
.unwrap();
assert_eq!(handler.submissions(), 2);
ordering
.register_subsystem(owner, Some(first_id), handler.clone())
.unwrap();
assert_eq!(handler.submissions(), 4);
}
#[test]
fn remote_submission_reorganizes_and_callback_errors_are_non_authoritative() {
let root = TempRoot::new("remote");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let (a, b) = (subsystem(b'a'), subsystem(b'b'));
let (handler_a, handler_b) = (Arc::new(Recorder::new()), Arc::new(Recorder::new()));
ordering
.register_subsystem(a, None, handler_a.clone())
.unwrap();
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
let first = transaction(GENESIS_PARENT, 20, 1, a, b"first");
let first_id = TxId::for_transaction(&first);
ordering.submit_txn(&first).unwrap();
ordering.submit_txn(&first).unwrap();
assert_eq!(handler_a.submissions(), 1);
let second = transaction(first_id, 50, 2, b, b"second");
let second_id = TxId::for_transaction(&second);
ordering.submit_txn(&second).unwrap();
let third = transaction(second_id, 50, 3, a, b"third");
let third_id = TxId::for_transaction(&third);
ordering.submit_txn(&third).unwrap();
handler_a.fail_reorg.store(true, Ordering::Relaxed);
let replacement = transaction(first_id, 1, 99, a, b"replacement");
let replacement_id = TxId::for_transaction(&replacement);
ordering.submit_txn(&replacement).unwrap();
assert!(ordering.contains(replacement_id));
assert!(!ordering.contains(second_id));
assert!(!ordering.contains(third_id));
assert_eq!(handler_a.reorgs.load(Ordering::Relaxed), 1);
assert_eq!(handler_b.reorgs.load(Ordering::Relaxed), 1);
assert_eq!(handler_a.submissions(), 2);
ordering
.register_subsystem(a, Some(first_id), handler_a.clone())
.unwrap();
assert_eq!(handler_a.submissions(), 3);
handler_b.fail_submit.store(true, Ordering::Relaxed);
ordering
.register_subsystem(b, None, handler_b.clone())
.unwrap();
let extension = transaction(replacement_id, 1, 100, b, b"extension");
assert!(ordering.submit_txn(&extension).is_ok());
assert!(ordering.contains(TxId::for_transaction(&extension)));
}
#[test]
fn callback_blocks_writers_but_not_post_commit_queries() {
let root = TempRoot::new("blocking");
let ordering = Arc::new(K1TxnOrdering::open(&root.0).unwrap());
let owner = subsystem(b'c');
let blocker = Arc::new(Blocker::new());
ordering
.register_subsystem(owner, None, blocker.clone())
.unwrap();
let first_ordering = ordering.clone();
let first = thread::spawn(move || {
first_ordering
.submit_local_txn(1, [1; 32], owner, b"first", |_| Ok([1; 64]))
.unwrap()
});
blocker.wait_until_entered();
let tip = ordering.tip().unwrap();
assert!(ordering.contains(tip));
let (sent, received) = mpsc::channel();
let second_ordering = ordering.clone();
thread::spawn(move || {
let result =
second_ordering.submit_local_txn(2, [1; 32], owner, b"second", |_| Ok([2; 64]));
sent.send(result).unwrap();
});
assert!(matches!(
received.recv_timeout(Duration::from_millis(50)),
Err(mpsc::RecvTimeoutError::Timeout)
));
blocker.release();
first.join().unwrap();
assert!(
received
.recv_timeout(Duration::from_secs(5))
.unwrap()
.is_ok()
);
}
#[test]
fn registration_replay_failure_is_retryable_and_queries_delegate() {
let root = TempRoot::new("replay");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let owner = subsystem(b'd');
let first = transaction(GENESIS_PARENT, 1, 1, owner, b"first");
let first_id = TxId::for_transaction(&first);
ordering.submit_txn(&first).unwrap();
let handler = Arc::new(Recorder::new());
handler.fail_submit.store(true, Ordering::Relaxed);
assert!(
ordering
.register_subsystem(owner, None, handler.clone())
.is_err()
);
handler.fail_submit.store(false, Ordering::Relaxed);
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
assert_eq!(handler.submissions(), 2);
assert_eq!(ordering.tip(), Some(first_id));
assert_eq!(ordering.get_txn(first_id).unwrap(), Some(first));
assert!(
ordering
.between_txids(GENESIS_PARENT, first_id)
.unwrap()
.is_empty()
);
assert!(ordering.register_subsystem(owner, None, handler).is_err());
}
#[test]
fn callback_panics_poison_writers_without_hiding_committed_chain() {
let root = TempRoot::new("submit-panic");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let owner = subsystem(b'e');
let handler = Arc::new(Recorder::new());
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
handler.panic_submit.store(true, Ordering::Relaxed);
let panicking_txn = transaction(GENESIS_PARENT, 1, 1, owner, b"panic");
let id = TxId::for_transaction(&panicking_txn);
assert!(catch_unwind(AssertUnwindSafe(|| ordering.submit_txn(&panicking_txn))).is_err());
assert!(ordering.contains(id));
assert!(catch_unwind(AssertUnwindSafe(|| ordering.submit_txn(&panicking_txn))).is_err());
let root = TempRoot::new("reorg-panic");
let ordering = K1TxnOrdering::open(&root.0).unwrap();
let handler = Arc::new(Recorder::new());
ordering
.register_subsystem(owner, None, handler.clone())
.unwrap();
let first = transaction(GENESIS_PARENT, 20, 1, owner, b"first");
let first_id = TxId::for_transaction(&first);
ordering.submit_txn(&first).unwrap();
let incumbent = transaction(first_id, 50, 2, owner, b"incumbent");
ordering.submit_txn(&incumbent).unwrap();
handler.panic_reorg.store(true, Ordering::Relaxed);
let replacement = transaction(first_id, 1, 3, owner, b"replacement");
let replacement_id = TxId::for_transaction(&replacement);
assert!(catch_unwind(AssertUnwindSafe(|| ordering.submit_txn(&replacement))).is_err());
assert!(ordering.contains(replacement_id));
assert!(
catch_unwind(AssertUnwindSafe(|| {
ordering.register_subsystem(owner, Some(first_id), handler)
}))
.is_err()
);
}
}