use std::collections::{BTreeMap, HashSet};
use std::sync::OnceLock;
use std::sync::atomic::{AtomicU64, Ordering};
use fsqlite_types::{CommitSeq, ObjectId, PageNumber, TxnToken, WitnessKey};
use tracing::{debug, info, warn};
use crate::observability;
use crate::ssi_abort_policy::{
SsiDecisionCard, SsiDecisionCardDraft, SsiDecisionQuery, SsiDecisionType, SsiEvidenceLedger,
};
use crate::witness_objects::{
AbortPolicy, AbortReason, AbortWitness, DependencyEdgeKind, EcsCommitProof, EcsDependencyEdge,
EdgeKeyBasis, KeySummary,
};
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SsiBusySnapshot {
pub txn: TxnToken,
pub reason: SsiAbortReason,
pub witness: AbortWitness,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum SsiAbortReason {
Pivot,
CommittedPivot,
MarkedForAbort,
}
static FSQLITE_EVIDENCE_RECORDS_TOTAL_COMMIT: AtomicU64 = AtomicU64::new(0);
static FSQLITE_EVIDENCE_RECORDS_TOTAL_ABORT: AtomicU64 = AtomicU64::new(0);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct EvidenceRecordMetricsSnapshot {
pub fsqlite_evidence_records_total_commit: u64,
pub fsqlite_evidence_records_total_abort: u64,
}
impl EvidenceRecordMetricsSnapshot {
#[must_use]
pub fn fsqlite_evidence_records_total(self) -> u64 {
self.fsqlite_evidence_records_total_commit + self.fsqlite_evidence_records_total_abort
}
}
fn ssi_evidence_ledger() -> &'static SsiEvidenceLedger {
static LEDGER: OnceLock<SsiEvidenceLedger> = OnceLock::new();
LEDGER.get_or_init(SsiEvidenceLedger::default)
}
#[must_use]
pub fn ssi_evidence_snapshot() -> Vec<SsiDecisionCard> {
ssi_evidence_ledger().snapshot()
}
#[must_use]
pub fn ssi_evidence_query(query: &SsiDecisionQuery) -> Vec<SsiDecisionCard> {
ssi_evidence_ledger().query(query)
}
#[must_use]
pub fn ssi_evidence_metrics_snapshot() -> EvidenceRecordMetricsSnapshot {
EvidenceRecordMetricsSnapshot {
fsqlite_evidence_records_total_commit: FSQLITE_EVIDENCE_RECORDS_TOTAL_COMMIT
.load(Ordering::Relaxed),
fsqlite_evidence_records_total_abort: FSQLITE_EVIDENCE_RECORDS_TOTAL_ABORT
.load(Ordering::Relaxed),
}
}
pub fn reset_ssi_evidence_metrics() {
FSQLITE_EVIDENCE_RECORDS_TOTAL_COMMIT.store(0, Ordering::Relaxed);
FSQLITE_EVIDENCE_RECORDS_TOTAL_ABORT.store(0, Ordering::Relaxed);
}
#[derive(Debug, Clone)]
pub struct SsiState {
pub txn: TxnToken,
pub begin_seq: CommitSeq,
pub has_in_rw: bool,
pub has_out_rw: bool,
pub rw_in_from: HashSet<TxnToken>,
pub rw_out_to: HashSet<TxnToken>,
pub edges_emitted: Vec<ObjectId>,
pub marked_for_abort: bool,
}
impl SsiState {
#[must_use]
pub fn new(txn: TxnToken, begin_seq: CommitSeq) -> Self {
Self {
txn,
begin_seq,
has_in_rw: false,
has_out_rw: false,
rw_in_from: HashSet::new(),
rw_out_to: HashSet::new(),
edges_emitted: Vec::new(),
marked_for_abort: false,
}
}
#[must_use]
pub fn has_dangerous_structure(&self) -> bool {
self.has_in_rw && self.has_out_rw
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct DiscoveredEdge {
pub from: TxnToken,
pub to: TxnToken,
pub overlap_key: WitnessKey,
pub source_is_active: bool,
pub source_has_in_rw: bool,
}
pub trait ActiveTxnView {
fn token(&self) -> TxnToken;
fn begin_seq(&self) -> CommitSeq;
fn is_active(&self) -> bool;
fn read_keys(&self) -> &[WitnessKey];
fn write_keys(&self) -> &[WitnessKey];
fn has_in_rw(&self) -> bool;
fn has_out_rw(&self) -> bool;
fn set_has_out_rw(&self, val: bool);
fn set_has_in_rw(&self, val: bool);
fn set_marked_for_abort(&self, val: bool);
}
#[derive(Debug, Clone)]
pub struct CommittedReaderInfo {
pub token: TxnToken,
pub begin_seq: CommitSeq,
pub commit_seq: CommitSeq,
pub had_in_rw: bool,
pub pages: Vec<PageNumber>,
}
#[derive(Debug, Clone)]
pub struct CommittedWriterInfo {
pub token: TxnToken,
pub commit_seq: CommitSeq,
pub had_out_rw: bool,
pub pages: Vec<PageNumber>,
}
#[derive(Debug, Clone, Copy)]
struct IndexedTxnRecord<'a> {
token: TxnToken,
begin_seq: CommitSeq,
commit_seq: Option<CommitSeq>,
source_is_active: bool,
source_has_in_rw: bool,
key: Option<&'a WitnessKey>,
}
fn build_reader_index<'a>(
committing_txn: TxnToken,
active_readers: &[&'a dyn ActiveTxnView],
committed_readers: &[CommittedReaderInfo],
target_pages: &HashSet<u32>,
) -> BTreeMap<u32, Vec<IndexedTxnRecord<'a>>> {
let mut index: BTreeMap<u32, Vec<IndexedTxnRecord<'a>>> = BTreeMap::new();
for reader in active_readers {
if reader.token() == committing_txn || !reader.is_active() {
continue;
}
for key in reader.read_keys() {
let page = witness_key_page(key);
if !target_pages.contains(&page) {
continue;
}
index.entry(page).or_default().push(IndexedTxnRecord {
token: reader.token(),
begin_seq: reader.begin_seq(),
commit_seq: None,
source_is_active: true,
source_has_in_rw: reader.has_in_rw(),
key: Some(key),
});
}
}
for reader in committed_readers {
if reader.token == committing_txn {
continue;
}
for page in &reader.pages {
let page_no = page.get();
if !target_pages.contains(&page_no) {
continue;
}
index.entry(page_no).or_default().push(IndexedTxnRecord {
token: reader.token,
begin_seq: reader.begin_seq,
commit_seq: Some(reader.commit_seq),
source_is_active: false,
source_has_in_rw: reader.had_in_rw,
key: None,
});
}
}
index
}
fn build_writer_index<'a>(
committing_txn: TxnToken,
active_writers: &[&'a dyn ActiveTxnView],
committed_writers: &[CommittedWriterInfo],
target_pages: &HashSet<u32>,
) -> BTreeMap<u32, Vec<IndexedTxnRecord<'a>>> {
let mut index: BTreeMap<u32, Vec<IndexedTxnRecord<'a>>> = BTreeMap::new();
for writer in active_writers {
if writer.token() == committing_txn || !writer.is_active() {
continue;
}
for key in writer.write_keys() {
let page = witness_key_page(key);
if !target_pages.contains(&page) {
continue;
}
index.entry(page).or_default().push(IndexedTxnRecord {
token: writer.token(),
begin_seq: writer.begin_seq(),
commit_seq: None,
source_is_active: true,
source_has_in_rw: writer.has_out_rw(),
key: Some(key),
});
}
}
for writer in committed_writers {
if writer.token == committing_txn {
continue;
}
for page in &writer.pages {
let page_no = page.get();
if !target_pages.contains(&page_no) {
continue;
}
index.entry(page_no).or_default().push(IndexedTxnRecord {
token: writer.token,
begin_seq: CommitSeq::ZERO,
commit_seq: Some(writer.commit_seq),
source_is_active: false,
source_has_in_rw: writer.had_out_rw,
key: None,
});
}
}
index
}
pub fn discover_incoming_edges(
committing_txn: TxnToken,
committing_begin_seq: CommitSeq,
committing_commit_seq: CommitSeq,
write_keys: &[WitnessKey],
active_readers: &[&dyn ActiveTxnView],
committed_readers: &[CommittedReaderInfo],
) -> Vec<DiscoveredEdge> {
let mut edges = Vec::new();
if write_keys.is_empty() {
return edges;
}
let target_pages: HashSet<u32> = write_keys.iter().map(witness_key_page).collect();
let index = build_reader_index(
committing_txn,
active_readers,
committed_readers,
&target_pages,
);
let committing_begin = committing_begin_seq.get();
let committing_end = committing_commit_seq.get();
let mut seen_sources = HashSet::new();
for write_key in write_keys {
let page = witness_key_page(write_key);
let Some(candidates) = index.get(&page) else {
continue;
};
for candidate in candidates {
if let Some(reader_key) = candidate.key {
if let (WitnessKey::Cell { tag: t1, .. }, WitnessKey::Cell { tag: t2, .. }) =
(write_key, reader_key)
{
if t1 != t2 {
continue;
}
}
}
let candidate_begin = candidate.begin_seq.get();
let candidate_end = candidate.commit_seq.map_or(u64::MAX, CommitSeq::get);
let overlaps = committing_begin < candidate_end && candidate_begin < committing_end;
if !overlaps || !seen_sources.insert(candidate.token) {
continue;
}
let source = if candidate.source_is_active {
"hot_plane_index"
} else {
"rcri_index"
};
debug!(
bead_id = "bd-31bo",
from = ?candidate.token,
to = ?committing_txn,
key = ?write_key,
source,
"discovered incoming rw-antidependency edge"
);
edges.push(DiscoveredEdge {
from: candidate.token,
to: committing_txn,
overlap_key: write_key.clone(),
source_is_active: candidate.source_is_active,
source_has_in_rw: candidate.source_has_in_rw,
});
}
}
edges
}
pub fn discover_outgoing_edges(
committing_txn: TxnToken,
committing_begin_seq: CommitSeq,
committing_commit_seq: CommitSeq,
read_keys: &[WitnessKey],
active_writers: &[&dyn ActiveTxnView],
committed_writers: &[CommittedWriterInfo],
) -> Vec<DiscoveredEdge> {
let mut edges = Vec::new();
if read_keys.is_empty() {
return edges;
}
let target_pages: HashSet<u32> = read_keys.iter().map(witness_key_page).collect();
let index = build_writer_index(
committing_txn,
active_writers,
committed_writers,
&target_pages,
);
let committing_begin = committing_begin_seq.get();
let committing_end = committing_commit_seq.get();
let mut seen_targets = HashSet::new();
for read_key in read_keys {
let page = witness_key_page(read_key);
let Some(candidates) = index.get(&page) else {
continue;
};
for candidate in candidates {
if let Some(writer_key) = candidate.key {
if let (WitnessKey::Cell { tag: t1, .. }, WitnessKey::Cell { tag: t2, .. }) =
(read_key, writer_key)
{
if t1 != t2 {
continue;
}
}
}
let candidate_begin = candidate.begin_seq.get();
let candidate_end = candidate.commit_seq.map_or(u64::MAX, CommitSeq::get);
let overlaps = committing_begin < candidate_end && candidate_begin < committing_end;
if !overlaps || !seen_targets.insert(candidate.token) {
continue;
}
let source = if candidate.source_is_active {
"hot_plane_index"
} else {
"commit_log_index"
};
debug!(
bead_id = "bd-31bo",
from = ?committing_txn,
to = ?candidate.token,
key = ?read_key,
source,
"discovered outgoing rw-antidependency edge"
);
edges.push(DiscoveredEdge {
from: committing_txn,
to: candidate.token,
overlap_key: read_key.clone(),
source_is_active: candidate.source_is_active,
source_has_in_rw: candidate.source_has_in_rw,
});
}
}
edges
}
pub(crate) fn witness_key_page(key: &WitnessKey) -> u32 {
match key {
WitnessKey::Page(p) => p.get(),
WitnessKey::Cell { btree_root, .. } | WitnessKey::KeyRange { btree_root, .. } => {
btree_root.get()
}
WitnessKey::ByteRange { page, .. } => page.get(),
WitnessKey::Custom { .. } => 0,
}
}
#[derive(Debug, Clone)]
pub struct SsiValidationOk {
pub edges: Vec<EcsDependencyEdge>,
pub edge_ids: Vec<ObjectId>,
pub commit_proof: EcsCommitProof,
pub ssi_state: SsiState,
}
#[allow(clippy::too_many_lines, clippy::too_many_arguments)]
pub fn ssi_validate_and_publish(
txn: TxnToken,
begin_seq: CommitSeq,
commit_seq: CommitSeq,
read_keys: &[WitnessKey],
write_keys: &[WitnessKey],
active_readers: &[&dyn ActiveTxnView],
active_writers: &[&dyn ActiveTxnView],
committed_readers: &[CommittedReaderInfo],
committed_writers: &[CommittedWriterInfo],
marked_for_abort: bool,
) -> Result<SsiValidationOk, SsiBusySnapshot> {
let mut state = SsiState::new(txn, begin_seq);
state.marked_for_abort = marked_for_abort;
let span = tracing::span!(
tracing::Level::INFO,
"ssi_validate",
txn_id = txn.id.get(),
read_set_size = read_keys.len(),
write_set_size = write_keys.len(),
conflict_detected = tracing::field::Empty,
decision_reason = tracing::field::Empty,
);
let _guard = span.enter();
info!(
bead_id = "bd-31bo",
txn = ?txn,
read_keys = read_keys.len(),
write_keys = write_keys.len(),
marked_for_abort,
"ssi_validate_and_publish: starting"
);
if write_keys.is_empty() {
record_evidence_decision(
SsiDecisionType::CommitAllowed,
txn,
begin_seq,
Some(commit_seq),
read_keys,
write_keys,
&[],
"read_only_fast_path",
);
span.record("conflict_detected", false);
span.record("decision_reason", "read_only_fast_path");
debug!(
bead_id = "bd-31bo",
txn = ?txn,
"ssi_validate: read-only fast path, skipping SSI"
);
let proof = build_commit_proof(txn, begin_seq, commit_seq, &state, &[], &[]);
observability::record_ssi_commit();
return Ok(SsiValidationOk {
edges: Vec::new(),
edge_ids: Vec::new(),
commit_proof: proof,
ssi_state: state,
});
}
if marked_for_abort {
record_evidence_decision(
SsiDecisionType::AbortCycle,
txn,
begin_seq,
Some(commit_seq),
read_keys,
write_keys,
&[],
"marked_for_abort",
);
span.record("conflict_detected", true);
span.record("decision_reason", "marked_for_abort");
warn!(
bead_id = "bd-31bo",
txn = ?txn,
"ssi_validate: transaction marked for abort by another committer"
);
observability::record_ssi_abort(fsqlite_observability::SsiAbortCategory::MarkedForAbort);
let witness = AbortWitness {
txn,
begin_seq,
abort_seq: commit_seq,
reason: AbortReason::SsiPivot,
edges_observed: Vec::new(),
};
return Err(SsiBusySnapshot {
txn,
reason: SsiAbortReason::MarkedForAbort,
witness,
});
}
let in_edges = discover_incoming_edges(
txn,
begin_seq,
commit_seq,
write_keys,
active_readers,
committed_readers,
);
let out_edges = discover_outgoing_edges(
txn,
begin_seq,
commit_seq,
read_keys,
active_writers,
committed_writers,
);
state.has_in_rw = !in_edges.is_empty();
state.has_out_rw = !out_edges.is_empty();
for edge in &in_edges {
state.rw_in_from.insert(edge.from);
}
for edge in &out_edges {
state.rw_out_to.insert(edge.to);
}
info!(
bead_id = "bd-31bo",
txn = ?txn,
incoming = in_edges.len(),
outgoing = out_edges.len(),
has_in_rw = state.has_in_rw,
has_out_rw = state.has_out_rw,
"ssi_validate: edge discovery complete"
);
if state.has_in_rw && state.has_out_rw {
let all_edges = build_dependency_edges(&in_edges, &out_edges, txn, commit_seq);
let discovered_edges: Vec<DiscoveredEdge> = in_edges
.iter()
.cloned()
.chain(out_edges.iter().cloned())
.collect();
record_evidence_decision(
SsiDecisionType::AbortWriteSkew,
txn,
begin_seq,
Some(commit_seq),
read_keys,
write_keys,
&discovered_edges,
"pivot_abort_dangerous_structure",
);
span.record("conflict_detected", true);
span.record("decision_reason", "pivot_abort");
warn!(
bead_id = "bd-31bo",
txn = ?txn,
in_sources = ?state.rw_in_from,
out_targets = ?state.rw_out_to,
"ssi_validate: PIVOT ABORT — dangerous structure detected"
);
observability::record_ssi_abort(fsqlite_observability::SsiAbortCategory::Pivot);
let witness = AbortWitness {
txn,
begin_seq,
abort_seq: commit_seq,
reason: AbortReason::SsiPivot,
edges_observed: all_edges,
};
return Err(SsiBusySnapshot {
txn,
reason: SsiAbortReason::Pivot,
witness,
});
}
for edge in &in_edges {
if edge.source_is_active {
for reader in active_readers {
if reader.token() == edge.from {
reader.set_has_out_rw(true);
if reader.has_in_rw() {
debug!(
bead_id = "bd-31bo",
pivot = ?edge.from,
"T3 rule: active reader is pivot, marking for abort"
);
reader.set_marked_for_abort(true);
}
break;
}
}
} else {
if edge.source_has_in_rw {
let discovered_edges: Vec<DiscoveredEdge> = in_edges
.iter()
.cloned()
.chain(out_edges.iter().cloned())
.collect();
record_evidence_decision(
SsiDecisionType::AbortCycle,
txn,
begin_seq,
Some(commit_seq),
read_keys,
write_keys,
&discovered_edges,
"committed_pivot_abort",
);
span.record("conflict_detected", true);
span.record("decision_reason", "committed_pivot_abort");
warn!(
bead_id = "bd-31bo",
txn = ?txn,
committed_pivot = ?edge.from,
"T3 rule: committed reader was pivot, T must abort"
);
observability::record_ssi_abort(
fsqlite_observability::SsiAbortCategory::CommittedPivot,
);
let all_edges = build_dependency_edges(&in_edges, &out_edges, txn, commit_seq);
let witness = AbortWitness {
txn,
begin_seq,
abort_seq: commit_seq,
reason: AbortReason::SsiPivot,
edges_observed: all_edges,
};
return Err(SsiBusySnapshot {
txn,
reason: SsiAbortReason::CommittedPivot,
witness,
});
}
}
}
for edge in &out_edges {
if edge.source_is_active {
for writer in active_writers {
if writer.token() == edge.to {
writer.set_has_in_rw(true);
if writer.has_out_rw() {
debug!(
bead_id = "bd-31bo",
pivot = ?edge.to,
"T3 rule: active writer is pivot, marking for abort"
);
writer.set_marked_for_abort(true);
}
break;
}
}
} else if edge.source_has_in_rw {
let discovered_edges: Vec<DiscoveredEdge> = in_edges
.iter()
.cloned()
.chain(out_edges.iter().cloned())
.collect();
record_evidence_decision(
SsiDecisionType::AbortCycle,
txn,
begin_seq,
Some(commit_seq),
read_keys,
write_keys,
&discovered_edges,
"committed_writer_pivot_abort",
);
span.record("conflict_detected", true);
span.record("decision_reason", "committed_writer_pivot_abort");
warn!(
bead_id = "bd-31bo",
txn = ?txn,
committed_pivot = ?edge.to,
"T3 rule: committed writer was pivot, T must abort"
);
observability::record_ssi_abort(
fsqlite_observability::SsiAbortCategory::CommittedPivot,
);
let all_edges = build_dependency_edges(&in_edges, &out_edges, txn, commit_seq);
let witness = AbortWitness {
txn,
begin_seq,
abort_seq: commit_seq,
reason: AbortReason::SsiPivot,
edges_observed: all_edges,
};
return Err(SsiBusySnapshot {
txn,
reason: SsiAbortReason::CommittedPivot,
witness,
});
}
}
if !out_edges.is_empty() {
debug!(
bead_id = "bd-31bo",
outgoing_edges = out_edges.len(),
"ssi_validate: outgoing edge propagation complete"
);
}
let all_edges = build_dependency_edges(&in_edges, &out_edges, txn, commit_seq);
let discovered_edges: Vec<DiscoveredEdge> = in_edges
.iter()
.cloned()
.chain(out_edges.iter().cloned())
.collect();
record_evidence_decision(
SsiDecisionType::CommitAllowed,
txn,
begin_seq,
Some(commit_seq),
read_keys,
write_keys,
&discovered_edges,
"commit_approved",
);
let edge_ids: Vec<ObjectId> = all_edges
.iter()
.enumerate()
.map(|(i, _)| {
let mut bytes = [0u8; 16];
bytes[..8].copy_from_slice(&txn.id.get().to_le_bytes());
bytes[8..12].copy_from_slice(&commit_seq.get().to_le_bytes()[..4]);
#[allow(clippy::cast_possible_truncation)]
let idx = i as u32;
bytes[12..16].copy_from_slice(&idx.to_le_bytes());
ObjectId::from_bytes(bytes)
})
.collect();
state.edges_emitted.clone_from(&edge_ids);
let proof = build_commit_proof(txn, begin_seq, commit_seq, &state, &edge_ids, &[]);
span.record("conflict_detected", false);
span.record("decision_reason", "commit_approved");
info!(
bead_id = "bd-31bo",
txn = ?txn,
edges_emitted = all_edges.len(),
"ssi_validate: commit approved, evidence published"
);
observability::record_ssi_commit();
Ok(SsiValidationOk {
edges: all_edges,
edge_ids,
commit_proof: proof,
ssi_state: state,
})
}
fn build_dependency_edges(
in_edges: &[DiscoveredEdge],
out_edges: &[DiscoveredEdge],
observer: TxnToken,
observation_seq: CommitSeq,
) -> Vec<EcsDependencyEdge> {
let mut result = Vec::with_capacity(in_edges.len() + out_edges.len());
for edge in in_edges.iter().chain(out_edges.iter()) {
result.push(EcsDependencyEdge {
kind: DependencyEdgeKind::RwAntiDependency,
from: edge.from,
to: edge.to,
key_basis: EdgeKeyBasis {
level: 0,
range_prefix: witness_key_page(&edge.overlap_key),
refinement: Some(KeySummary::ExactKeys(vec![edge.overlap_key.clone()])),
},
observed_by: observer,
observation_seq,
});
}
result
}
fn build_commit_proof(
txn: TxnToken,
begin_seq: CommitSeq,
commit_seq: CommitSeq,
state: &SsiState,
edge_ids: &[ObjectId],
merge_witnesses: &[ObjectId],
) -> EcsCommitProof {
EcsCommitProof {
txn,
begin_seq,
commit_seq,
has_in_rw: state.has_in_rw,
has_out_rw: state.has_out_rw,
read_witness_refs: Vec::new(),
write_witness_refs: Vec::new(),
index_segments_used: Vec::new(),
edges_emitted: edge_ids.to_vec(),
merge_witnesses: merge_witnesses.to_vec(),
abort_policy: AbortPolicy::AbortPivot,
}
}
fn decision_outcome(decision_type: SsiDecisionType) -> &'static str {
match decision_type {
SsiDecisionType::CommitAllowed => "commit",
SsiDecisionType::AbortWriteSkew
| SsiDecisionType::AbortPhantom
| SsiDecisionType::AbortCycle => "abort",
}
}
fn witness_keys_to_pages(keys: &[WitnessKey]) -> Vec<PageNumber> {
let mut pages: Vec<PageNumber> = keys
.iter()
.filter_map(|key| PageNumber::new(witness_key_page(key)))
.collect();
pages.sort_by_key(|page| page.get());
pages.dedup();
pages
}
fn edge_conflicting_txns(txn: TxnToken, edges: &[DiscoveredEdge]) -> Vec<TxnToken> {
let mut txns = Vec::new();
for edge in edges {
if edge.from != txn {
txns.push(edge.from);
}
if edge.to != txn {
txns.push(edge.to);
}
}
txns.sort_by(|left, right| {
left.id
.get()
.cmp(&right.id.get())
.then_with(|| left.epoch.get().cmp(&right.epoch.get()))
});
txns.dedup();
txns
}
fn edge_conflict_pages(edges: &[DiscoveredEdge]) -> Vec<PageNumber> {
let mut pages: Vec<PageNumber> = edges
.iter()
.filter_map(|edge| PageNumber::new(witness_key_page(&edge.overlap_key)))
.collect();
pages.sort_by_key(|page| page.get());
pages.dedup();
pages
}
fn estimate_evidence_size_bytes(
read_pages: &[PageNumber],
write_pages: &[PageNumber],
conflict_pages: &[PageNumber],
conflicting_txns: &[TxnToken],
rationale: &str,
) -> u64 {
#[allow(clippy::cast_possible_truncation)]
{
let words = read_pages.len()
+ write_pages.len()
+ conflict_pages.len()
+ conflicting_txns.len() * 2
+ rationale.len();
(words * std::mem::size_of::<u64>()) as u64
}
}
#[allow(clippy::too_many_arguments)]
fn record_evidence_decision(
decision_type: SsiDecisionType,
txn: TxnToken,
begin_seq: CommitSeq,
commit_seq: Option<CommitSeq>,
read_keys: &[WitnessKey],
write_keys: &[WitnessKey],
edges: &[DiscoveredEdge],
rationale: &str,
) {
let read_pages = witness_keys_to_pages(read_keys);
let write_pages = witness_keys_to_pages(write_keys);
let conflict_pages = edge_conflict_pages(edges);
let conflicting_txns = edge_conflicting_txns(txn, edges);
let evidence_size_bytes = estimate_evidence_size_bytes(
&read_pages,
&write_pages,
&conflict_pages,
&conflicting_txns,
rationale,
);
let outcome = decision_outcome(decision_type);
let decision_id = fsqlite_observability::next_decision_id();
let span = tracing::span!(
target: "fsqlite.evidence",
tracing::Level::INFO,
"evidence_record",
decision_id,
outcome,
evidence_size_bytes
);
let _guard = span.enter();
let mut draft = SsiDecisionCardDraft::new(
decision_type,
txn,
begin_seq,
conflicting_txns.clone(),
conflict_pages.clone(),
read_pages.clone(),
write_pages.clone(),
rationale,
)
.with_decision_id(decision_id);
if let Some(seq) = commit_seq {
draft = draft.with_commit_seq(seq);
}
ssi_evidence_ledger().record_async(draft);
if matches!(decision_type, SsiDecisionType::CommitAllowed) {
FSQLITE_EVIDENCE_RECORDS_TOTAL_COMMIT.fetch_add(1, Ordering::Relaxed);
} else {
FSQLITE_EVIDENCE_RECORDS_TOTAL_ABORT.fetch_add(1, Ordering::Relaxed);
}
info!(
decision_id,
outcome,
decision_type = %decision_type,
txn_id = txn.id.get(),
read_pages = read_pages.len(),
write_pages = write_pages.len(),
conflicting_txns = conflicting_txns.len(),
conflict_pages = conflict_pages.len(),
"ssi decision evidence recorded"
);
debug!(
decision_id,
decision_type = %decision_type,
txn = ?txn,
read_pages = ?read_pages,
write_pages = ?write_pages,
conflicting_txns = ?conflicting_txns,
conflict_pages = ?conflict_pages,
rationale,
"ssi decision evidence details"
);
}
#[cfg(test)]
mod tests {
use std::collections::HashSet;
use super::*;
use fsqlite_types::{TxnEpoch, TxnId};
use std::cell::Cell;
struct MockActiveTxn {
token: TxnToken,
begin_seq: CommitSeq,
active: bool,
reads: Vec<WitnessKey>,
writes: Vec<WitnessKey>,
has_in: Cell<bool>,
has_out: Cell<bool>,
marked: Cell<bool>,
}
impl MockActiveTxn {
fn new(id: u64, epoch: u32, begin_seq: u64) -> Self {
Self {
token: TxnToken::new(TxnId::new(id).unwrap(), TxnEpoch::new(epoch)),
begin_seq: CommitSeq::new(begin_seq),
active: true,
reads: Vec::new(),
writes: Vec::new(),
has_in: Cell::new(false),
has_out: Cell::new(false),
marked: Cell::new(false),
}
}
fn with_reads(mut self, keys: Vec<WitnessKey>) -> Self {
self.reads = keys;
self
}
fn with_writes(mut self, keys: Vec<WitnessKey>) -> Self {
self.writes = keys;
self
}
fn with_has_in_rw(self, val: bool) -> Self {
self.has_in.set(val);
self
}
#[allow(dead_code)]
fn committed(mut self) -> Self {
self.active = false;
self
}
}
impl ActiveTxnView for MockActiveTxn {
fn token(&self) -> TxnToken {
self.token
}
fn begin_seq(&self) -> CommitSeq {
self.begin_seq
}
fn is_active(&self) -> bool {
self.active
}
fn read_keys(&self) -> &[WitnessKey] {
&self.reads
}
fn write_keys(&self) -> &[WitnessKey] {
&self.writes
}
fn has_in_rw(&self) -> bool {
self.has_in.get()
}
fn has_out_rw(&self) -> bool {
self.has_out.get()
}
fn set_has_out_rw(&self, val: bool) {
self.has_out.set(val);
}
fn set_has_in_rw(&self, val: bool) {
self.has_in.set(val);
}
fn set_marked_for_abort(&self, val: bool) {
self.marked.set(val);
}
}
fn page_key(pgno: u32) -> WitnessKey {
WitnessKey::Page(PageNumber::new(pgno).unwrap())
}
fn keys_to_pages(keys: &[WitnessKey]) -> Vec<PageNumber> {
let mut pages = Vec::new();
for key in keys {
if let WitnessKey::Page(page) = key {
pages.push(*page);
}
}
pages
}
#[test]
fn test_ssi_read_only_skip() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(10)], &[], &[],
&[],
&[],
&[],
false,
);
let ok = result.expect("read-only txn should commit");
assert!(ok.edges.is_empty(), "no edges for read-only");
assert!(!ok.ssi_state.has_in_rw);
assert!(!ok.ssi_state.has_out_rw);
}
#[test]
fn test_ssi_no_edges_commit() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(10)],
&[page_key(20)], &[],
&[],
&[],
&[],
false,
);
let ok = result.expect("no overlap → no edges → commit");
assert!(ok.edges.is_empty());
assert!(!ok.ssi_state.has_in_rw);
assert!(!ok.ssi_state.has_out_rw);
}
#[test]
fn test_safe_snapshot_shortcut_no_conflict_commit() {
let txn = TxnToken::new(TxnId::new(8_001).unwrap(), TxnEpoch::new(0));
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(10),
CommitSeq::new(11),
&[page_key(1_000)],
&[page_key(2_000)],
&[],
&[],
&[],
&[],
false,
)
.expect("safe snapshot should commit without conflict work");
assert!(
result.edges.is_empty(),
"safe snapshot should emit no edges"
);
assert!(
result.edge_ids.is_empty(),
"safe snapshot should emit no edge identifiers"
);
assert!(!result.ssi_state.has_in_rw);
assert!(!result.ssi_state.has_out_rw);
}
#[test]
fn test_ssi_only_incoming_edge_commit() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[], &[page_key(5)], &readers, &[],
&[],
&[],
false,
);
let ok = result.expect("only incoming edge → commit allowed");
assert!(ok.ssi_state.has_in_rw);
assert!(!ok.ssi_state.has_out_rw);
assert!(!ok.edges.is_empty());
}
#[test]
fn test_ssi_only_outgoing_edge_commit() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let writer = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(7)]);
let writers: Vec<&dyn ActiveTxnView> = vec![&writer];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)], &[page_key(20)], &[],
&writers,
&[],
&[],
false,
);
let ok = result.expect("only outgoing edge → commit allowed");
assert!(!ok.ssi_state.has_in_rw);
assert!(ok.ssi_state.has_out_rw);
}
#[test]
fn test_ssi_pivot_both_edges_abort() {
let txn = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(1, 0, 1).with_reads(vec![page_key(5)]);
let writer = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(7)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let writers: Vec<&dyn ActiveTxnView> = vec![&writer];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)], &[page_key(5)], &readers,
&writers,
&[],
&[],
false,
);
let err = result.expect_err("both in + out rw → MUST abort");
assert_eq!(err.reason, SsiAbortReason::Pivot);
assert_eq!(err.witness.reason, AbortReason::SsiPivot);
}
#[test]
fn test_ssi_dangerous_structure_detection() {
let t2 = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let t1 = MockActiveTxn::new(1, 0, 1).with_reads(vec![page_key(10)]);
let t3 = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(20)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&t1];
let writers: Vec<&dyn ActiveTxnView> = vec![&t3];
let result = ssi_validate_and_publish(
t2,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(20)], &[page_key(10)], &readers,
&writers,
&[],
&[],
false,
);
assert!(result.is_err(), "dangerous structure → abort");
let err = result.unwrap_err();
assert_eq!(err.reason, SsiAbortReason::Pivot);
}
#[test]
fn test_discover_incoming_from_hot_plane() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let edges = discover_incoming_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(5)],
&readers,
&[],
);
assert_eq!(edges.len(), 1);
assert!(edges[0].source_is_active);
assert_eq!(edges[0].from.id.get(), 2);
}
#[test]
fn test_discover_outgoing_from_hot_plane() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let writer = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(7)]);
let writers: Vec<&dyn ActiveTxnView> = vec![&writer];
let edges = discover_outgoing_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)],
&writers,
&[],
);
assert_eq!(edges.len(), 1);
assert!(edges[0].source_is_active);
assert_eq!(edges[0].to.id.get(), 3);
}
#[test]
fn test_discover_outgoing_from_commit_index() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let committed_w = CommittedWriterInfo {
token: TxnToken::new(TxnId::new(3).unwrap(), TxnEpoch::new(0)),
commit_seq: CommitSeq::new(3),
had_out_rw: false,
pages: vec![PageNumber::new(7).unwrap()],
};
let edges = discover_outgoing_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)],
&[],
&[committed_w],
);
assert_eq!(edges.len(), 1);
assert!(!edges[0].source_is_active);
}
#[test]
fn test_discover_incoming_from_recently_committed() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let committed_r = CommittedReaderInfo {
token: TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0)),
begin_seq: CommitSeq::new(0),
commit_seq: CommitSeq::new(3),
had_in_rw: false,
pages: vec![PageNumber::new(5).unwrap()],
};
let edges = discover_incoming_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(5)],
&[],
&[committed_r],
);
assert_eq!(edges.len(), 1);
assert!(!edges[0].source_is_active);
}
#[test]
fn test_edge_gap_without_commit_index() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let committed_w = CommittedWriterInfo {
token: TxnToken::new(TxnId::new(3).unwrap(), TxnEpoch::new(0)),
commit_seq: CommitSeq::new(3),
had_out_rw: false,
pages: vec![PageNumber::new(7).unwrap()],
};
let edges_hot_only = discover_outgoing_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)],
&[], &[], );
assert!(
edges_hot_only.is_empty(),
"hot-plane only misses committed writer"
);
let edges_full = discover_outgoing_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)],
&[],
&[committed_w],
);
assert_eq!(edges_full.len(), 1, "commit index catches the edge");
}
#[test]
fn test_edge_gap_without_recently_committed() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let committed_r = CommittedReaderInfo {
token: TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0)),
begin_seq: CommitSeq::new(0),
commit_seq: CommitSeq::new(3),
had_in_rw: false,
pages: vec![PageNumber::new(5).unwrap()],
};
let edges_hot_only = discover_incoming_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(5)],
&[],
&[], );
assert!(
edges_hot_only.is_empty(),
"hot-plane only misses committed reader"
);
let edges_full = discover_incoming_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(5)],
&[],
&[committed_r],
);
assert_eq!(edges_full.len(), 1, "RCRI catches the edge");
}
#[test]
fn test_interval_overlap_excludes_future_active_reader() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let late_reader = MockActiveTxn::new(2, 0, 9).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&late_reader];
let edges = discover_incoming_edges(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(5)],
&readers,
&[],
);
assert!(
edges.is_empty(),
"reader interval [9,+inf) does not overlap [1,5]"
);
}
#[test]
fn test_interval_overlap_excludes_stale_committed_writer() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let stale_writer = CommittedWriterInfo {
token: TxnToken::new(TxnId::new(3).unwrap(), TxnEpoch::new(0)),
commit_seq: CommitSeq::new(4),
had_out_rw: false,
pages: vec![PageNumber::new(7).unwrap()],
};
let edges = discover_outgoing_edges(
txn,
CommitSeq::new(5),
CommitSeq::new(8),
&[page_key(7)],
&[],
&[stale_writer],
);
assert!(
edges.is_empty(),
"writer interval (-inf,4] does not overlap [5,8]"
);
}
#[test]
fn test_bank_transfer_write_skew_prevented() {
let t1 = TxnToken::new(TxnId::new(11).unwrap(), TxnEpoch::new(0));
let t2 = TxnToken::new(TxnId::new(12).unwrap(), TxnEpoch::new(0));
let t2_active = MockActiveTxn::new(12, 0, 1).with_reads(vec![page_key(100), page_key(200)]);
let readers_for_t1: Vec<&dyn ActiveTxnView> = vec![&t2_active];
let t1_reads = [page_key(100), page_key(200)];
let t1_writes = [page_key(100)];
let t1_commit = ssi_validate_and_publish(
t1,
CommitSeq::new(1),
CommitSeq::new(2),
&t1_reads,
&t1_writes,
&readers_for_t1,
&[],
&[],
&[],
false,
)
.expect("first transfer leg should commit");
let committed_reader_t1 = CommittedReaderInfo {
token: t1,
begin_seq: CommitSeq::new(1),
commit_seq: CommitSeq::new(2),
had_in_rw: t1_commit.ssi_state.has_in_rw,
pages: vec![PageNumber::new(100).unwrap(), PageNumber::new(200).unwrap()],
};
let committed_writer_t1 = CommittedWriterInfo {
token: t1,
commit_seq: CommitSeq::new(2),
had_out_rw: t1_commit.ssi_state.has_out_rw,
pages: vec![PageNumber::new(100).unwrap()],
};
let t2_reads = [page_key(100), page_key(200)];
let t2_writes = [page_key(200)];
let t2_result = ssi_validate_and_publish(
t2,
CommitSeq::new(1),
CommitSeq::new(3),
&t2_reads,
&t2_writes,
&[],
&[],
&[committed_reader_t1],
&[committed_writer_t1],
false,
);
assert!(
t2_result.is_err(),
"second transfer leg must abort to prevent write skew"
);
}
#[test]
fn test_doctor_on_call_write_skew_prevented() {
let d1 = TxnToken::new(TxnId::new(21).unwrap(), TxnEpoch::new(0));
let d2 = TxnToken::new(TxnId::new(22).unwrap(), TxnEpoch::new(0));
let d2_active = MockActiveTxn::new(22, 0, 1).with_reads(vec![page_key(310), page_key(311)]);
let readers_for_d1: Vec<&dyn ActiveTxnView> = vec![&d2_active];
let d1_reads = [page_key(310), page_key(311)];
let d1_writes = [page_key(310)];
let d1_commit = ssi_validate_and_publish(
d1,
CommitSeq::new(1),
CommitSeq::new(2),
&d1_reads,
&d1_writes,
&readers_for_d1,
&[],
&[],
&[],
false,
)
.expect("first doctor update should commit");
let committed_reader_d1 = CommittedReaderInfo {
token: d1,
begin_seq: CommitSeq::new(1),
commit_seq: CommitSeq::new(2),
had_in_rw: d1_commit.ssi_state.has_in_rw,
pages: vec![PageNumber::new(310).unwrap(), PageNumber::new(311).unwrap()],
};
let committed_writer_d1 = CommittedWriterInfo {
token: d1,
commit_seq: CommitSeq::new(2),
had_out_rw: d1_commit.ssi_state.has_out_rw,
pages: vec![PageNumber::new(310).unwrap()],
};
let d2_reads = [page_key(310), page_key(311)];
let d2_writes = [page_key(311)];
let d2_result = ssi_validate_and_publish(
d2,
CommitSeq::new(1),
CommitSeq::new(3),
&d2_reads,
&d2_writes,
&[],
&[],
&[committed_reader_d1],
&[committed_writer_d1],
false,
);
assert!(
d2_result.is_err(),
"second doctor update must abort to preserve on-call invariant"
);
}
#[test]
fn test_t3_rule_active_pivot_marked() {
let t = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let r = MockActiveTxn::new(2, 0, 1)
.with_reads(vec![page_key(5)])
.with_has_in_rw(true);
let readers: Vec<&dyn ActiveTxnView> = vec![&r];
let result = ssi_validate_and_publish(
t,
CommitSeq::new(1),
CommitSeq::new(5),
&[], &[page_key(5)], &readers,
&[],
&[],
&[],
false,
);
result.expect("T has only incoming edge, should commit");
assert!(r.has_out.get(), "R.has_out_rw should be set to true");
assert!(r.marked.get(), "R should be marked for abort (T3 rule)");
}
#[test]
fn test_t3_rule_committed_pivot_forces_abort() {
let t = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let committed_r = CommittedReaderInfo {
token: TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0)),
begin_seq: CommitSeq::new(0),
commit_seq: CommitSeq::new(3),
had_in_rw: true, pages: vec![PageNumber::new(5).unwrap()],
};
let result = ssi_validate_and_publish(
t,
CommitSeq::new(1),
CommitSeq::new(5),
&[],
&[page_key(5)],
&[],
&[],
&[committed_r],
&[],
false,
);
let err = result.expect_err("committed pivot → T must abort");
assert_eq!(err.reason, SsiAbortReason::CommittedPivot);
}
#[test]
fn test_t3_rule_active_no_in_rw_no_mark() {
let t = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let r = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&r];
let result = ssi_validate_and_publish(
t,
CommitSeq::new(1),
CommitSeq::new(5),
&[],
&[page_key(5)],
&readers,
&[],
&[],
&[],
false,
);
result.expect("T should commit");
assert!(r.has_out.get(), "R.has_out_rw should be set");
assert!(!r.marked.get(), "R should NOT be marked (no in_rw)");
}
#[test]
fn test_refinement_eliminates_false_edge() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(10)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(30)],
&[page_key(20)],
&readers,
&[],
&[],
&[],
false,
);
let ok = result.expect("no overlap → commit");
assert!(!ok.ssi_state.has_in_rw);
}
#[test]
fn test_skip_refinement_safe() {
let t = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let t1 = MockActiveTxn::new(1, 0, 1).with_reads(vec![page_key(5)]);
let t3 = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&t1];
let writers: Vec<&dyn ActiveTxnView> = vec![&t3];
let result = ssi_validate_and_publish(
t,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(5)],
&[page_key(5)],
&readers,
&writers,
&[],
&[],
false,
);
assert!(result.is_err(), "without refinement, overlap → abort");
}
#[test]
fn test_dependency_edge_published() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[],
&[page_key(5)],
&readers,
&[],
&[],
&[],
false,
);
let ok = result.unwrap();
assert!(!ok.edges.is_empty(), "edge must be published");
assert_eq!(ok.edges[0].kind, DependencyEdgeKind::RwAntiDependency);
assert_eq!(ok.edges[0].from.id.get(), 2); assert_eq!(ok.edges[0].to.id.get(), 1); }
#[test]
fn test_commit_proof_published() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(10)],
&[page_key(20)],
&[],
&[],
&[],
&[],
false,
);
let ok = result.unwrap();
assert_eq!(ok.commit_proof.txn, txn);
assert_eq!(ok.commit_proof.begin_seq.get(), 1);
assert_eq!(ok.commit_proof.commit_seq.get(), 5);
assert_eq!(ok.commit_proof.abort_policy, AbortPolicy::AbortPivot);
}
#[test]
fn test_abort_witness_published() {
let txn = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(1, 0, 1).with_reads(vec![page_key(5)]);
let writer = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(7)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let writers: Vec<&dyn ActiveTxnView> = vec![&writer];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)],
&[page_key(5)],
&readers,
&writers,
&[],
&[],
false,
);
let err = result.unwrap_err();
assert_eq!(err.witness.txn, txn);
assert_eq!(err.witness.reason, AbortReason::SsiPivot);
assert!(
!err.witness.edges_observed.is_empty(),
"abort witness must contain edges"
);
}
#[test]
fn test_ssi_state_has_in_rw_flag() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[],
&[page_key(5)],
&readers,
&[],
&[],
&[],
false,
);
let ok = result.unwrap();
assert!(ok.ssi_state.has_in_rw, "incoming edge must set has_in_rw");
}
#[test]
fn test_ssi_state_has_out_rw_flag() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let writer = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(7)]);
let writers: Vec<&dyn ActiveTxnView> = vec![&writer];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(7)],
&[page_key(20)],
&[],
&writers,
&[],
&[],
false,
);
let ok = result.unwrap();
assert!(ok.ssi_state.has_out_rw, "outgoing edge must set has_out_rw");
}
#[test]
fn test_ssi_state_marked_for_abort() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(10)],
&[page_key(20)],
&[],
&[],
&[],
&[],
true, );
let err = result.expect_err("marked_for_abort → must abort");
assert_eq!(err.reason, SsiAbortReason::MarkedForAbort);
}
#[test]
fn test_ssi_state_edges_emitted_tracking() {
let txn = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_key(5)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(5),
&[],
&[page_key(5)],
&readers,
&[],
&[],
&[],
false,
);
let ok = result.unwrap();
assert_eq!(ok.edge_ids.len(), ok.edges.len());
assert_eq!(ok.ssi_state.edges_emitted.len(), ok.edges.len());
}
#[test]
fn test_conservative_pivot_rule() {
let t2 = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let t1 = MockActiveTxn::new(1, 0, 1).with_reads(vec![page_key(10)]);
let t3 = MockActiveTxn::new(3, 0, 1).with_writes(vec![page_key(20)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&t1];
let writers: Vec<&dyn ActiveTxnView> = vec![&t3];
let result = ssi_validate_and_publish(
t2,
CommitSeq::new(1),
CommitSeq::new(5),
&[page_key(20)],
&[page_key(10)],
&readers,
&writers,
&[],
&[],
false,
);
assert!(
result.is_err(),
"conservative rule: abort even with all active"
);
}
#[test]
fn test_false_positive_bounded() {
let mut commits = 0_u32;
let mut aborts = 0_u32;
for i in 0..100_u64 {
let txn = TxnToken::new(TxnId::new(i + 1).unwrap(), TxnEpoch::new(0));
#[allow(clippy::cast_possible_truncation)]
let pg = (i as u32) * 2 + 1;
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(i + 2),
&[page_key(pg)],
&[page_key(pg + 1)],
&[], &[], &[],
&[],
false,
);
match result {
Ok(_) => commits += 1,
Err(_) => aborts += 1,
}
}
assert_eq!(aborts, 0, "no false positives with non-overlapping writes");
assert_eq!(commits, 100);
}
#[test]
#[allow(clippy::redundant_clone, clippy::cloned_ref_to_slice_refs)]
fn test_write_skew_prevented() {
let t1_token = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let t2_token = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let page_a = page_key(10);
let page_b = page_key(20);
let _t1_view = MockActiveTxn::new(1, 0, 1).with_reads(vec![page_a.clone(), page_b.clone()]);
let t2_view = MockActiveTxn::new(2, 0, 1).with_reads(vec![page_a.clone(), page_b.clone()]);
let t2_readers: Vec<&dyn ActiveTxnView> = vec![&t2_view];
let result_t1 = ssi_validate_and_publish(
t1_token,
CommitSeq::new(1),
CommitSeq::new(2),
&[page_a.clone(), page_b.clone()], &[page_a.clone()], &t2_readers, &[],
&[],
&[],
false,
);
let ok_t1 = result_t1.expect("T1 should commit (only incoming)");
assert!(ok_t1.ssi_state.has_in_rw);
let reader_t1 = CommittedReaderInfo {
token: t1_token,
begin_seq: CommitSeq::new(1),
commit_seq: CommitSeq::new(2),
had_in_rw: ok_t1.ssi_state.has_in_rw,
pages: vec![PageNumber::new(10).unwrap(), PageNumber::new(20).unwrap()],
};
let writer_t1 = CommittedWriterInfo {
token: t1_token,
commit_seq: CommitSeq::new(2),
had_out_rw: ok_t1.ssi_state.has_out_rw,
pages: vec![PageNumber::new(10).unwrap()],
};
let result_t2 = ssi_validate_and_publish(
t2_token,
CommitSeq::new(1),
CommitSeq::new(3),
&[page_a.clone(), page_b.clone()], &[page_b], &[],
&[],
&[reader_t1], &[writer_t1], false,
);
assert!(result_t2.is_err(), "write skew must be prevented");
}
#[test]
fn test_concurrent_inserts_different_pages_no_abort() {
let t1 = TxnToken::new(TxnId::new(1).unwrap(), TxnEpoch::new(0));
let _t2 = TxnToken::new(TxnId::new(2).unwrap(), TxnEpoch::new(0));
let t2_view = MockActiveTxn::new(2, 0, 1)
.with_reads(vec![page_key(10)])
.with_writes(vec![page_key(10)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&t2_view];
let writers: Vec<&dyn ActiveTxnView> = vec![&t2_view];
let result = ssi_validate_and_publish(
t1,
CommitSeq::new(1),
CommitSeq::new(2),
&[page_key(5)],
&[page_key(5)],
&readers,
&writers,
&[],
&[],
false,
);
result.expect("different pages → no conflict → both commit");
}
#[test]
fn test_phantom_batch_insert_scan_conflict_prevented() {
let t_scan = TxnToken::new(TxnId::new(301).unwrap(), TxnEpoch::new(0));
let t_insert = TxnToken::new(TxnId::new(302).unwrap(), TxnEpoch::new(0));
let range_witness = page_key(900);
let aggregate_page = page_key(901);
let active_insert = MockActiveTxn::new(302, 0, 1)
.with_reads(vec![aggregate_page.clone()])
.with_writes(vec![range_witness.clone()]);
let readers_for_scan: Vec<&dyn ActiveTxnView> = vec![&active_insert];
let writers_for_scan: Vec<&dyn ActiveTxnView> = vec![&active_insert];
let scan_result = ssi_validate_and_publish(
t_scan,
CommitSeq::new(1),
CommitSeq::new(2),
std::slice::from_ref(&range_witness),
std::slice::from_ref(&aggregate_page),
&readers_for_scan,
&writers_for_scan,
&[],
&[],
false,
);
assert!(
scan_result.is_err(),
"scan+aggregate transaction should abort under phantom-style cycle"
);
let insert_result = ssi_validate_and_publish(
t_insert,
CommitSeq::new(1),
CommitSeq::new(3),
std::slice::from_ref(&aggregate_page),
std::slice::from_ref(&range_witness),
&[],
&[],
&[],
&[],
false,
);
assert!(
insert_result.is_ok(),
"one side of the phantom-style cycle must still commit"
);
}
#[test]
#[allow(clippy::too_many_lines)]
fn test_adversarial_three_way_cycle_breaks_with_single_abort() {
let begin_seq = CommitSeq::new(1);
let t1 = TxnToken::new(TxnId::new(401).unwrap(), TxnEpoch::new(0));
let t2 = TxnToken::new(TxnId::new(402).unwrap(), TxnEpoch::new(0));
let t3 = TxnToken::new(TxnId::new(403).unwrap(), TxnEpoch::new(0));
let v1 = MockActiveTxn::new(401, 0, begin_seq.get())
.with_reads(vec![page_key(1000)])
.with_writes(vec![page_key(1001)]);
let v2 = MockActiveTxn::new(402, 0, begin_seq.get())
.with_reads(vec![page_key(1001)])
.with_writes(vec![page_key(1002)]);
let v3 = MockActiveTxn::new(403, 0, begin_seq.get())
.with_reads(vec![page_key(1002)])
.with_writes(vec![page_key(1000)]);
let mut committed_readers = Vec::new();
let mut committed_writers = Vec::new();
let mut commits = 0_u32;
let mut aborts = 0_u32;
let readers_t1: Vec<&dyn ActiveTxnView> = vec![&v2];
let writers_t1: Vec<&dyn ActiveTxnView> = vec![&v3];
let t1_res = ssi_validate_and_publish(
t1,
begin_seq,
CommitSeq::new(2),
&v1.reads,
&v1.writes,
&readers_t1,
&writers_t1,
&committed_readers,
&committed_writers,
false,
);
match t1_res {
Ok(ok) => {
commits += 1;
committed_readers.push(CommittedReaderInfo {
token: t1,
begin_seq,
commit_seq: CommitSeq::new(2),
had_in_rw: ok.ssi_state.has_in_rw,
pages: keys_to_pages(&v1.reads),
});
committed_writers.push(CommittedWriterInfo {
token: t1,
commit_seq: CommitSeq::new(2),
had_out_rw: ok.ssi_state.has_out_rw,
pages: keys_to_pages(&v1.writes),
});
}
Err(_) => aborts += 1,
}
let readers_t2: Vec<&dyn ActiveTxnView> = vec![&v3];
let writers_t2: Vec<&dyn ActiveTxnView> = vec![&v3];
let t2_res = ssi_validate_and_publish(
t2,
begin_seq,
CommitSeq::new(3),
&v2.reads,
&v2.writes,
&readers_t2,
&writers_t2,
&committed_readers,
&committed_writers,
false,
);
match t2_res {
Ok(ok) => {
commits += 1;
committed_readers.push(CommittedReaderInfo {
token: t2,
begin_seq,
commit_seq: CommitSeq::new(3),
had_in_rw: ok.ssi_state.has_in_rw,
pages: keys_to_pages(&v2.reads),
});
committed_writers.push(CommittedWriterInfo {
token: t2,
commit_seq: CommitSeq::new(3),
had_out_rw: ok.ssi_state.has_out_rw,
pages: keys_to_pages(&v2.writes),
});
}
Err(_) => aborts += 1,
}
let t3_res = ssi_validate_and_publish(
t3,
begin_seq,
CommitSeq::new(4),
&v3.reads,
&v3.writes,
&[],
&[],
&committed_readers,
&committed_writers,
false,
);
match t3_res {
Ok(_) => commits += 1,
Err(_) => aborts += 1,
}
assert_eq!(commits + aborts, 3);
assert_eq!(aborts, 1, "exactly one abort should break the 3-cycle");
}
#[test]
#[allow(clippy::too_many_lines)]
fn test_100_writer_adversarial_schedule_with_serialization_checker() {
struct StressTxn {
token_id: u64,
token: TxnToken,
reads: Vec<WitnessKey>,
writes: Vec<WitnessKey>,
view: MockActiveTxn,
}
let begin_seq = CommitSeq::new(1);
let mut txns = Vec::new();
let mut next_id = 1_u64;
for pair in 0..10_u32 {
let base = 2000_u32 + pair * 10;
let reads = vec![page_key(base), page_key(base + 1)];
let token_a = TxnToken::new(TxnId::new(next_id).unwrap(), TxnEpoch::new(0));
let view_a = MockActiveTxn::new(next_id, 0, begin_seq.get())
.with_reads(reads.clone())
.with_writes(vec![page_key(base)]);
txns.push(StressTxn {
token_id: next_id,
token: token_a,
reads: reads.clone(),
writes: vec![page_key(base)],
view: view_a,
});
next_id += 1;
let token_b = TxnToken::new(TxnId::new(next_id).unwrap(), TxnEpoch::new(0));
let view_b = MockActiveTxn::new(next_id, 0, begin_seq.get())
.with_reads(reads.clone())
.with_writes(vec![page_key(base + 1)]);
txns.push(StressTxn {
token_id: next_id,
token: token_b,
reads,
writes: vec![page_key(base + 1)],
view: view_b,
});
next_id += 1;
}
while txns.len() < 100 {
let disjoint = 5000_u32 + u32::try_from(txns.len()).unwrap();
let token = TxnToken::new(TxnId::new(next_id).unwrap(), TxnEpoch::new(0));
let reads = vec![page_key(disjoint)];
let writes = vec![page_key(disjoint + 10_000)];
let view = MockActiveTxn::new(next_id, 0, begin_seq.get())
.with_reads(reads.clone())
.with_writes(writes.clone());
txns.push(StressTxn {
token_id: next_id,
token,
reads,
writes,
view,
});
next_id += 1;
}
let mut committed_ids = HashSet::new();
let mut committed_readers = Vec::new();
let mut committed_writers = Vec::new();
let mut abort_count = 0_u32;
for idx in 0..txns.len() {
let current = &txns[idx];
let active_tail = &txns[idx + 1..];
let active_readers: Vec<&dyn ActiveTxnView> = active_tail
.iter()
.map(|txn| &txn.view as &dyn ActiveTxnView)
.collect();
let active_writers: Vec<&dyn ActiveTxnView> = active_tail
.iter()
.map(|txn| &txn.view as &dyn ActiveTxnView)
.collect();
let commit_seq = CommitSeq::new(u64::try_from(idx).unwrap() + 2);
let result = ssi_validate_and_publish(
current.token,
begin_seq,
commit_seq,
¤t.reads,
¤t.writes,
&active_readers,
&active_writers,
&committed_readers,
&committed_writers,
false,
);
match result {
Ok(ok) => {
committed_ids.insert(current.token_id);
committed_readers.push(CommittedReaderInfo {
token: current.token,
begin_seq,
commit_seq,
had_in_rw: ok.ssi_state.has_in_rw,
pages: keys_to_pages(¤t.reads),
});
committed_writers.push(CommittedWriterInfo {
token: current.token,
commit_seq,
had_out_rw: ok.ssi_state.has_out_rw,
pages: keys_to_pages(¤t.writes),
});
}
Err(_) => abort_count += 1,
}
}
let total = u32::try_from(txns.len()).unwrap();
assert_eq!(
u32::try_from(committed_ids.len()).unwrap() + abort_count,
total,
"all 100 commit attempts must complete (no deadlock/livelock)"
);
let mut mandatory_aborts = 0_u32;
for pair in 0..10_u64 {
let tx_a = pair * 2 + 1;
let tx_b = pair * 2 + 2;
let a_committed = committed_ids.contains(&tx_a);
let b_committed = committed_ids.contains(&tx_b);
assert!(
a_committed ^ b_committed,
"conflict pair ({tx_a},{tx_b}) must commit exactly one member"
);
mandatory_aborts += 1;
}
for id in 21_u64..=100_u64 {
assert!(
committed_ids.contains(&id),
"disjoint writer {id} should commit"
);
}
let false_positive_aborts = abort_count.saturating_sub(mandatory_aborts);
assert!(
false_positive_aborts <= 5,
"false positive aborts must stay under 5%: {false_positive_aborts}/100"
);
}
#[test]
fn test_long_running_reader_stable_snapshot_under_writer_churn() {
let long_reader =
MockActiveTxn::new(9_001, 0, 1).with_reads(vec![page_key(700), page_key(701)]);
let active_readers: Vec<&dyn ActiveTxnView> = vec![&long_reader];
let mut commits = 0_u32;
for i in 0..200_u64 {
let writer = TxnToken::new(TxnId::new(9_100 + i).unwrap(), TxnEpoch::new(0));
let write_key = if i % 2 == 0 {
page_key(700)
} else {
page_key(701)
};
let result = ssi_validate_and_publish(
writer,
CommitSeq::new(1),
CommitSeq::new(i + 2),
&[],
&[write_key],
&active_readers,
&[],
&[],
&[],
false,
);
assert!(
result.is_ok(),
"writer churn should not deadlock/abort readers"
);
commits += 1;
}
assert_eq!(commits, 200);
assert!(
!long_reader.marked.get(),
"long-running read-only snapshot must not be marked for abort"
);
}
#[test]
fn test_evidence_metrics_count_by_outcome() {
let before = ssi_evidence_metrics_snapshot();
let commit_txn = TxnToken::new(TxnId::new(90_001).unwrap(), TxnEpoch::new(0));
let commit_result = ssi_validate_and_publish(
commit_txn,
CommitSeq::new(1),
CommitSeq::new(2),
&[page_key(500)],
&[page_key(600)],
&[],
&[],
&[],
&[],
false,
);
commit_result.expect("commit decision should be recorded");
let abort_txn = TxnToken::new(TxnId::new(90_002).unwrap(), TxnEpoch::new(0));
let reader = MockActiveTxn::new(90_003, 0, 1).with_reads(vec![page_key(700)]);
let writer = MockActiveTxn::new(90_004, 0, 1).with_writes(vec![page_key(800)]);
let readers: Vec<&dyn ActiveTxnView> = vec![&reader];
let writers: Vec<&dyn ActiveTxnView> = vec![&writer];
let abort_result = ssi_validate_and_publish(
abort_txn,
CommitSeq::new(1),
CommitSeq::new(3),
&[page_key(800)],
&[page_key(700)],
&readers,
&writers,
&[],
&[],
false,
);
abort_result.expect_err("pivot abort should be recorded");
let after = ssi_evidence_metrics_snapshot();
assert!(
after.fsqlite_evidence_records_total_commit
> before.fsqlite_evidence_records_total_commit
);
assert!(
after.fsqlite_evidence_records_total_abort
> before.fsqlite_evidence_records_total_abort
);
assert!(
after.fsqlite_evidence_records_total() >= before.fsqlite_evidence_records_total() + 2
);
}
#[test]
fn test_evidence_store_queryable_by_txn_id() {
let txn = TxnToken::new(TxnId::new(90_101).unwrap(), TxnEpoch::new(0));
let result = ssi_validate_and_publish(
txn,
CommitSeq::new(1),
CommitSeq::new(2),
&[page_key(901)],
&[page_key(902)],
&[],
&[],
&[],
&[],
false,
);
result.expect("commit should succeed");
let rows = ssi_evidence_query(&SsiDecisionQuery {
txn_id: Some(txn.id.get()),
..SsiDecisionQuery::default()
});
assert!(
!rows.is_empty(),
"evidence ledger should contain row for txn_id"
);
let last = rows.last().unwrap();
assert_eq!(last.txn.id.get(), txn.id.get());
assert!(last.decision_id > 0, "decision_id must be populated");
}
}